diff --git a/README.md b/README.md index 2cfd4b150..6dcf3a17c 100644 --- a/README.md +++ b/README.md @@ -494,7 +494,7 @@ The current scope split is: | Discovery and health | `hosts`, `apps`, `source_ips`, `status`, `stats`, `ingest_rate`, `silent_hosts`, `clock_skew` | | Analytics and correlation | `timeline`, `patterns`, `anomalies`, `compare`, `correlate`, `topic_correlate`, `similar_incidents`, `recurring_error_comparison`, `incident_context` | | Fleet and topology | `map`, `host_state`, `fleet_state`, `correlate_state`, `graph`, `compose_status`, `compose_doctor` | -| AI sessions and scoped evidence | `sessions`, `search_sessions`, `evidence_scope`, `abuse`, `abuse_incidents`, `abuse_investigate`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_tools`, `list_ai_projects` | +| AI sessions and scoped evidence | `sessions`, `search_sessions`, `session_investigate`, `evidence_scope`, `abuse`, `abuse_incidents`, `abuse_investigate`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_tools`, `list_ai_projects` | | AI operational events | `skill_events`, `skill_incidents`, `skill_investigate`, `mcp_events`, `mcp_incidents`, `mcp_investigate`, `hook_events`, `hook_incidents`, `hook_investigate` | | Errors and administration | `unaddressed_errors`, `ack_error`, `unack_error`, `notifications_recent`, `notifications_test`, `file_tails`, `llm_invocations`, `artifact_evidence_record` | | Reference | `help` | diff --git a/docs/INVENTORY.md b/docs/INVENTORY.md index 8caff63c4..4ae18ecaa 100644 --- a/docs/INVENTORY.md +++ b/docs/INVENTORY.md @@ -1,7 +1,7 @@ --- title: "Component Inventory -- cortex" created: 2026-04-04 -updated: 2026-09-27 +updated: 2026-09-30 --- # Component Inventory -- cortex @@ -86,6 +86,7 @@ that registry by `src/mcp/schemas.rs::tool_definitions()`. | `status` | Lightweight runtime status: DB health, queue/backpressure state, listener/writer counters, OTLP counters | no | | `sessions` | AI transcript sessions grouped by project/tool/session/host | no | | `search_sessions` | Ranked grouped session search | no | +| `session_investigate` | Bounded evidence bundle rooted at one AI session | no | | `evidence_scope` | Historical Agent Observatory evidence for a Git branch or worktree | no | | `abuse` | Abuse-term detector with same-session context | no | | `abuse_incidents` | Groups abuse hits into scored incident candidates | no | diff --git a/docs/mcp/SCHEMA.md b/docs/mcp/SCHEMA.md index 906ad90c0..432bf83d7 100644 --- a/docs/mcp/SCHEMA.md +++ b/docs/mcp/SCHEMA.md @@ -1,7 +1,7 @@ --- title: "Tool Schema Documentation -- cortex" created: "2026-07-30" -updated: "2026-08-04" +updated: 2026-09-30 --- # Tool Schema Documentation -- cortex @@ -45,6 +45,7 @@ selects one of the actions below. The mechanically generated current count is in | `apps` | `cortex:read` | cheap | Distinct application names with counts | | `sessions` | `cortex:read` | cheap | AI transcript session inventory | | `search_sessions` | `cortex:read` | cheap | FTS5 search over AI transcript sessions | +| `session_investigate` | `cortex:read` | expensive | Bounded evidence bundle rooted at one AI session | | `evidence_scope` | `cortex:read` | moderate | Historical Agent Observatory evidence for a Git branch or worktree | | `abuse` | `cortex:read` | moderate | Abuse-term hits with same-session context | | `abuse_incidents` | `cortex:read` | moderate | Grouped abuse incident candidates | @@ -175,12 +176,13 @@ boundary. | `query` | `search`, `search_sessions`, `correlate`, `similar_incidents` | | `hostname` | `search`, `filter`, `tail`, `correlate`, `host_state`, `ai_correlate`, `apps`, `sessions`, `timeline`, `patterns`, `context`, `similar_incidents`, `incident_context` | | `host_id` | Authoritative heartbeat identity for `host_state` | -| `host` | Optional host_id-or-hostname filter for `correlate_state` | +| `host` | Optional host_id-or-hostname filter for `correlate_state`; exact session identity qualifier for `session_investigate` | | `reference_time` | Required window center for `correlate_state`; for `correlate`, required unless `query` is given (then derived from an AI-session search) | | `source_ip` | `search`, `filter`, `tail`, `correlate`, `ai_correlate` | | `source_kind` | `filter` only; aliases Docker, file-tail, command-history, shell-history, transcript, and AI-tool rows | | `project` | `filter`, `sessions`, `search_sessions`, `abuse`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_tools` | | `tool` | `filter`, `sessions`, `search_sessions`, `abuse`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_projects` | +| `session_id` | Required for `session_investigate`; optional exact native identity filter for `sessions` | | `branch`, `worktree` | `evidence_scope`; at least one is required by service validation | | `kinds`, `include_payload`, `after_id` | `evidence_scope` filtering and durable pagination | | `session_id` | `filter`, `ai_correlate` | diff --git a/docs/mcp/TESTS.md b/docs/mcp/TESTS.md index 3879f60af..eb88b4093 100644 --- a/docs/mcp/TESTS.md +++ b/docs/mcp/TESTS.md @@ -1,7 +1,7 @@ --- title: "Testing Guide -- cortex" created: "2026-07-30" -updated: 2026-09-27 +updated: 2026-09-30 --- # Testing Guide -- cortex @@ -77,7 +77,7 @@ the repo-local debug binary at `target/debug/cortex`, so repo-local builds do not require an installed shell binary. Action registry covered by live/script references: `search`, `filter`, `tail`, `errors`, -`hosts`, `map`, `host_state`, `fleet_state`, `correlate_state`, `topic_correlate`, `sessions`, `search_sessions`, `evidence_scope`, `abuse`, `abuse_incidents`, `abuse_investigate`, `ai_correlate`, `usage_blocks`, `project_context`, +`hosts`, `map`, `host_state`, `fleet_state`, `correlate_state`, `topic_correlate`, `sessions`, `search_sessions`, `session_investigate`, `evidence_scope`, `abuse`, `abuse_incidents`, `abuse_investigate`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_tools`, `list_ai_projects`, `correlate`, `stats`, `status`, `apps`, `source_ips`, `timeline`, `patterns`, `context`, `get`, `ingest_rate`, `silent_hosts`, `clock_skew`, `anomalies`, `compare`, `compose_status`, @@ -195,7 +195,7 @@ curl -s -X POST http://localhost:3100/mcp \ ## Testing checklist - [ ] **All actions return expected shape** -- cortex search, cortex tail, cortex errors, cortex hosts, cortex host_state, cortex sessions, cortex correlate, cortex stats, cortex status, cortex help -- [ ] **AI session analytics and scoped lifecycle evidence return expected shape and seeded rows** -- cortex search_sessions, cortex evidence_scope, cortex abuse, cortex sessions_correlate, cortex usage_blocks, cortex project_context, cortex list_ai_tools, cortex list_ai_projects +- [ ] **AI session analytics and scoped lifecycle evidence return expected shape and seeded rows** -- cortex search_sessions, cortex session_investigate, cortex evidence_scope, cortex abuse, cortex sessions_correlate, cortex usage_blocks, cortex project_context, cortex list_ai_tools, cortex list_ai_projects - [ ] **Auth: valid token** -- 200 with correct Bearer token - [ ] **Auth: invalid token** -- 401 Unauthorized - [ ] **Auth: no token when required** -- 401 Unauthorized diff --git a/docs/mcp/TOOLS.md b/docs/mcp/TOOLS.md index efda909c6..c0cdf01c1 100644 --- a/docs/mcp/TOOLS.md +++ b/docs/mcp/TOOLS.md @@ -1,7 +1,7 @@ --- title: "MCP Tools Reference -- cortex" created: "2026-07-30" -updated: "2026-08-04" +updated: 2026-09-30 --- # MCP Tools Reference -- cortex @@ -24,6 +24,7 @@ cortex exposes one MCP tool named `cortex`. The required | `correlate_state` | Correlate logs with heartbeat window summaries around a reference time | | `sessions` | AI transcript sessions by project | | `search_sessions` | Ranked grouped session search | +| `session_investigate` | Bounded evidence bundle rooted at one AI session | | `evidence_scope` | Historical Agent Observatory evidence for a Git branch or worktree | | `abuse` | Abuse hits in AI transcripts with same-session context | | `abuse_incidents` | Groups abuse hits into scored incident candidates | @@ -182,6 +183,21 @@ Required arguments: `action = "search_sessions"`, `query` Optional arguments: `project`, `tool`, `from`, `to`, `limit`. +## cortex session_investigate + +Build an evidence bundle for one exact AI session. Required arguments: +`action = "session_investigate"`, `session_id`. Optional arguments: `tool`, +`project`, `host`, `limit` (1–200), and `severity_min`. Provide identity qualifiers +when a native session ID is shared by several transcripts. + +The response contains `metadata` and `result`, including rendered transcript, +scoped graph correlation, Observatory lineage, tools/skills/hooks, artifacts, +notifications, incidents, and retained deletion lineage. Each section reports +truncation; `partial_reasons` records omitted evidence. The complete JSON envelope +is capped at 64 KiB and wall time at two seconds; exhausted wall time returns a +retryable busy error. Transcript references are passive claims with +`verified = false`. + ## cortex evidence_scope Return a bounded historical page of Agent Observatory events associated with an diff --git a/packages/cortex-rmcp/README.md b/packages/cortex-rmcp/README.md index 2cfd4b150..6dcf3a17c 100644 --- a/packages/cortex-rmcp/README.md +++ b/packages/cortex-rmcp/README.md @@ -494,7 +494,7 @@ The current scope split is: | Discovery and health | `hosts`, `apps`, `source_ips`, `status`, `stats`, `ingest_rate`, `silent_hosts`, `clock_skew` | | Analytics and correlation | `timeline`, `patterns`, `anomalies`, `compare`, `correlate`, `topic_correlate`, `similar_incidents`, `recurring_error_comparison`, `incident_context` | | Fleet and topology | `map`, `host_state`, `fleet_state`, `correlate_state`, `graph`, `compose_status`, `compose_doctor` | -| AI sessions and scoped evidence | `sessions`, `search_sessions`, `evidence_scope`, `abuse`, `abuse_incidents`, `abuse_investigate`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_tools`, `list_ai_projects` | +| AI sessions and scoped evidence | `sessions`, `search_sessions`, `session_investigate`, `evidence_scope`, `abuse`, `abuse_incidents`, `abuse_investigate`, `ai_correlate`, `usage_blocks`, `project_context`, `list_ai_tools`, `list_ai_projects` | | AI operational events | `skill_events`, `skill_incidents`, `skill_investigate`, `mcp_events`, `mcp_incidents`, `mcp_investigate`, `hook_events`, `hook_incidents`, `hook_investigate` | | Errors and administration | `unaddressed_errors`, `ack_error`, `unack_error`, `notifications_recent`, `notifications_test`, `file_tails`, `llm_invocations`, `artifact_evidence_record` | | Reference | `help` | diff --git a/plugins/cortex/skills/using-cortex/references/operations.md b/plugins/cortex/skills/using-cortex/references/operations.md index cc2fd7676..3393bbc49 100644 --- a/plugins/cortex/skills/using-cortex/references/operations.md +++ b/plugins/cortex/skills/using-cortex/references/operations.md @@ -23,6 +23,7 @@ A single MCP tool, `mcp__cortex__cortex`, dispatches on a required `action` argu | `correlate_state` | Correlate logs with heartbeat summaries around a reference time | | `sessions` | AI transcript sessions by project | | `search_sessions` | Ranked grouped session search | +| `session_investigate` | Bounded evidence bundle rooted at one AI session | | `evidence_scope` | Historical Agent Observatory evidence for a Git branch or worktree | | `abuse` | Abuse hits in AI transcripts with same-session context | | `abuse_incidents` | Groups abuse hits into scored incident candidates | diff --git a/scripts/test-agent-memory-symlinks.sh b/scripts/test-agent-memory-symlinks.sh index 55341b0f5..f7718a769 100644 --- a/scripts/test-agent-memory-symlinks.sh +++ b/scripts/test-agent-memory-symlinks.sh @@ -3,6 +3,11 @@ set -euo pipefail script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" checker="$script_dir/check-agent-memory-symlinks.sh" +# Hooks export repository-local Git variables. A fixture's `git -C ... init` +# must select its own directory rather than reinitialize the live repository. +while IFS= read -r variable; do + unset "$variable" +done < <(git rev-parse --local-env-vars) fixture="$(mktemp -d)" trap 'rm -rf "$fixture"' EXIT git -C "$fixture" init -q @@ -98,4 +103,21 @@ expect_fail "force-staged private override" git -C "$fixture" rm --cached -q AGENTS.override.md expect_pass "private override removed from index" +if [[ "${CORTEX_INSTRUCTION_FIXTURE_CHILD:-}" != "1" ]]; then + hook_fixture="$fixture/hook-environment" + mkdir "$hook_fixture" + git -C "$hook_fixture" init -q + printf 'Preserve the hook repository index\n' > "$hook_fixture/sentinel" + git -C "$hook_fixture" add sentinel + cp "$hook_fixture/.git/config" "$fixture/config-before" + cp "$hook_fixture/.git/index" "$fixture/index-before" + env CORTEX_INSTRUCTION_FIXTURE_CHILD=1 \ + GIT_DIR="$hook_fixture/.git" GIT_WORK_TREE="$hook_fixture" \ + GIT_INDEX_FILE="$hook_fixture/.git/index" \ + bash "$script_dir/test-agent-memory-symlinks.sh" + cmp "$fixture/config-before" "$hook_fixture/.git/config" + cmp "$fixture/index-before" "$hook_fixture/.git/index" + checks=$((checks + 1)) +fi + printf "[agent-memory-test] OK — %s regression cases\n" "$checks" diff --git a/src/api.rs b/src/api.rs index b9206c72f..58caf2e10 100644 --- a/src/api.rs +++ b/src/api.rs @@ -736,6 +736,7 @@ async fn observatory_runs( since: query.since, until: query.until, active_only: query.active_only.unwrap_or(false), + ..Default::default() }, query.cursor, query.limit.unwrap_or(50), diff --git a/src/app.rs b/src/app.rs index c95421440..4f1572f7b 100644 --- a/src/app.rs +++ b/src/app.rs @@ -237,6 +237,11 @@ pub use models::{ ServiceJournalEntry, ServiceLogsRequest, ServiceLogsResponse, + SessionExternalReference, + SessionExternalReferenceKind, + SessionInvestigateRequest, + SessionInvestigateResponse, + SessionObservatoryEvidence, SeverityCount, SilentHostsRequest, SilentHostsResponse, diff --git a/src/app/models/ai_incidents.rs b/src/app/models/ai_incidents.rs index cf9fed88a..3d9e1b50b 100644 --- a/src/app/models/ai_incidents.rs +++ b/src/app/models/ai_incidents.rs @@ -419,6 +419,7 @@ pub struct GraphSessionCorrelation { /// discover related hosts; `false` for the time-windowed fallback (session /// not yet projected into the graph). pub used_graph: bool, + pub session_entity_keys: Vec, pub discovered_hosts: Vec, pub discovered_entities: Vec, pub logs: Vec, diff --git a/src/app/models/ai_sessions.rs b/src/app/models/ai_sessions.rs index 6f1448a4f..ad1fc23d1 100644 --- a/src/app/models/ai_sessions.rs +++ b/src/app/models/ai_sessions.rs @@ -5,6 +5,7 @@ use super::*; pub struct ListSessionsRequest { pub project: Option, pub tool: Option, + pub session_id: Option, pub host: Option, pub since: Option, pub until: Option, diff --git a/src/app/models/investigation.rs b/src/app/models/investigation.rs index 7eb736ffe..2e2b4d162 100644 --- a/src/app/models/investigation.rs +++ b/src/app/models/investigation.rs @@ -177,6 +177,109 @@ pub struct AskInvestigationResponse { pub logs: Vec, } +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct SessionInvestigateRequest { + pub session_id: String, + pub tool: Option, + pub project: Option, + pub host: Option, + pub limit: Option, + pub severity_min: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)] +#[serde(rename_all = "snake_case")] +pub enum SessionExternalReferenceKind { + GithubPullRequest, + GithubIssue, + GithubReference, + LinearIssue, + Url, + CommitSha, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)] +pub struct SessionExternalReference { + pub kind: SessionExternalReferenceKind, + pub value: String, + pub source_position: Option, + pub evidence_kind: String, + pub trust_level: String, + pub verified: bool, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SessionRetentionLineageEntry { + pub log_id: i64, + pub hostname: String, + pub app_name: Option, + pub severity: String, + pub ai_project: Option, + pub ai_tool: Option, + pub ai_session_id: Option, + pub retained_until_epoch: i64, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct SessionSourceEvidenceSummary { + pub count: usize, + pub log_ids: Vec, + pub truncated: bool, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct SessionObservatoryEvidence { + pub runs: Vec, + pub ambiguous_run: bool, + pub parent_run: Option, + pub previous_run: Option, + pub actors: Vec, + pub actors_truncated: bool, + pub related_runs: Vec, + pub related_runs_truncated: bool, + pub repository: Option, + pub worktree: Option, + pub commits: Vec, + pub commits_truncated: bool, + pub events: Vec, + pub spans: Vec, + pub metrics: Vec, + pub events_truncated: bool, + pub spans_truncated: bool, + pub metrics_truncated: bool, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SessionInvestigateResponse { + pub session: AiSessionEntry, + pub transcript: Vec, + pub transcript_has_more: bool, + pub correlation: Option, + pub graph_neighborhood: Option, + pub skill_events: Vec, + pub skill_events_truncated: bool, + pub mcp_events: Vec, + pub mcp_events_truncated: bool, + pub hook_events: Vec, + pub hook_events_truncated: bool, + pub artifact_evidence: Vec, + pub artifact_evidence_truncated: bool, + pub observatory: SessionObservatoryEvidence, + pub incident_context: IncidentContextResponse, + pub notifications: Vec, + pub notifications_truncated: bool, + pub retention_lineage: Vec, + pub retention_lineage_truncated: bool, + pub related_sessions: Vec, + pub related_sessions_truncated: bool, + pub external_references: Vec, + pub external_references_truncated: bool, + pub source_counts: BTreeMap, + pub source_evidence: BTreeMap, + pub partial_reasons: Vec, +} + pub fn app_entity_summary(entity: &GraphEntity) -> AppEntitySummary { AppEntitySummary { id: entity.id, @@ -288,22 +391,22 @@ pub fn app_graph_from_around_response(around: &GraphAroundResponse) -> AppGraphR } pub fn safe_passive_text(input: &str, max_chars: usize) -> String { - let mut out = input + static BEARER_VALUE: std::sync::LazyLock = std::sync::LazyLock::new(|| { + regex::Regex::new(r"(?i)\bBearer\s+[^\s,;]+").expect("static bearer redaction pattern") + }); + let scrubbed = crate::receiver::enrichment::scrub_ai_message(input, None); + let scrubbed = BEARER_VALUE.replace_all(&scrubbed, "[REDACTED]"); + let mut out = crate::assessment::redact_secrets(&scrubbed) .chars() .filter(|ch| !ch.is_control() || matches!(ch, '\n' | '\t')) .collect::(); - for marker in [ - "sk-proj-", - "Bearer ", - "password=", - "token=", - "CORTEX_API_TOKEN=", - ] { - out = out.replace(marker, "[redacted]"); - } if out.chars().count() > max_chars { out = out.chars().take(max_chars).collect::(); out.push_str("..."); } out } + +#[cfg(test)] +#[path = "investigation_tests.rs"] +mod tests; diff --git a/src/app/models/investigation_tests.rs b/src/app/models/investigation_tests.rs new file mode 100644 index 000000000..8a98fd5b5 --- /dev/null +++ b/src/app/models/investigation_tests.rs @@ -0,0 +1,25 @@ +use super::*; + +#[test] +fn passive_text_redacts_secret_values_before_truncation() { + for text in [ + "password=passive-secret-value", + "Authorization: Bearer passive-secret-value", + "Bearer passive-secret-value", + "CORTEX_API_TOKEN=passive-secret-value", + "token=passive-secret-value", + "sk-proj-passive-secret-value-0123456789abcdefgh", + ] { + let output = safe_passive_text(text, 500); + assert!( + !output.contains("passive-secret-value"), + "secret value survived redaction" + ); + let short = safe_passive_text(text, 16); + assert!( + !short.contains("passive-secret"), + "truncation exposed a partial secret" + ); + } + assert_eq!(safe_passive_text("safe\u{1b} text", 500), "safe text"); +} diff --git a/src/app/services.rs b/src/app/services.rs index 947151ce2..5d3ea59c1 100644 --- a/src/app/services.rs +++ b/src/app/services.rs @@ -126,6 +126,11 @@ mod mcp_backfill; mod mcp_events; mod mcp_incidents; mod rag; +mod session_graph_correlation; +mod session_investigation; +mod session_investigation_budget; +mod session_investigation_sections; +mod session_investigation_support; mod session_pages; mod skill_assessment; mod skill_backfill; diff --git a/src/app/services/ai.rs b/src/app/services/ai.rs index 77691ecc3..5ba28f3c4 100644 --- a/src/app/services/ai.rs +++ b/src/app/services/ai.rs @@ -1,67 +1,6 @@ use super::*; -/// Log fan-out cap for the graph-anchored session lane of `ai_correlate`. -/// Clamped again to `[1, 1000]` inside `db::correlate_session_graph`. -const GRAPH_SESSION_LOG_LIMIT: usize = 500; - -/// Shape the DB-layer `SessionGraphInputs` into the API response, classifying -/// each row into a source lane (`agent_command` / `shell_history` / -/// `graph:host:`) and counting the agent-command and shell-history lanes. -/// Heartbeat summaries are filtered to the discovered hosts. Returns `None` when -/// the session has no rows at all (empty bounds). -fn build_graph_session_correlation( - session_id: String, - inputs: db::SessionGraphInputs, - summaries: Vec, -) -> Option { - let (session_start, session_end) = inputs.bounds?; - - let truncated = inputs.logs.len() >= GRAPH_SESSION_LOG_LIMIT; - let mut agent_command_count = 0usize; - let mut shell_history_count = 0usize; - let logs: Vec = inputs - .logs - .into_iter() - .map(|entry| { - let source_kind = row_source_kind(&entry); - let discovery = if entry.source_ip.starts_with("agent-command://") { - agent_command_count += 1; - "agent_command".to_string() - } else if source_kind.as_deref() == Some("shell-history") { - shell_history_count += 1; - "shell_history".to_string() - } else { - format!("graph:host:{}", entry.hostname) - }; - CorrelatedLogRow { - entry: entry.into(), - source_kind, - discovery, - } - }) - .collect(); - - let discovered: std::collections::HashSet<&str> = - inputs.discovered_hosts.iter().map(String::as_str).collect(); - let heartbeat_summaries: Vec = summaries - .into_iter() - .filter(|s| discovered.contains(s.hostname.as_str())) - .collect(); - - Some(GraphSessionCorrelation { - session_id, - session_start, - session_end, - used_graph: inputs.used_graph, - discovered_hosts: inputs.discovered_hosts, - discovered_entities: inputs.discovered_entities, - logs, - agent_command_count, - shell_history_count, - heartbeat_summaries, - truncated, - }) -} +use super::session_graph_correlation::{GRAPH_SESSION_LOG_LIMIT, build_graph_session_correlation}; impl CortexService { pub async fn list_sessions( @@ -73,10 +12,11 @@ impl CortexService { // The unbounded (no time-window) path reads from the periodically // refreshed rollup; expose its staleness so callers know the `as_of`. // Time-windowed queries run live, so no staleness applies. - let unbounded = from.is_none() && to.is_none(); + let unbounded = from.is_none() && to.is_none() && req.session_id.is_none(); let params = db::ListAiSessionsParams { ai_project: req.project, ai_tool: req.tool, + ai_session_id: req.session_id, host: req.host, since: from, until: to, diff --git a/src/app/services/ai_correlate_tests.rs b/src/app/services/ai_correlate_tests.rs index ceedc63fb..f0b347cff 100644 --- a/src/app/services/ai_correlate_tests.rs +++ b/src/app/services/ai_correlate_tests.rs @@ -59,9 +59,11 @@ fn row_source_kind_parses_metadata() { fn build_graph_session_correlation_classifies_lanes_and_filters_heartbeats() { let inputs = db::SessionGraphInputs { bounds: Some(("2026-01-01T00:00:00Z".into(), "2026-01-01T00:10:00Z".into())), + session_entity_keys: vec!["cortex:claude:s1".into()], discovered_hosts: vec!["devhost".into()], discovered_entities: vec!["devhost".into(), "cortex".into()], used_graph: true, + source_fields_truncated: false, logs: vec![ db_log( "agent-command://devhost/claude/s1", diff --git a/src/app/services/session_graph_correlation.rs b/src/app/services/session_graph_correlation.rs new file mode 100644 index 000000000..3bf810ffd --- /dev/null +++ b/src/app/services/session_graph_correlation.rs @@ -0,0 +1,103 @@ +use super::*; + +/// Log fan-out cap for the graph-anchored session lane of `ai_correlate`. +/// Clamped again to `[1, 1000]` inside `db::correlate_session_graph`. +pub(super) const GRAPH_SESSION_LOG_LIMIT: usize = 500; + +/// Shape the DB-layer `SessionGraphInputs` into the API response, classifying +/// each row into a source lane (`agent_command` / `shell_history` / +/// `graph:host:`) and counting the agent-command and shell-history lanes. +/// Heartbeat summaries are filtered to the discovered hosts. Returns `None` when +/// the session has no rows at all (empty bounds). +pub(super) fn build_graph_session_correlation( + session_id: String, + inputs: db::SessionGraphInputs, + summaries: Vec, +) -> Option { + let (session_start, session_end) = inputs.bounds?; + + let truncated = inputs.source_fields_truncated || inputs.logs.len() >= GRAPH_SESSION_LOG_LIMIT; + let mut agent_command_count = 0usize; + let mut shell_history_count = 0usize; + let logs: Vec = inputs + .logs + .into_iter() + .map(|entry| { + let source_kind = row_source_kind(&entry); + let discovery = if entry.source_ip.starts_with("agent-command://") { + agent_command_count += 1; + "agent_command".to_string() + } else if source_kind.as_deref() == Some("shell-history") { + shell_history_count += 1; + "shell_history".to_string() + } else { + format!("graph:host:{}", entry.hostname) + }; + CorrelatedLogRow { + entry: entry.into(), + source_kind, + discovery, + } + }) + .collect(); + + let discovered: std::collections::HashSet<&str> = + inputs.discovered_hosts.iter().map(String::as_str).collect(); + let heartbeat_summaries: Vec = summaries + .into_iter() + .filter(|s| discovered.contains(s.hostname.as_str())) + .collect(); + + Some(GraphSessionCorrelation { + session_id, + session_start, + session_end, + used_graph: inputs.used_graph, + session_entity_keys: inputs.session_entity_keys, + discovered_hosts: inputs.discovered_hosts, + discovered_entities: inputs.discovered_entities, + logs, + agent_command_count, + shell_history_count, + heartbeat_summaries, + truncated, + }) +} + +impl CortexService { + pub(super) async fn session_graph_correlation( + &self, + session: &AiSessionEntry, + severity_min: Option<&str>, + ) -> ServiceResult> { + let levels = severity_at_or_above(severity_min.unwrap_or("info"))?; + let session_id = session.session_id.clone(); + let sid = session_id.clone(); + let scope = db::SessionGraphScope { + project: session.project.clone(), + tool: session.tool.clone(), + host: session.hostname.clone(), + severity_in: levels, + }; + let (inputs, summaries) = self + .run_db("session_investigation_graph", move |pool| { + let inputs = db::correlate_session_graph_scoped( + pool, + &sid, + &scope, + GRAPH_SESSION_LOG_LIMIT, + )?; + let summaries = match &inputs.bounds { + Some((start, end)) if inputs.used_graph => { + db::heartbeat_window_summaries(pool, start, end, None)? + } + _ => Vec::new(), + }; + Ok((inputs, summaries)) + }) + .await?; + Ok(build_graph_session_correlation( + session_id, inputs, summaries, + )) + } +} diff --git a/src/app/services/session_investigation.rs b/src/app/services/session_investigation.rs new file mode 100644 index 000000000..f4d2fa5d4 --- /dev/null +++ b/src/app/services/session_investigation.rs @@ -0,0 +1,391 @@ +use std::time::Instant; + +use super::*; + +use super::session_investigation_support::{ + build_session_source_counts, extract_session_external_references, + related_sessions_for_investigation, session_investigation_metadata, + summarize_session_source_evidence, +}; + +impl CortexService { + pub async fn session_investigate( + &self, + req: models::SessionInvestigateRequest, + ) -> ServiceResult> { + let budget = InvestigationBudget::default(); + super::session_investigation_budget::within_session_budget( + budget.max_wall_time_ms, + self.session_investigate_inner(req), + ) + .await + } + + async fn session_investigate_inner( + &self, + req: models::SessionInvestigateRequest, + ) -> ServiceResult> { + let started = Instant::now(); + let mut budget = InvestigationBudget::default(); + let session_id = req.session_id.trim().to_string(); + if session_id.is_empty() { + return Err(ServiceError::InvalidInput( + "session_id must not be empty".to_string(), + )); + } + let section_limit = req.limit.unwrap_or(100).clamp(1, 200); + // The composition has multiple independently bounded evidence lanes. + budget.max_log_rows = 1_000 + section_limit; + budget.max_evidence_rows = section_limit * 18 + 512; + + let sessions = self + .list_sessions(ListSessionsRequest { + project: req.project.clone(), + tool: req.tool.clone(), + session_id: Some(session_id.clone()), + host: req.host.clone(), + since: None, + until: None, + limit: Some(20), + offset: None, + }) + .await? + .sessions; + let mut exact = sessions + .into_iter() + .filter(|session| session.session_id == session_id) + .collect::>(); + if exact.is_empty() { + return Err(ServiceError::NotFound(format!( + "AI session not found: {session_id}" + ))); + } + if exact.len() > 1 { + return Err(ServiceError::InvalidInput( + "session_id is ambiguous; provide tool, project, and/or host".to_string(), + )); + } + let session = exact.remove(0); + + let transcript_page = self + .rendered_session_page(models::RenderedSessionPageRequest { + project: session.project.clone(), + tool: session.tool.clone(), + session_id: session.session_id.clone(), + host: session.hostname.clone(), + cursor: None, + limit: Some(section_limit), + }) + .await?; + + let correlation = self + .session_graph_correlation(&session, req.severity_min.as_deref()) + .await?; + let session_entity_keys = correlation + .as_ref() + .map(|value| value.session_entity_keys.clone()) + .unwrap_or_default(); + let session_graph_entity_ambiguous = session_entity_keys.len() > 1; + let mut graph_neighborhood_truncated = false; + let graph_neighborhood = match session_entity_keys.as_slice() { + [session_key] => { + let around = self + .graph_around(GraphAroundRequest { + mode: Some("around".to_string()), + entity_type: Some("ai_session".to_string()), + key: Some(session_key.clone()), + depth: Some(1), + limit: Some(section_limit.min(100)), + evidence_sample_limit: Some(3), + payload_budget: Some(32_768), + ..Default::default() + }) + .await?; + graph_neighborhood_truncated = around.metadata.truncated; + Some(models::app_graph_from_around_response(&around)) + } + _ => None, + }; + + let incident_context = self + .incident_context(models::IncidentContextRequest { + since: Some(session.first_seen.clone()), + until: Some(session.last_seen.clone()), + host: Some(session.hostname.clone()), + app: None, + query: None, + severity_min: req.severity_min.clone(), + limit: Some(section_limit.min(200)), + }) + .await?; + + let mut notifications = super::session_investigation_sections::session_notifications( + self, + &session, + section_limit, + ) + .await?; + let notifications_truncated = notifications.len() > section_limit as usize; + notifications.truncate(section_limit as usize); + + let retention_project = session.project.clone(); + let retention_tool = session.tool.clone(); + let retention_session_id = session.session_id.clone(); + let retention_host = session.hostname.clone(); + let mut retention_lineage = self + .run_db("session_investigate_retention_lineage", move |pool| { + let conn = pool.get()?; + let mut stmt = conn.prepare( + "SELECT id, hostname, app_name, severity, ai_project, ai_tool, ai_session_id, + deleted_at + 900 + FROM stream_deleted_log_lineage + WHERE ai_project = ?1 + AND lower(ai_tool) = lower(?2) + AND ai_session_id = ?3 + AND hostname = ?4 + ORDER BY id DESC + LIMIT ?5", + )?; + let rows = stmt.query_map( + rusqlite::params![ + retention_project, + retention_tool, + retention_session_id, + retention_host, + i64::from(section_limit) + 1 + ], + |row| { + Ok(models::SessionRetentionLineageEntry { + log_id: row.get(0)?, + hostname: row.get(1)?, + app_name: row.get(2)?, + severity: row.get(3)?, + ai_project: row.get(4)?, + ai_tool: row.get(5)?, + ai_session_id: row.get(6)?, + retained_until_epoch: row.get(7)?, + }) + }, + )?; + Ok(rows.collect::, rusqlite::Error>>()?) + }) + .await?; + let retention_lineage_truncated = retention_lineage.len() > section_limit as usize; + retention_lineage.truncate(section_limit as usize); + + let (skill_events, mcp_events, hook_events) = tokio::try_join!( + self.list_skill_events(models::ListSkillEventsRequest { + tool: Some(session.tool.clone()), + project: Some(session.project.clone()), + session_id: Some(session.session_id.clone()), + hostname: Some(session.hostname.clone()), + limit: Some(section_limit), + ..Default::default() + }), + self.list_mcp_events(models::ListMcpEventsRequest { + tool: Some(session.tool.clone()), + project: Some(session.project.clone()), + session_id: Some(session.session_id.clone()), + hostname: Some(session.hostname.clone()), + limit: Some(section_limit), + ..Default::default() + }), + self.list_hook_events(models::ListHookEventsRequest { + tool: Some(session.tool.clone()), + project: Some(session.project.clone()), + session_id: Some(session.session_id.clone()), + hostname: Some(session.hostname.clone()), + limit: Some(section_limit), + ..Default::default() + }) + )?; + + let (artifact_by_correlation, artifact_by_request) = tokio::try_join!( + self.list_artifact_evidence(models::ListArtifactEvidenceRequest { + correlation_id: Some(session.session_id.clone()), + from: Some(session.first_seen.clone()), + to: Some(session.last_seen.clone()), + limit: Some(section_limit), + ..Default::default() + }), + self.list_artifact_evidence(models::ListArtifactEvidenceRequest { + request_id: Some(session.session_id.clone()), + from: Some(session.first_seen.clone()), + to: Some(session.last_seen.clone()), + limit: Some(section_limit), + ..Default::default() + }) + )?; + let (artifact_evidence, artifact_evidence_truncated) = + super::session_investigation_sections::merge_artifact_evidence( + artifact_by_correlation, + artifact_by_request, + section_limit as usize, + ); + + let observatory = super::session_investigation_sections::session_observatory( + self, + &session, + section_limit, + ) + .await?; + + let (related_sessions, related_sessions_truncated) = + related_sessions_for_investigation(self, &session).await?; + + let (external_references, external_references_truncated) = + extract_session_external_references(&transcript_page.events, section_limit as usize); + let source_evidence = + summarize_session_source_evidence(correlation.as_ref(), section_limit as usize); + let graph_counts = graph_neighborhood.as_ref().map(|graph| { + ( + graph.entities.len(), + graph.relationships.len(), + graph.evidence.len(), + ) + }); + let heartbeat_summary_count = correlation + .as_ref() + .map_or(0, |value| value.heartbeat_summaries.len()); + let source_counts = build_session_source_counts( + &source_evidence, + &observatory, + skill_events.events.len(), + mcp_events.events.len(), + hook_events.events.len(), + artifact_evidence.len(), + graph_counts, + heartbeat_summary_count, + incident_context.error_logs.len(), + notifications.len(), + retention_lineage.len(), + ); + + let mut partial_reasons = Vec::new(); + if external_references_truncated { + partial_reasons.push("external_references_truncated".to_string()); + } + if transcript_page.has_more { + partial_reasons.push("transcript_truncated".to_string()); + } + if correlation.as_ref().is_some_and(|value| value.truncated) { + partial_reasons.push("correlation_truncated".to_string()); + } + if graph_neighborhood_truncated { + partial_reasons.push("graph_neighborhood_truncated".to_string()); + } + if related_sessions_truncated { + partial_reasons.push("related_sessions_truncated".to_string()); + } + if observatory.commits_truncated { + partial_reasons.push("observatory_commits_truncated".to_string()); + } + if session_graph_entity_ambiguous { + partial_reasons.push("session_graph_entity_ambiguous".to_string()); + } + if skill_events.truncated { + partial_reasons.push("skill_events_truncated".to_string()); + } + if mcp_events.truncated { + partial_reasons.push("mcp_events_truncated".to_string()); + } + if hook_events.truncated { + partial_reasons.push("hook_events_truncated".to_string()); + } + if artifact_evidence_truncated { + partial_reasons.push("artifact_evidence_truncated".to_string()); + } + if observatory.ambiguous_run { + partial_reasons.push("observatory_run_ambiguous".to_string()); + } + if observatory.actors_truncated { + partial_reasons.push("observatory_actors_truncated".to_string()); + } + if observatory.related_runs_truncated { + partial_reasons.push("observatory_related_runs_truncated".to_string()); + } + if observatory.events_truncated { + partial_reasons.push("observatory_events_truncated".to_string()); + } + if observatory.spans_truncated { + partial_reasons.push("observatory_spans_truncated".to_string()); + } + if observatory.metrics_truncated { + partial_reasons.push("observatory_metrics_truncated".to_string()); + } + if incident_context.error_logs_truncated { + partial_reasons.push("incident_context_errors_truncated".to_string()); + } + if notifications_truncated { + partial_reasons.push("notifications_truncated".to_string()); + } + if retention_lineage_truncated { + partial_reasons.push("retention_lineage_truncated".to_string()); + } + + let metadata = session_investigation_metadata( + &budget, + started, + correlation.as_ref(), + graph_neighborhood.is_some(), + &partial_reasons, + transcript_page.events.len(), + observatory.events.len(), + ); + let envelope = InvestigationEnvelope { + metadata, + result: models::SessionInvestigateResponse { + session, + transcript: transcript_page.events, + transcript_has_more: transcript_page.has_more, + correlation, + graph_neighborhood, + skill_events: skill_events.events, + skill_events_truncated: skill_events.truncated, + mcp_events: mcp_events.events, + mcp_events_truncated: mcp_events.truncated, + hook_events: hook_events.events, + hook_events_truncated: hook_events.truncated, + artifact_evidence, + artifact_evidence_truncated, + observatory, + incident_context, + notifications, + notifications_truncated, + retention_lineage, + retention_lineage_truncated, + related_sessions, + related_sessions_truncated, + external_references, + external_references_truncated, + source_counts, + source_evidence, + partial_reasons, + }, + }; + self.run_db("session_investigate_finalize", move |_| { + let mut envelope = + super::session_investigation_budget::finalize_session_envelope(envelope) + .map_err(anyhow::Error::from)?; + envelope.metadata.budget_used.wall_time_ms = + started.elapsed().as_millis().min(u32::MAX as u128) as u32; + if started.elapsed() + > std::time::Duration::from_millis(u64::from( + envelope.metadata.budget.max_wall_time_ms, + )) + { + return Err(anyhow::Error::from(ServiceError::Busy( + "session_investigation_wall_time_budget_exceeded".into(), + ))); + } + // Updating elapsed usage can change the escaped envelope byte count. + super::session_investigation_budget::finalize_session_envelope(envelope) + .map_err(anyhow::Error::from) + }) + .await + } +} + +#[cfg(test)] +#[path = "session_investigation_tests.rs"] +mod tests; diff --git a/src/app/services/session_investigation_budget.rs b/src/app/services/session_investigation_budget.rs new file mode 100644 index 000000000..b5945af86 --- /dev/null +++ b/src/app/services/session_investigation_budget.rs @@ -0,0 +1,454 @@ +use super::session_investigation_support::{ + build_session_source_counts, summarize_session_source_evidence, +}; +use super::*; + +pub(super) async fn within_session_budget( + wall_time_ms: u32, + work: impl std::future::Future>, +) -> ServiceResult { + let started = std::time::Instant::now(); + let deadline = std::time::Duration::from_millis(u64::from(wall_time_ms)); + let result = tokio::time::timeout(deadline, work).await.map_err(|_| { + ServiceError::Busy("session_investigation_wall_time_budget_exceeded".into()) + })?; + // A future that blocks within a ready poll can overshoot a Tokio timer. + if started.elapsed() > deadline { + return Err(ServiceError::Busy( + "session_investigation_wall_time_budget_exceeded".into(), + )); + } + result +} + +type SessionEnvelope = InvestigationEnvelope; + +/// Bound the complete escaped JSON envelope, including its metadata. Preserve +/// section omission reasons and rebuild counts after removing returned rows. +pub(super) fn finalize_session_envelope( + mut envelope: SessionEnvelope, +) -> ServiceResult { + loop { + refresh_counts(&mut envelope); + // The byte count contributes digits to its own serialized envelope. + // Iterate to a fixed point before comparing with the advertised limit. + loop { + let size = serde_json::to_vec(&envelope) + .map_err(|error| ServiceError::Internal(error.into()))? + .len(); + if envelope.metadata.budget_used.payload_bytes as usize == size { + break; + } + envelope.metadata.budget_used.payload_bytes = u32::try_from(size).unwrap_or(u32::MAX); + } + if envelope.metadata.budget_used.payload_bytes <= envelope.metadata.payload_limit_bytes { + return Ok(envelope); + } + if !prune_largest_section(&mut envelope.result)? { + return Err(ServiceError::Busy( + "session_investigation_identity_exceeds_payload_budget".into(), + )); + } + let reason = "payload_budget_truncated".to_string(); + if !envelope.result.partial_reasons.contains(&reason) { + envelope.result.partial_reasons.push(reason); + } + envelope.metadata.partial = true; + envelope.metadata.truncated = true; + envelope.metadata.partial_reasons = envelope.result.partial_reasons.clone(); + envelope.metadata.truncation_reasons = envelope + .result + .partial_reasons + .iter() + .filter(|reason| reason.contains("truncated")) + .cloned() + .collect(); + } +} + +fn refresh_counts(envelope: &mut SessionEnvelope) { + let result = &mut envelope.result; + let source_id_limit = result + .source_evidence + .values() + .map(|summary| summary.log_ids.len()) + .max() + .unwrap_or(50) + .max(1); + result.source_evidence = + summarize_session_source_evidence(result.correlation.as_ref(), source_id_limit); + result.source_counts = build_session_source_counts( + &result.source_evidence, + &result.observatory, + result.skill_events.len(), + result.mcp_events.len(), + result.hook_events.len(), + result.artifact_evidence.len(), + result.graph_neighborhood.as_ref().map(|graph| { + ( + graph.entities.len(), + graph.relationships.len(), + graph.evidence.len(), + ) + }), + result + .correlation + .as_ref() + .map_or(0, |correlation| correlation.heartbeat_summaries.len()), + result.incident_context.error_logs.len(), + result.notifications.len(), + result.retention_lineage.len(), + ); + result + .source_counts + .insert("transcript".into(), result.transcript.len()); + result + .source_counts + .insert("related_session".into(), result.related_sessions.len()); + result.source_counts.insert( + "external_reference".into(), + result.external_references.len(), + ); + result + .source_counts + .insert("attributed_commit".into(), result.observatory.commits.len()); + result + .source_counts + .insert("observatory_run".into(), result.observatory.runs.len()); + envelope.metadata.budget_used.log_rows = (result.incident_context.error_logs.len() + + result + .correlation + .as_ref() + .map_or(0, |correlation| correlation.logs.len())) + as u32; + let observatory = &result.observatory; + let graph_rows = result.graph_neighborhood.as_ref().map_or(0, |graph| { + graph.entities.len() + graph.relationships.len() + graph.evidence.len() + }); + envelope.metadata.budget_used.evidence_rows = (result.transcript.len() + + result.skill_events.len() + + result.mcp_events.len() + + result.hook_events.len() + + result.artifact_evidence.len() + + result.notifications.len() + + result.retention_lineage.len() + + result.related_sessions.len() + + result.external_references.len() + + result.incident_context.ai_sessions.len() + + observatory.runs.len() + + observatory.actors.len() + + observatory.related_runs.len() + + observatory.commits.len() + + observatory.events.len() + + observatory.spans.len() + + observatory.metrics.len() + + graph_rows + + usize::from(observatory.parent_run.is_some()) + + usize::from(observatory.previous_run.is_some()) + + usize::from(observatory.repository.is_some()) + + usize::from(observatory.worktree.is_some()) + + result + .correlation + .as_ref() + .map_or(0, |correlation| correlation.heartbeat_summaries.len())) + as u32; +} + +fn prune_largest_section(result: &mut models::SessionInvestigateResponse) -> ServiceResult { + let mut largest = (0, 0); + macro_rules! candidate { + ($id:expr, $value:expr, $present:expr) => { + if $present { + let bytes = serde_json::to_vec(&$value) + .map_err(|error| ServiceError::Internal(error.into()))? + .len(); + if bytes > largest.1 { + largest = ($id, bytes); + } + } + }; + } + candidate!(1, result.transcript, !result.transcript.is_empty()); + candidate!(2, result.correlation, result.correlation.is_some()); + candidate!( + 3, + result.graph_neighborhood, + result.graph_neighborhood.is_some() + ); + candidate!(4, result.skill_events, !result.skill_events.is_empty()); + candidate!(5, result.mcp_events, !result.mcp_events.is_empty()); + candidate!(6, result.hook_events, !result.hook_events.is_empty()); + candidate!( + 7, + result.artifact_evidence, + !result.artifact_evidence.is_empty() + ); + candidate!( + 8, + result.observatory.runs, + !result.observatory.runs.is_empty() + ); + candidate!( + 9, + result.observatory.actors, + !result.observatory.actors.is_empty() + ); + candidate!( + 10, + result.observatory.related_runs, + !result.observatory.related_runs.is_empty() + ); + candidate!( + 11, + result.observatory.commits, + !result.observatory.commits.is_empty() + ); + candidate!( + 12, + result.observatory.events, + !result.observatory.events.is_empty() + ); + candidate!( + 13, + result.observatory.spans, + !result.observatory.spans.is_empty() + ); + candidate!( + 14, + result.observatory.metrics, + !result.observatory.metrics.is_empty() + ); + candidate!( + 15, + result.incident_context.error_logs, + !result.incident_context.error_logs.is_empty() + ); + candidate!(16, result.notifications, !result.notifications.is_empty()); + candidate!( + 17, + result.retention_lineage, + !result.retention_lineage.is_empty() + ); + candidate!( + 18, + result.related_sessions, + !result.related_sessions.is_empty() + ); + candidate!( + 19, + result.external_references, + !result.external_references.is_empty() + ); + candidate!( + 20, + result.observatory.parent_run, + result.observatory.parent_run.is_some() + ); + candidate!( + 21, + result.observatory.previous_run, + result.observatory.previous_run.is_some() + ); + candidate!( + 22, + result.observatory.repository, + result.observatory.repository.is_some() + ); + candidate!( + 23, + result.observatory.worktree, + result.observatory.worktree.is_some() + ); + candidate!( + 24, + result.incident_context.ai_sessions, + !result.incident_context.ai_sessions.is_empty() + ); + candidate!( + 25, + result.incident_context.by_app, + !result.incident_context.by_app.is_empty() + ); + candidate!(26, result.session.title, result.session.title.is_some()); + candidate!( + 27, + result.session.transcript_path, + result.session.transcript_path.is_some() + ); + let section = match largest.0 { + 0 => return Ok(false), + 1 => { + result.transcript.truncate(result.transcript.len() / 2); + result.transcript_has_more = true; + "transcript" + } + 2 => { + result.correlation = None; + "correlation" + } + 3 => { + result.graph_neighborhood = None; + "graph_neighborhood" + } + 4 => { + result.skill_events.truncate(result.skill_events.len() / 2); + result.skill_events_truncated = true; + "skill_events" + } + 5 => { + result.mcp_events.truncate(result.mcp_events.len() / 2); + result.mcp_events_truncated = true; + "mcp_events" + } + 6 => { + result.hook_events.truncate(result.hook_events.len() / 2); + result.hook_events_truncated = true; + "hook_events" + } + 7 => { + result + .artifact_evidence + .truncate(result.artifact_evidence.len() / 2); + result.artifact_evidence_truncated = true; + "artifact_evidence" + } + 8 => { + result + .observatory + .runs + .truncate(result.observatory.runs.len() / 2); + "observatory_runs" + } + 9 => { + result + .observatory + .actors + .truncate(result.observatory.actors.len() / 2); + result.observatory.actors_truncated = true; + "observatory_actors" + } + 10 => { + result + .observatory + .related_runs + .truncate(result.observatory.related_runs.len() / 2); + result.observatory.related_runs_truncated = true; + "observatory_related_runs" + } + 11 => { + result + .observatory + .commits + .truncate(result.observatory.commits.len() / 2); + result.observatory.commits_truncated = true; + "observatory_commits" + } + 12 => { + result + .observatory + .events + .truncate(result.observatory.events.len() / 2); + result.observatory.events_truncated = true; + "observatory_events" + } + 13 => { + result + .observatory + .spans + .truncate(result.observatory.spans.len() / 2); + result.observatory.spans_truncated = true; + "observatory_spans" + } + 14 => { + result + .observatory + .metrics + .truncate(result.observatory.metrics.len() / 2); + result.observatory.metrics_truncated = true; + "observatory_metrics" + } + 15 => { + result + .incident_context + .error_logs + .truncate(result.incident_context.error_logs.len() / 2); + result.incident_context.error_logs_truncated = true; + "incident_context_errors" + } + 16 => { + result + .notifications + .truncate(result.notifications.len() / 2); + result.notifications_truncated = true; + "notifications" + } + 17 => { + result + .retention_lineage + .truncate(result.retention_lineage.len() / 2); + result.retention_lineage_truncated = true; + "retention_lineage" + } + 18 => { + result + .related_sessions + .truncate(result.related_sessions.len() / 2); + result.related_sessions_truncated = true; + "related_sessions" + } + 19 => { + result + .external_references + .truncate(result.external_references.len() / 2); + result.external_references_truncated = true; + "external_references" + } + 20 => { + result.observatory.parent_run = None; + "parent_run" + } + 21 => { + result.observatory.previous_run = None; + "previous_run" + } + 22 => { + result.observatory.repository = None; + "repository" + } + 23 => { + result.observatory.worktree = None; + "worktree" + } + 24 => { + result + .incident_context + .ai_sessions + .truncate(result.incident_context.ai_sessions.len() / 2); + "incident_ai_sessions" + } + 25 => { + result + .incident_context + .by_app + .truncate(result.incident_context.by_app.len() / 2); + "incident_apps" + } + 26 => { + result.session.title = None; + "session_title" + } + _ => { + result.session.transcript_path = None; + "transcript_path" + } + }; + let reason = format!("{section}_payload_truncated"); + if !result.partial_reasons.contains(&reason) { + result.partial_reasons.push(reason); + } + Ok(true) +} + +#[cfg(test)] +#[path = "session_investigation_budget_tests.rs"] +mod tests; diff --git a/src/app/services/session_investigation_budget_tests.rs b/src/app/services/session_investigation_budget_tests.rs new file mode 100644 index 000000000..a323b9b1a --- /dev/null +++ b/src/app/services/session_investigation_budget_tests.rs @@ -0,0 +1,145 @@ +use super::*; + +fn minimal_envelope() -> SessionEnvelope { + let session = AiSessionEntry { + session_key: "h:codex:p:s".into(), + project: "p".into(), + tool: "codex".into(), + session_id: "s".into(), + hostname: "h".into(), + transcript_path: None, + first_seen: "2026-09-21T00:00:00Z".into(), + last_seen: "2026-09-21T00:10:00Z".into(), + event_count: 1, + title: None, + title_provenance: None, + }; + SessionEnvelope { + metadata: super::super::session_investigation_support::session_investigation_metadata( + &InvestigationBudget::default(), + std::time::Instant::now(), + None, + false, + &[], + 1, + 0, + ), + result: models::SessionInvestigateResponse { + session, + transcript: vec![], + transcript_has_more: false, + correlation: None, + graph_neighborhood: None, + skill_events: vec![], + skill_events_truncated: false, + mcp_events: vec![], + mcp_events_truncated: false, + hook_events: vec![], + hook_events_truncated: false, + artifact_evidence: vec![], + artifact_evidence_truncated: false, + observatory: Default::default(), + incident_context: IncidentContextResponse { + window_from: "2026-09-21T00:00:00Z".into(), + window_to: "2026-09-21T00:10:00Z".into(), + total_logs: 0, + by_severity: vec![], + by_app: vec![], + error_logs: vec![], + error_logs_truncated: false, + ai_sessions: vec![], + }, + notifications: vec![], + notifications_truncated: false, + retention_lineage: vec![], + retention_lineage_truncated: false, + related_sessions: vec![], + related_sessions_truncated: false, + external_references: vec![], + external_references_truncated: false, + source_counts: Default::default(), + source_evidence: Default::default(), + partial_reasons: vec![], + }, + } +} + +#[test] +fn escaped_payload_budget_includes_metadata_and_refreshes_counts() { + let mut envelope = minimal_envelope(); + envelope + .result + .transcript + .push(models::RenderedSessionEvent { + position: 1, + timestamp: "2026-09-21T00:00:00Z".into(), + kind: models::RenderedSessionEventKind::Assistant, + text: "\"🦀".repeat(40_000), + redacted: false, + parse_warning: None, + }); + envelope.result.source_counts.insert("transcript".into(), 1); + let result = finalize_session_envelope(envelope).unwrap(); + let size = serde_json::to_vec(&result).unwrap().len(); + assert!(size <= result.metadata.payload_limit_bytes as usize); + assert_eq!(size, result.metadata.budget_used.payload_bytes as usize); + assert!(result.metadata.partial && result.metadata.truncated); + assert!(result.result.transcript_has_more); + assert_eq!( + result.result.source_counts["transcript"], + result.result.transcript.len() + ); + assert!( + result + .result + .partial_reasons + .iter() + .any(|reason| reason == "transcript_payload_truncated") + ); +} + +#[test] +fn oversized_required_identity_fails_closed() { + let mut envelope = minimal_envelope(); + envelope.result.session.project = "p".repeat(100_000); + assert!(matches!( + finalize_session_envelope(envelope), + Err(ServiceError::Busy(_)) + )); +} + +#[test] +fn ambiguity_is_partial_without_claiming_truncation() { + let reasons = vec!["observatory_run_ambiguous".to_string()]; + let metadata = super::super::session_investigation_support::session_investigation_metadata( + &InvestigationBudget::default(), + std::time::Instant::now(), + None, + false, + &reasons, + 0, + 0, + ); + assert!(metadata.partial); + assert!(!metadata.truncated); + assert!(metadata.truncation_reasons.is_empty()); + assert_eq!(metadata.auth_state, "unknown"); +} + +#[tokio::test] +async fn elapsed_budget_returns_retryable_busy() { + let result: ServiceResult<()> = within_session_budget(1, std::future::pending()).await; + assert!( + matches!(result, Err(ServiceError::Busy(message)) if message == "session_investigation_wall_time_budget_exceeded") + ); +} + +#[tokio::test] +async fn ready_future_that_overshoots_wall_budget_fails_closed() { + let result = within_session_budget(1, async { + std::thread::sleep(std::time::Duration::from_millis(5)); + Ok(()) + }) + .await; + assert!(matches!(result, Err(ServiceError::Busy(_)))); +} diff --git a/src/app/services/session_investigation_sections.rs b/src/app/services/session_investigation_sections.rs new file mode 100644 index 000000000..d753927cb --- /dev/null +++ b/src/app/services/session_investigation_sections.rs @@ -0,0 +1,188 @@ +use super::*; + +pub(super) async fn session_notifications( + service: &CortexService, + session: &AiSessionEntry, + limit: u32, +) -> ServiceResult> { + let host = session.hostname.clone(); + let start = session.first_seen.clone(); + let end = session.last_seen.clone(); + service + .run_db("session_investigate_notifications", move |pool| { + let conn = pool.get()?; + let mut stmt = conn.prepare( + "SELECT id,outbox_id,rule_id,hostname,fired_at,status_code + FROM notification_firings + WHERE hostname=?1 AND julianday(fired_at)>=julianday(?2) + AND julianday(fired_at)<=julianday(?3) + ORDER BY fired_at DESC,id DESC LIMIT ?4", + )?; + Ok(stmt + .query_map( + rusqlite::params![host, start, end, i64::from(limit) + 1], + |row| { + Ok(db::notifications::FiringRow { + id: row.get(0)?, + outbox_id: row.get(1)?, + rule_id: row.get(2)?, + hostname: row.get(3)?, + fired_at: row.get(4)?, + status_code: row.get(5)?, + }) + }, + )? + .collect::>>()?) + }) + .await +} + +pub(super) fn merge_artifact_evidence( + correlation: models::ListArtifactEvidenceResponse, + request: models::ListArtifactEvidenceResponse, + limit: usize, +) -> (Vec, bool) { + let source_truncated = correlation.truncated || request.truncated; + let mut events = BTreeMap::new(); + for event in correlation.events.into_iter().chain(request.events) { + events.insert(event.cortex_log_id, event); + } + let truncated = source_truncated || events.len() > limit; + (events.into_values().rev().take(limit).collect(), truncated) +} + +pub(super) async fn session_observatory( + service: &CortexService, + session: &AiSessionEntry, + section_limit: u32, +) -> ServiceResult { + let observatory_session_id = session.session_id.clone(); + let observatory_tool = session.tool.clone(); + let observatory_host = session.hostname.clone(); + service + .run_db("session_investigate_observatory", move |pool| { + let query = db::agent_observatory::AgentRunQuery { + tools: vec![observatory_tool.clone()], + host: Some(observatory_host.clone()), + native_session_id: Some(observatory_session_id.clone()), + ..Default::default() + }; + let mut runs = + db::agent_observatory::list_observatory_runs(pool, &query, None, 50, i64::MAX)?; + runs.retain(|run| { + run.native_session_id == observatory_session_id + && run.tool.eq_ignore_ascii_case(&observatory_tool) + && run.hostname == observatory_host + }); + let ambiguous_run = runs.len() > 1; + let Some(run) = runs.first().cloned().filter(|_| !ambiguous_run) else { + return Ok(models::SessionObservatoryEvidence { + runs, + ambiguous_run, + ..Default::default() + }); + }; + let resolved = db::agent_observatory::resolve_observatory_run(pool, &run.run_key)?; + let (run_id, identity) = resolved.ok_or_else(|| { + anyhow::anyhow!("Agent Observatory run disappeared during session investigation") + })?; + let parent_run = match run.parent_run_id { + Some(id) => db::agent_observatory::resolve_observatory_run_row(pool, id)?, + None => None, + }; + let previous_run = match run.previous_run_id { + Some(id) => db::agent_observatory::resolve_observatory_run_row(pool, id)?, + None => None, + }; + let mut actors = db::agent_observatory::list_observatory_run_actors( + pool, + run_id, + section_limit as usize, + )?; + let actors_truncated = actors.len() > section_limit as usize; + actors.truncate(section_limit as usize); + let worktree = match run.primary_worktree_id { + Some(id) => db::agent_observatory::resolve_observatory_worktree(pool, id)?, + None => None, + }; + let repository = match worktree.as_ref() { + Some(worktree) => db::agent_observatory::resolve_observatory_repository( + pool, + worktree.repository_id, + )?, + None => None, + }; + let mut commits = db::agent_observatory::list_agent_run_attributed_commits( + pool, + run_id, + section_limit as usize, + )?; + let commits_truncated = commits.len() > section_limit as usize; + commits.truncate(section_limit as usize); + let mut related_runs = match run.primary_worktree_id { + Some(worktree_id) => db::agent_observatory::list_observatory_runs( + pool, + &db::agent_observatory::AgentRunQuery { + worktree_id: Some(worktree_id), + exclude_run_id: Some(run_id), + ..Default::default() + }, + None, + section_limit as usize, + i64::MAX, + )?, + None => Vec::new(), + }; + related_runs.retain(|candidate| candidate.id != run_id); + let related_runs_truncated = related_runs.len() > section_limit as usize; + related_runs.truncate(section_limit as usize); + let events = db::agent_observatory::list_observatory_events( + pool, + &run.run_key, + &db::agent_observatory::AgentEventQuery::default(), + None, + section_limit as usize, + true, + i64::MAX, + )?; + let spans = db::agent_observatory::list_observatory_spans( + pool, + run_id, + &identity, + &db::agent_observatory::TelemetryQuery::default(), + None, + section_limit as usize, + i64::MAX, + )?; + let metrics = db::agent_observatory::list_observatory_metrics( + pool, + run_id, + &identity, + &db::agent_observatory::TelemetryQuery::default(), + None, + section_limit as usize, + i64::MAX, + )?; + Ok(models::SessionObservatoryEvidence { + runs, + ambiguous_run, + parent_run, + previous_run, + actors, + actors_truncated, + related_runs, + related_runs_truncated, + repository, + worktree, + commits, + commits_truncated, + events_truncated: events.len() > section_limit as usize, + spans_truncated: spans.len() > section_limit as usize, + metrics_truncated: metrics.len() > section_limit as usize, + events: events.into_iter().take(section_limit as usize).collect(), + spans: spans.into_iter().take(section_limit as usize).collect(), + metrics: metrics.into_iter().take(section_limit as usize).collect(), + }) + }) + .await +} diff --git a/src/app/services/session_investigation_support.rs b/src/app/services/session_investigation_support.rs new file mode 100644 index 000000000..3ac697c86 --- /dev/null +++ b/src/app/services/session_investigation_support.rs @@ -0,0 +1,249 @@ +use std::time::Instant; + +use super::*; + +pub(super) async fn related_sessions_for_investigation( + service: &CortexService, + session: &AiSessionEntry, +) -> ServiceResult<(Vec, bool)> { + let mut sessions = service + .list_sessions(ListSessionsRequest { + project: Some(session.project.clone()), + tool: None, + session_id: None, + host: None, + since: Some(session.first_seen.clone()), + until: Some(session.last_seen.clone()), + limit: Some(22), + offset: None, + }) + .await? + .sessions + .into_iter() + .filter(|candidate| candidate.session_key != session.session_key) + .collect::>(); + let truncated = sessions.len() > 20; + sessions.truncate(20); + Ok((sessions, truncated)) +} + +#[allow(clippy::too_many_arguments)] +pub(super) fn build_session_source_counts( + source_evidence: &std::collections::BTreeMap, + observatory: &models::SessionObservatoryEvidence, + skill_event_count: usize, + mcp_event_count: usize, + hook_event_count: usize, + artifact_evidence_count: usize, + graph_counts: Option<(usize, usize, usize)>, + heartbeat_summary_count: usize, + incident_error_log_count: usize, + notification_count: usize, + retention_lineage_count: usize, +) -> std::collections::BTreeMap { + let mut source_counts = source_evidence + .iter() + .map(|(kind, summary)| (kind.clone(), summary.count)) + .collect::>(); + for event in &observatory.events { + *source_counts + .entry(format!("observatory:{}", event.source_kind)) + .or_default() += 1; + } + if skill_event_count > 0 { + source_counts.insert("skill_event".to_string(), skill_event_count); + } + if mcp_event_count > 0 { + source_counts.insert("mcp_event".to_string(), mcp_event_count); + } + if hook_event_count > 0 { + source_counts.insert("hook_event".to_string(), hook_event_count); + } + if artifact_evidence_count > 0 { + source_counts.insert("artifact_evidence".to_string(), artifact_evidence_count); + } + if !observatory.spans.is_empty() { + source_counts.insert("otlp_span".to_string(), observatory.spans.len()); + } + if !observatory.metrics.is_empty() { + source_counts.insert("otlp_metric".to_string(), observatory.metrics.len()); + } + if let Some((entities, relationships, evidence)) = graph_counts { + source_counts.insert("graph_entity".to_string(), entities); + source_counts.insert("graph_relationship".to_string(), relationships); + source_counts.insert("graph_evidence".to_string(), evidence); + } + if heartbeat_summary_count > 0 { + source_counts.insert("heartbeat_summary".to_string(), heartbeat_summary_count); + } + if incident_error_log_count > 0 { + source_counts.insert("incident_error_log".to_string(), incident_error_log_count); + } + if notification_count > 0 { + source_counts.insert("notification".to_string(), notification_count); + } + if !observatory.actors.is_empty() { + source_counts.insert("observatory_actor".to_string(), observatory.actors.len()); + } + if observatory.parent_run.is_some() { + source_counts.insert("parent_run".to_string(), 1); + } + if observatory.previous_run.is_some() { + source_counts.insert("previous_run".to_string(), 1); + } + if !observatory.related_runs.is_empty() { + source_counts.insert( + "same_worktree_run".to_string(), + observatory.related_runs.len(), + ); + } + if retention_lineage_count > 0 { + source_counts.insert("retention_lineage".to_string(), retention_lineage_count); + } + source_counts +} + +pub(super) fn session_investigation_metadata( + budget: &InvestigationBudget, + started: Instant, + correlation: Option<&GraphSessionCorrelation>, + has_graph_neighborhood: bool, + partial_reasons: &[String], + transcript_rows: usize, + evidence_rows: usize, +) -> InvestigationMetadata { + let graph_calls = u32::from(correlation.is_some()) + u32::from(has_graph_neighborhood); + let log_rows = correlation.map_or(0, |value| value.logs.len() as u32); + let partial = !partial_reasons.is_empty(); + InvestigationMetadata { + server_version: env!("CARGO_PKG_VERSION").to_string(), + schema_version: INVESTIGATION_UI_VERSION.to_string(), + graph_projection_status: correlation + .map(|value| if value.used_graph { "used" } else { "fallback" }.to_string()), + source_watermark: None, + degraded_reasons: correlation + .filter(|value| !value.used_graph) + .map(|_| vec!["session_graph_entity_unavailable".to_string()]) + .unwrap_or_default(), + truncated: partial_reasons + .iter() + .any(|reason| reason.contains("truncated")), + truncation_reasons: partial_reasons + .iter() + .filter(|reason| reason.contains("truncated")) + .cloned() + .collect(), + partial, + partial_reasons: partial_reasons.to_vec(), + auth_state: "unknown".to_string(), + budget: budget.clone(), + budget_used: InvestigationBudgetUsed { + graph_calls, + log_rows, + evidence_rows: evidence_rows.min(u32::MAX as usize) as u32, + candidate_explanations: 0, + wall_time_ms: started.elapsed().as_millis().min(u32::MAX as u128) as u32, + payload_bytes: 0, + }, + payload_limit_bytes: budget.max_payload_bytes, + version_skew: (transcript_rows > 200).then(|| "transcript_limit_exceeded".to_string()), + } +} + +pub(super) fn summarize_session_source_evidence( + correlation: Option<&GraphSessionCorrelation>, + requested_limit: usize, +) -> BTreeMap { + let id_limit = requested_limit.clamp(1, 50); + let mut summaries = BTreeMap::::new(); + let Some(correlation) = correlation else { + return summaries; + }; + + for row in &correlation.logs { + let kind = row + .source_kind + .clone() + .unwrap_or_else(|| "unknown".to_string()); + let summary = summaries.entry(kind).or_default(); + summary.count += 1; + if summary.log_ids.len() < id_limit { + summary.log_ids.push(row.entry.id); + } else { + summary.truncated = true; + } + } + + summaries +} + +pub(super) fn extract_session_external_references( + events: &[models::RenderedSessionEvent], + limit: usize, +) -> (Vec, bool) { + let mut references = BTreeMap::new(); + for event in events { + for raw in event.text.split_whitespace() { + let token = raw.trim_matches(|ch: char| { + matches!( + ch, + ',' | '.' | ';' | ':' | ')' | '(' | ']' | '[' | '}' | '{' | '"' | '\'' + ) + }); + if token.is_empty() { + continue; + } + let kind = if token.starts_with("http://") || token.starts_with("https://") { + if token.contains("github.com/") && token.contains("/pull/") { + models::SessionExternalReferenceKind::GithubPullRequest + } else if token.contains("github.com/") && token.contains("/issues/") { + models::SessionExternalReferenceKind::GithubIssue + } else { + models::SessionExternalReferenceKind::Url + } + } else if looks_like_linear_identifier(token) { + models::SessionExternalReferenceKind::LinearIssue + } else if token.strip_prefix('#').is_some_and(|number| { + !number.is_empty() && number.chars().all(|ch| ch.is_ascii_digit()) + }) { + models::SessionExternalReferenceKind::GithubReference + } else if looks_like_commit_sha(token) { + models::SessionExternalReferenceKind::CommitSha + } else { + continue; + }; + references + .entry((kind.clone(), safe_passive_text(token, 500))) + .or_insert(models::SessionExternalReference { + kind, + value: safe_passive_text(token, 500), + source_position: Some(event.position), + evidence_kind: "transcript_text".to_string(), + trust_level: "claimed".to_string(), + verified: false, + }); + if references.len() > limit { + return (references.into_values().take(limit).collect(), true); + } + } + } + (references.into_values().collect(), false) +} + +pub(super) fn looks_like_linear_identifier(token: &str) -> bool { + let Some((prefix, suffix)) = token.split_once('-') else { + return false; + }; + (2..=12).contains(&prefix.len()) + && prefix + .chars() + .all(|ch| ch.is_ascii_uppercase() || ch.is_ascii_digit()) + && !suffix.is_empty() + && suffix.chars().all(|ch| ch.is_ascii_digit()) +} + +pub(super) fn looks_like_commit_sha(token: &str) -> bool { + (7..=40).contains(&token.len()) + && token.chars().all(|ch| ch.is_ascii_hexdigit()) + && token.chars().any(|ch| ch.is_ascii_alphabetic()) +} diff --git a/src/app/services/session_investigation_tests.rs b/src/app/services/session_investigation_tests.rs new file mode 100644 index 000000000..8fcdc1c33 --- /dev/null +++ b/src/app/services/session_investigation_tests.rs @@ -0,0 +1,229 @@ +use super::super::session_investigation_support::looks_like_linear_identifier; +use super::*; + +fn event(position: i64, text: &str) -> models::RenderedSessionEvent { + models::RenderedSessionEvent { + position, + timestamp: "2026-09-21T00:00:00Z".to_string(), + kind: models::RenderedSessionEventKind::Assistant, + text: text.to_string(), + redacted: false, + parse_warning: None, + } +} + +#[test] +fn external_references_are_classified_and_deduplicated() { + let events = vec![ + event( + 1, + "Worked on CLD-1149 and https://github.com/dinglebear-ai/labby/pull/717 with commit abc1234.", + ), + event( + 2, + "Follow-up #731 https://github.com/dinglebear-ai/labby/issues/723 https://linear.app/lime-technology/issue/CLD-1149/foo", + ), + ]; + + let (refs, _) = extract_session_external_references(&events, 100); + + assert!(refs.iter().any(|item| { + item.kind == models::SessionExternalReferenceKind::LinearIssue && item.value == "CLD-1149" + })); + assert!(refs.iter().any(|item| { + item.kind == models::SessionExternalReferenceKind::GithubPullRequest + && item.value.contains("/pull/717") + })); + assert!(refs.iter().any(|item| { + item.kind == models::SessionExternalReferenceKind::GithubIssue + && item.value.contains("/issues/723") + })); + assert!(refs.iter().any(|item| { + item.kind == models::SessionExternalReferenceKind::GithubReference && item.value == "#731" + })); + assert!(refs.iter().any(|item| { + item.kind == models::SessionExternalReferenceKind::CommitSha && item.value == "abc1234" + })); +} + +#[test] +fn linear_identifier_parser_rejects_ordinary_hyphenated_text() { + assert!(looks_like_linear_identifier("U8-1122")); + assert!(looks_like_linear_identifier("CLD-1149")); + assert!(!looks_like_linear_identifier("session-investigate")); + assert!(!looks_like_linear_identifier("abc-123")); + assert!(!looks_like_linear_identifier("CLD-next")); +} + +fn correlated_log(id: i64, source_kind: Option<&str>) -> CorrelatedLogRow { + CorrelatedLogRow { + entry: models::LogEntry { + id, + timestamp: "2026-09-21T00:00:00Z".to_string(), + hostname: "macpoo".to_string(), + facility: None, + severity: "info".to_string(), + app_name: Some("test".to_string()), + process_id: None, + message: format!("log-{id}"), + received_at: "2026-09-21T00:00:00Z".to_string(), + source_ip: "test://source".to_string(), + ai_tool: None, + ai_project: None, + ai_session_id: None, + ai_transcript_path: None, + metadata_json: None, + }, + source_kind: source_kind.map(str::to_string), + discovery: "test".to_string(), + } +} + +#[test] +fn source_evidence_groups_source_kinds_and_bounds_log_ids() { + let correlation = GraphSessionCorrelation { + session_id: "session-1".to_string(), + session_start: "2026-09-21T00:00:00Z".to_string(), + session_end: "2026-09-21T00:10:00Z".to_string(), + used_graph: true, + session_entity_keys: vec!["project:codex:session-1".to_string()], + discovered_hosts: Vec::new(), + discovered_entities: Vec::new(), + logs: vec![ + correlated_log(1, Some("docker-stream")), + correlated_log(2, Some("docker-stream")), + correlated_log(3, Some("agent-command")), + correlated_log(4, None), + ], + agent_command_count: 1, + shell_history_count: 0, + heartbeat_summaries: Vec::new(), + truncated: false, + }; + + let summaries = summarize_session_source_evidence(Some(&correlation), 1); + + assert_eq!(summaries["docker-stream"].count, 2); + assert_eq!(summaries["docker-stream"].log_ids, vec![1]); + assert!(summaries["docker-stream"].truncated); + assert_eq!(summaries["agent-command"].log_ids, vec![3]); + assert_eq!(summaries["unknown"].log_ids, vec![4]); +} + +fn fixture_service() -> (tempfile::TempDir, CortexService) { + let dir = tempfile::tempdir().unwrap(); + let storage = crate::config::StorageConfig::for_test(dir.path().join("investigate.db")); + let pool = std::sync::Arc::new(db::init_pool(&storage).unwrap()); + (dir, CortexService::new(pool, storage)) +} + +fn fixture_session() -> AiSessionEntry { + db::AiSessionEntry { + ai_project: "p".into(), + ai_tool: "codex".into(), + ai_session_id: "s".into(), + ai_transcript_path: None, + hostname: "h".into(), + first_seen: "2026-09-21T00:00:00Z".into(), + last_seen: "2026-09-21T00:10:00Z".into(), + event_count: 1, + title: None, + title_provenance: None, + } + .into() +} + +#[tokio::test] +async fn notification_identity_and_end_window_are_applied_before_limit() { + let (_dir, service) = fixture_service(); + let conn = service.pool.get().unwrap(); + for (host, time) in [ + ("h", "2026-09-21T00:00:00Z"), + ("h", "2026-09-21T00:05:00Z"), + ("h", "2026-09-21T00:10:00Z"), + ("h", "2026-09-21T01:00:00Z"), + ] { + conn.execute("INSERT INTO notification_firings(outbox_id,rule_id,severity,hostname,fired_at) VALUES(1,'r','err',?1,?2)", rusqlite::params![host,time]).unwrap(); + } + for _ in 0..501 { + conn.execute("INSERT INTO notification_firings(outbox_id,rule_id,severity,hostname,fired_at) VALUES(1,'r','err','other','2026-09-21T00:09:00Z')", []).unwrap(); + } + drop(conn); + let rows = super::super::session_investigation_sections::session_notifications( + &service, + &fixture_session(), + 2, + ) + .await + .unwrap(); + assert_eq!( + rows.len(), + 3, + "sentinel detects truncation among matching rows" + ); + assert!(rows.iter().all(|row| row.hostname == "h")); + assert!( + rows.iter() + .all(|row| row.fired_at.as_str() <= "2026-09-21T00:10:00Z") + ); +} + +#[test] +fn artifact_union_is_deduplicated_bounded_and_signals_overflow() { + let entry = |id| models::ArtifactEvidenceEntry { + cortex_log_id: id, + event: serde_json::from_value(serde_json::json!({ + "schemaVersion": "dinglebear.cortex-artifact-evidence/v1", "eventId": format!("e-{id}"), + "eventKind": "installed", "sourceSystem": "test", "sourceIssuer": "fixture", + "observedAt": "2026-09-21T00:05:00Z" + })) + .unwrap(), + }; + let (events, truncated) = super::super::session_investigation_sections::merge_artifact_evidence( + models::ListArtifactEvidenceResponse { + events: vec![entry(1), entry(2)], + truncated: false, + }, + models::ListArtifactEvidenceResponse { + events: vec![entry(2), entry(3)], + truncated: false, + }, + 2, + ); + assert_eq!( + events + .iter() + .map(|item| item.cortex_log_id) + .collect::>(), + vec![3, 2] + ); + assert!(truncated); +} + +#[test] +fn external_references_deduplicate_across_positions_and_cap_unique_values() { + let events = vec![event(1, "#123 #456"), event(2, "#123 #789")]; + let (refs, truncated) = extract_session_external_references(&events, 2); + assert!(truncated); + assert_eq!(refs.len(), 2); + assert_eq!(refs.iter().filter(|item| item.value == "#123").count(), 1); + assert!( + refs.iter() + .all(|item| !item.verified && item.trust_level == "claimed") + ); +} + +#[tokio::test] +async fn related_sessions_signal_truncation() { + let (_dir, service) = fixture_service(); + let conn = service.pool.get().unwrap(); + for id in 0..23 { + conn.execute("INSERT INTO logs(timestamp,hostname,severity,message,raw,source_ip,ai_tool,ai_project,ai_session_id) VALUES('2026-09-21T00:05:00Z','h','info','message','','fixture','codex','p',?1)", [format!("session-{id}")]).unwrap(); + } + drop(conn); + let (rows, truncated) = related_sessions_for_investigation(&service, &fixture_session()) + .await + .unwrap(); + assert_eq!(rows.len(), 20); + assert!(truncated); +} diff --git a/src/cli/dispatch.rs b/src/cli/dispatch.rs index 45d9edb64..a65f55180 100644 --- a/src/cli/dispatch.rs +++ b/src/cli/dispatch.rs @@ -149,6 +149,7 @@ impl SessionsArgs { ListSessionsRequest { project: self.project, tool: self.tool, + session_id: None, host: self.host, since: self.since, until: self.until, diff --git a/src/cli/dispatch_tests.rs b/src/cli/dispatch_tests.rs index 75ec4b1cd..839a5f37d 100644 --- a/src/cli/dispatch_tests.rs +++ b/src/cli/dispatch_tests.rs @@ -247,7 +247,7 @@ fn sessions_args_into_request_snapshot() { let req = args.into_request(); assert_eq!( format!("{req:?}"), - "ListSessionsRequest { project: Some(\"/home/me/proj\"), tool: Some(\"claude\"), host: None, since: None, until: None, limit: Some(20), offset: None }" + "ListSessionsRequest { project: Some(\"/home/me/proj\"), tool: Some(\"claude\"), session_id: None, host: None, since: None, until: None, limit: Some(20), offset: None }" ); } diff --git a/src/db.rs b/src/db.rs index 7f416d9c4..3e6440330 100644 --- a/src/db.rs +++ b/src/db.rs @@ -136,12 +136,13 @@ pub use pool::{ }; pub(crate) use pool::{WriteConnBusy, is_pool_acquire_failure, try_write_conn_for}; pub use queries::{ - RollupRefresh, SEVERITY_LEVELS, ai_session_rollup_status, correlate_session_graph, - durable_stream_page, get_error_summary, get_stats, incident_context_summary, - investigate_ai_incidents, list_ai_projects, list_ai_sessions, list_ai_tools, list_hosts, - prune_expired_stream_lineage, prune_timeline_rollup, refresh_ai_session_rollup_if_stale, - refresh_timeline_rollup, rendered_session_page, search_ai_abuse, search_ai_anchors, - search_ai_incidents, search_ai_related_logs, search_ai_sessions, search_logs, severity_to_num, + RollupRefresh, SEVERITY_LEVELS, SessionGraphScope, ai_session_rollup_status, + correlate_session_graph, correlate_session_graph_scoped, durable_stream_page, + get_error_summary, get_stats, incident_context_summary, investigate_ai_incidents, + list_ai_projects, list_ai_sessions, list_ai_tools, list_hosts, prune_expired_stream_lineage, + prune_timeline_rollup, refresh_ai_session_rollup_if_stale, refresh_timeline_rollup, + rendered_session_page, search_ai_abuse, search_ai_anchors, search_ai_incidents, + search_ai_related_logs, search_ai_sessions, search_logs, severity_to_num, similar_incidents_clusters, tail_logs, timeline_rollup_status, topic_correlate_inputs, }; pub(crate) use skill_events::insert_skill_events_in_tx; diff --git a/src/db/agent_observatory.rs b/src/db/agent_observatory.rs index e0218b618..3cdeb2937 100644 --- a/src/db/agent_observatory.rs +++ b/src/db/agent_observatory.rs @@ -16,8 +16,8 @@ mod run_commits; #[cfg(test)] pub use run_commits::list_agent_run_commits; pub use run_commits::{ - AgentRunCommitUpsert, commit_attribution_evidence, git_commit_by_repository_sha, - upsert_agent_run_commit, + AgentRunAttributedCommit, AgentRunCommitUpsert, commit_attribution_evidence, + git_commit_by_repository_sha, list_agent_run_attributed_commits, upsert_agent_run_commit, }; #[path = "agent_observatory_sources.rs"] mod sources; @@ -61,12 +61,14 @@ pub use queries::{get_worktree_by_key, reconcile_repository}; #[path = "agent_observatory_read.rs"] mod read; pub use read::{ - AgentEventQuery, AgentRunQuery, EvidenceScopePage, EvidenceScopeQuery, ObservatoryEventRow, - ObservatoryMetricRow, ObservatoryRepositoryRow, ObservatoryRunRow, ObservatorySpanRow, - ObservatoryWorktreeRow, RepositoryQuery, RunTelemetryIdentity, TelemetryQuery, - list_observatory_events, list_observatory_metrics, list_observatory_repositories, - list_observatory_runs, list_observatory_spans, list_observatory_worktrees, - resolve_observatory_run, scoped_evidence_events, + AgentEventQuery, AgentRunQuery, EvidenceScopePage, EvidenceScopeQuery, ObservatoryActorRow, + ObservatoryEventRow, ObservatoryMetricRow, ObservatoryRepositoryRow, ObservatoryRunRow, + ObservatorySpanRow, ObservatoryWorktreeRow, RepositoryQuery, RunTelemetryIdentity, + TelemetryQuery, list_observatory_events, list_observatory_metrics, + list_observatory_repositories, list_observatory_run_actors, list_observatory_runs, + list_observatory_spans, list_observatory_worktrees, resolve_observatory_repository, + resolve_observatory_run, resolve_observatory_run_row, resolve_observatory_worktree, + scoped_evidence_events, }; use crate::db::pool::DbPool; diff --git a/src/db/agent_observatory_read.rs b/src/db/agent_observatory_read.rs index fbb86ad3c..eb320caca 100644 --- a/src/db/agent_observatory_read.rs +++ b/src/db/agent_observatory_read.rs @@ -2,7 +2,7 @@ use super::DbPool; use anyhow::{Context, Result}; -use rusqlite::{Row, params_from_iter, types::Value}; +use rusqlite::{OptionalExtension, Row, params_from_iter, types::Value}; #[path = "agent_observatory_read_cursor.rs"] mod cursor; use cursor::{bounded_limit, int_cursor, push_filter, text_cursor}; @@ -73,6 +73,74 @@ pub fn list_observatory_repositories( .context("list observatory repositories") } +pub fn resolve_observatory_repository( + pool: &DbPool, + repository_id: i64, +) -> Result> { + let conn = pool.get()?; + conn.query_row( + "SELECT r.id,r.repository_key,r.hostname,r.primary_path,r.display_name,r.first_seen_at,r.last_seen_at,r.removed_at,(SELECT COUNT(*) FROM repository_worktrees w WHERE w.repository_id=r.id),(SELECT COUNT(*) FROM agent_runs a JOIN agent_run_worktrees rw ON rw.run_id=a.id JOIN repository_worktrees w ON w.id=rw.worktree_id WHERE w.repository_id=r.id AND a.status IN ('starting','active','waiting','idle')) FROM repositories r WHERE r.id=?1", + [repository_id], + |r| { + Ok(ObservatoryRepositoryRow { + id: r.get(0)?, + key: r.get(1)?, + hostname: r.get(2)?, + primary_path: r.get(3)?, + name: r.get(4)?, + first_seen_at: r.get(5)?, + last_seen_at: r.get(6)?, + removed_at: r.get(7)?, + worktree_count: r.get(8)?, + active_run_count: r.get(9)?, + }) + }, + ) + .optional() + .context("resolve observatory repository") +} + +pub fn resolve_observatory_worktree( + pool: &DbPool, + worktree_id: i64, +) -> Result> { + let conn = pool.get()?; + conn.query_row( + "SELECT id,worktree_key,repository_id,hostname,path,branch_ref,branch_name,head_sha,upstream_ref,detached,bare,locked,lock_reason,prunable,prune_reason,dirty,staged_count,unstaged_count,untracked_count,ahead,behind,first_seen_at,last_seen_at,removed_at FROM repository_worktrees WHERE id=?1", + [worktree_id], + |r| { + Ok(ObservatoryWorktreeRow { + id: r.get(0)?, + key: r.get(1)?, + repository_id: r.get(2)?, + hostname: r.get(3)?, + path: r.get(4)?, + branch_ref: r.get(5)?, + branch: r.get(6)?, + head_sha: r.get(7)?, + upstream_ref: r.get(8)?, + detached: r.get(9)?, + bare: r.get(10)?, + locked: r.get(11)?, + lock_reason: r.get(12)?, + prunable: r.get(13)?, + prune_reason: r.get(14)?, + dirty: r.get(15)?, + staged: r.get(16)?, + unstaged: r.get(17)?, + untracked: r.get(18)?, + ahead: r.get(19)?, + behind: r.get(20)?, + first_seen_at: r.get(21)?, + last_seen_at: r.get(22)?, + removed_at: r.get(23)?, + }) + }, + ) + .optional() + .context("resolve observatory worktree") +} + /// Historical and follow-safe evidence projection for one branch/worktree. /// The keyset watermark is the durable `agent_run_events.id`; callers can /// replay a bounded page and then poll with `after_id` without a race window. @@ -233,8 +301,19 @@ pub fn list_observatory_runs( ) -> Result> { let conn = pool.get()?; let mut values = Vec::new(); - let mut sql="SELECT DISTINCT a.id,a.run_key,a.native_session_id,a.tool,a.provider_tool,a.hostname,a.status,a.status_reason,a.status_observed_at,a.started_at,a.last_activity_at,a.ended_at,a.transcript_path,a.primary_worktree_id,a.primary_branch,a.start_head_sha,a.current_head_sha,a.event_count,a.error_count,a.freshness_json FROM agent_runs a WHERE 1=1".to_string(); + let mut sql="SELECT DISTINCT a.id,a.run_key,a.native_session_id,a.tool,a.provider_tool,a.hostname,a.parent_run_id,a.previous_run_id,a.status,a.status_reason,a.status_observed_at,a.started_at,a.last_activity_at,a.ended_at,a.transcript_path,a.primary_worktree_id,a.primary_branch,a.start_head_sha,a.current_head_sha,a.event_count,a.error_count,a.freshness_json FROM agent_runs a WHERE 1=1".to_string(); push_filter(&mut sql, &mut values, "a.id <= ?", high_water); + if let Some(session_id) = &q.native_session_id { + push_filter( + &mut sql, + &mut values, + "a.native_session_id = ?", + session_id.clone(), + ); + } + if let Some(run_id) = q.exclude_run_id { + push_filter(&mut sql, &mut values, "a.id != ?", run_id); + } if let Some(id) = q.worktree_id { push_filter( &mut sql, @@ -265,11 +344,22 @@ pub fn list_observatory_runs( values.extend(q.statuses.iter().cloned().map(Value::from)); } if !q.tools.is_empty() { + let tool_column = if q.native_session_id.is_some() { + "lower(a.tool)" + } else { + "a.tool" + }; sql.push_str(&format!( - " AND a.tool IN ({})", + " AND {tool_column} IN ({})", vec!["?"; q.tools.len()].join(",") )); - values.extend(q.tools.iter().cloned().map(Value::from)); + values.extend(q.tools.iter().map(|tool| { + Value::from(if q.native_session_id.is_some() { + tool.to_ascii_lowercase() + } else { + tool.clone() + }) + })); } if q.active_only { sql.push_str(" AND a.status IN ('starting','active','waiting','idle')"); @@ -304,6 +394,54 @@ pub fn list_observatory_runs( .collect::>() .context("list observatory runs") } +pub fn resolve_observatory_run_row( + pool: &DbPool, + run_id: i64, +) -> Result> { + pool.get()? + .query_row( + "SELECT a.id,a.run_key,a.native_session_id,a.tool,a.provider_tool,a.hostname,a.parent_run_id,a.previous_run_id,a.status,a.status_reason,a.status_observed_at,a.started_at,a.last_activity_at,a.ended_at,a.transcript_path,a.primary_worktree_id,a.primary_branch,a.start_head_sha,a.current_head_sha,a.event_count,a.error_count,a.freshness_json FROM agent_runs a WHERE a.id=?1", + [run_id], + run_row, + ) + .optional() + .context("resolve observatory run row") +} + +pub fn list_observatory_run_actors( + pool: &DbPool, + run_id: i64, + limit: usize, +) -> Result> { + let conn = pool.get()?; + let mut stmt = conn.prepare( + "SELECT id,actor_key,run_id,native_actor_id,actor_type,display_name,started_at,last_activity_at,ended_at,metadata_json + FROM agent_run_actors + WHERE run_id=?1 + ORDER BY COALESCE(last_activity_at,started_at,'') DESC,id DESC + LIMIT ?2", + )?; + stmt.query_map( + rusqlite::params![run_id, (bounded_limit(limit, 200) + 1) as i64], + |r| { + Ok(ObservatoryActorRow { + id: r.get(0)?, + actor_key: r.get(1)?, + run_id: r.get(2)?, + native_actor_id: r.get(3)?, + actor_type: r.get(4)?, + display_name: r.get(5)?, + started_at: r.get(6)?, + last_activity_at: r.get(7)?, + ended_at: r.get(8)?, + metadata_json: r.get(9)?, + }) + }, + )? + .collect::>() + .context("list observatory run actors") +} + fn run_row(r: &Row<'_>) -> rusqlite::Result { Ok(ObservatoryRunRow { id: r.get(0)?, @@ -312,20 +450,22 @@ fn run_row(r: &Row<'_>) -> rusqlite::Result { tool: r.get(3)?, provider_tool: r.get(4)?, hostname: r.get(5)?, - status: r.get(6)?, - status_reason: r.get(7)?, - status_observed_at: r.get(8)?, - started_at: r.get(9)?, - last_activity_at: r.get(10)?, - ended_at: r.get(11)?, - transcript_path: r.get(12)?, - primary_worktree_id: r.get(13)?, - primary_branch: r.get(14)?, - start_head_sha: r.get(15)?, - current_head_sha: r.get(16)?, - event_count: r.get(17)?, - error_count: r.get(18)?, - freshness_json: r.get(19)?, + parent_run_id: r.get(6)?, + previous_run_id: r.get(7)?, + status: r.get(8)?, + status_reason: r.get(9)?, + status_observed_at: r.get(10)?, + started_at: r.get(11)?, + last_activity_at: r.get(12)?, + ended_at: r.get(13)?, + transcript_path: r.get(14)?, + primary_worktree_id: r.get(15)?, + primary_branch: r.get(16)?, + start_head_sha: r.get(17)?, + current_head_sha: r.get(18)?, + event_count: r.get(19)?, + error_count: r.get(20)?, + freshness_json: r.get(21)?, }) } diff --git a/src/db/agent_observatory_read_models.rs b/src/db/agent_observatory_read_models.rs index 4d25b276b..f04b24686 100644 --- a/src/db/agent_observatory_read_models.rs +++ b/src/db/agent_observatory_read_models.rs @@ -27,6 +27,8 @@ pub struct RepositoryQuery { } #[derive(Debug, Clone, Default, Serialize)] pub struct AgentRunQuery { + pub native_session_id: Option, + pub exclude_run_id: Option, pub repository_id: Option, pub worktree_id: Option, pub branch: Option, @@ -132,6 +134,10 @@ pub struct ObservatoryRunRow { pub tool: String, pub provider_tool: Option, pub hostname: String, + #[serde(serialize_with = "serialize_optional_id")] + pub parent_run_id: Option, + #[serde(serialize_with = "serialize_optional_id")] + pub previous_run_id: Option, pub status: String, pub status_reason: String, pub status_observed_at: String, @@ -148,6 +154,22 @@ pub struct ObservatoryRunRow { pub error_count: i64, pub freshness_json: String, } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct ObservatoryActorRow { + #[serde(serialize_with = "serialize_id")] + pub id: i64, + pub actor_key: String, + #[serde(serialize_with = "serialize_id")] + pub run_id: i64, + pub native_actor_id: String, + pub actor_type: Option, + pub display_name: Option, + pub started_at: Option, + pub last_activity_at: Option, + pub ended_at: Option, + pub metadata_json: String, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct ObservatoryEventRow { #[serde(serialize_with = "serialize_id")] diff --git a/src/db/agent_observatory_read_tests.rs b/src/db/agent_observatory_read_tests.rs index 77a2ae69b..53cae81f5 100644 --- a/src/db/agent_observatory_read_tests.rs +++ b/src/db/agent_observatory_read_tests.rs @@ -311,3 +311,108 @@ fn run_resolution_reports_unknown_runs_as_none() { assert_eq!(identity.provider_tool.as_deref(), Some("openai")); assert_eq!(identity.native_session_id, "session"); } + +#[test] +fn run_reads_expose_lineage_and_bounded_actors() { + let dir = tempfile::tempdir().unwrap(); + let pool = init_pool(&StorageConfig::for_test(dir.path().join("run-lineage.db"))).unwrap(); + let conn = pool.get().unwrap(); + conn.execute("INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,status,status_observed_at,started_at,last_activity_at) VALUES('parent','parent-session','codex','host','completed','2026-08-21T09:00:00Z','2026-08-21T09:00:00Z','2026-08-21T09:30:00Z')", []).unwrap(); + let parent_id = conn.last_insert_rowid(); + conn.execute("INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,status,status_observed_at,started_at,last_activity_at) VALUES('previous','previous-session','codex','host','completed','2026-08-21T09:30:00Z','2026-08-21T09:30:00Z','2026-08-21T09:45:00Z')", []).unwrap(); + let previous_id = conn.last_insert_rowid(); + conn.execute( + "INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,parent_run_id,previous_run_id,status,status_observed_at,started_at,last_activity_at) VALUES('current','session','codex','host',?1,?2,'active','2026-08-21T10:00:00Z','2026-08-21T10:00:00Z','2026-08-21T10:10:00Z')", + rusqlite::params![parent_id, previous_id], + ).unwrap(); + let run_id = conn.last_insert_rowid(); + for (key, native, activity) in [ + ("actor-a", "subagent-a", "2026-08-21T10:05:00Z"), + ("actor-b", "subagent-b", "2026-08-21T10:06:00Z"), + ] { + conn.execute( + "INSERT INTO agent_run_actors(actor_key,run_id,native_actor_id,actor_type,display_name,started_at,last_activity_at,metadata_json) VALUES(?1,?2,?3,'subagent',?3,'2026-08-21T10:00:00Z',?4,'{}')", + rusqlite::params![key, run_id, native, activity], + ).unwrap(); + } + drop(conn); + + let current = resolve_observatory_run_row(&pool, run_id).unwrap().unwrap(); + assert_eq!(current.parent_run_id, Some(parent_id)); + assert_eq!(current.previous_run_id, Some(previous_id)); + assert_eq!( + resolve_observatory_run_row(&pool, parent_id) + .unwrap() + .unwrap() + .run_key, + "parent" + ); + assert_eq!( + resolve_observatory_run_row(&pool, previous_id) + .unwrap() + .unwrap() + .run_key, + "previous" + ); + + let actors = list_observatory_run_actors(&pool, run_id, 1).unwrap(); + assert_eq!( + actors.len(), + 2, + "bounded reads return limit + 1 for truncation detection" + ); + assert_eq!(actors[0].native_actor_id, "subagent-b"); +} + +#[test] +fn exact_native_identity_is_filtered_before_run_limit() { + let dir = tempfile::tempdir().unwrap(); + let pool = init_pool(&StorageConfig::for_test(dir.path().join("exact.db"))).unwrap(); + let conn = pool.get().unwrap(); + conn.execute("INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,status,status_observed_at,started_at,last_activity_at) VALUES('exact','s','codex','h','completed','2026-09-21T00:00:00Z','2026-09-21T00:00:00Z','2026-09-21T00:00:00Z')", []).unwrap(); + for id in 0..60 { + conn.execute("INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,status,status_observed_at,started_at,last_activity_at) VALUES(?1,?1,'codex','h','active','2026-09-21T01:00:00Z','2026-09-21T01:00:00Z','2026-09-21T01:00:00Z')", [format!("prefix-s-{id}")]).unwrap(); + } + drop(conn); + let query = AgentRunQuery { + native_session_id: Some("s".into()), + tools: vec!["Codex".into()], + host: Some("h".into()), + ..Default::default() + }; + let rows = list_observatory_runs(&pool, &query, None, 1, i64::MAX).unwrap(); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0].run_key, "exact"); + pool.get().unwrap().execute("INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,status,status_observed_at,started_at,last_activity_at) VALUES('case-variant','s','CODEX','h','active','2026-09-21T01:00:00Z','2026-09-21T01:00:00Z','2026-09-21T01:00:00Z')", []).unwrap(); + assert_eq!( + list_observatory_runs(&pool, &query, None, 1, i64::MAX) + .unwrap() + .len(), + 2, + "sentinel retains exact ambiguity" + ); +} + +#[test] +fn attributed_commit_join_is_bounded_with_sentinel_and_preserves_fields() { + let dir = tempfile::tempdir().unwrap(); + let pool = init_pool(&StorageConfig::for_test(dir.path().join("commits.db"))).unwrap(); + let conn = pool.get().unwrap(); + conn.execute("INSERT INTO repositories(repository_key,hostname,common_git_dir,primary_path,display_name,first_seen_at,last_seen_at) VALUES('repo','h','/git','/p','p','2026-09-21T00:00:00Z','2026-09-21T00:00:00Z')", []).unwrap(); + let repository_id = conn.last_insert_rowid(); + conn.execute("INSERT INTO agent_runs(run_key,native_session_id,tool,hostname,status,status_observed_at,started_at,last_activity_at) VALUES('run','s','codex','h','active','2026-09-21T00:00:00Z','2026-09-21T00:00:00Z','2026-09-21T00:00:00Z')", []).unwrap(); + let run_id = conn.last_insert_rowid(); + for id in 0..3 { + conn.execute("INSERT INTO git_commits(repository_id,sha,subject,first_observed_at,last_observed_at) VALUES(?1,?2,?3,'2026-09-21T00:00:00Z','2026-09-21T00:00:00Z')", rusqlite::params![repository_id,format!("sha-{id}"),format!("subject-{id}")]).unwrap(); + let commit_id = conn.last_insert_rowid(); + conn.execute("INSERT INTO agent_run_commits(relation_key,run_id,commit_id,evidence_kind,evidence_source,trust_level,confidence,first_seen_at,last_seen_at) VALUES(?1,?2,?3,'attributed','test','claimed',0.5,'2026-09-21T00:00:00Z','2026-09-21T00:00:00Z')", rusqlite::params![format!("relation-{id}"),run_id,commit_id]).unwrap(); + } + drop(conn); + let commits = + crate::db::agent_observatory::list_agent_run_attributed_commits(&pool, run_id, 1).unwrap(); + assert_eq!(commits.len(), 2); + assert_eq!(commits[0].commit.sha, "sha-0"); + assert_eq!(commits[0].commit.subject, "subject-0"); + assert_eq!(commits[0].relation.commit_id, commits[0].commit.id); + assert_eq!(commits[0].relation.confidence, 0.5); +} diff --git a/src/db/agent_observatory_run_commits.rs b/src/db/agent_observatory_run_commits.rs index 1dd2c833e..3e58f8f86 100644 --- a/src/db/agent_observatory_run_commits.rs +++ b/src/db/agent_observatory_run_commits.rs @@ -29,6 +29,12 @@ pub struct AgentRunCommitRow { pub metadata_json: String, } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct AgentRunAttributedCommit { + pub relation: AgentRunCommitRow, + pub commit: super::GitCommitRow, +} + #[derive(Debug, Clone, PartialEq)] pub struct AgentRunCommitUpsert { pub run_id: i64, @@ -190,6 +196,53 @@ pub fn list_agent_run_commits(pool: &DbPool, run_id: i64) -> Result>>()?) } +pub fn list_agent_run_attributed_commits( + pool: &DbPool, + run_id: i64, + limit: usize, +) -> Result> { + if run_id <= 0 { + bail!("run_id must be positive"); + } + let connection = pool.get().context("acquire database connection")?; + let mut statement = connection.prepare( + "SELECT r.id,r.relation_key,r.run_id,r.commit_id,r.worktree_id,r.evidence_kind, + r.evidence_source,r.trust_level,r.confidence,r.first_seen_at,r.last_seen_at,r.metadata_json, + c.id,c.repository_id,c.sha,c.parent_shas_json,c.author_name,c.author_email_hash, + c.authored_at,c.committed_at,c.subject,c.changed_files,c.insertions,c.deletions, + c.changed_paths_json,c.first_observed_at,c.last_observed_at,c.reachable,c.metadata_json + FROM agent_run_commits r JOIN git_commits c ON c.id=r.commit_id + WHERE r.run_id=?1 ORDER BY r.first_seen_at,r.id LIMIT ?2", + )?; + statement + .query_map(params![run_id, (limit.clamp(1, 200) + 1) as i64], |item| { + Ok(AgentRunAttributedCommit { + relation: row(item)?, + commit: super::GitCommitRow { + id: item.get(12)?, + repository_id: item.get(13)?, + sha: item.get(14)?, + parent_shas_json: item.get(15)?, + author_name: item.get(16)?, + author_email_hash: item.get(17)?, + authored_at: item.get(18)?, + committed_at: item.get(19)?, + subject: item.get(20)?, + changed_files: item.get(21)?, + insertions: item.get(22)?, + deletions: item.get(23)?, + changed_paths_json: item.get(24)?, + first_observed_at: item.get(25)?, + last_observed_at: item.get(26)?, + reachable: item.get(27)?, + metadata_json: item.get(28)?, + }, + }) + })? + .collect::>() + .context("list attributed git commits") +} + pub fn commit_attribution_evidence( pool: &DbPool, worktree_id: i64, diff --git a/src/db/models.rs b/src/db/models.rs index 41850f8a3..9572f5a78 100644 --- a/src/db/models.rs +++ b/src/db/models.rs @@ -65,6 +65,7 @@ pub struct DockerCheckpoint { pub struct ListAiSessionsParams { pub ai_project: Option, pub ai_tool: Option, + pub ai_session_id: Option, pub host: Option, pub since: Option, pub until: Option, @@ -281,10 +282,12 @@ pub struct AiRelatedWindow { #[derive(Debug, Clone, Default)] pub struct SessionGraphInputs { pub bounds: Option<(String, String)>, + pub session_entity_keys: Vec, pub discovered_hosts: Vec, pub discovered_entities: Vec, pub used_graph: bool, pub logs: Vec, + pub source_fields_truncated: bool, } /// A graph entity matched while resolving a topic string, with how it matched diff --git a/src/db/queries.rs b/src/db/queries.rs index 3b24194fd..48a12a3a8 100644 --- a/src/db/queries.rs +++ b/src/db/queries.rs @@ -37,6 +37,10 @@ use super::models::{ use super::pool::DbPool; use super::queries_service_instances; +#[path = "queries_session_graph.rs"] +mod session_graph; +pub use session_graph::{SessionGraphScope, correlate_session_graph_scoped}; + const SEARCH_FTS_CANDIDATE_CAP: usize = 10_000; const SIMILAR_INCIDENT_FTS_CANDIDATE_CAP: usize = 5_000; /// Cap on the FTS match-set materialization in the fast (index-led) search @@ -584,7 +588,7 @@ pub fn list_ai_sessions( params: &ListAiSessionsParams, ) -> Result> { let time_filtered = params.since.is_some() || params.until.is_some(); - if !time_filtered && ai_session_rollup_is_populated(pool)? { + if !time_filtered && params.ai_session_id.is_none() && ai_session_rollup_is_populated(pool)? { return list_ai_sessions_from_rollup(pool, params); } list_ai_sessions_live(pool, params) @@ -833,6 +837,11 @@ pub fn list_ai_sessions_live( bindings.push(rusqlite::types::Value::Text(tool.clone())); idx += 1; } + if let Some(session_id) = ¶ms.ai_session_id { + sql.push_str(&format!(" AND ai_session_id = ?{idx}")); + bindings.push(rusqlite::types::Value::Text(session_id.clone())); + idx += 1; + } if let Some(hostname) = ¶ms.host { sql.push_str(&format!(" AND hostname = ?{idx}")); bindings.push(rusqlite::types::Value::Text(hostname.clone())); @@ -930,6 +939,11 @@ fn list_ai_sessions_from_rollup( bindings.push(rusqlite::types::Value::Text(tool.clone())); idx += 1; } + if let Some(session_id) = ¶ms.ai_session_id { + sql.push_str(&format!(" AND ai_session_id = ?{idx}")); + bindings.push(rusqlite::types::Value::Text(session_id.clone())); + idx += 1; + } if let Some(hostname) = ¶ms.host { sql.push_str(&format!(" AND hostname = ?{idx}")); bindings.push(rusqlite::types::Value::Text(hostname.clone())); @@ -2163,10 +2177,12 @@ pub fn correlate_session_graph( Ok(SessionGraphInputs { bounds: Some((start, end)), + session_entity_keys: session_keys, discovered_hosts, discovered_entities, used_graph, logs, + source_fields_truncated: false, }) } diff --git a/src/db/queries_session_graph.rs b/src/db/queries_session_graph.rs new file mode 100644 index 000000000..7438dc3b9 --- /dev/null +++ b/src/db/queries_session_graph.rs @@ -0,0 +1,168 @@ +use super::*; + +#[derive(Debug, Clone)] +pub struct SessionGraphScope { + pub project: String, + pub tool: String, + pub host: String, + pub severity_in: Vec, +} + +/// Correlate an already resolved session without treating its native ID as a +/// globally unique identity. Graph keys omit the host, so shared keys fall back +/// to the selected session's own evidence instead of attributing other hosts. +pub fn correlate_session_graph_scoped( + pool: &DbPool, + session_id: &str, + scope: &SessionGraphScope, + limit: usize, +) -> Result { + let limit = limit.clamp(1, 1000); + let conn = pool.get()?; + let bounds = conn.query_row( + "SELECT MIN(timestamp), MAX(timestamp) FROM logs + WHERE ai_session_id=?1 AND ai_project=?2 AND lower(ai_tool)=lower(?3) + AND hostname=?4", + params![session_id, scope.project, scope.tool, scope.host], + |row| { + Ok(( + row.get::<_, Option>(0)?, + row.get::<_, Option>(1)?, + )) + }, + )?; + let (Some(start), Some(end)) = bounds else { + return Ok(SessionGraphInputs::default()); + }; + let key = format!( + "{}:{}:{}", + scope.project.trim().to_ascii_lowercase(), + scope.tool.trim().to_ascii_lowercase(), + session_id + ); + let shared_key: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM logs WHERE ai_session_id=?1 + AND lower(trim(ai_project))=lower(trim(?2)) AND lower(trim(ai_tool))=lower(trim(?3)) + AND hostname<>?4)", + params![session_id, scope.project, scope.tool, scope.host], + |row| row.get(0), + )?; + let exists: bool = conn.query_row( + "SELECT EXISTS(SELECT 1 FROM graph_entities WHERE entity_type='ai_session' AND canonical_key=?1)", + [&key], |row| row.get(0), + )?; + let used_graph = exists && !shared_key; + let mut discovered_hosts = Vec::new(); + let mut discovered_entities = Vec::new(); + if used_graph { + for entity in super::super::graph::graph_walk_n_hops(&conn, std::slice::from_ref(&key), 2)? + { + match entity.entity_type.as_str() { + super::super::graph::ENTITY_TYPE_HOST => { + discovered_hosts.push(entity.canonical_key.clone()) + } + super::super::graph::ENTITY_TYPE_CONTAINER => { + if let Some(host) = + super::super::entity_resolution::container_key_host(&entity.canonical_key) + { + discovered_hosts.push(host.to_string()); + } + } + super::super::graph::ENTITY_TYPE_SERVICE_INSTANCE => { + if let Some((host, _)) = + super::super::entity_resolution::split_service_instance_key( + &entity.canonical_key, + ) + { + discovered_hosts.push(host.to_string()); + } + } + _ => {} + } + discovered_entities.push(entity.canonical_key); + } + discovered_hosts.sort(); + discovered_hosts.dedup(); + discovered_entities.sort(); + discovered_entities.dedup(); + } + let mut bindings: Vec = vec![ + session_id.to_string().into(), + scope.project.clone().into(), + scope.tool.clone().into(), + scope.host.clone().into(), + start.clone().into(), + end.clone().into(), + ]; + let own = + "(l.ai_session_id=?1 AND l.ai_project=?2 AND lower(l.ai_tool)=lower(?3) AND l.hostname=?4)"; + let mut selection = own.to_string(); + if used_graph && !discovered_hosts.is_empty() { + let first_host = bindings.len() + 1; + bindings.extend( + discovered_hosts + .iter() + .cloned() + .map(rusqlite::types::Value::Text), + ); + let hosts = (first_host..=bindings.len()) + .map(|index| format!("?{index}")) + .collect::>() + .join(","); + // Host evidence is a temporal correlation, never another AI session's + // transcript presented as part of the selected identity. + selection = format!("({own} OR (l.ai_session_id IS NULL AND l.hostname IN ({hosts})))"); + } + let severity = if scope.severity_in.is_empty() { + String::new() + } else { + let first = bindings.len() + 1; + bindings.extend( + scope + .severity_in + .iter() + .cloned() + .map(rusqlite::types::Value::Text), + ); + let placeholders = (first..=bindings.len()) + .map(|index| format!("?{index}")) + .collect::>() + .join(","); + format!(" AND l.severity IN ({placeholders})") + }; + bindings.push((limit as i64).into()); + // Cap text before materializing rows. JSON metadata is omitted intact + // rather than returning an invalid JSON prefix; propagate the omission. + let columns = FTS_SELECT_COLS + .replace("l.message", "substr(l.message,1,2048)") + .replace( + "l.metadata_json", + "CASE WHEN length(l.metadata_json)>4096 THEN NULL ELSE l.metadata_json END", + ); + let sql = format!( + "SELECT {columns},(length(l.message)>2048 OR COALESCE(length(l.metadata_json)>4096,0)) FROM logs l WHERE {selection} + AND l.timestamp>=?5 AND l.timestamp<=?6{severity} ORDER BY l.timestamp DESC,l.id DESC LIMIT ?{}", + bindings.len() + ); + let rows = conn + .prepare(&sql)? + .query_map(rusqlite::params_from_iter(bindings.iter()), |row| { + Ok((map_row(row)?, row.get::<_, bool>(15)?)) + })? + .collect::>>()?; + let source_fields_truncated = rows.iter().any(|(_, truncated)| *truncated); + let logs = rows.into_iter().map(|(entry, _)| entry).collect(); + Ok(SessionGraphInputs { + bounds: Some((start, end)), + session_entity_keys: if used_graph { vec![key] } else { Vec::new() }, + discovered_hosts, + discovered_entities, + used_graph, + logs, + source_fields_truncated, + }) +} + +#[cfg(test)] +#[path = "queries_session_graph_tests.rs"] +mod tests; diff --git a/src/db/queries_session_graph_tests.rs b/src/db/queries_session_graph_tests.rs new file mode 100644 index 000000000..b57726105 --- /dev/null +++ b/src/db/queries_session_graph_tests.rs @@ -0,0 +1,218 @@ +use super::*; +use crate::db::{LogBatchEntry, init_pool, insert_logs_batch}; + +fn fixture() -> (DbPool, tempfile::TempDir) { + let dir = tempfile::tempdir().unwrap(); + let pool = init_pool(&StorageConfig::for_test(dir.path().join("scope.db"))).unwrap(); + (pool, dir) +} + +fn row(time: &str, host: &str, project: &str, tool: &str, session: &str) -> LogBatchEntry { + LogBatchEntry { + timestamp: time.into(), + hostname: host.into(), + severity: "info".into(), + message: format!("{project}/{tool}/{host}"), + raw: "fixture".into(), + source_ip: "fixture://scope".into(), + ai_project: Some(project.into()), + ai_tool: Some(tool.into()), + ai_session_id: Some(session.into()), + facility: None, + app_name: None, + process_id: None, + docker_checkpoint: None, + ai_transcript_path: None, + metadata_json: None, + http_status: None, + auth_outcome: None, + dns_blocked: None, + event_action: None, + parse_error: None, + } +} + +fn scope() -> SessionGraphScope { + SessionGraphScope { + project: "chosen".into(), + tool: "codex".into(), + host: "host-a".into(), + severity_in: vec!["info".into()], + } +} + +fn graph_entity(conn: &rusqlite::Connection, kind: &str, key: &str) -> i64 { + conn.execute( + "INSERT INTO graph_entities(entity_type,canonical_key,display_label,trust_level) + VALUES(?1,?2,?2,'verified')", + params![kind, key], + ) + .unwrap(); + conn.last_insert_rowid() +} + +#[test] +fn fallback_reused_native_ids_preserve_selected_project_tool_host_and_bounds() { + let (pool, _dir) = fixture(); + insert_logs_batch( + &pool, + &[ + row( + "2026-09-21T08:00:00Z", + "host-a", + "chosen", + "CODEX", + "reused_%", + ), + row( + "2026-09-21T09:00:00Z", + "host-b", + "chosen", + "codex", + "reused_%", + ), + row( + "2026-09-21T10:00:00Z", + "host-a", + "other", + "codex", + "reused_%", + ), + row( + "2026-09-21T11:00:00Z", + "host-a", + "chosen", + "claude", + "reused_%", + ), + ], + ) + .unwrap(); + let result = correlate_session_graph_scoped(&pool, "reused_%", &scope(), 10).unwrap(); + assert!(!result.used_graph); + assert_eq!(result.logs.len(), 1); + assert_eq!(result.logs[0].message, "chosen/CODEX/host-a"); + assert_eq!( + result.bounds, + Some(("2026-09-21T08:00:00Z".into(), "2026-09-21T08:00:00Z".into())) + ); +} + +#[test] +fn graph_uses_exact_seed_and_keeps_host_context_without_other_session_transcripts() { + let (pool, _dir) = fixture(); + let mut context = row( + "2026-09-21T08:01:00Z", + "host-a", + "unused", + "unused", + "unused", + ); + context.ai_session_id = None; + context.ai_project = None; + context.ai_tool = None; + context.message = "host context".into(); + insert_logs_batch( + &pool, + &[ + row( + "2026-09-21T08:00:00Z", + "host-a", + "chosen", + "codex", + "reused_%", + ), + row( + "2026-09-21T08:02:00Z", + "host-a", + "chosen", + "codex", + "reused_%", + ), + row( + "2026-09-21T08:01:00Z", + "host-a", + "other", + "codex", + "reused_%", + ), + context, + ], + ) + .unwrap(); + { + let conn = pool.get().unwrap(); + let selected = graph_entity(&conn, "ai_session", "chosen:codex:reused_%"); + graph_entity(&conn, "ai_session", "other:codex:reused_%"); + let host = graph_entity(&conn, "host", "host-a"); + conn.execute("INSERT INTO graph_relationships(relationship_key,src_entity_id,dst_entity_id, + relationship_type,reason_code,trust_level,confidence,last_seen_at) + VALUES('scope-link',?1,?2,'runs_on','log_app_name','inferred',0.5,'2026-09-21T08:00:00Z')", params![selected,host]).unwrap(); + } + let result = correlate_session_graph_scoped(&pool, "reused_%", &scope(), 10).unwrap(); + assert!(result.used_graph); + assert_eq!(result.session_entity_keys, ["chosen:codex:reused_%"]); + assert_eq!(result.logs.len(), 3); + assert!(result.logs.iter().any(|row| row.message == "host context")); + assert!( + result + .logs + .iter() + .all(|row| row.ai_project.as_deref() != Some("other")) + ); +} + +#[test] +fn graph_key_shared_by_hosts_falls_back_without_claiming_other_host_evidence() { + let (pool, _dir) = fixture(); + insert_logs_batch( + &pool, + &[ + row("2026-09-21T08:00:00Z", "host-a", "chosen", "codex", "same"), + row("2026-09-21T08:01:00Z", "host-b", "chosen", "codex", "same"), + ], + ) + .unwrap(); + graph_entity(&pool.get().unwrap(), "ai_session", "chosen:codex:same"); + let result = correlate_session_graph_scoped(&pool, "same", &scope(), 10).unwrap(); + assert!(!result.used_graph); + assert!(result.session_entity_keys.is_empty()); + assert!(result.discovered_hosts.is_empty()); + assert_eq!(result.logs.len(), 1); + assert_eq!(result.logs[0].hostname, "host-a"); +} + +#[test] +fn severity_filter_precedes_the_correlation_row_limit() { + let (pool, _dir) = fixture(); + let mut error = row("2026-09-21T08:00:00Z", "host-a", "chosen", "codex", "same"); + error.severity = "err".into(); + insert_logs_batch( + &pool, + &[ + error, + row("2026-09-21T08:01:00Z", "host-a", "chosen", "codex", "same"), + row("2026-09-21T08:02:00Z", "host-a", "chosen", "codex", "same"), + ], + ) + .unwrap(); + let mut scope = scope(); + scope.severity_in = vec!["err".into()]; + let result = correlate_session_graph_scoped(&pool, "same", &scope, 1).unwrap(); + assert_eq!(result.logs.len(), 1); + assert_eq!(result.logs[0].severity, "err"); +} + +#[test] +fn oversized_source_fields_are_bounded_with_explicit_partial_evidence() { + let (pool, _dir) = fixture(); + let mut source = row("2026-09-21T08:00:00Z", "host-a", "chosen", "codex", "large"); + source.message = "😀".repeat(10_000); + source.metadata_json = + Some(serde_json::json!({"source_kind":"fixture", "large":"x".repeat(10_000)}).to_string()); + insert_logs_batch(&pool, &[source]).unwrap(); + let result = correlate_session_graph_scoped(&pool, "large", &scope(), 10).unwrap(); + assert!(result.source_fields_truncated); + assert_eq!(result.logs[0].message.chars().count(), 2048); + assert!(result.logs[0].metadata_json.is_none()); +} diff --git a/src/db/queries_tests.rs b/src/db/queries_tests.rs index 3af337009..bdce8441e 100644 --- a/src/db/queries_tests.rs +++ b/src/db/queries_tests.rs @@ -2054,6 +2054,7 @@ fn ai_session_queries_respect_filters() { &ListAiSessionsParams { ai_project: Some("/tmp/a".into()), ai_tool: Some("claude".into()), + ai_session_id: Some("s1".into()), host: Some("host-a".into()), since: Some("2026-01-01T00:00:00Z".into()), until: Some("2026-01-01T23:59:59Z".into()), @@ -2065,6 +2066,19 @@ fn ai_session_queries_respect_filters() { assert_eq!(listed.len(), 1); assert_eq!(listed[0].ai_session_id, "s1"); + let by_session = list_ai_sessions( + &pool, + &ListAiSessionsParams { + ai_session_id: Some("s2".into()), + limit: Some(10), + ..Default::default() + }, + ) + .unwrap(); + assert_eq!(by_session.len(), 1); + assert_eq!(by_session[0].ai_session_id, "s2"); + assert_eq!(by_session[0].ai_project, "/tmp/b"); + let searched = search_ai_sessions( &pool, &SearchAiSessionsParams { @@ -2390,6 +2404,7 @@ fn default_session_params() -> ListAiSessionsParams { ListAiSessionsParams { ai_project: None, ai_tool: None, + ai_session_id: None, host: None, since: None, until: None, @@ -3491,6 +3506,7 @@ fn bench_stats_and_sessions() { let params = ListAiSessionsParams { ai_project: None, ai_tool: None, + ai_session_id: None, host: None, since: None, until: None, @@ -3730,3 +3746,43 @@ fn lint_flags_unquoted_hyphen_term_alongside_a_quoted_phrase() { .to_string(); assert!(err.contains("NOT operator"), "{err}"); } + +#[test] +fn exact_session_lookup_finds_new_evidence_after_rollup_refresh() { + let (pool, _dir) = test_pool(); + insert_logs_batch( + &pool, + &[make_ai_entry( + "2026-01-01T00:00:00Z", + "host-a", + "claude", + "/tmp/project", + "old", + "old event", + )], + ) + .unwrap(); + refresh_ai_session_rollup(&pool).unwrap(); + insert_logs_batch( + &pool, + &[make_ai_entry( + "2026-01-01T00:01:00Z", + "host-a", + "claude", + "/tmp/project", + "fresh", + "fresh event", + )], + ) + .unwrap(); + let result = list_ai_sessions( + &pool, + &ListAiSessionsParams { + ai_session_id: Some("fresh".into()), + ..Default::default() + }, + ) + .unwrap(); + assert_eq!(result.len(), 1); + assert_eq!(result[0].ai_session_id, "fresh"); +} diff --git a/src/filetail/supervisor_tests.rs b/src/filetail/supervisor_tests.rs index fa036b39e..c3d7e9785 100644 --- a/src/filetail/supervisor_tests.rs +++ b/src/filetail/supervisor_tests.rs @@ -257,6 +257,7 @@ async fn supervisor_coalesces_checkpoint_writes_for_a_burst() { configured.start_at_end = false; registry.upsert(configured).unwrap(); let baseline_writes = registry.write_count(); + let burst_started = std::time::Instant::now(); supervisor.reconcile().await.unwrap(); let mut writer = tokio::fs::OpenOptions::new() @@ -291,8 +292,14 @@ async fn supervisor_coalesces_checkpoint_writes_for_a_burst() { .await .unwrap(); + // The supervisor may legitimately flush every 250ms while the fixture's + // asynchronous writes and reads are delayed by other tests. Allow those + // elapsed intervals plus initialization/final persistence, while rejecting + // a checkpoint write for each line. + let elapsed_flushes = burst_started.elapsed().as_millis().div_ceil(250) as usize; + let allowed_writes = (elapsed_flushes + 2).min(199); assert!( - registry.write_count() - baseline_writes <= 4, + registry.write_count() - baseline_writes <= allowed_writes, "checkpoint persistence scaled with line count" ); token.cancel(); diff --git a/src/mcp/actions.rs b/src/mcp/actions.rs index 71a913904..4fc22f403 100644 --- a/src/mcp/actions.rs +++ b/src/mcp/actions.rs @@ -72,6 +72,7 @@ pub(super) enum ActionHandler { ListApps, ListSessions, SearchSessions, + SessionInvestigate, EvidenceScope, SearchAbuse, AbuseIncidents, @@ -381,6 +382,17 @@ pub(super) const ACTION_SPECS: &[ActionSpec] = &[ Cheap, SearchSessions ), + action_spec!( + "session_investigate", + Read, + "Build a bounded evidence bundle rooted at one AI session", + Expensive, + SessionInvestigate, + exact: { + allowed: &["session_id", "tool", "project", "host", "limit", "severity_min"], + required: &["session_id"] + } + ), action_spec!( "evidence_scope", Read, diff --git a/src/mcp/rmcp_server_tests.rs b/src/mcp/rmcp_server_tests.rs index 8354d3283..491d826ea 100644 --- a/src/mcp/rmcp_server_tests.rs +++ b/src/mcp/rmcp_server_tests.rs @@ -159,6 +159,7 @@ fn minimal_args_for_action(action: &str) -> Value { match action { "correlate" => json!({"action": action, "reference_time": "2026-01-01T00:00:00Z"}), "search_sessions" => json!({"action": action, "query": "mounted"}), + "session_investigate" => json!({"action": action, "session_id": "mounted-session"}), "ai_correlate" => json!({"action": action, "project": "/tmp/project"}), "project_context" => json!({"action": action, "project": "/tmp/project"}), "context" => json!({"action": action, "log_id": 1}), diff --git a/src/mcp/tools.rs b/src/mcp/tools.rs index ccc041458..f6ec873bd 100644 --- a/src/mcp/tools.rs +++ b/src/mcp/tools.rs @@ -24,8 +24,9 @@ use crate::app::{ ListAiToolsRequest, ListAppsRequest, ListArtifactEvidenceRequest, ListHookEventsRequest, ListMcpEventsRequest, ListSessionsRequest, ListSkillEventsRequest, ListSourceIpsRequest, LlmInvocationsRequest, NotificationsRecentRequest, PatternsRequest, ProjectContextRequest, - SearchLogsRequest, SearchSessionsRequest, SilentHostsRequest, TailLogsRequest, TimelineRequest, - TopicCorrelateRequest, UnaddressedErrorsRequest, UsageBlocksRequest, + SearchLogsRequest, SearchSessionsRequest, SessionInvestigateRequest, SilentHostsRequest, + TailLogsRequest, TimelineRequest, TopicCorrelateRequest, UnaddressedErrorsRequest, + UsageBlocksRequest, }; use crate::artifact_evidence::ArtifactEvidenceInput; @@ -96,6 +97,7 @@ async fn dispatch_cortex_action( H::ListApps => tool_list_apps(state, args).await, H::ListSessions => tool_list_sessions(state, args).await, H::SearchSessions => tool_search_sessions(state, args).await, + H::SessionInvestigate => tool_session_investigate(state, args).await, H::EvidenceScope => tool_evidence_scope(state, args).await, H::SearchAbuse => tool_search_abuse(state, args).await, H::AbuseIncidents => tool_abuse_incidents(state, args).await, @@ -262,6 +264,12 @@ async fn tool_search_sessions(state: &AppState, args: Value) -> anyhow::Result anyhow::Result { + let req: SessionInvestigateRequest = action_payload(args, "session_investigate")?; + let response = state.service.session_investigate(req).await?; + Ok(serde_json::to_value(response)?) +} + async fn tool_search_abuse(state: &AppState, args: Value) -> anyhow::Result { let req: AbuseSearchRequest = action_payload(args, "abuse")?; let response = state.service.search_abuse(req).await?; diff --git a/src/mcp/tools_tests.rs b/src/mcp/tools_tests.rs index 3b15fafe8..8f6b13af9 100644 --- a/src/mcp/tools_tests.rs +++ b/src/mcp/tools_tests.rs @@ -1375,6 +1375,7 @@ fn sample_args_for_action(action: &str) -> Option { json!({"action": action, "reference_time": "2026-01-01T00:00:00Z"}) } "search_sessions" => json!({"action": action, "query": "schema"}), + "session_investigate" => json!({"action": action, "session_id": "schema-session"}), "evidence_scope" => json!({"action": action, "branch": "codex/schema-test"}), "ai_correlate" => json!({"action": action, "project": "/tmp/project"}), "topic_correlate" => json!({"action": action, "topic": "schema-test"}), @@ -1484,6 +1485,7 @@ fn typed_unknown_field_samples() -> Vec { "apps", "sessions", "search_sessions", + "session_investigate", "abuse", "abuse_incidents", "abuse_investigate", @@ -1606,6 +1608,35 @@ async fn schema_actions_are_dispatchable() { ) .unwrap(); } + db::insert_logs_batch( + &h.pool, + &[db::LogBatchEntry { + timestamp: "2026-01-01T00:00:00Z".to_string(), + hostname: "schema-session-host".to_string(), + facility: Some("agent".to_string()), + severity: "info".to_string(), + app_name: Some("codex".to_string()), + process_id: None, + message: "schema session transcript".to_string(), + raw: "schema session transcript".to_string(), + source_ip: "agent-command://schema-session-host/codex/schema-session".to_string(), + docker_checkpoint: None, + ai_tool: Some("codex".to_string()), + ai_project: Some("/schema/project".to_string()), + ai_session_id: Some("schema-session".to_string()), + ai_transcript_path: Some("/schema/transcript.jsonl".to_string()), + metadata_json: Some( + r#"{"source_kind":"agent-command","agent_command":{"cwd":"/schema/project"}}"# + .to_string(), + ), + http_status: None, + auth_outcome: None, + dns_blocked: None, + event_action: Some("command".to_string()), + parse_error: None, + }], + ) + .unwrap(); { let _guard = db::graph::GRAPH_TEST_LOCK.lock(); db::graph::refresh_graph_projection(&h.pool).unwrap(); @@ -1953,3 +1984,53 @@ async fn help_action_dispatch_returns_the_tool_reference() { assert!(reference.contains("## cortex search\n")); assert!(reference.contains("## cortex ack_error\n")); } + +#[tokio::test] +async fn session_investigate_rejects_unsupported_window_minutes() { + let harness = TestHarness::new(); + let error = execute_tool( + &harness.state, + "cortex", + json!({ + "action": "session_investigate", "session_id": "s", "window_minutes": 5 + }), + None, + ) + .await + .unwrap_err(); + assert!(error.to_string().contains("unknown field `window_minutes`")); +} + +#[tokio::test] +async fn session_investigate_dispatch_preserves_selected_identity_and_bounded_metadata() { + let harness = TestHarness::new(); + let conn = harness.pool.get().unwrap(); + for host in ["selected", "other"] { + conn.execute("INSERT INTO logs(timestamp,hostname,severity,message,raw,source_ip,ai_tool,ai_project,ai_session_id) VALUES('2026-09-21T00:00:00Z',?1,'info','user message','','fixture','codex','p','shared')", [host]).unwrap(); + } + drop(conn); + let response = execute_tool(&harness.state, "cortex", json!({ + "action": "session_investigate", "session_id": "shared", "host": "selected", "project": "p", "tool": "codex" + }), None).await.unwrap(); + assert_eq!(response["result"]["session"]["hostname"], "selected"); + assert!( + response["result"]["correlation"]["logs"] + .as_array() + .unwrap() + .iter() + .all(|log| log["entry"]["hostname"] == "selected") + ); + assert_eq!(response["metadata"]["auth_state"], "unknown"); + assert!( + response["metadata"]["budget_used"]["payload_bytes"] + .as_u64() + .unwrap() + <= 65_536 + ); + assert_eq!( + response["metadata"]["budget_used"]["payload_bytes"] + .as_u64() + .unwrap() as usize, + serde_json::to_vec(&response).unwrap().len() + ); +} diff --git a/src/surfaces.rs b/src/surfaces.rs index bbba6d7f2..40ab2ec90 100644 --- a/src/surfaces.rs +++ b/src/surfaces.rs @@ -332,6 +332,7 @@ pub const SURFACE_SPECS: &[SurfaceSpec] = &[ mcp!("abuse_incidents", Sessions, Canonical, Read), mcp!("abuse_investigate", Sessions, Canonical, Read), mcp!("ai_correlate", Sessions, Canonical, Read), + mcp!("session_investigate", Sessions, Canonical, Read), mcp!( "topic_correlate", Correlate, diff --git a/tests/TEST_COVERAGE.md b/tests/TEST_COVERAGE.md index d49643383..0241852c3 100644 --- a/tests/TEST_COVERAGE.md +++ b/tests/TEST_COVERAGE.md @@ -35,13 +35,13 @@ This table is generated from the compiled `SurfaceContract` and `profiles.json`; | Inventory | Count | |---|---:| -| mcp surfaces | 60 | +| mcp surfaces | 61 | | rest surfaces | 98 | | cli surfaces | 180 | | ingest surfaces | 32 | | artifact surfaces | 0 | | browser surfaces | 0 | -| all surfaces | 370 | +| all surfaces | 371 | | runnable profiles | 21 | Profiles: `agent`, `artifacts`, `auth`, `compose-isolated`, `docker-boundary-full`, `docker-boundary-reduced`, `fleet-mutating`, `fleet-read-only`, `full`, `isolated`, `legacy-central-pull`, `mcp`, `mutation`, `noop`, `notifications`, `security`, `smoke`, `soak`, `stateful`, `storage`, `upgrade` diff --git a/tests/live/phases/mcp/run.sh b/tests/live/phases/mcp/run.sh index 7f078ec61..161ceca65 100644 --- a/tests/live/phases/mcp/run.sh +++ b/tests/live/phases/mcp/run.sh @@ -214,6 +214,7 @@ mcp_phase_run() { artifact_evidence_record) jq -e '.result.structuredContent.inserted==true and .result.structuredContent.event.eventId=="mcp-live-event" and .result.structuredContent.event.artifactId=="mcp-live-artifact"' "$output" >/dev/null || result=fail ;; artifact_evidence) jq -e 'any(.result.structuredContent.events[]?;.eventId=="mcp-live-event" and .artifactId=="mcp-live-artifact")' "$output" >/dev/null || result=fail ;; search_sessions) jq -e --arg s "$MCP_LIVE_SESSION" 'any(.result.structuredContent.sessions[]?;.session_id==$s and .event_count>0)' "$output" >/dev/null || result=fail ;; + session_investigate) jq -e --arg s "$MCP_LIVE_SESSION" '.result.structuredContent.result.session.session_id==$s and (.result.structuredContent.result.transcript|length)>0 and .result.structuredContent.metadata.budget_used.payload_bytes<=65536' "$output" >/dev/null || result=fail ;; sessions) jq -e --arg s "$MCP_LIVE_SESSION" 'any(.result.structuredContent.sessions[]?;.session_id==$s and .event_count>0)' "$output" >/dev/null || result=fail ;; correlate) jq -e --arg h "$MCP_LIVE_ERROR_HOST" '.result.structuredContent.total_events>0 and any(.result.structuredContent.hosts[]?;.hostname==$h and (.events|length)>0)' "$output" >/dev/null || result=fail ;; correlate_state) jq -e --arg h "$MCP_LIVE_TOPIC_HOST" 'any(.result.structuredContent.hosts[]?;.hostname==$h and .heartbeat_summary!=null and (.logs|length)>0)' "$output" >/dev/null || result=fail ;; diff --git a/tests/live/phases/mcp/scenarios.json b/tests/live/phases/mcp/scenarios.json index cdd2e6f91..e32dd00a4 100644 --- a/tests/live/phases/mcp/scenarios.json +++ b/tests/live/phases/mcp/scenarios.json @@ -3,6 +3,7 @@ "arguments": { "search": {"query":"cortex"}, "evidence_scope": {"branch":"main"}, + "session_investigate": {"session_id":"mcp-live-session"}, "host_state": {"host":"missing-live-host"}, "correlate": {"query":"cortex","reference_time":"2026-08-27T12:00:00Z"}, "correlate_state": {"reference_time":"2026-08-27T12:00:00Z"}, @@ -43,7 +44,7 @@ "hook_investigate":"evidence","hosts":"hosts","incident_context":"total_logs","ingest_rate":"buckets","list_ai_projects":"projects", "list_ai_tools":"tools","llm_invocations":"$array","map":"schema","mcp_events":"events","mcp_incidents":"incidents","mcp_investigate":"evidence", "notifications_recent":"$array","patterns":"patterns","project_context":"project","search":"logs","search_sessions":"sessions","sessions":"sessions", - "evidence_scope":"items","recurring_error_comparison":"comparisons","silent_hosts":"hosts","similar_incidents":"clusters","skill_events":"events","skill_incidents":"incidents","skill_investigate":"evidence", + "session_investigate":"result","evidence_scope":"items","recurring_error_comparison":"comparisons","silent_hosts":"hosts","similar_incidents":"clusters","skill_events":"events","skill_incidents":"incidents","skill_investigate":"evidence", "source_ips":"source_ips","stats":"total_logs","status":"status","tail":"logs","timeline":"points","topic_correlate":"topic", "unaddressed_errors":"signatures","usage_blocks":"blocks","ack_error":"signature_hash","unack_error":"signature_hash","host_state":"host_id","compose_doctor":"diagnostics","notifications_test":"result","graph":"resolved_entity" }, diff --git a/xtask/src/pre_push.rs b/xtask/src/pre_push.rs index 2dcdebc8c..73785a081 100644 --- a/xtask/src/pre_push.rs +++ b/xtask/src/pre_push.rs @@ -359,10 +359,21 @@ fn configure_shell_command(command: &mut Command, step: &str) { command.arg("-c").arg(step); } +fn remove_repository_git_environment(root: &Path, command: &mut Command) -> Result<()> { + // Git hooks export repository-local variables. Validation creates temporary + // repositories, so inheriting these can redirect fixture writes into this + // checkout even when the fixture uses `git -C`. + for key in git_output(root, &["rev-parse", "--local-env-vars"])?.lines() { + command.env_remove(key); + } + Ok(()) +} + fn run_command(root: &Path, step: &PlanStep) -> Result<()> { println!("\n==> {}\n{}", step.name, step.command); let mut command = Command::new("bash"); configure_shell_command(&mut command, step.command); + remove_repository_git_environment(root, &mut command)?; command.current_dir(root); for (key, _) in std::env::vars() { if key.starts_with("CARGO_PROFILE_") { diff --git a/xtask/src/pre_push_tests.rs b/xtask/src/pre_push_tests.rs index 6285f0f48..52b4cdb04 100644 --- a/xtask/src/pre_push_tests.rs +++ b/xtask/src/pre_push_tests.rs @@ -76,6 +76,37 @@ fn pre_push_steps_preserve_the_callers_toolchain_environment() { assert_eq!(args, vec!["-c", "cargo --version"]); } +#[test] +fn validation_fixtures_cannot_modify_the_hook_repository() { + let temp = tempfile::tempdir().unwrap(); + let protected = temp.path().join("protected"); + let fixture = temp.path().join("fixture"); + std::fs::create_dir(&protected).unwrap(); + std::fs::create_dir(&fixture).unwrap(); + let mut init = Command::new("git"); + init.arg("-C").arg(&protected).arg("init"); + remove_repository_git_environment(Path::new("."), &mut init).unwrap(); + assert!(init.output().unwrap().status.success()); + let config_path = protected.join(".git/config"); + let original_config = std::fs::read(&config_path).unwrap(); + + let mut command = Command::new("bash"); + configure_shell_command( + &mut command, + "git -C \"$FIXTURE_PATH\" init && git -C \"$FIXTURE_PATH\" config core.bare true", + ); + command + .env("FIXTURE_PATH", &fixture) + .env("GIT_DIR", protected.join(".git")) + .env("GIT_WORK_TREE", &protected) + .env("GIT_INDEX_FILE", protected.join(".git/index")); + remove_repository_git_environment(&protected, &mut command).unwrap(); + assert!(command.output().unwrap().status.success()); + assert_eq!(std::fs::read(config_path).unwrap(), original_config); + assert!(fixture.join(".git/config").is_file()); + assert!(!protected.join(".git/index").exists()); +} + #[test] fn full_mode_keeps_the_old_expensive_suite_available() { let plan = plan_for(&["README.md"], true);