From 6380f6a2345914255c322c37637a69063c008d34 Mon Sep 17 00:00:00 2001 From: Jake Magar Date: Thu, 24 Sep 2026 11:20:49 -0400 Subject: [PATCH 1/4] feat(api): expose bounded graph discovery and changes --- README.md | 4 +- docs/api.md | 11 +- docs/architecture.md | 4 +- src/api.rs | 37 +++- src/api_tests.rs | 273 ++++++++++++++++++++++- src/app.rs | 8 + src/app/error.rs | 5 + src/app/models/graph.rs | 72 +++++++ src/app/services.rs | 39 ++-- src/app/services/graph_discovery.rs | 324 ++++++++++++++++++++++++++++ src/app/services/graph_safety.rs | 1 + src/app/services/graph_support.rs | 1 + src/db.rs | 1 + src/db/graph.rs | 82 ++++++- src/db/graph_discovery.rs | 306 ++++++++++++++++++++++++++ src/db/graph_inventory.rs | 1 + src/db/pool.rs | 89 +++++++- src/db/pool_tests.rs | 6 +- src/mcp/rmcp_server.rs | 5 +- src/surfaces/api.rs | 3 + 20 files changed, 1236 insertions(+), 36 deletions(-) create mode 100644 src/app/services/graph_discovery.rs create mode 100644 src/db/graph_discovery.rs diff --git a/README.md b/README.md index c73f682b5..0a390468c 100644 --- a/README.md +++ b/README.md @@ -17,7 +17,7 @@ Cortex began as a syslog receiver. It now covers network logs, Docker, managed f | Area | What Cortex provides | | --- | --- | | Ingest | UDP/TCP syslog, OTLP/HTTP logs, metrics, and traces, Docker logs and events, managed file tails, host heartbeats, AI transcripts, shell history, agent command records, and fleet inventory | -| Storage | SQLite in WAL mode, FTS5 full-text search, bounded metadata, retention, storage budgets, maintenance jobs, checkpoints, and 61 sequential schema migrations | +| Storage | SQLite in WAL mode, FTS5 full-text search, bounded metadata, retention, storage budgets, maintenance jobs, checkpoints, and 62 sequential schema migrations | | Investigation | Search, filtering, context, timelines, patterns, anomaly comparison, cross-source correlation, recurring error signatures, deterministic incident bundles, and graph explanations | | Fleet intelligence | SSH and API inventory collectors, host state, service topology, container and route relationships, redacted evidence, and rebuildable graph projections | | AI operations | Claude, Codex, Gemini CLI, and Antigravity session indexing; skill, MCP, and hook event extraction where each provider exposes them; incident clustering; and guarded local LLM assessments | @@ -701,7 +701,7 @@ Cortex uses SQLite with: - Online backup support - Integrity checks, checkpoints, and vacuum workflows -The current schema history contains 61 sequential migrations. CI derives this denominator from `KNOWN_SCHEMA_VERSION` and the migration registry. Migration 60 introduced the transcript-only `ai_logs_fts` projection. Migration 61 rebuilds that derived index with provider scope and adds an AI-only `(timestamp, ai_tool)` index so time-bounded session searches can derive a safe rowid floor before FTS walks common terms. On large transcript histories these one-time derived-index rebuilds can hold the startup write transaction while they populate. Subsequent transcript inserts and deletes maintain `ai_logs_fts` incrementally, while global `logs_fts` remains the full-corpus search index. Server forwarding receipts have a seven-day replay horizon and are removed when their canonical evidence is deleted. The sender spool is intentionally shorter and bounded: an individual source retains at most 1,024 records or 1 MiB and evicts records older than one day; the aggregate spool retains at most 4,096 records or 4 MiB. Eviction removes the original payload and retains a pending gap marker so evidence loss remains visible. A retry after the server receipt horizon is a new ingestion attempt. +The current schema history contains 62 sequential migrations. CI derives this denominator from `KNOWN_SCHEMA_VERSION` and the migration registry. Migration 60 introduced the transcript-only `ai_logs_fts` projection. Migration 61 rebuilds that derived index with provider scope and adds an AI-only `(timestamp, ai_tool)` index so time-bounded session searches can derive a safe rowid floor before FTS walks common terms. On large transcript histories these one-time derived-index rebuilds can hold the startup write transaction while they populate. Subsequent transcript inserts and deletes maintain `ai_logs_fts` incrementally, while global `logs_fts` remains the full-corpus search index. Server forwarding receipts have a seven-day replay horizon and are removed when their canonical evidence is deleted. The sender spool is intentionally shorter and bounded: an individual source retains at most 1,024 records or 1 MiB and evicts records older than one day; the aggregate spool retains at most 4,096 records or 4 MiB. Eviction removes the original payload and retains a pending gap marker so evidence loss remains visible. A retry after the server receipt horizon is a new ingestion attempt. ### Authoritative and derived data diff --git a/docs/api.md b/docs/api.md index bb7c981f2..921ea5a96 100644 --- a/docs/api.md +++ b/docs/api.md @@ -24,7 +24,7 @@ updated: 2026-07-30 ## Endpoint matrix -93 method/path bindings total. Scope is `read` (mounted via `axum::routing::get`, +96 method/path bindings total. Scope is `read` (mounted via `axum::routing::get`, hits read-side `db_permits`) or `admin`. Database maintenance and integrity checks share one process-wide maintenance gate; concurrent attempts receive a busy response. Admin mutations are audited before the service call. @@ -138,16 +138,21 @@ compatibility routes. | GET | `/api/compose/status` | read | (none) | `ComposeMcpStatus { container_name, ownership, runtime_state, health?, published_ports, diagnostics }` | 200, 401, 500 | Y | Redacted read-only projection. If the container cannot run Docker inspection, this still returns 200 with `runtime_state="docker_unavailable"` and diagnostic code `docker_unavailable`. | | GET | `/api/compose/doctor` | read | (none) | `ComposeMcpStatus { container_name, ownership, runtime_state, health?, published_ports, diagnostics }` | 200, 401, **503**, 500 | Y | Strict readiness check. Healthy Compose-owned deployment returns 200; Docker/ownership/runtime unready states return 503 with the same structured projection, not a generic error envelope. | -### Investigation graph queries (4) +### Investigation graph queries (7) | Method | Path | Scope | Request | Response (top-level) | Status codes | Idempotent | Notes | | --- | --- | --- | --- | --- | --- | --- | --- | | GET | `/api/graph/entity` | read | query: `entity_id?` or `entity_type` + `key` or `alias_type` + `alias_key`; `payload_budget?` | `GraphEntityLookupResponse { resolved_entity?, candidates, metadata }` | 200, 400, 401, 404, 503, 500 | Y | Resolves one graph entity by id, canonical key, or alias without rebuilding the projection. | +| GET | `/api/graph/entities` | read | query: `limit?` (1–25), `cursor?`, `entity_type?`, `source_kind?`, `trust_level?`, `query?` | `{ entities, next_cursor?, snapshot_cursor, metadata }` | 200, 400, 401, 409, 503, 500 | Y | Stable ID-ordered inventory page with exact filters and bounded label/key search. | +| GET | `/api/graph/relationships` | read | query: `limit?` (1–25), `cursor?`, `snapshot_cursor?` | `{ relationships, next_cursor?, snapshot_cursor, metadata }` | 200, 400, 401, 409, 503, 500 | Y | Relationship bootstrap from the same projection as entity pages. Includes up to 3 evidence IDs and their source kinds per relationship. | +| GET | `/api/graph/changes` | read | query: `cursor` (required), `limit?` (1–25) | `{ changes, next_cursor, metadata }` | 200, 400, 401, 410, 503, 500 | Y | Ordered entity/relationship upserts and tombstones after snapshot bootstrap. | | GET | `/api/graph/around` | read | query: entity selector, `depth?` (1 only), `limit?`, `evidence_sample_limit?`, `payload_budget?` | `GraphAroundResponse { resolved_entity, entities, relationships, evidence, metadata }` | 200, 400, 401, 404, 503, 500 | Y | Bounded one-hop neighborhood with allowlisted evidence samples. | | GET | `/api/graph/explain` | read | query: entity selector, `depth?` (clamped to 3), `beam_width?`, `max_chains?`, `evidence_sample_limit?`, `payload_budget?` | `GraphExplainResponse { resolved_entity, chains, narrative, open_questions, missing_evidence, next_queries, metadata }` | 200, 400, 401, 404, 503, 500 | Y | Deterministic evidence-backed explanation; weak evidence becomes open questions, not causal claims. | | GET | `/api/graph/evidence` | read | query: `evidence_id` (REQUIRED, minimum 1), `payload_budget?` | `GraphEvidenceLookupResponse { evidence, relationship, src_entity, dst_entity, source_log_summary?, missing_source_reason?, metadata }` | 200, 400, 401, 404, 503, 500 | Y | Proof lookup for one evidence row. Source summaries are redacted/truncated and exclude raw frames and raw metadata. | -**Total: 93 method/path bindings** (current surface registry, including syslog, +Inventory clients start at `/api/graph/entities`, follow `next_cursor` to exhaustion, then page `/api/graph/relationships` with the returned `snapshot_cursor`. A 409 means the projection changed between pages: discard that bootstrap and restart. After both inventories complete, poll `/api/graph/changes` with `snapshot_cursor`, applying events in sequence and storing each `next_cursor`. A 410 means the retained journal no longer covers the cursor: discard the local copy and bootstrap again. Change upserts reflect the latest committed projection; an object removed before its event is read becomes a tombstone. Responses report projection status, source watermark, completion time, degradation, and truncation. The journal retains at most seven days or 200,000 events, whichever is less; consumers should poll about once per minute and recover from expiry. This is a read-only surface under the same bearer authorization and heavy-read admission as other graph queries. Inventory exposes canonical `(entity_type, canonical_key)` identities, not aliases or raw logs/session content. Host, app, and service-instance names must not be merged by display label; resolve aliases with `/api/graph/entity`. Relationship keys are projection-scoped; apply tombstones and upserts when rebuilding changes endpoint IDs. + +**Total: 96 method/path bindings** (current surface registry, including syslog, surface-parity, AI, graph, compose, notification, error-ack, and DB routes; includes the 3 hook routes above, added alongside the `ai_hook_events` subsystem). diff --git a/docs/architecture.md b/docs/architecture.md index d05f7949b..3a9eaa709 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -32,14 +32,14 @@ SQLite database and a service layer: | `config.rs` | all | Layered config: defaults → `config.toml` → `~/.cortex/.env` → process env; startup validation (non-loopback auth gate) | | `runtime.rs` + `runtime/` | all | `RuntimeCore`: wires pool, ingest, auth policy; spawns the maintenance tasks below | | `app/` | core | `CortexService` service layer — shared limits/validation for MCP, REST, and CLI | -| `db/` | core | SQLite pool + 61 sequential migrations, FTS5 queries, retention and storage-budget maintenance | +| `db/` | core | SQLite pool + 62 sequential migrations, FTS5 queries, retention and storage-budget maintenance | | `receiver/` + `receiver.rs` | core | UDP + TCP listeners (supervised with restart + backoff), RFC 3164/5424 + CEF parsing | | `ingest.rs` | core | mpsc channel + batch writer (one pool connection reserved for this writer) | | `otlp.rs` + `otlp/` | core | OTLP/HTTP protobuf ingest: `POST /v1/logs` (4 MiB cap), `POST /v1/metrics` and `POST /v1/traces` (8 MiB cap); all use `CORTEX_TOKEN` auth | | `agent/`, `heartbeat_agent.rs` | inventory | Host-local cortex agent, including Docker log streaming from the local socket | | `docker_ingest/` | core | Legacy central pull container stdout/stderr + lifecycle events via explicit remote Docker Engine HTTP endpoints | | `mcp/` | core | RMCP Streamable HTTP server, `ACTION_SPECS` registry, scope gates, `/health` + `/health/full` | -| `api.rs` | core | Always-on `/api/*` REST surface (93 method/path bindings), bearer-token gated | +| `api.rs` | core | Always-on `/api/*` REST surface (96 method/path bindings), bearer-token gated | | `scanner/`, `sessions_watch.rs` | core | AI transcript scanning/scrubbing and the host-side watch daemon | | `inventory/` | inventory | Collectors (SSH, Docker, UniFi/Unraid/media APIs), redaction, normalized cache | | `heartbeat.rs` / `heartbeat_agent.rs` | inventory | `POST /v1/heartbeats` ingest + host-local heartbeat agent | diff --git a/src/api.rs b/src/api.rs index 3461bbb21..8c77e1bdb 100644 --- a/src/api.rs +++ b/src/api.rs @@ -36,7 +36,8 @@ use crate::app::{ CorrelateEventsRequest, CorrelateStateRequest, CortexService, DbBackupRequest, DbCheckpointRequest, DbIntegrityRequest, DbVacuumRequest, FeedLogsRequest, FileTailRequest, FilterLogsRequest, FleetStateRequest, GetErrorsRequest, GetLogRequest, GraphAroundRequest, - GraphEntityLookupRequest, GraphEvidenceLookupRequest, GraphExplainRequest, HostStateRequest, + GraphChangesRequest, GraphEntitiesRequest, GraphEntityLookupRequest, + GraphEvidenceLookupRequest, GraphExplainRequest, GraphRelationshipsRequest, HostStateRequest, IncidentContextRequest, IngestRateRequest, ListAiProjectsRequest, ListAiToolsRequest, ListAppsRequest, ListArtifactEvidenceRequest, ListHookEventsRequest, ListMcpEventsRequest, ListSessionsRequest, ListSkillEventsRequest, ListSourceIpsRequest, LlmInvocationsRequest, @@ -354,6 +355,9 @@ pub fn router(state: ApiState) -> anyhow::Result { ) .contract_route("GET /api/incident-context", get(incident_context)) .contract_route("GET /api/graph/entity", get(graph_entity)) + .contract_route("GET /api/graph/entities", get(graph_entities)) + .contract_route("GET /api/graph/relationships", get(graph_relationships)) + .contract_route("GET /api/graph/changes", get(graph_changes)) .contract_route("GET /api/graph/around", get(graph_around)) .contract_route("GET /api/graph/explain", get(graph_explain)) .contract_route("GET /api/graph/evidence", get(graph_evidence)) @@ -1758,6 +1762,27 @@ async fn graph_entity( respond(state.service.graph_entity_lookup(q).await) } +async fn graph_entities( + State(state): State, + Query(q): Query, +) -> axum::response::Response { + respond(state.service.graph_entities(q).await) +} + +async fn graph_relationships( + State(state): State, + Query(q): Query, +) -> axum::response::Response { + respond(state.service.graph_relationships(q).await) +} + +async fn graph_changes( + State(state): State, + Query(q): Query, +) -> axum::response::Response { + respond(state.service.graph_changes(q).await) +} + async fn graph_around( State(state): State, Query(q): Query, @@ -2371,6 +2396,16 @@ async fn ai_prune_checkpoints( fn respond(result: crate::app::ServiceResult) -> axum::response::Response { match result { Ok(value) => Json(value).into_response(), + Err(crate::app::ServiceError::Conflict(msg)) => ( + StatusCode::CONFLICT, + Json(json!({"error": msg, "recovery": "restart_snapshot"})), + ) + .into_response(), + Err(crate::app::ServiceError::Gone(msg)) => ( + StatusCode::GONE, + Json(json!({"error": msg, "recovery": "restart_snapshot"})), + ) + .into_response(), Err(crate::app::ServiceError::InvalidInput(msg)) => { (StatusCode::BAD_REQUEST, Json(json!({"error": msg}))).into_response() } diff --git a/src/api_tests.rs b/src/api_tests.rs index 46ea70a4a..0b8424980 100644 --- a/src/api_tests.rs +++ b/src/api_tests.rs @@ -15,9 +15,9 @@ use super::*; #[test] fn documented_rest_route_count_matches_router_registrations() { let binding_count = crate::surfaces::api_bindings().count(); - assert_eq!(binding_count, 93, "update the documented API denominator"); - assert!(include_str!("../docs/api.md").contains("93 method/path bindings total")); - assert!(include_str!("../docs/architecture.md").contains("(93 method/path bindings)")); + assert_eq!(binding_count, 96, "update the documented API denominator"); + assert!(include_str!("../docs/api.md").contains("96 method/path bindings total")); + assert!(include_str!("../docs/architecture.md").contains("(96 method/path bindings)")); } /// Build the router for a test, layering a `MockConnectInfo` so handlers @@ -3442,6 +3442,273 @@ async fn fleet_state_accepts_include_ok_and_sort_params() { // ─── /api/graph ───────────────────────────────────────────────────────────── +#[tokio::test] +async fn graph_inventory_lists_a_bounded_authenticated_snapshot() { + let (state, pool, _dir) = test_state(Some("secret".into())); + db::insert_logs_batch( + &pool, + &[entry( + "2026-01-01T00:00:00.000Z", + "graph-api-host", + "info", + "graph api seed", + "10.0.0.8:514", + )], + ) + .unwrap(); + { + let _guard = db::graph::GRAPH_TEST_LOCK.lock(); + db::graph::refresh_graph_projection(&pool).unwrap(); + } + let app = test_router(state); + let (unauthorized, _) = get_json(app.clone(), "/api/graph/entities?limit=1", None).await; + assert_eq!(unauthorized, axum::http::StatusCode::UNAUTHORIZED); + let (status, value) = get_json(app, "/api/graph/entities?limit=1", Some("secret")).await; + assert_eq!(status, axum::http::StatusCode::OK); + assert_eq!(value["entities"].as_array().unwrap().len(), 1); + assert!(value["snapshot_cursor"].is_string()); +} + +#[tokio::test] +async fn graph_inventory_cursor_restarts_when_projection_changes() { + let (state, pool, _dir) = test_state(Some("secret".into())); + db::insert_logs_batch( + &pool, + &[ + entry( + "2026-01-01T00:00:00.000Z", + "first-host", + "info", + "first", + "10.0.0.1:514", + ), + entry( + "2026-01-01T00:00:01.000Z", + "second-host", + "info", + "second", + "10.0.0.2:514", + ), + ], + ) + .unwrap(); + { + let _guard = db::graph::GRAPH_TEST_LOCK.lock(); + db::graph::refresh_graph_projection(&pool).unwrap(); + } + let app = test_router(state); + let (status, first) = + get_json(app.clone(), "/api/graph/entities?limit=1", Some("secret")).await; + assert_eq!(status, axum::http::StatusCode::OK); + let cursor = first["next_cursor"].as_str().unwrap(); + let (status, second) = get_json( + app.clone(), + &format!("/api/graph/entities?limit=1&cursor={cursor}"), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK); + assert_ne!(first["entities"][0]["id"], second["entities"][0]["id"]); + + pool.get() + .unwrap() + .execute( + "INSERT INTO graph_entities (entity_type, canonical_key, display_label, trust_level) + VALUES ('host', 'third-host', 'third-host', 'claimed')", + [], + ) + .unwrap(); + let (status, changed) = get_json( + app, + &format!("/api/graph/entities?limit=1&cursor={cursor}"), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::CONFLICT); + assert_eq!(changed["recovery"], "restart_snapshot"); +} + +#[tokio::test] +async fn graph_change_feed_reports_upserts_deletions_and_expiry() { + let (state, pool, _dir) = test_state(Some("secret".into())); + let app = test_router(state); + let (_, snapshot) = get_json(app.clone(), "/api/graph/entities?limit=1", Some("secret")).await; + let cursor = snapshot["snapshot_cursor"].as_str().unwrap(); + let id: i64 = pool + .get() + .unwrap() + .query_row( + "INSERT INTO graph_entities (entity_type, canonical_key, display_label, trust_level) + VALUES ('host', 'new-host', 'new-host', 'claimed') RETURNING id", + [], + |row| row.get(0), + ) + .unwrap(); + let (status, upsert) = get_json( + app.clone(), + &format!("/api/graph/changes?cursor={cursor}&limit=1"), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK); + assert_eq!(upsert["changes"][0]["operation"], "upsert"); + assert_eq!(upsert["changes"][0]["entity"]["canonical_key"], "new-host"); + let cursor = upsert["next_cursor"].as_str().unwrap(); + pool.get() + .unwrap() + .execute("DELETE FROM graph_entities WHERE id = ?1", [id]) + .unwrap(); + let (status, deletion) = get_json( + app.clone(), + &format!("/api/graph/changes?cursor={cursor}"), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK); + assert_eq!(deletion["changes"][0]["operation"], "delete"); + assert!(deletion["changes"][0]["entity"].is_null()); + pool.get() + .unwrap() + .execute("DELETE FROM graph_change_events WHERE seq = 1", []) + .unwrap(); + let (status, expired) = get_json( + app, + &format!( + "/api/graph/changes?cursor={}", + snapshot["snapshot_cursor"].as_str().unwrap() + ), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::GONE); + assert_eq!(expired["recovery"], "restart_snapshot"); +} + +#[tokio::test] +async fn graph_relationship_inventory_bootstraps_with_matching_snapshot_and_provenance() { + let (state, pool, _dir) = test_state(Some("secret".into())); + db::insert_logs_batch( + &pool, + &[entry( + "2026-01-01T00:00:00.000Z", + "graph-api-host", + "info", + "graph api seed", + "10.0.0.8:514", + )], + ) + .unwrap(); + { + let _guard = db::graph::GRAPH_TEST_LOCK.lock(); + db::graph::refresh_graph_projection(&pool).unwrap(); + } + let app = test_router(state); + let (_, snapshot) = get_json( + app.clone(), + "/api/graph/entities?entity_type=host&limit=1", + Some("secret"), + ) + .await; + assert_eq!(snapshot["entities"][0]["entity_type"], "host"); + let snapshot_cursor = snapshot["snapshot_cursor"].as_str().unwrap(); + let (status, page) = get_json( + app.clone(), + &format!("/api/graph/relationships?limit=1&snapshot_cursor={snapshot_cursor}"), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK); + assert_eq!(page["relationships"].as_array().unwrap().len(), 1); + assert!( + !page["relationships"][0]["evidence_ids"] + .as_array() + .unwrap() + .is_empty() + ); + assert!( + !page["relationships"][0]["source_kinds"] + .as_array() + .unwrap() + .is_empty() + ); + pool.get() + .unwrap() + .execute( + "UPDATE graph_relationship_evidence SET source_kind = 'heartbeat' + WHERE id = (SELECT MIN(id) FROM graph_relationship_evidence)", + [], + ) + .unwrap(); + let (change_status, changes) = get_json( + app.clone(), + &format!("/api/graph/changes?cursor={snapshot_cursor}"), + Some("secret"), + ) + .await; + assert_eq!(change_status, axum::http::StatusCode::OK); + assert!(changes["changes"].as_array().unwrap().iter().any(|change| { + change["object_kind"] == "relationship" && change["operation"] == "upsert" + })); + let (bad_status, _) = get_json(app, "/api/graph/entities?limit=26", Some("secret")).await; + assert_eq!(bad_status, axum::http::StatusCode::BAD_REQUEST); +} + +#[tokio::test] +async fn graph_full_rebuild_emits_tombstones_without_replaying_unchanged_rows() { + let (state, pool, _dir) = test_state(Some("secret".into())); + db::insert_logs_batch( + &pool, + &[entry( + "2026-01-01T00:00:00.000Z", + "kept-host", + "info", + "graph seed", + "10.0.0.8:514", + )], + ) + .unwrap(); + { + let _guard = db::graph::GRAPH_TEST_LOCK.lock(); + db::graph::refresh_graph_projection(&pool).unwrap(); + } + let app = test_router(state); + let (_, snapshot) = get_json(app.clone(), "/api/graph/entities?limit=1", Some("secret")).await; + let cursor = snapshot["snapshot_cursor"].as_str().unwrap(); + pool.get() + .unwrap() + .execute( + "INSERT INTO graph_entities (entity_type, canonical_key, display_label, trust_level) + VALUES ('host', 'temporary-host', 'temporary-host', 'claimed')", + [], + ) + .unwrap(); + { + let _guard = db::graph::GRAPH_TEST_LOCK.lock(); + db::graph::refresh_graph_projection(&pool).unwrap(); + } + let (status, page) = get_json( + app.clone(), + &format!("/api/graph/changes?cursor={cursor}&limit=25"), + Some("secret"), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK); + let events = page["changes"].as_array().unwrap(); + assert!(events.iter().any(|event| { + event["operation"] == "delete" + && event["item_key"] + .as_str() + .unwrap() + .contains("temporary-host") + })); + assert!( + !events + .iter() + .any(|event| event["item_key"].as_str().unwrap().contains("kept-host")), + "unchanged entities must not be replayed during a full rebuild" + ); +} + #[tokio::test] async fn graph_routes_return_shared_service_payloads() { let (state, pool, _dir) = test_state(Some("secret".into())); diff --git a/src/app.rs b/src/app.rs index 34139842d..0c92e6cb9 100644 --- a/src/app.rs +++ b/src/app.rs @@ -120,6 +120,12 @@ pub use models::{ GetLogResponse, GraphAroundRequest, GraphAroundResponse, + GraphChange, + GraphChangesRequest, + GraphChangesResponse, + GraphDiscoveryMetadata, + GraphEntitiesRequest, + GraphEntitiesResponse, GraphEntity, GraphEntityCandidate, GraphEntityLookupRequest, @@ -137,6 +143,8 @@ pub use models::{ GraphRebuildResponse, GraphRebuildStatsResponse, GraphRelationship, + GraphRelationshipsRequest, + GraphRelationshipsResponse, GraphResponseMetadata, GraphSourceLogSummary, HomelabMapRequest, diff --git a/src/app/error.rs b/src/app/error.rs index 4cf267257..f40394e2d 100644 --- a/src/app/error.rs +++ b/src/app/error.rs @@ -8,6 +8,11 @@ use thiserror::Error; /// result chains until all call sites are migrated to explicit `map_err`. #[derive(Debug, Error)] pub enum ServiceError { + #[error("{0}")] + Conflict(String), + + #[error("{0}")] + Gone(String), /// Caller-supplied argument was invalid. #[error("{0}")] InvalidInput(String), diff --git a/src/app/models/graph.rs b/src/app/models/graph.rs index 3d1b67905..651ade635 100644 --- a/src/app/models/graph.rs +++ b/src/app/models/graph.rs @@ -63,6 +63,76 @@ pub struct GraphExplainRequest { pub payload_budget: Option, } +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct GraphEntitiesRequest { + pub cursor: Option, + pub limit: Option, + pub entity_type: Option, + pub source_kind: Option, + pub trust_level: Option, + pub query: Option, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct GraphRelationshipsRequest { + pub cursor: Option, + pub snapshot_cursor: Option, + pub limit: Option, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct GraphChangesRequest { + pub cursor: String, + pub limit: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct GraphDiscoveryMetadata { + pub projection_status: String, + pub source_watermark: String, + pub last_completed_at: Option, + pub is_degraded: bool, + pub truncated: bool, + pub recovery: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct GraphEntitiesResponse { + pub entities: Vec, + pub next_cursor: Option, + pub snapshot_cursor: String, + pub metadata: GraphDiscoveryMetadata, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct GraphRelationshipsResponse { + pub relationships: Vec, + pub next_cursor: Option, + pub snapshot_cursor: String, + pub metadata: GraphDiscoveryMetadata, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct GraphChange { + pub seq: i64, + pub object_kind: String, + pub operation: String, + pub item_key: String, + pub occurred_at: String, + pub entity: Option, + pub relationship: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct GraphChangesResponse { + pub changes: Vec, + pub next_cursor: String, + pub metadata: GraphDiscoveryMetadata, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct GraphProjectionStatusResponse { pub projection_status: String, @@ -161,6 +231,8 @@ pub struct GraphRelationship { pub confidence: f64, pub evidence_count: i64, pub evidence_ids: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub source_kinds: Vec, pub first_seen_at: Option, pub last_seen_at: Option, } diff --git a/src/app/services.rs b/src/app/services.rs index 777a3fb93..d858a9128 100644 --- a/src/app/services.rs +++ b/src/app/services.rs @@ -29,24 +29,26 @@ use super::models::{ DbIntegrityResult, DbMaintenanceStatus, DbStats, DbVacuumRequest, DbVacuumResult, FeedLogsRequest, FeedLogsResponse, FilterLogsRequest, FleetStateHostRow, FleetStateRequest, FleetStateResponse, FleetStateSummary, GetErrorsRequest, GetErrorsResponse, GetLogRequest, - GetLogResponse, GraphAroundRequest, GraphAroundResponse, GraphEntity, GraphEntityCandidate, - GraphEntityLookupRequest, GraphEntityLookupResponse, GraphEntitySummary, GraphEvidence, - GraphEvidenceLookupRequest, GraphEvidenceLookupResponse, GraphExplainRequest, - GraphExplainResponse, GraphIncidentNarrative, GraphNarrativeChain, GraphNextQuery, - GraphProjectionStatusResponse, GraphRebuildResponse, GraphRebuildStatsResponse, - GraphRelationship, GraphResponseMetadata, GraphSessionCorrelation, GraphSourceLogSummary, - HomelabMapAnswerRow, HomelabMapAnswerTruncation, HomelabMapGraphAnswer, HomelabMapGraphTarget, - HomelabMapNextQuery, HomelabMapNode, HomelabMapProofQuery, HomelabMapRequest, - HomelabMapResponse, HomelabMapSummary, HookIncidentEvidence, HookIncidentSummary, - INVESTIGATION_UI_VERSION, IncidentContextRequest, IncidentContextResponse, IncidentEvent, - IncidentRequest, IncidentResponse, IngestRateRequest, IngestRateResponse, InvestigationBudget, - InvestigationBudgetUsed, InvestigationClaim, InvestigationClaimType, InvestigationEnvelope, - InvestigationMetadata, ListAiProjectsRequest, ListAiProjectsResponse, ListAiToolsRequest, - ListAiToolsResponse, ListAppsRequest, ListAppsResponse, ListHostsResponse, ListSessionsRequest, - ListSessionsResponse, ListSourceIpsRequest, ListSourceIpsResponse, LlmInvocationsRequest, - LogEntry, MaintenanceJobStatus, McpIncidentEvidence, McpIncidentSummary, - NotificationsRecentRequest, PatternsRequest, PatternsResponse, ProjectContextRequest, - ProjectContextResponse, RecurringErrorComparisonEntry, RecurringErrorComparisonRequest, + GetLogResponse, GraphAroundRequest, GraphAroundResponse, GraphChange, GraphChangesRequest, + GraphChangesResponse, GraphDiscoveryMetadata, GraphEntitiesRequest, GraphEntitiesResponse, + GraphEntity, GraphEntityCandidate, GraphEntityLookupRequest, GraphEntityLookupResponse, + GraphEntitySummary, GraphEvidence, GraphEvidenceLookupRequest, GraphEvidenceLookupResponse, + GraphExplainRequest, GraphExplainResponse, GraphIncidentNarrative, GraphNarrativeChain, + GraphNextQuery, GraphProjectionStatusResponse, GraphRebuildResponse, GraphRebuildStatsResponse, + GraphRelationship, GraphRelationshipsRequest, GraphRelationshipsResponse, + GraphResponseMetadata, GraphSessionCorrelation, GraphSourceLogSummary, HomelabMapAnswerRow, + HomelabMapAnswerTruncation, HomelabMapGraphAnswer, HomelabMapGraphTarget, HomelabMapNextQuery, + HomelabMapNode, HomelabMapProofQuery, HomelabMapRequest, HomelabMapResponse, HomelabMapSummary, + HookIncidentEvidence, HookIncidentSummary, INVESTIGATION_UI_VERSION, IncidentContextRequest, + IncidentContextResponse, IncidentEvent, IncidentRequest, IncidentResponse, IngestRateRequest, + IngestRateResponse, InvestigationBudget, InvestigationBudgetUsed, InvestigationClaim, + InvestigationClaimType, InvestigationEnvelope, InvestigationMetadata, ListAiProjectsRequest, + ListAiProjectsResponse, ListAiToolsRequest, ListAiToolsResponse, ListAppsRequest, + ListAppsResponse, ListHostsResponse, ListSessionsRequest, ListSessionsResponse, + ListSourceIpsRequest, ListSourceIpsResponse, LlmInvocationsRequest, LogEntry, + MaintenanceJobStatus, McpIncidentEvidence, McpIncidentSummary, NotificationsRecentRequest, + PatternsRequest, PatternsResponse, ProjectContextRequest, ProjectContextResponse, + RecurringErrorComparisonEntry, RecurringErrorComparisonRequest, RecurringErrorComparisonResponse, RecurringErrorEvidenceBundle, RecurringErrorNextQuery, RequestActor, ResolvedTopicEntity, SearchLogsRequest, SearchLogsResponse, SearchSessionsRequest, SearchSessionsResponse, SearchedSessionEntry, ServiceJournalEntry, @@ -101,6 +103,7 @@ mod error_detection; mod file_tails; mod filters; mod graph; +mod graph_discovery; mod graph_limits; mod graph_safety; mod graph_support; diff --git a/src/app/services/graph_discovery.rs b/src/app/services/graph_discovery.rs new file mode 100644 index 000000000..71d35f35d --- /dev/null +++ b/src/app/services/graph_discovery.rs @@ -0,0 +1,324 @@ +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; + +use super::graph_safety::{graph_entity_safe, graph_relationship_safe, redact_graph_text}; +use super::graph_support::graph_relationship_to_model; +use super::*; + +const MAX_PAGE: u32 = 25; +const MAX_RESPONSE_BYTES: usize = 65_536; + +#[derive(Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Cursor { + version: u8, + kind: String, + anchor: i64, + after: i64, + filter: String, +} + +fn encode_cursor(kind: &str, anchor: i64, after: i64, filter: &str) -> String { + hex::encode( + serde_json::to_vec(&Cursor { + version: 1, + kind: kind.into(), + anchor, + after, + filter: filter.into(), + }) + .expect("cursor is serializable"), + ) +} + +fn decode_cursor(encoded: &str, kind: &str, filter: &str) -> ServiceResult { + if encoded.len() > 512 { + return Err(ServiceError::InvalidInput("graph cursor too long".into())); + } + let raw = hex::decode(encoded) + .map_err(|_| ServiceError::InvalidInput("invalid graph cursor".into()))?; + let cursor: Cursor = serde_json::from_slice(&raw) + .map_err(|_| ServiceError::InvalidInput("invalid graph cursor".into()))?; + if cursor.version != 1 + || cursor.kind != kind + || cursor.filter != filter + || cursor.anchor < 0 + || cursor.after < 0 + { + return Err(ServiceError::InvalidInput( + "graph cursor does not match this request".into(), + )); + } + Ok(cursor) +} + +fn page_limit(limit: Option) -> ServiceResult { + let limit = limit.unwrap_or(MAX_PAGE); + if !(1..=MAX_PAGE).contains(&limit) { + return Err(ServiceError::InvalidInput( + "graph page limit must be 1..=25".into(), + )); + } + Ok(limit) +} + +fn metadata( + projection_status: String, + source_watermark: String, + last_completed_at: Option, + is_degraded: bool, + truncated: bool, +) -> GraphDiscoveryMetadata { + GraphDiscoveryMetadata { + projection_status, + source_watermark, + last_completed_at, + is_degraded, + truncated, + recovery: None, + } +} + +fn relation(item: db::graph_discovery::RelationshipItem) -> GraphRelationship { + let mut relationship = graph_relationship_to_model(item.row, None, None, item.evidence_ids); + relationship.source_kinds = item.source_kinds; + graph_relationship_safe(relationship) +} + +impl CortexService { + pub async fn graph_entities( + &self, + req: GraphEntitiesRequest, + ) -> ServiceResult { + let limit = page_limit(req.limit)?; + let kind = req.entity_type.unwrap_or_default(); + let source = req.source_kind.unwrap_or_default(); + let trust = req.trust_level.unwrap_or_default(); + let query = req.query.unwrap_or_default().trim().to_owned(); + if (!kind.is_empty() && !db::graph::is_known_entity_type(&kind)) + || (!trust.is_empty() + && !["verified", "claimed", "inferred", "correlated", "refuted"] + .contains(&trust.as_str())) + || source.len() > 64 + || !source + .chars() + .all(|c| c.is_ascii_alphanumeric() || "_:-.".contains(c)) + || query.len() > 80 + || query.chars().any(char::is_control) + { + return Err(ServiceError::InvalidInput( + "invalid graph inventory filter".into(), + )); + } + let fingerprint = hex::encode(Sha256::digest( + format!("{kind}\0{source}\0{trust}\0{query}").as_bytes(), + )); + let cursor = req + .cursor + .as_deref() + .map(|value| decode_cursor(value, "entities", &fingerprint)) + .transpose()?; + let after = cursor.as_ref().map_or(0, |cursor| cursor.after); + let expected_anchor = cursor.as_ref().map(|cursor| cursor.anchor); + self.with_heavy_read_permit("graph.entities", || async move { + let page = self + .run_db("graph.entities", move |pool| { + db::graph_discovery::entities( + pool, + after, + expected_anchor, + limit, + db::graph_discovery::EntityFilters { + entity_type: &kind, + source_kind: &source, + trust_level: &trust, + query: &query, + }, + ) + }) + .await?; + if page.changed { + return Err(ServiceError::Conflict( + "graph snapshot changed during pagination".into(), + )); + } + let next_cursor = page.has_more.then(|| { + encode_cursor( + "entities", + page.anchor, + page.rows.last().expect("nonempty when more").id, + &fingerprint, + ) + }); + let mut response = GraphEntitiesResponse { + entities: page + .rows + .into_iter() + .map(|row| graph_entity_safe(row.into())) + .collect(), + next_cursor, + snapshot_cursor: encode_cursor("changes", page.anchor, page.anchor, ""), + metadata: metadata( + page.projection_status, + page.source_watermark, + page.last_completed_at, + page.is_degraded, + page.has_more, + ), + }; + while serde_json::to_vec(&response).is_ok_and(|bytes| bytes.len() > MAX_RESPONSE_BYTES) + { + response.entities.pop(); + let last = response.entities.last().ok_or_else(|| { + ServiceError::Internal(anyhow::anyhow!("graph entity exceeds response budget")) + })?; + response.next_cursor = Some(encode_cursor( + "entities", + page.anchor, + last.id, + &fingerprint, + )); + response.metadata.truncated = true; + } + Ok(response) + }) + .await + } + + pub async fn graph_relationships( + &self, + req: GraphRelationshipsRequest, + ) -> ServiceResult { + let limit = page_limit(req.limit)?; + let cursor = req + .cursor + .as_deref() + .map(|value| decode_cursor(value, "relationships", "")) + .transpose()?; + let snapshot = req + .snapshot_cursor + .as_deref() + .map(|value| decode_cursor(value, "changes", "")) + .transpose()?; + if let (Some(cursor), Some(snapshot)) = (&cursor, &snapshot) + && cursor.anchor != snapshot.anchor + { + return Err(ServiceError::InvalidInput( + "graph snapshot cursors disagree".into(), + )); + } + let after = cursor.as_ref().map_or(0, |cursor| cursor.after); + let expected_anchor = cursor + .as_ref() + .map(|cursor| cursor.anchor) + .or_else(|| snapshot.as_ref().map(|cursor| cursor.anchor)); + self.with_heavy_read_permit("graph.relationships", || async move { + let page = self + .run_db("graph.relationships", move |pool| { + db::graph_discovery::relationships(pool, after, expected_anchor, limit) + }) + .await?; + if page.changed { + return Err(ServiceError::Conflict( + "graph snapshot changed during pagination".into(), + )); + } + let next_cursor = page.has_more.then(|| { + encode_cursor( + "relationships", + page.anchor, + page.rows.last().expect("nonempty when more").row.id, + "", + ) + }); + let mut response = GraphRelationshipsResponse { + relationships: page.rows.into_iter().map(relation).collect(), + next_cursor, + snapshot_cursor: encode_cursor("changes", page.anchor, page.anchor, ""), + metadata: metadata( + page.projection_status, + page.source_watermark, + page.last_completed_at, + page.is_degraded, + page.has_more, + ), + }; + while serde_json::to_vec(&response).is_ok_and(|bytes| bytes.len() > MAX_RESPONSE_BYTES) + { + response.relationships.pop(); + let last = response.relationships.last().ok_or_else(|| { + ServiceError::Internal(anyhow::anyhow!( + "graph relationship exceeds response budget" + )) + })?; + response.next_cursor = + Some(encode_cursor("relationships", page.anchor, last.id, "")); + response.metadata.truncated = true; + } + Ok(response) + }) + .await + } + + pub async fn graph_changes( + &self, + req: GraphChangesRequest, + ) -> ServiceResult { + let limit = page_limit(req.limit)?; + let cursor = decode_cursor(&req.cursor, "changes", "")?; + if cursor.after != cursor.anchor { + return Err(ServiceError::InvalidInput("invalid change cursor".into())); + } + self.with_heavy_read_permit("graph.changes", || async move { + let page = self + .run_db("graph.changes", move |pool| { + db::graph_discovery::changes(pool, cursor.after, limit) + }) + .await?; + if page.expired { + return Err(ServiceError::Gone("graph change cursor expired".into())); + } + let next = page.rows.last().map_or(page.latest, |event| event.seq); + let mut response = GraphChangesResponse { + changes: page + .rows + .into_iter() + .map(|event| GraphChange { + seq: event.seq, + object_kind: event.object_kind, + operation: event.operation, + item_key: redact_graph_text(event.item_key), + occurred_at: event.occurred_at, + entity: event.entity.map(|row| graph_entity_safe(row.into())), + relationship: event.relationship.map(|row| { + relation(db::graph_discovery::RelationshipItem { + row, + evidence_ids: event.evidence_ids, + source_kinds: event.source_kinds, + }) + }), + }) + .collect(), + next_cursor: encode_cursor("changes", next, next, ""), + metadata: metadata( + page.projection_status, + page.source_watermark, + page.last_completed_at, + page.is_degraded, + page.has_more, + ), + }; + while serde_json::to_vec(&response).is_ok_and(|bytes| bytes.len() > MAX_RESPONSE_BYTES) + { + response.changes.pop(); + let last = response.changes.last().ok_or_else(|| { + ServiceError::Internal(anyhow::anyhow!("graph change exceeds response budget")) + })?; + response.next_cursor = encode_cursor("changes", last.seq, last.seq, ""); + response.metadata.truncated = true; + } + Ok(response) + }) + .await + } +} diff --git a/src/app/services/graph_safety.rs b/src/app/services/graph_safety.rs index 3cd66768d..e360ad19e 100644 --- a/src/app/services/graph_safety.rs +++ b/src/app/services/graph_safety.rs @@ -57,6 +57,7 @@ pub(super) fn graph_entity_safe(entity: GraphEntity) -> GraphEntity { GraphEntity { canonical_key: redact_graph_text(entity.canonical_key), display_label: redact_graph_text(entity.display_label), + source_kind: redact_graph_text(entity.source_kind), source_id: redact_graph_text(entity.source_id), ..entity } diff --git a/src/app/services/graph_support.rs b/src/app/services/graph_support.rs index 147942cad..530c9eb54 100644 --- a/src/app/services/graph_support.rs +++ b/src/app/services/graph_support.rs @@ -70,6 +70,7 @@ pub(super) fn graph_relationship_to_model( confidence: row.confidence, evidence_count: row.evidence_count, evidence_ids, + source_kinds: Vec::new(), first_seen_at: row.first_seen_at, last_seen_at: row.last_seen_at, } diff --git a/src/db.rs b/src/db.rs index 0a0435a6c..d36b40995 100644 --- a/src/db.rs +++ b/src/db.rs @@ -9,6 +9,7 @@ pub mod entity_resolution; pub(crate) mod error_signatures; pub mod graph; pub(crate) mod graph_confidence; +pub mod graph_discovery; pub mod graph_findings; pub mod graph_inventory; mod graph_resolver_projection; diff --git a/src/db/graph.rs b/src/db/graph.rs index 7429890a4..7060d388e 100644 --- a/src/db/graph.rs +++ b/src/db/graph.rs @@ -961,7 +961,7 @@ fn graph_evidence_for_relationships( Ok(rows) } -fn graph_entity_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result { +pub(crate) fn graph_entity_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result { Ok(GraphEntityRow { id: row.get(0)?, entity_type: row.get(1)?, @@ -975,7 +975,9 @@ fn graph_entity_from_row(row: &rusqlite::Row<'_>) -> rusqlite::Result) -> rusqlite::Result { +pub(crate) fn graph_relationship_from_row( + row: &rusqlite::Row<'_>, +) -> rusqlite::Result { Ok(GraphRelationshipRow { id: row.get(0)?, relationship_key: row.get(1)?, @@ -1374,6 +1376,7 @@ fn merge_graph_delta( )?; tx.execute("DROP TABLE IF EXISTS _graph_entity_idmap", [])?; tx.execute("DROP TABLE IF EXISTS _graph_rel_idmap", [])?; + prune_graph_change_events(&tx)?; tx.commit()?; Ok(GraphRebuildStats { @@ -2975,6 +2978,66 @@ fn swap_graph_projection( // caller's connection, so the lock is taken second by construction. let _guard = graph_staging_write_lock(conn); let tx = conn.transaction()?; + // Emit the logical diff before replacing the projection. The ordinary + // row triggers are paused inside this transaction so unchanged rows do + // not flood the incremental feed during a full rebuild. + tx.execute( + "UPDATE graph_change_capture SET enabled = 0 WHERE id = 1", + [], + )?; + tx.execute_batch( + "INSERT INTO graph_change_events (object_kind, operation, item_key) + SELECT 'relationship', 'delete', old.relationship_key + FROM graph_relationships old + LEFT JOIN _graph_relationships_staging next ON next.relationship_key = old.relationship_key + WHERE next.id IS NULL; + INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type) + SELECT 'entity', 'delete', old.entity_type || char(31) || old.canonical_key, old.entity_type + FROM graph_entities old + LEFT JOIN _graph_entities_staging next + ON next.entity_type = old.entity_type AND next.canonical_key = old.canonical_key + WHERE next.id IS NULL; + INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type, entity_id) + SELECT 'entity', 'upsert', next.entity_type || char(31) || next.canonical_key, + next.entity_type, next.id + FROM _graph_entities_staging next + LEFT JOIN graph_entities old + ON old.entity_type = next.entity_type AND old.canonical_key = next.canonical_key + WHERE old.id IS NULL OR old.id IS NOT next.id OR old.display_label IS NOT next.display_label + OR old.source_kind IS NOT next.source_kind OR old.source_id IS NOT next.source_id + OR old.trust_level IS NOT next.trust_level OR old.first_seen_at IS NOT next.first_seen_at + OR old.last_seen_at IS NOT next.last_seen_at; + INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) + SELECT 'relationship', 'upsert', next.relationship_key, next.id + FROM _graph_relationships_staging next + LEFT JOIN graph_relationships old ON old.relationship_key = next.relationship_key + WHERE old.id IS NULL OR old.id IS NOT next.id + OR old.src_entity_id IS NOT next.src_entity_id OR old.dst_entity_id IS NOT next.dst_entity_id + OR old.relationship_type IS NOT next.relationship_type OR old.reason_code IS NOT next.reason_code + OR old.trust_level IS NOT next.trust_level OR old.confidence IS NOT next.confidence + OR old.evidence_count IS NOT next.evidence_count OR old.first_seen_at IS NOT next.first_seen_at + OR old.last_seen_at IS NOT next.last_seen_at + OR EXISTS ( + SELECT 1 FROM _graph_evidence_staging staged_evidence + LEFT JOIN graph_relationship_evidence old_evidence + ON old_evidence.relationship_id = old.id + AND old_evidence.evidence_key = staged_evidence.evidence_key + WHERE staged_evidence.relationship_id = next.id + AND (old_evidence.id IS NULL + OR old_evidence.id IS NOT staged_evidence.id + OR old_evidence.source_kind IS NOT staged_evidence.source_kind + OR old_evidence.source_id IS NOT staged_evidence.source_id + OR old_evidence.observed_at IS NOT staged_evidence.observed_at + OR old_evidence.trust_level IS NOT staged_evidence.trust_level) + ) + OR EXISTS ( + SELECT 1 FROM graph_relationship_evidence old_evidence + LEFT JOIN _graph_evidence_staging staged_evidence + ON staged_evidence.relationship_id = next.id + AND staged_evidence.evidence_key = old_evidence.evidence_key + WHERE old_evidence.relationship_id = old.id AND staged_evidence.id IS NULL + );", + )?; tx.execute("DELETE FROM graph_relationship_evidence", [])?; tx.execute("DELETE FROM graph_relationships", [])?; tx.execute("DELETE FROM graph_entity_aliases", [])?; @@ -3046,6 +3109,11 @@ fn swap_graph_projection( chunk_count ], )?; + tx.execute( + "UPDATE graph_change_capture SET enabled = 1 WHERE id = 1", + [], + )?; + prune_graph_change_events(&tx)?; tx.commit()?; Ok(GraphRebuildStats { source_row_count, @@ -3058,6 +3126,16 @@ fn swap_graph_projection( }) } +pub(crate) fn prune_graph_change_events(conn: &rusqlite::Connection) -> Result<()> { + conn.execute( + "DELETE FROM graph_change_events + WHERE occurred_at < strftime('%Y-%m-%dT%H:%M:%fZ', 'now', '-7 days') + OR seq <= COALESCE((SELECT MAX(seq) - 200000 FROM graph_change_events), 0)", + [], + )?; + Ok(()) +} + fn graph_source_watermark(conn: &rusqlite::Connection) -> Result { let max_log_id: i64 = conn.query_row("SELECT COALESCE(MAX(id), 0) FROM logs", [], |r| r.get(0))?; diff --git a/src/db/graph_discovery.rs b/src/db/graph_discovery.rs new file mode 100644 index 000000000..77a236e7b --- /dev/null +++ b/src/db/graph_discovery.rs @@ -0,0 +1,306 @@ +//! Stable, bounded pages over the committed graph projection and its change journal. +//! Every page is read in one SQLite transaction so a concurrent rebuild cannot +//! combine a cursor from one projection with rows from another. + +use anyhow::Result; +use rusqlite::{OptionalExtension, params}; + +use super::DbPool; +use super::graph::{self, GraphEntityRow, GraphRelationshipRow}; + +#[derive(Debug)] +pub struct SnapshotPage { + pub rows: Vec, + pub anchor: i64, + pub has_more: bool, + pub changed: bool, + pub projection_status: String, + pub source_watermark: String, + pub last_completed_at: Option, + pub is_degraded: bool, +} + +#[derive(Debug)] +pub struct ChangeRow { + pub seq: i64, + pub object_kind: String, + pub operation: String, + pub item_key: String, + pub occurred_at: String, + pub entity: Option, + pub relationship: Option, + pub evidence_ids: Vec, + pub source_kinds: Vec, +} + +#[derive(Debug)] +pub struct RelationshipItem { + pub row: GraphRelationshipRow, + pub evidence_ids: Vec, + pub source_kinds: Vec, +} + +pub struct EntityFilters<'a> { + pub entity_type: &'a str, + pub source_kind: &'a str, + pub trust_level: &'a str, + pub query: &'a str, +} + +#[derive(Debug)] +pub struct ChangePage { + pub rows: Vec, + pub latest: i64, + pub has_more: bool, + pub expired: bool, + pub projection_status: String, + pub source_watermark: String, + pub last_completed_at: Option, + pub is_degraded: bool, +} + +fn sequence(conn: &rusqlite::Connection) -> Result { + Ok(conn.query_row( + "SELECT COALESCE((SELECT seq FROM sqlite_sequence WHERE name = 'graph_change_events'), 0)", + [], + |row| row.get(0), + )?) +} + +fn metadata(conn: &rusqlite::Connection) -> Result<(String, String, Option, bool)> { + Ok(conn.query_row( + "SELECT projection_status, source_watermark, last_completed_at, is_degraded + FROM graph_projection_meta WHERE id = 1", + [], + |row| { + Ok(( + row.get(0)?, + row.get(1)?, + row.get(2)?, + row.get::<_, i64>(3)? != 0, + )) + }, + )?) +} + +pub fn entities( + pool: &DbPool, + after: i64, + expected_anchor: Option, + limit: u32, + filters: EntityFilters<'_>, +) -> Result> { + let mut conn = pool.get()?; + let tx = conn.transaction()?; + let anchor = sequence(&tx)?; + let (projection_status, source_watermark, last_completed_at, is_degraded) = metadata(&tx)?; + if expected_anchor.is_some_and(|expected| expected != anchor) { + return Ok(SnapshotPage { + rows: vec![], + anchor, + has_more: false, + changed: true, + projection_status, + source_watermark, + last_completed_at, + is_degraded, + }); + } + let mut stmt = tx.prepare( + "SELECT id, entity_type, canonical_key, display_label, source_kind, source_id, + trust_level, first_seen_at, last_seen_at + FROM graph_entities + WHERE id > ?1 AND (?2 = '' OR entity_type = ?2) + AND (?3 = '' OR source_kind = ?3) AND (?4 = '' OR trust_level = ?4) + AND (?5 = '' OR instr(lower(display_label), lower(?5)) > 0 + OR instr(lower(canonical_key), lower(?5)) > 0) + ORDER BY id LIMIT ?6", + )?; + let mut rows = stmt + .query_map( + params![ + after, + filters.entity_type, + filters.source_kind, + filters.trust_level, + filters.query, + i64::from(limit) + 1 + ], + graph::graph_entity_from_row, + )? + .collect::>>()?; + let has_more = rows.len() > limit as usize; + rows.truncate(limit as usize); + Ok(SnapshotPage { + rows, + anchor, + has_more, + changed: false, + projection_status, + source_watermark, + last_completed_at, + is_degraded, + }) +} + +pub fn relationships( + pool: &DbPool, + after: i64, + expected_anchor: Option, + limit: u32, +) -> Result> { + let mut conn = pool.get()?; + let tx = conn.transaction()?; + let anchor = sequence(&tx)?; + let (projection_status, source_watermark, last_completed_at, is_degraded) = metadata(&tx)?; + if expected_anchor.is_some_and(|expected| expected != anchor) { + return Ok(SnapshotPage { + rows: vec![], + anchor, + has_more: false, + changed: true, + projection_status, + source_watermark, + last_completed_at, + is_degraded, + }); + } + let mut stmt = tx.prepare( + "SELECT id, relationship_key, src_entity_id, dst_entity_id, relationship_type, + reason_code, trust_level, confidence, evidence_count, first_seen_at, last_seen_at + FROM graph_relationships WHERE id > ?1 ORDER BY id LIMIT ?2", + )?; + let mut rows = stmt + .query_map( + params![after, i64::from(limit) + 1], + graph::graph_relationship_from_row, + )? + .collect::>>()?; + let has_more = rows.len() > limit as usize; + rows.truncate(limit as usize); + let rows = rows + .into_iter() + .map(|row| { + let (evidence_ids, source_kinds) = relationship_provenance(&tx, row.id)?; + Ok(RelationshipItem { + row, + evidence_ids, + source_kinds, + }) + }) + .collect::>>()?; + Ok(SnapshotPage { + rows, + anchor, + has_more, + changed: false, + projection_status, + source_watermark, + last_completed_at, + is_degraded, + }) +} + +pub fn changes(pool: &DbPool, after: i64, limit: u32) -> Result { + let mut conn = pool.get()?; + let tx = conn.transaction()?; + let latest = sequence(&tx)?; + let oldest: Option = + tx.query_row("SELECT MIN(seq) FROM graph_change_events", [], |row| { + row.get(0) + })?; + let (projection_status, source_watermark, last_completed_at, is_degraded) = metadata(&tx)?; + let expired = after < oldest.unwrap_or(latest + 1) - 1 || after > latest; + if expired { + return Ok(ChangePage { + rows: vec![], + latest, + has_more: false, + expired, + projection_status, + source_watermark, + last_completed_at, + is_degraded, + }); + } + let mut stmt = tx.prepare( + "SELECT seq, object_kind, operation, item_key, occurred_at + FROM graph_change_events WHERE seq > ?1 ORDER BY seq LIMIT ?2", + )?; + let mut rows = stmt + .query_map(params![after, i64::from(limit) + 1], |row| { + Ok(ChangeRow { + seq: row.get(0)?, + object_kind: row.get(1)?, + operation: row.get(2)?, + item_key: row.get(3)?, + occurred_at: row.get(4)?, + entity: None, + relationship: None, + evidence_ids: Vec::new(), + source_kinds: Vec::new(), + }) + })? + .collect::>>()?; + let has_more = rows.len() > limit as usize; + rows.truncate(limit as usize); + for event in &mut rows { + if event.operation != "upsert" { + continue; + } + if event.object_kind == "entity" { + if let Some((kind, key)) = event.item_key.split_once('\u{1f}') { + event.entity = tx.query_row( + "SELECT id, entity_type, canonical_key, display_label, source_kind, source_id, + trust_level, first_seen_at, last_seen_at + FROM graph_entities WHERE entity_type = ?1 AND canonical_key = ?2", + params![kind, key], graph::graph_entity_from_row, + ).optional()?; + } + } else { + event.relationship = tx.query_row( + "SELECT id, relationship_key, src_entity_id, dst_entity_id, relationship_type, + reason_code, trust_level, confidence, evidence_count, first_seen_at, last_seen_at + FROM graph_relationships WHERE relationship_key = ?1", + [&event.item_key], graph::graph_relationship_from_row, + ).optional()?; + if let Some(relationship) = &event.relationship { + (event.evidence_ids, event.source_kinds) = + relationship_provenance(&tx, relationship.id)?; + } + } + if event.entity.is_none() && event.relationship.is_none() { + event.operation = "delete".into(); + } + } + Ok(ChangePage { + rows, + latest, + has_more, + expired: false, + projection_status, + source_watermark, + last_completed_at, + is_degraded, + }) +} + +fn relationship_provenance( + conn: &rusqlite::Connection, + id: i64, +) -> Result<(Vec, Vec)> { + let mut stmt = conn.prepare( + "SELECT id, source_kind FROM graph_relationship_evidence + WHERE relationship_id = ?1 ORDER BY observed_at DESC, id DESC LIMIT 3", + )?; + let rows = stmt + .query_map([id], |row| { + Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)) + })? + .collect::>>()?; + let ids = rows.iter().map(|row| row.0).collect(); + let mut kinds: Vec = rows.into_iter().map(|row| row.1).collect(); + kinds.sort(); + kinds.dedup(); + Ok((ids, kinds)) +} diff --git a/src/db/graph_inventory.rs b/src/db/graph_inventory.rs index ef79bfd19..14364718a 100644 --- a/src/db/graph_inventory.rs +++ b/src/db/graph_inventory.rs @@ -134,6 +134,7 @@ fn apply_projection_plan( let stats = graph_counts(&tx)?; update_projection_meta(&tx, &stats)?; + graph::prune_graph_change_events(&tx)?; tx.commit().context("commit inventory graph projection")?; Ok(stats) } diff --git a/src/db/pool.rs b/src/db/pool.rs index cbb677ac7..072da9c5d 100644 --- a/src/db/pool.rs +++ b/src/db/pool.rs @@ -258,7 +258,7 @@ pub(crate) fn try_write_conn_for( } } -pub const KNOWN_SCHEMA_VERSION: i64 = 61; +pub const KNOWN_SCHEMA_VERSION: i64 = 62; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct SchemaVersionInfo { @@ -3632,6 +3632,93 @@ pub fn init_pool(config: &StorageConfig) -> Result { ); } + // Migration 62: bounded read-only graph discovery and durable change receipts. + // Triggers cover incremental and inventory writers; a full projection swap + // records its logical diff explicitly while capture is disabled in its tx. + if !migration_applied(&conn, 62)? { + conn.execute_batch( + "BEGIN IMMEDIATE; + CREATE TABLE graph_change_capture ( + id INTEGER PRIMARY KEY CHECK (id = 1), + enabled INTEGER NOT NULL CHECK (enabled IN (0, 1)) + ); + INSERT INTO graph_change_capture VALUES (1, 1); + CREATE TABLE graph_change_events ( + seq INTEGER PRIMARY KEY AUTOINCREMENT, + object_kind TEXT NOT NULL CHECK (object_kind IN ('entity', 'relationship')), + operation TEXT NOT NULL CHECK (operation IN ('upsert', 'delete')), + item_key TEXT NOT NULL, + entity_type TEXT, + entity_id INTEGER, + relationship_id INTEGER, + occurred_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) + ); + CREATE INDEX idx_graph_changes_time ON graph_change_events(occurred_at, seq); + CREATE TRIGGER graph_entity_insert_change AFTER INSERT ON graph_entities + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type, entity_id) + VALUES ('entity', 'upsert', new.entity_type || char(31) || new.canonical_key, new.entity_type, new.id); + END; + CREATE TRIGGER graph_entity_update_change AFTER UPDATE OF + display_label, source_kind, source_id, trust_level, first_seen_at, last_seen_at + ON graph_entities WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 + AND (old.display_label IS NOT new.display_label OR old.source_kind IS NOT new.source_kind + OR old.source_id IS NOT new.source_id OR old.trust_level IS NOT new.trust_level + OR old.first_seen_at IS NOT new.first_seen_at OR old.last_seen_at IS NOT new.last_seen_at) BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type, entity_id) + VALUES ('entity', 'upsert', new.entity_type || char(31) || new.canonical_key, new.entity_type, new.id); + END; + CREATE TRIGGER graph_entity_delete_change AFTER DELETE ON graph_entities + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type) + VALUES ('entity', 'delete', old.entity_type || char(31) || old.canonical_key, old.entity_type); + END; + CREATE TRIGGER graph_relationship_insert_change AFTER INSERT ON graph_relationships + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) + VALUES ('relationship', 'upsert', new.relationship_key, new.id); + END; + CREATE TRIGGER graph_relationship_update_change AFTER UPDATE OF + src_entity_id, dst_entity_id, relationship_type, reason_code, trust_level, + confidence, evidence_count, first_seen_at, last_seen_at + ON graph_relationships WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 + AND (old.src_entity_id IS NOT new.src_entity_id OR old.dst_entity_id IS NOT new.dst_entity_id + OR old.relationship_type IS NOT new.relationship_type OR old.reason_code IS NOT new.reason_code + OR old.trust_level IS NOT new.trust_level OR old.confidence IS NOT new.confidence + OR old.evidence_count IS NOT new.evidence_count OR old.first_seen_at IS NOT new.first_seen_at + OR old.last_seen_at IS NOT new.last_seen_at) BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) + VALUES ('relationship', 'upsert', new.relationship_key, new.id); + END; + CREATE TRIGGER graph_relationship_delete_change AFTER DELETE ON graph_relationships + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key) + VALUES ('relationship', 'delete', old.relationship_key); + END; + CREATE TRIGGER graph_evidence_insert_change AFTER INSERT ON graph_relationship_evidence + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) + SELECT 'relationship', 'upsert', relationship_key, id + FROM graph_relationships WHERE id = new.relationship_id; + END; + CREATE TRIGGER graph_evidence_update_change AFTER UPDATE ON graph_relationship_evidence + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) + SELECT 'relationship', 'upsert', relationship_key, id + FROM graph_relationships WHERE id = new.relationship_id; + END; + CREATE TRIGGER graph_evidence_delete_change AFTER DELETE ON graph_relationship_evidence + WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN + INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) + SELECT 'relationship', 'upsert', relationship_key, id + FROM graph_relationships WHERE id = old.relationship_id; + END; + INSERT INTO schema_migrations (version) VALUES (62); + COMMIT;", + )?; + tracing::info!("Migration 62: graph discovery change journal ready"); + } + if table_exists(&conn, "host_heartbeats")? && table_exists(&conn, "host_heartbeats_latest")? { let deleted_heartbeat_latest = conn.execute( "DELETE FROM host_heartbeats_latest diff --git a/src/db/pool_tests.rs b/src/db/pool_tests.rs index 9615f9270..d6f328576 100644 --- a/src/db/pool_tests.rs +++ b/src/db/pool_tests.rs @@ -10,9 +10,9 @@ use rusqlite::OptionalExtension; #[test] fn documented_schema_count_matches_known_version() { - assert_eq!(KNOWN_SCHEMA_VERSION, 61); - assert!(include_str!("../../README.md").contains("61 sequential schema migrations")); - assert!(include_str!("../../docs/architecture.md").contains("61 sequential migrations")); + assert_eq!(KNOWN_SCHEMA_VERSION, 62); + assert!(include_str!("../../README.md").contains("62 sequential schema migrations")); + assert!(include_str!("../../docs/architecture.md").contains("62 sequential migrations")); } #[test] diff --git a/src/mcp/rmcp_server.rs b/src/mcp/rmcp_server.rs index 05f9a64bf..d98ac026c 100644 --- a/src/mcp/rmcp_server.rs +++ b/src/mcp/rmcp_server.rs @@ -566,7 +566,10 @@ fn classify_tool_error(error: &anyhow::Error) -> ToolErrorClass { Some(ServiceError::NotFound(_)) | Some(ServiceError::RowNotFound) => { ToolErrorClass::NotFound } - Some(ServiceError::ConstraintViolation { .. }) => ToolErrorClass::Conflict, + Some(ServiceError::ConstraintViolation { .. }) | Some(ServiceError::Conflict(_)) => { + ToolErrorClass::Conflict + } + Some(ServiceError::Gone(_)) => ToolErrorClass::NotFound, Some(ServiceError::Internal(_)) | None => ToolErrorClass::Internal, } } diff --git a/src/surfaces/api.rs b/src/surfaces/api.rs index a0716ba16..19e9e9d67 100644 --- a/src/surfaces/api.rs +++ b/src/surfaces/api.rs @@ -271,6 +271,9 @@ pub(super) const API_SURFACE_SPECS: &[SurfaceSpec] = &[ Read ), api!("/api/graph/entity", Graph, Canonical, Read), + api!("/api/graph/entities", Graph, Canonical, Read), + api!("/api/graph/relationships", Graph, Canonical, Read), + api!("/api/graph/changes", Graph, Canonical, Read), api!("/api/graph/around", Graph, Canonical, Read), api!("/api/graph/explain", Graph, Canonical, Read), api!("/api/graph/evidence", Graph, Canonical, Read), From 41b2581bb154b3cf93ac9db67d7c66b8d22cc89e Mon Sep 17 00:00:00 2001 From: Jake Magar Date: Thu, 24 Sep 2026 11:22:07 -0400 Subject: [PATCH 2/4] docs(npm): sync graph discovery schema count --- packages/cortex-rmcp/README.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/cortex-rmcp/README.md b/packages/cortex-rmcp/README.md index c73f682b5..0a390468c 100644 --- a/packages/cortex-rmcp/README.md +++ b/packages/cortex-rmcp/README.md @@ -17,7 +17,7 @@ Cortex began as a syslog receiver. It now covers network logs, Docker, managed f | Area | What Cortex provides | | --- | --- | | Ingest | UDP/TCP syslog, OTLP/HTTP logs, metrics, and traces, Docker logs and events, managed file tails, host heartbeats, AI transcripts, shell history, agent command records, and fleet inventory | -| Storage | SQLite in WAL mode, FTS5 full-text search, bounded metadata, retention, storage budgets, maintenance jobs, checkpoints, and 61 sequential schema migrations | +| Storage | SQLite in WAL mode, FTS5 full-text search, bounded metadata, retention, storage budgets, maintenance jobs, checkpoints, and 62 sequential schema migrations | | Investigation | Search, filtering, context, timelines, patterns, anomaly comparison, cross-source correlation, recurring error signatures, deterministic incident bundles, and graph explanations | | Fleet intelligence | SSH and API inventory collectors, host state, service topology, container and route relationships, redacted evidence, and rebuildable graph projections | | AI operations | Claude, Codex, Gemini CLI, and Antigravity session indexing; skill, MCP, and hook event extraction where each provider exposes them; incident clustering; and guarded local LLM assessments | @@ -701,7 +701,7 @@ Cortex uses SQLite with: - Online backup support - Integrity checks, checkpoints, and vacuum workflows -The current schema history contains 61 sequential migrations. CI derives this denominator from `KNOWN_SCHEMA_VERSION` and the migration registry. Migration 60 introduced the transcript-only `ai_logs_fts` projection. Migration 61 rebuilds that derived index with provider scope and adds an AI-only `(timestamp, ai_tool)` index so time-bounded session searches can derive a safe rowid floor before FTS walks common terms. On large transcript histories these one-time derived-index rebuilds can hold the startup write transaction while they populate. Subsequent transcript inserts and deletes maintain `ai_logs_fts` incrementally, while global `logs_fts` remains the full-corpus search index. Server forwarding receipts have a seven-day replay horizon and are removed when their canonical evidence is deleted. The sender spool is intentionally shorter and bounded: an individual source retains at most 1,024 records or 1 MiB and evicts records older than one day; the aggregate spool retains at most 4,096 records or 4 MiB. Eviction removes the original payload and retains a pending gap marker so evidence loss remains visible. A retry after the server receipt horizon is a new ingestion attempt. +The current schema history contains 62 sequential migrations. CI derives this denominator from `KNOWN_SCHEMA_VERSION` and the migration registry. Migration 60 introduced the transcript-only `ai_logs_fts` projection. Migration 61 rebuilds that derived index with provider scope and adds an AI-only `(timestamp, ai_tool)` index so time-bounded session searches can derive a safe rowid floor before FTS walks common terms. On large transcript histories these one-time derived-index rebuilds can hold the startup write transaction while they populate. Subsequent transcript inserts and deletes maintain `ai_logs_fts` incrementally, while global `logs_fts` remains the full-corpus search index. Server forwarding receipts have a seven-day replay horizon and are removed when their canonical evidence is deleted. The sender spool is intentionally shorter and bounded: an individual source retains at most 1,024 records or 1 MiB and evicts records older than one day; the aggregate spool retains at most 4,096 records or 4 MiB. Eviction removes the original payload and retains a pending gap marker so evidence loss remains visible. A retry after the server receipt horizon is a new ingestion attempt. ### Authoritative and derived data From ced16185db828c26c55936c0163793cad520aeb1 Mon Sep 17 00:00:00 2001 From: Jake Magar Date: Thu, 24 Sep 2026 11:31:36 -0400 Subject: [PATCH 3/4] fix(db): support graph journal during schema repair --- src/cli/output/graph_tests.rs | 2 ++ src/db/pool.rs | 44 +++++++++++++++++++++++++---------- 2 files changed, 34 insertions(+), 12 deletions(-) diff --git a/src/cli/output/graph_tests.rs b/src/cli/output/graph_tests.rs index dece9fecf..baacf97fb 100644 --- a/src/cli/output/graph_tests.rs +++ b/src/cli/output/graph_tests.rs @@ -71,6 +71,7 @@ fn graph_relationship() -> cortex::app::GraphRelationship { confidence: 0.9, evidence_count: 1, evidence_ids: vec![9], + source_kinds: vec!["log".into()], first_seen_at: Some("2026-06-13T00:00:00Z".into()), last_seen_at: Some("2026-06-13T00:01:00Z".into()), } @@ -287,6 +288,7 @@ fn graph_evidence_json_output_accepts_safe_response() { confidence: 0.9, evidence_count: 1, evidence_ids: vec![9], + source_kinds: vec!["log".into()], first_seen_at: None, last_seen_at: None, }, diff --git a/src/db/pool.rs b/src/db/pool.rs index 072da9c5d..305c0ae45 100644 --- a/src/db/pool.rs +++ b/src/db/pool.rs @@ -2511,6 +2511,14 @@ pub fn init_pool(config: &StorageConfig) -> Result { conn.execute_batch(&format!( "BEGIN IMMEDIATE; + -- A repaired database can replay this table rebuild after the + -- graph journal was installed. Evidence triggers reference the + -- relationship table while it is replaced; recreate them after + -- all migrations complete. + DROP TRIGGER IF EXISTS graph_evidence_insert_change; + DROP TRIGGER IF EXISTS graph_evidence_update_change; + DROP TRIGGER IF EXISTS graph_evidence_delete_change; + CREATE TABLE graph_entities_new ( id INTEGER PRIMARY KEY AUTOINCREMENT, entity_type TEXT NOT NULL CHECK (entity_type IN ( @@ -3654,12 +3662,26 @@ pub fn init_pool(config: &StorageConfig) -> Result { occurred_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) ); CREATE INDEX idx_graph_changes_time ON graph_change_events(occurred_at, seq); - CREATE TRIGGER graph_entity_insert_change AFTER INSERT ON graph_entities + INSERT INTO schema_migrations (version) VALUES (62); + COMMIT;", + )?; + tracing::info!("Migration 62: graph discovery change journal ready"); + } + + // Historical repair tests can replay a sparse pre-graph schema. Install + // capture triggers only once all graph tables exist; a later pool open + // recreates triggers after any table-rebuild migration has dropped them. + if table_exists(&conn, "graph_entities")? + && table_exists(&conn, "graph_relationships")? + && table_exists(&conn, "graph_relationship_evidence")? + { + conn.execute_batch( + " CREATE TRIGGER IF NOT EXISTS graph_entity_insert_change AFTER INSERT ON graph_entities WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type, entity_id) VALUES ('entity', 'upsert', new.entity_type || char(31) || new.canonical_key, new.entity_type, new.id); END; - CREATE TRIGGER graph_entity_update_change AFTER UPDATE OF + CREATE TRIGGER IF NOT EXISTS graph_entity_update_change AFTER UPDATE OF display_label, source_kind, source_id, trust_level, first_seen_at, last_seen_at ON graph_entities WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 AND (old.display_label IS NOT new.display_label OR old.source_kind IS NOT new.source_kind @@ -3668,17 +3690,17 @@ pub fn init_pool(config: &StorageConfig) -> Result { INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type, entity_id) VALUES ('entity', 'upsert', new.entity_type || char(31) || new.canonical_key, new.entity_type, new.id); END; - CREATE TRIGGER graph_entity_delete_change AFTER DELETE ON graph_entities + CREATE TRIGGER IF NOT EXISTS graph_entity_delete_change AFTER DELETE ON graph_entities WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key, entity_type) VALUES ('entity', 'delete', old.entity_type || char(31) || old.canonical_key, old.entity_type); END; - CREATE TRIGGER graph_relationship_insert_change AFTER INSERT ON graph_relationships + CREATE TRIGGER IF NOT EXISTS graph_relationship_insert_change AFTER INSERT ON graph_relationships WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) VALUES ('relationship', 'upsert', new.relationship_key, new.id); END; - CREATE TRIGGER graph_relationship_update_change AFTER UPDATE OF + CREATE TRIGGER IF NOT EXISTS graph_relationship_update_change AFTER UPDATE OF src_entity_id, dst_entity_id, relationship_type, reason_code, trust_level, confidence, evidence_count, first_seen_at, last_seen_at ON graph_relationships WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 @@ -3690,33 +3712,31 @@ pub fn init_pool(config: &StorageConfig) -> Result { INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) VALUES ('relationship', 'upsert', new.relationship_key, new.id); END; - CREATE TRIGGER graph_relationship_delete_change AFTER DELETE ON graph_relationships + CREATE TRIGGER IF NOT EXISTS graph_relationship_delete_change AFTER DELETE ON graph_relationships WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key) VALUES ('relationship', 'delete', old.relationship_key); END; - CREATE TRIGGER graph_evidence_insert_change AFTER INSERT ON graph_relationship_evidence + CREATE TRIGGER IF NOT EXISTS graph_evidence_insert_change AFTER INSERT ON graph_relationship_evidence WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) SELECT 'relationship', 'upsert', relationship_key, id FROM graph_relationships WHERE id = new.relationship_id; END; - CREATE TRIGGER graph_evidence_update_change AFTER UPDATE ON graph_relationship_evidence + CREATE TRIGGER IF NOT EXISTS graph_evidence_update_change AFTER UPDATE ON graph_relationship_evidence WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) SELECT 'relationship', 'upsert', relationship_key, id FROM graph_relationships WHERE id = new.relationship_id; END; - CREATE TRIGGER graph_evidence_delete_change AFTER DELETE ON graph_relationship_evidence + CREATE TRIGGER IF NOT EXISTS graph_evidence_delete_change AFTER DELETE ON graph_relationship_evidence WHEN (SELECT enabled FROM graph_change_capture WHERE id = 1) = 1 BEGIN INSERT INTO graph_change_events (object_kind, operation, item_key, relationship_id) SELECT 'relationship', 'upsert', relationship_key, id FROM graph_relationships WHERE id = old.relationship_id; END; - INSERT INTO schema_migrations (version) VALUES (62); - COMMIT;", +" )?; - tracing::info!("Migration 62: graph discovery change journal ready"); } if table_exists(&conn, "host_heartbeats")? && table_exists(&conn, "host_heartbeats_latest")? { From c95897daffd417534b65da7aec3803ab9565bc70 Mon Sep 17 00:00:00 2001 From: Jake Magar Date: Thu, 24 Sep 2026 11:58:29 -0400 Subject: [PATCH 4/4] test(cortex): qualify graph inventory routes in live sweep --- tests/live/phases/surfaces/rest_sweep.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/tests/live/phases/surfaces/rest_sweep.py b/tests/live/phases/surfaces/rest_sweep.py index 0c6cb616f..423a9d213 100755 --- a/tests/live/phases/surfaces/rest_sweep.py +++ b/tests/live/phases/surfaces/rest_sweep.py @@ -46,6 +46,9 @@ def request(base: str, method: str, path: str, token: str | None, admin: str | N "/api/correlate": {"reference_time": "2026-08-27T00:00:00Z", "limit": "5"}, "/api/correlate-state": {"reference_time": "2026-08-27T00:00:00Z", "limit": "5"}, "/api/graph/entity": {"entity_type": "host", "key": "cortex-live"}, + "/api/graph/entities": {"limit": "5"}, + "/api/graph/relationships": {"limit": "5"}, + "/api/graph/changes": {"limit": "5"}, "/api/graph/around": {"entity_type": "host", "key": "cortex-live", "depth": "1", "limit": "5"}, "/api/graph/explain": {"entity_type": "host", "key": "cortex-live", "depth": "1", "max_chains": "5"}, "/api/graph/evidence": {"evidence_id": "1"}, @@ -167,6 +170,9 @@ def evidence(status: int, body: bytes, headers: dict[str, str]) -> dict: "GET /api/fleet-state": ("object", "hosts summary"), "GET /api/get": ("object", "log"), "GET /api/graph/around": ("object", "candidates entities evidence metadata next_queries relationships resolved_entity"), + "GET /api/graph/entities": ("object", "entities metadata next_cursor snapshot_cursor"), + "GET /api/graph/relationships": ("object", "metadata next_cursor relationships snapshot_cursor"), + "GET /api/graph/changes": ("object", "changes metadata next_cursor"), "GET /api/graph/entity": ("object", "candidates metadata resolved_entity"), "GET /api/graph/evidence": ("object", "dst_entity evidence metadata missing_source_reason relationship source_log_summary src_entity"), "GET /api/graph/explain": ("object", "candidates chains evidence metadata missing_evidence narrative next_queries open_questions resolved_entity"), @@ -407,6 +413,14 @@ def main() -> int: raise RuntimeError("run-owned session identity was not discoverable") from error QUERY["/api/sessions/rendered"] = {**session_query, "limit": "5"} QUERY["/api/streams/sessions"] = session_query + graph_status, graph_payload, _ = request(base, "GET", "/api/graph/entities?limit=1", read_token, None) + try: + graph_cursor = json.loads(graph_payload)["snapshot_cursor"] if graph_status == 200 else None + except (KeyError, TypeError, json.JSONDecodeError): + graph_cursor = None + if not isinstance(graph_cursor, str) or not graph_cursor: + raise RuntimeError(f"graph inventory change cursor was not available: status={graph_status}") + QUERY["/api/graph/changes"]["cursor"] = graph_cursor for graph_path in ("/api/graph/entity", "/api/graph/around", "/api/graph/explain", "/api/v1/graph/entity", "/api/v1/graph/around", "/api/v1/graph/explain"): QUERY[graph_path]["key"] = fixture_host