Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down Expand Up @@ -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

Expand Down
11 changes: 8 additions & 3 deletions docs/api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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).
Expand Down
4 changes: 2 additions & 2 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
4 changes: 2 additions & 2 deletions packages/cortex-rmcp/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down Expand Up @@ -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

Expand Down
37 changes: 36 additions & 1 deletion src/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -354,6 +355,9 @@ pub fn router(state: ApiState) -> anyhow::Result<Router> {
)
.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))
Expand Down Expand Up @@ -1758,6 +1762,27 @@ async fn graph_entity(
respond(state.service.graph_entity_lookup(q).await)
}

async fn graph_entities(
State(state): State<ApiState>,
Query(q): Query<GraphEntitiesRequest>,
) -> axum::response::Response {
respond(state.service.graph_entities(q).await)
}

async fn graph_relationships(
State(state): State<ApiState>,
Query(q): Query<GraphRelationshipsRequest>,
) -> axum::response::Response {
respond(state.service.graph_relationships(q).await)
}

async fn graph_changes(
State(state): State<ApiState>,
Query(q): Query<GraphChangesRequest>,
) -> axum::response::Response {
respond(state.service.graph_changes(q).await)
}

async fn graph_around(
State(state): State<ApiState>,
Query(q): Query<GraphAroundRequest>,
Expand Down Expand Up @@ -2371,6 +2396,16 @@ async fn ai_prune_checkpoints(
fn respond<T: serde::Serialize>(result: crate::app::ServiceResult<T>) -> 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()
}
Expand Down
Loading
Loading