diff --git a/Justfile b/Justfile index fc7d5dc5..869ca075 100644 --- a/Justfile +++ b/Justfile @@ -23,6 +23,17 @@ check: cargo check bash scripts/check-rust-module-size.sh --limit 500 +# AO-001: parse/compile/validate the Agent Observatory planning contracts +# (JSON, SQL, Rust, TypeScript) and fail on unresolved placeholders. +check-agent-observatory-contracts: + bash scripts/check-agent-observatory-contracts.sh + +# ENV-004: confirm the deprecated CORTEX_AGENT_AI_TRANSCRIPTS env var only +# appears in the approved allowlist locations. +validate-transcript-forward-env-rename: + bash scripts/validate-transcript-forward-env-rename.sh + bash scripts/test-validate-transcript-forward-env-rename.sh + lint: cargo clippy -- -D warnings @@ -42,9 +53,6 @@ coverage-html: test-doc: cargo test --doc -check-agent-observatory-contracts: - bash scripts/check-agent-observatory-contracts.sh - docker-build: docker build -f config/Dockerfile -t cortex . diff --git a/docs/plans/agent-observatory/proof/PROOF.md b/docs/plans/agent-observatory/proof/PROOF.md new file mode 100644 index 00000000..d92e499d --- /dev/null +++ b/docs/plans/agent-observatory/proof/PROOF.md @@ -0,0 +1,140 @@ +# Agent Observatory implementation proof + +## AO-001 Add a planning-contract verification script +commit/worktree SHA: 4b84b406 (task started) +RED: `just --justfile Justfile --working-directory . check-agent-observatory-contracts` +RED result: exit 1, `justfile does not contain recipe check-agent-observatory-contracts` +GREEN: `just --justfile Justfile --working-directory . check-agent-observatory-contracts` +GREEN result: exit 0; JSON contracts ok; SQL integrity ok; Rust contract tests 2 passed; TypeScript 5.9.3 ok; placeholder audit ok +REGRESSION: `bash -n scripts/check-agent-observatory-contracts.sh && git diff --check` +REGRESSION result: shell syntax valid and diff whitespace clean +FILES: `scripts/check-agent-observatory-contracts.sh`, `Justfile` +NOTES: TypeScript resolution is network-free and requires exact 5.9.3 from TSC, future web/node_modules, PATH, or npm's offline cache. + +## AO-002 Lock schema and projection version constants +commit/worktree SHA: 06b15c04 (task started) +RED: `cargo test --manifest-path Cargo.toml --locked agent_observatory::tests --lib` +RED result: E0432 unresolved imports for `AGENT_OBSERVATORY_SCHEMA_VERSION` and `AGENT_OBSERVATORY_PROJECTION_VERSION` +GREEN: `cargo --config 'build.rustc-wrapper=""' test --manifest-path Cargo.toml --locked agent_observatory::tests --lib` +GREEN result: 2 passed; target schema 47 and projection version 1 locked; runtime schema remains below target until migrations land +REGRESSION: `cargo test --manifest-path Cargo.toml --locked known_schema_version_matches_migration_head --lib && cargo fmt --all -- --check` +REGRESSION result: runtime schema-head test passed; rustfmt and diff checks clean +FILES: `src/agent_observatory.rs`, `src/agent_observatory_tests.rs`, `src/lib.rs` +NOTES: The planned target constant is intentionally distinct from `db::KNOWN_SCHEMA_VERSION`; runtime version advances only with implemented migrations. + +## AO-003 Implement migration 44 repository table +commit/worktree SHA: ef73d297 (task started) +RED: `cargo --config 'build.rustc-wrapper=""' test --manifest-path Cargo.toml --locked init_pool_creates_agent_observatory_repository_schema_scaffold --lib` +RED result: expected repository columns, got an empty list because `repositories` did not exist +GREEN: `cargo --config 'build.rustc-wrapper=""' test --manifest-path Cargo.toml --locked init_pool_creates_agent_observatory_repository_schema_scaffold --lib` +GREEN result: 1 passed; exact columns, indexes, uniqueness, reopen preservation, and no premature migration-44 marker verified +REGRESSION: `cargo test --manifest-path Cargo.toml --locked known_schema_version_matches_migration_head --lib && cargo fmt --all -- --check && git diff --check` +REGRESSION result: runtime schema remains 43; focused formatted test passed; diff clean +FILES: `src/db/pool.rs`, `src/db/pool_tests.rs` +NOTES: The repository DDL is an idempotent migration-44 scaffold. AO-007 will mark migration 44 only after worktrees, observations, and commits are complete. + +## AO-004 Add repository worktree table +commit/worktree SHA: 9a7d6c25 (task started) +RED: `env CARGO_TARGET_DIR=.cache/cargo cargo --config 'build.rustc-wrapper=""' test --manifest-path Cargo.toml --locked init_pool_creates_agent_observatory_worktree_schema_scaffold --lib` +RED result: expected worktree columns, got an empty list because `repository_worktrees` did not exist +GREEN: same focused command with the worktree DDL implemented +GREEN result: 1 passed; exact columns, branch/HEAD state, host/path uniqueness, repository cascade, and empty `PRAGMA foreign_key_check` verified +REGRESSION: pinned-target `known_schema_version_matches_migration_head`, `cargo fmt --all -- --check`, and `git diff --check` +REGRESSION result: runtime schema remains 43; formatting and diff checks clean +FILES: `src/db/pool.rs`, `src/db/pool_tests.rs` +NOTES: The worktree DDL remains part of the unmarked migration-44 scaffold; all Cargo proof commands now pin `CARGO_TARGET_DIR` to this worktree to avoid cross-worktree lock pollution. + +## AO-005 Add repository observations table +commit/worktree SHA: fd3d8342 (task started) +RED: pinned-target `init_pool_creates_agent_observatory_observation_schema_scaffold` +RED result: expected observation columns, got an empty list because `repository_observations` did not exist +GREEN: same focused command with observation DDL and indexes implemented +GREEN result: 1 passed; exact columns, unique observation keys, invalid JSON rejection, deterministic `(observed_at DESC, id DESC)` ordering, and named repository/worktree indexes verified +REGRESSION: pinned-target `known_schema_version_matches_migration_head`, `cargo fmt --all -- --check`, and `git diff --check` +REGRESSION result: runtime schema remains 43; formatting and diff checks clean +FILES: `src/db/pool.rs`, `src/db/pool_tests.rs` +NOTES: Query-plan proof requires `idx_repository_observations_repo_time`; migration 44 remains deliberately unmarked. + +## AO-006 Add exact Git commit table +commit/worktree SHA: 58c180f8 (task started) +RED: pinned-target init_pool_creates_agent_observatory_git_commit_schema_scaffold +RED result: expected exact-commit columns, got an empty list because git_commits did not exist +GREEN: same focused command with exact-commit DDL implemented +GREEN result: 1 passed; exact columns, per-repository SHA uniqueness, cross-repository SHA reuse, JSON rejection, reachability update, and metadata-only storage verified +REGRESSION: pinned-target known_schema_version_matches_migration_head, cargo fmt --all -- --check, and git diff --check +REGRESSION result: runtime schema remains 43; formatting and diff checks clean +FILES: src/db/pool.rs, src/db/pool_tests.rs +NOTES: No patch, diff, blob, or plaintext author-email column exists; migration 44 remains unmarked until AO-007. + +## AO-007 Complete migration 44 version bookkeeping +commit/worktree SHA: f50abffc (task started) +RED: pinned-target migration_44_applies_from_schema_43_and_is_idempotent +RED result: schema migration head remained 43 after the simulated schema-43 database reopened +GREEN: same focused command after wrapping all migration-44 DDL and the version marker in one BEGIN IMMEDIATE transaction +GREEN result: 1 passed; schema 43 upgraded to 44, all four tables created, legacy stream state preserved, foreign-key/integrity checks clean, and repeated reopen kept one marker +REGRESSION: pinned-target init_pool_creates_agent_observatory_ suite, known_schema_version_matches_migration_head, cargo fmt, and git diff --check +REGRESSION result: 4 topology tests passed; schema-head test passed at 44; formatting and diff checks clean +FILES: src/db/pool.rs, src/db/pool_tests.rs +NOTES: KNOWN_SCHEMA_VERSION now advances truthfully to 44 only after the complete topology migration commits atomically. + +## AO-008 Superseded: partial migration 45 agent_runs scaffold removed +commit/worktree SHA: 696d60c5 (task started) +RED: isolated pinned-target init_pool_creates_agent_observatory_run_schema_scaffold +RED result: expected run columns, got an empty list because agent_runs did not exist +GREEN: same focused command with the agent-runs table and four indexes implemented +GREEN result: 1 passed; lifecycle status constraints, host/tool/native-session identity, nullable primary worktree, required indexes, and active-run query-plan use verified +REGRESSION: pinned-target known_schema_version_matches_migration_head, cargo fmt, and git diff --check +REGRESSION result: runtime schema remains 44; formatting and diff checks clean +FILES: src/db/pool.rs, src/db/pool_tests.rs +NOTES: Adversarial review found that creating unversioned migration-45 tables exposed a partial schema. The scaffold was removed; migration 45 must land atomically in its complete implementation PR. + +## AO-009 Superseded: partial actors and run/worktree evidence scaffold removed +commit/worktree SHA: cb6f7f77 (task started) +RED: isolated pinned-target init_pool_creates_agent_observatory_actor_and_worktree_evidence_schema +RED result: expected actor columns, got an empty list because agent_run_actors and agent_run_worktrees did not exist +GREEN: same focused command with both tables and three indexes implemented +GREEN result: 1 passed; actor identity and JSON checks, confidence/trust constraints, evidence tuple dedupe, deterministic primary ordering, multiple-worktree history, and all indexes verified +REGRESSION: pinned-target known_schema_version_matches_migration_head, cargo fmt --all -- --check, git diff --check +REGRESSION result: runtime schema remains 44; formatting and diff clean +FILES: src/db/pool.rs, src/db/pool_tests.rs +NOTES: These unversioned tables were removed with AO-008. `init_pool_does_not_create_partial_agent_observatory_migration_45` now proves all three partial tables remain absent. + +## ENV-001 Add the new resolver and compatibility alias +commit/worktree SHA: 2e22dc19 (task started) +RED: pinned-target env_new_only_true_enables_forwarding +RED result: new CORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true was ignored and forwarding remained false +GREEN: pinned-target transcript_forward_env_ test filter +GREEN result: 4 passed; precedence matrix, warning codes, authoritative replacement, legacy-only compatibility, and local sessions-watch independence verified +REGRESSION: pinned-target heartbeat_agent::tests plus cargo fmt and git diff --check +REGRESSION result: 46 passed; formatting and diff clean +FILES: src/heartbeat_agent.rs, src/heartbeat_agent_tests.rs +NOTES: The deprecated name is centralized in AI_TRANSCRIPT_FORWARD_LEGACY_ENV; from_env emits exactly one warning selected by the pure resolver. + +## ENV-002 Switch all generated and deployed configuration to the new name +commit/worktree SHA: 06c64c09 (task started) +RED: setup generation, persisted env resolution, and Linux deployment fixtures expecting CORTEX_AGENT_AI_TRANSCRIPT_FORWARD only +RED result: all three failed because output preserved CORTEX_AGENT_AI_TRANSCRIPTS or omitted the replacement +GREEN: the same three focused tests after setup/deploy normalization +GREEN result: 3 passed; legacy-only persisted/process environment values are emitted under CORTEX_AGENT_AI_TRANSCRIPT_FORWARD and the legacy key is absent +REGRESSION: setup::heartbeat_agent::tests, agent_deploy::tests, cargo fmt, git diff --check, production occurrence audit +REGRESSION result: 12 setup tests and 32 deployment tests passed; the only non-test source occurrence of CORTEX_AGENT_AI_TRANSCRIPTS is the compatibility constant +FILES: src/setup/heartbeat_agent.rs, src/setup/heartbeat_agent_tests.rs, src/agent_deploy.rs, src/agent_deploy_tests.rs +NOTES: Replacement values are authoritative; legacy persisted values are normalized instead of copied verbatim, and generated files never contain both names. + +## PR-173 adversarial remediation +RED: review beads `syslog-mcp-2axwk`, `syslog-mcp-e5o51`, `syslog-mcp-ftiu1`, and `syslog-mcp-34jai` +RED result: startup created an unversioned migration-45 subset; doctor inspected only the final duplicate, replaced symlinks, and could overwrite a concurrently changed file; validation allowlisted whole files and ignored extensionless tracked text. +GREEN: focused database and doctor tests plus both transcript-forward validator scripts +GREEN result: migration-45 tables remain absent; ambiguous legacy/duplicate configuration fails closed for manual editing; conflicting duplicates and symlinks fail without mutation; tracked occurrence validation rejects both extensionless fixtures and invalid occurrences inside an otherwise approved file. +FILES: `src/db/pool.rs`, `src/db/pool_tests.rs`, `src/setup/doctor.rs`, `src/setup/doctor_tests.rs`, `scripts/validate-transcript-forward-env-rename.sh`, `scripts/test-validate-transcript-forward-env-rename.sh`, `Justfile` +NOTES: Contract, OpenAPI, schema, Rust/TypeScript type, architecture, research, specification, validator, and golden-fixture artifacts were resolved to the versions already validated and merged through PR #172. + +## PR-173 independent re-review remediation +RED: review beads `syslog-mcp-afx03` and `syslog-mcp-ya4a7` +RED result: reread-before-rename still allowed a noncooperative edit or symlink swap in the final compare/replace window; documentation context accepted the substring `red` inside unrelated words such as `configured`. +GREEN: fail-closed doctor migration plus injected post-read mutation tests; whole-word documentation context grammar plus an executable-assignment rejection and code-fence negative fixture. +GREEN result: no automatic rewrite path remains, so forced noncooperative edits and symlink swaps are preserved; executable legacy assignments are rejected even in allowlisted documentation, and unrelated substrings grant no exception. +REGRESSION: post-merge privacy scanner and adversarial fixtures, render-template hostile-input test, Agent Observatory contracts, transcript validator fixtures, doctor tests, workflow/Kache contracts, full clippy, and full nextest. +REGRESSION result: all hermetic gates passed; full nextest ran 2,754 tests with 2 skipped and no failures. The live deployment-host check requires a deployment-local `hosts.env`, which is intentionally absent from this worktree. +FILES: `src/setup/doctor.rs`, `src/setup/doctor_tests.rs`, `src/setup/doctor_transcript_forward_tests.rs`, `scripts/validate-transcript-forward-env-rename.sh`, `scripts/test-validate-transcript-forward-env-rename.sh` +NOTES: Current `origin/main` through PR #174 was merged without conflicts, preserving the privacy scanner/config/runbook changes and the validated PR #172 Agent Observatory artifacts. diff --git a/scripts/test-validate-transcript-forward-env-rename.sh b/scripts/test-validate-transcript-forward-env-rename.sh new file mode 100755 index 00000000..24afd373 --- /dev/null +++ b/scripts/test-validate-transcript-forward-env-rename.sh @@ -0,0 +1,59 @@ +#!/usr/bin/env bash +set -euo pipefail + +repo_root="$(git rev-parse --show-toplevel)" +validator="$repo_root/scripts/validate-transcript-forward-env-rename.sh" +tmp="$(mktemp -d)" +trap 'rm -rf "$tmp"' EXIT + +git -C "$tmp" init -q +mkdir -p "$tmp/scripts" "$tmp/src" "$tmp/docs/contracts" +cp "$validator" "$tmp/scripts/validate-transcript-forward-env-rename.sh" +cat > "$tmp/src/heartbeat_agent.rs" <<'EOF' +pub const AI_TRANSCRIPT_FORWARD_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD"; +pub const AI_TRANSCRIPT_FORWARD_LEGACY_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPTS"; +EOF +cat > "$tmp/docs/contracts/agent-observatory.md" <<'EOF' +Current: `CORTEX_AGENT_AI_TRANSCRIPT_FORWARD`; deprecated compatibility alias: `CORTEX_AGENT_AI_TRANSCRIPTS`. +EOF +git -C "$tmp" add . +(cd "$tmp" && bash scripts/validate-transcript-forward-env-rename.sh) + +cat > "$tmp/systemd-unit" <<'EOF' +Environment=CORTEX_AGENT_AI_TRANSCRIPTS=true +EOF +git -C "$tmp" add systemd-unit +if (cd "$tmp" && bash scripts/validate-transcript-forward-env-rename.sh >/dev/null 2>&1); then + echo "validator accepted a deprecated key in an extensionless systemd file" >&2 + exit 1 +fi + +rm "$tmp/systemd-unit" +git -C "$tmp" add -u +cat > "$tmp/src/heartbeat_agent.rs" <<'EOF' +pub const AI_TRANSCRIPT_FORWARD_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD"; +// CORTEX_AGENT_AI_TRANSCRIPTS is convenient shorthand. +EOF +git -C "$tmp" add src/heartbeat_agent.rs +if (cd "$tmp" && bash scripts/validate-transcript-forward-env-rename.sh >/dev/null 2>&1); then + echo "validator accepted an unapproved occurrence inside an allowlisted file" >&2 + exit 1 +fi + +cat > "$tmp/src/heartbeat_agent.rs" <<'EOF' +pub const AI_TRANSCRIPT_FORWARD_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD"; +pub const AI_TRANSCRIPT_FORWARD_LEGACY_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPTS"; +EOF +cat > "$tmp/docs/contracts/agent-observatory.md" <<'EOF' +This configured example must fail: +```ini +CORTEX_AGENT_AI_TRANSCRIPTS=true +``` +EOF +git -C "$tmp" add src/heartbeat_agent.rs docs/contracts/agent-observatory.md +if (cd "$tmp" && bash scripts/validate-transcript-forward-env-rename.sh >/dev/null 2>&1); then + echo "validator accepted executable legacy config in an allowlisted doc" >&2 + exit 1 +fi + +echo "transcript-forward rename validator negative fixtures passed" diff --git a/scripts/validate-transcript-forward-env-rename.sh b/scripts/validate-transcript-forward-env-rename.sh new file mode 100755 index 00000000..70197e30 --- /dev/null +++ b/scripts/validate-transcript-forward-env-rename.sh @@ -0,0 +1,109 @@ +#!/usr/bin/env bash +# Validate every tracked-text occurrence of the deprecated transcript-forward +# environment key. This intentionally uses `git grep` without extension filters +# so extensionless config, systemd units, and newly added file types are covered. +set -euo pipefail + +python3 - <<'PY' +from __future__ import annotations + +import re +import subprocess +import sys + +legacy = "CORTEX_AGENT_AI_TRANSCRIPTS" +new = "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD" +validator = "scripts/validate-transcript-forward-env-rename.sh" + +result = subprocess.run( + ["git", "grep", "-n", "-I", legacy, "--", ":!" + validator], + check=False, + stdout=subprocess.PIPE, + text=True, +) +if result.returncode not in (0, 1): + raise SystemExit(result.returncode) + +doc_paths = { + "docs/contracts/agent-observatory.md", + "docs/plans/2026-07-31-agent-observatory-implementation.md", + "docs/plans/agent-observatory/01a-transcript-forward-env-rename.md", + "docs/plans/agent-observatory/06-production-hardening-and-docs.md", + "docs/plans/agent-observatory/proof/PROOF.md", + "docs/research/2026-07-31-agent-observatory.md", + "docs/specs/agent-observatory.md", +} + +def allowed(path: str, line: str) -> bool: + if path == "Justfile": + return line.lstrip().startswith("# ENV-004:") + if path == "src/heartbeat_agent.rs": + return re.fullmatch( + r'pub const AI_TRANSCRIPT_FORWARD_LEGACY_ENV: &str = "' + legacy + r'";', + line.strip(), + ) is not None + if path in { + "src/agent_deploy_tests.rs", + "src/heartbeat_agent_tests.rs", + "src/setup/doctor_tests.rs", + "src/setup/doctor_transcript_forward_tests.rs", + "src/setup/heartbeat_agent_tests.rs", + }: + # Test occurrences must be string fixtures or assertions, never an env! + # assignment in executable workflow/config syntax. + stripped = line.strip() + return ('"' + legacy) in stripped or (legacy + '=') in stripped + if path == "scripts/test-validate-transcript-forward-env-rename.sh": + # This harness deliberately injects forbidden occurrences into an + # isolated repository. Only its fixture/assertion lines may name the + # legacy key; executable configuration in this repository stays banned. + stripped = line.strip() + return ( + ('"' + legacy) in stripped + or legacy + "=" in stripped + or "deprecated compatibility alias" in stripped + or "convenient shorthand" in stripped + ) + if path in doc_paths: + lowered = line.lower() + if re.search(rf"(?:^|[^A-Za-z0-9_]){re.escape(legacy)}\s*=", line): + return False + return re.search( + r"\b(?:deprecated|deprecation|legacy|compatibility|red|regression|removal|rename|generated)\b", + lowered, + ) is not None or "operationally misleading" in lowered + return False + +violations: list[str] = [] +occurrences = 0 +for record in result.stdout.splitlines(): + path, number, line = record.split(":", 2) + occurrences += line.count(legacy) + if not allowed(path, line): + violations.append(f"{path}:{number}:{line}") + +if violations: + print(f"error: unapproved {legacy} occurrence(s):", file=sys.stderr) + print("\n".join(violations), file=sys.stderr) + raise SystemExit(1) + +required = { + "src/heartbeat_agent.rs": new, + "docs/contracts/agent-observatory.md": new, +} +for path, token in required.items(): + text = open(path, encoding="utf-8").read() + if token not in text: + raise SystemExit(f"error: {path} must document/use {token}") + +env_example = subprocess.run( + ["git", "ls-files", "--error-unmatch", ".env.example"], + check=False, + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, +).returncode == 0 +if env_example and legacy in open(".env.example", encoding="utf-8").read(): + raise SystemExit(f"error: .env.example contains deprecated {legacy}") + +print(f"validated {occurrences} approved tracked-text occurrence(s) of {legacy}") +PY diff --git a/src/agent_deploy.rs b/src/agent_deploy.rs index 0e2fa552..8207ec9f 100644 --- a/src/agent_deploy.rs +++ b/src/agent_deploy.rs @@ -542,7 +542,22 @@ fn resolve_linux_agent_env( if let Some(syslog_file) = prev_get("CORTEX_AGENT_SYSLOG_FILE") { out.push(("CORTEX_AGENT_SYSLOG_FILE".to_string(), syslog_file)); } + let transcript_forward = prev_get(crate::heartbeat_agent::AI_TRANSCRIPT_FORWARD_ENV) + .or_else(|| prev_get(crate::heartbeat_agent::AI_TRANSCRIPT_FORWARD_LEGACY_ENV)); + if let Some(value) = transcript_forward { + out.push(( + crate::heartbeat_agent::AI_TRANSCRIPT_FORWARD_ENV.to_string(), + value, + )); + } for key in crate::heartbeat_agent::OPTIONAL_ENV_KEYS { + if matches!( + *key, + crate::heartbeat_agent::AI_TRANSCRIPT_FORWARD_ENV + | crate::heartbeat_agent::AI_TRANSCRIPT_FORWARD_LEGACY_ENV + ) { + continue; + } if let Some(value) = prev_get(key) { out.push(((*key).to_string(), value)); } diff --git a/src/agent_deploy_tests.rs b/src/agent_deploy_tests.rs index 139fd6ec..ec664152 100644 --- a/src/agent_deploy_tests.rs +++ b/src/agent_deploy_tests.rs @@ -325,7 +325,8 @@ printf 'ssh %s\n' "$*" >> "$CORTEX_TEST_AGENT_DEPLOY_LOG" assert!(log.contains("CORTEX_AGENT_JOURNALD='true'")); assert!(log.contains("CORTEX_SYSLOG_TARGET='old-syslog.example:1514'")); assert!(log.contains("CORTEX_AGENT_FILE_TAILS='/var/log/app.log:app'")); - assert!(log.contains("CORTEX_AGENT_AI_TRANSCRIPTS='true'")); + assert!(log.contains("CORTEX_AGENT_AI_TRANSCRIPT_FORWARD='true'")); + assert!(!log.contains("CORTEX_AGENT_AI_TRANSCRIPTS='")); assert!(log.contains("CORTEX_AGENT_COMMAND_FORWARD='true'")); assert!(log.contains("CORTEX_AGENT_SHELL_HISTORY_FORWARD='true'")); assert!(log.contains("CORTEX_AGENT_AUTO_UPDATE='false'")); @@ -661,7 +662,11 @@ fn resolve_linux_agent_env_preserves_auth_and_flags_without_defaults() { env_get(&env, "CORTEX_AGENT_FILE_TAILS"), Some("/var/log/app.log:app") ); - assert_eq!(env_get(&env, "CORTEX_AGENT_AI_TRANSCRIPTS"), Some("true")); + assert_eq!( + env_get(&env, "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD"), + Some("true") + ); + assert_eq!(env_get(&env, "CORTEX_AGENT_AI_TRANSCRIPTS"), None); assert_eq!(env_get(&env, "CORTEX_AGENT_COMMAND_FORWARD"), Some("true")); assert_eq!( env_get(&env, "CORTEX_AGENT_SHELL_HISTORY_FORWARD"), diff --git a/src/agent_observatory.rs b/src/agent_observatory.rs new file mode 100644 index 00000000..57fccc72 --- /dev/null +++ b/src/agent_observatory.rs @@ -0,0 +1,14 @@ +//! Agent Observatory domain constants and coordination entry points. + +/// Schema version reached after all planned Agent Observatory migrations. +/// +/// This is intentionally separate from `db::KNOWN_SCHEMA_VERSION` until +/// migrations 44 through 47 are implemented and verified. +pub const AGENT_OBSERVATORY_SCHEMA_VERSION: i64 = 47; + +/// Version of the durable Agent Observatory projection contract. +pub const AGENT_OBSERVATORY_PROJECTION_VERSION: u32 = 1; + +#[cfg(test)] +#[path = "agent_observatory_tests.rs"] +mod tests; diff --git a/src/agent_observatory_tests.rs b/src/agent_observatory_tests.rs new file mode 100644 index 00000000..143ebc9f --- /dev/null +++ b/src/agent_observatory_tests.rs @@ -0,0 +1,18 @@ +use super::{AGENT_OBSERVATORY_PROJECTION_VERSION, AGENT_OBSERVATORY_SCHEMA_VERSION}; + +#[test] +fn planned_schema_and_projection_versions_are_locked() { + assert_eq!(AGENT_OBSERVATORY_SCHEMA_VERSION, 47); + assert_eq!(AGENT_OBSERVATORY_PROJECTION_VERSION, 1); +} + +// Compile-time (not runtime) check: the runtime schema must never claim to have +// already applied Agent Observatory migrations that don't exist yet. Both +// constants are `const i64`, so this ordering is provable at compile time — +// asserting it in `const _` catches drift as soon as this `#[cfg(test)]` +// module is compiled (e.g. `cargo test`, `cargo build --tests`), rather than +// only when the `#[test]` above actually executes. NOTE: this module is +// `#[cfg(test)]`-gated, so a plain `cargo build`/`cargo build --release` +// (which never compiles the test target) does not evaluate this assertion — +// it is a test-compile-time guarantee, not a production-build one. +const _: () = assert!(crate::db::KNOWN_SCHEMA_VERSION < AGENT_OBSERVATORY_SCHEMA_VERSION); diff --git a/src/db/pool.rs b/src/db/pool.rs index c3088ab4..bbefa7c1 100644 --- a/src/db/pool.rs +++ b/src/db/pool.rs @@ -2,7 +2,7 @@ //! core. //! //! Owns the full schema: the `logs` table + FTS5 index, AI/graph/heartbeat -//! projections, and the **40 sequential migrations** tracked by +//! projections, and the **44 sequential migrations** tracked by //! `KNOWN_SCHEMA_VERSION`. Migrations run at startup; heavy ones log //! `Migration N: starting ...` lines, and the one-time //! `auto_vacuum=INCREMENTAL` conversion VACUUM is logged loudly (it can take @@ -39,7 +39,7 @@ pub fn write_lock() -> parking_lot::ReentrantMutexGuard<'static, ()> { WRITE_LOCK.lock() } -pub const KNOWN_SCHEMA_VERSION: i64 = 43; +pub const KNOWN_SCHEMA_VERSION: i64 = 44; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct SchemaVersionInfo { @@ -2509,6 +2509,120 @@ pub fn init_pool(config: &StorageConfig) -> Result { tracing::info!("Migration 43: stream_last_seen rollup for stream-silence alerting"); } + // Migration 44: Agent Observatory repository, worktree, observation, + // and exact-commit topology. The DDL and version marker share one + // transaction so startup never reports a partially applied migration. + if !migration_applied(&conn, 44)? { + conn.execute_batch( + "BEGIN IMMEDIATE; + CREATE TABLE IF NOT EXISTS repositories ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + repository_key TEXT NOT NULL UNIQUE, + hostname TEXT NOT NULL, + common_git_dir TEXT NOT NULL, + primary_path TEXT NOT NULL, + display_name TEXT NOT NULL, + remote_url_hash TEXT, + first_seen_at TEXT NOT NULL, + last_seen_at TEXT NOT NULL, + removed_at TEXT, + metadata_json TEXT NOT NULL DEFAULT '{}' CHECK (json_valid(metadata_json)), + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + UNIQUE(hostname, common_git_dir) + ); + CREATE INDEX IF NOT EXISTS idx_repositories_host_seen + ON repositories(hostname, last_seen_at DESC); + CREATE INDEX IF NOT EXISTS idx_repositories_display + ON repositories(display_name COLLATE NOCASE); + + CREATE TABLE IF NOT EXISTS repository_worktrees ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + worktree_key TEXT NOT NULL UNIQUE, + repository_id INTEGER NOT NULL REFERENCES repositories(id) ON DELETE CASCADE, + hostname TEXT NOT NULL, + path TEXT NOT NULL, + git_dir TEXT NOT NULL, + branch_ref TEXT, + branch_name TEXT, + head_sha TEXT, + upstream_ref TEXT, + detached INTEGER NOT NULL DEFAULT 0 CHECK (detached IN (0, 1)), + bare INTEGER NOT NULL DEFAULT 0 CHECK (bare IN (0, 1)), + locked INTEGER NOT NULL DEFAULT 0 CHECK (locked IN (0, 1)), + lock_reason TEXT, + prunable INTEGER NOT NULL DEFAULT 0 CHECK (prunable IN (0, 1)), + prune_reason TEXT, + dirty INTEGER NOT NULL DEFAULT 0 CHECK (dirty IN (0, 1)), + staged_count INTEGER NOT NULL DEFAULT 0 CHECK (staged_count >= 0), + unstaged_count INTEGER NOT NULL DEFAULT 0 CHECK (unstaged_count >= 0), + untracked_count INTEGER NOT NULL DEFAULT 0 CHECK (untracked_count >= 0), + ahead INTEGER CHECK (ahead IS NULL OR ahead >= 0), + behind INTEGER CHECK (behind IS NULL OR behind >= 0), + status_hash TEXT, + first_seen_at TEXT NOT NULL, + last_seen_at TEXT NOT NULL, + removed_at TEXT, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + updated_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + UNIQUE(hostname, path) + ); + CREATE INDEX IF NOT EXISTS idx_worktrees_repo_active + ON repository_worktrees(repository_id, removed_at, last_seen_at DESC); + CREATE INDEX IF NOT EXISTS idx_worktrees_branch + ON repository_worktrees(branch_name, last_seen_at DESC); + CREATE INDEX IF NOT EXISTS idx_worktrees_head + ON repository_worktrees(repository_id, head_sha); + + CREATE TABLE IF NOT EXISTS repository_observations ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + observation_key TEXT NOT NULL UNIQUE, + repository_id INTEGER NOT NULL REFERENCES repositories(id) ON DELETE CASCADE, + worktree_id INTEGER REFERENCES repository_worktrees(id) ON DELETE CASCADE, + observed_at TEXT NOT NULL, + observation_kind TEXT NOT NULL CHECK (observation_kind IN ( + 'discovered', 'status', 'head', 'branch', 'worktree_added', + 'worktree_removed', 'overflow_reconcile', 'periodic_reconcile', 'error' + )), + old_head_sha TEXT, + new_head_sha TEXT, + summary TEXT NOT NULL DEFAULT '', + payload_json TEXT NOT NULL DEFAULT '{}' CHECK (json_valid(payload_json)), + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')) + ); + CREATE INDEX IF NOT EXISTS idx_repository_observations_worktree_time + ON repository_observations(worktree_id, observed_at DESC, id DESC); + CREATE INDEX IF NOT EXISTS idx_repository_observations_repo_time + ON repository_observations(repository_id, observed_at DESC, id DESC); + + CREATE TABLE IF NOT EXISTS git_commits ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + repository_id INTEGER NOT NULL REFERENCES repositories(id) ON DELETE CASCADE, + sha TEXT NOT NULL, + parent_shas_json TEXT NOT NULL DEFAULT '[]' CHECK (json_valid(parent_shas_json)), + author_name TEXT, + author_email_hash TEXT, + authored_at TEXT, + committed_at TEXT, + subject TEXT NOT NULL DEFAULT '', + changed_files INTEGER CHECK (changed_files IS NULL OR changed_files >= 0), + insertions INTEGER CHECK (insertions IS NULL OR insertions >= 0), + deletions INTEGER CHECK (deletions IS NULL OR deletions >= 0), + changed_paths_json TEXT NOT NULL DEFAULT '[]' CHECK (json_valid(changed_paths_json)), + first_observed_at TEXT NOT NULL, + last_observed_at TEXT NOT NULL, + reachable INTEGER NOT NULL DEFAULT 1 CHECK (reachable IN (0, 1)), + metadata_json TEXT NOT NULL DEFAULT '{}' CHECK (json_valid(metadata_json)), + UNIQUE(repository_id, sha) + ); + CREATE INDEX IF NOT EXISTS idx_git_commits_repo_time + ON git_commits(repository_id, committed_at DESC, id DESC); + INSERT OR IGNORE INTO schema_migrations (version) VALUES (44); + COMMIT;", + )?; + tracing::info!("Migration 44: Agent Observatory repository topology"); + } + 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 3b2b55c2..1b42729f 100644 --- a/src/db/pool_tests.rs +++ b/src/db/pool_tests.rs @@ -5,6 +5,7 @@ use crate::db::{ TRUST_LEVELS, insert_logs_batch, is_known_entity_type, is_known_evidence_source_kind, is_known_reason_code, is_known_relationship_type, is_known_trust_level, }; +use rusqlite::OptionalExtension; fn test_storage_config(db_path: std::path::PathBuf) -> StorageConfig { StorageConfig::for_test(db_path) @@ -367,6 +368,693 @@ fn known_schema_version_matches_migration_head() { assert_eq!(info.known_version, KNOWN_SCHEMA_VERSION); } +#[test] +fn init_pool_creates_agent_observatory_repository_schema_scaffold() { + let dir = tempfile::tempdir().unwrap(); + let config = test_storage_config(dir.path().join("observatory-repositories.db")); + + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + + let columns: Vec = conn + .prepare("PRAGMA table_info(repositories)") + .unwrap() + .query_map([], |row| row.get(1)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!( + columns, + vec![ + "id", + "repository_key", + "hostname", + "common_git_dir", + "primary_path", + "display_name", + "remote_url_hash", + "first_seen_at", + "last_seen_at", + "removed_at", + "metadata_json", + "created_at", + "updated_at", + ] + ); + + let indexes: Vec = conn + .prepare( + "SELECT name FROM sqlite_master + WHERE type = 'index' AND tbl_name = 'repositories' + ORDER BY name", + ) + .unwrap() + .query_map([], |row| row.get(0)) + .unwrap() + .collect::>() + .unwrap(); + assert!( + indexes + .iter() + .any(|name| name == "idx_repositories_display") + ); + assert!( + indexes + .iter() + .any(|name| name == "idx_repositories_host_seen") + ); + + conn.execute( + "INSERT INTO repositories + (repository_key, hostname, common_git_dir, primary_path, display_name, + first_seen_at, last_seen_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6)", + rusqlite::params![ + "v1|6:devhost|20:/workspace/cortex/.git", + "devhost", + "/workspace/cortex/.git", + "/workspace/cortex", + "cortex", + "2026-07-31T23:00:00.000Z", + ], + ) + .unwrap(); + assert!( + conn.execute( + "INSERT INTO repositories + (repository_key, hostname, common_git_dir, primary_path, display_name, + first_seen_at, last_seen_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6)", + rusqlite::params![ + "v1|6:devhost|20:/workspace/cortex/.git", + "other-host", + "/workspace/other/.git", + "/workspace/other", + "other", + "2026-07-31T23:00:00.000Z", + ], + ) + .is_err(), + "repository_key must be globally unique" + ); + assert!( + conn.execute( + "INSERT INTO repositories + (repository_key, hostname, common_git_dir, primary_path, display_name, + first_seen_at, last_seen_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6)", + rusqlite::params![ + "different-key", + "devhost", + "/workspace/cortex/.git", + "/workspace/cortex-copy", + "cortex-copy", + "2026-07-31T23:00:00.000Z", + ], + ) + .is_err(), + "hostname/common_git_dir must identify one repository" + ); + drop(conn); + drop(pool); + + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + let row_count: i64 = conn + .query_row("SELECT COUNT(*) FROM repositories", [], |row| row.get(0)) + .unwrap(); + assert_eq!(row_count, 1, "reopening must preserve repository rows"); + let migration_44_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM schema_migrations WHERE version = 44", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!( + migration_44_count, 1, + "completed migration 44 must remain marked exactly once" + ); +} + +#[test] +fn init_pool_creates_agent_observatory_worktree_schema_scaffold() { + let dir = tempfile::tempdir().unwrap(); + let config = test_storage_config(dir.path().join("observatory-worktrees.db")); + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + + conn.execute( + "INSERT INTO repositories + (repository_key, hostname, common_git_dir, primary_path, display_name, + first_seen_at, last_seen_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6)", + rusqlite::params![ + "repo-key", + "devhost", + "/workspace/cortex/.git", + "/workspace/cortex", + "cortex", + "2026-08-01T01:00:00.000Z", + ], + ) + .unwrap(); + let repository_id = conn.last_insert_rowid(); + + let columns: Vec = conn + .prepare("PRAGMA table_info(repository_worktrees)") + .unwrap() + .query_map([], |row| row.get(1)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!( + columns, + vec![ + "id", + "worktree_key", + "repository_id", + "hostname", + "path", + "git_dir", + "branch_ref", + "branch_name", + "head_sha", + "upstream_ref", + "detached", + "bare", + "locked", + "lock_reason", + "prunable", + "prune_reason", + "dirty", + "staged_count", + "unstaged_count", + "untracked_count", + "ahead", + "behind", + "status_hash", + "first_seen_at", + "last_seen_at", + "removed_at", + "created_at", + "updated_at", + ] + ); + + conn.execute( + "INSERT INTO repository_worktrees + (worktree_key, repository_id, hostname, path, git_dir, branch_ref, + branch_name, head_sha, upstream_ref, dirty, staged_count, + unstaged_count, untracked_count, ahead, behind, first_seen_at, last_seen_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, 1, 2, 3, 4, 5, 6, ?10, ?10)", + rusqlite::params![ + "worktree-key", + repository_id, + "devhost", + "/workspace/cortex", + "/workspace/cortex/.git", + "refs/heads/feat/agent-observatory", + "feat/agent-observatory", + "0123456789012345678901234567890123456789", + "refs/remotes/origin/feat/agent-observatory", + "2026-08-01T01:00:00.000Z", + ], + ) + .unwrap(); + + assert!( + conn.execute( + "INSERT INTO repository_worktrees + (worktree_key, repository_id, hostname, path, git_dir, first_seen_at, last_seen_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6)", + rusqlite::params![ + "different-key", + repository_id, + "devhost", + "/workspace/cortex", + "/workspace/cortex/.git/worktrees/duplicate", + "2026-08-01T01:00:00.000Z", + ], + ) + .is_err(), + "hostname/path must identify one worktree" + ); + + let state: (String, String, i64, i64, i64, i64, i64) = conn + .query_row( + "SELECT branch_name, head_sha, dirty, staged_count, unstaged_count, + untracked_count, ahead + FROM repository_worktrees WHERE worktree_key = 'worktree-key'", + [], + |row| { + Ok(( + row.get(0)?, + row.get(1)?, + row.get(2)?, + row.get(3)?, + row.get(4)?, + row.get(5)?, + row.get(6)?, + )) + }, + ) + .unwrap(); + assert_eq!( + state, + ( + "feat/agent-observatory".to_string(), + "0123456789012345678901234567890123456789".to_string(), + 1, + 2, + 3, + 4, + 5, + ) + ); + + conn.execute("DELETE FROM repositories WHERE id = ?1", [repository_id]) + .unwrap(); + let remaining: i64 = conn + .query_row("SELECT COUNT(*) FROM repository_worktrees", [], |row| { + row.get(0) + }) + .unwrap(); + assert_eq!( + remaining, 0, + "repository deletion must cascade to worktrees" + ); + + let foreign_key_violation: Option = conn + .query_row("PRAGMA foreign_key_check", [], |row| row.get(0)) + .optional() + .unwrap(); + assert_eq!(foreign_key_violation, None); +} + +#[test] +fn init_pool_creates_agent_observatory_observation_schema_scaffold() { + let dir = tempfile::tempdir().unwrap(); + let config = test_storage_config(dir.path().join("observatory-observations.db")); + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + + conn.execute( + "INSERT INTO repositories + (repository_key, hostname, common_git_dir, primary_path, display_name, + first_seen_at, last_seen_at) + VALUES ('repo-key', 'devhost', '/workspace/cortex/.git', + '/workspace/cortex', 'cortex', ?1, ?1)", + ["2026-08-01T01:00:00.000Z"], + ) + .unwrap(); + let repository_id = conn.last_insert_rowid(); + conn.execute( + "INSERT INTO repository_worktrees + (worktree_key, repository_id, hostname, path, git_dir, first_seen_at, last_seen_at) + VALUES ('worktree-key', ?1, 'devhost', '/workspace/cortex', + '/workspace/cortex/.git', ?2, ?2)", + rusqlite::params![repository_id, "2026-08-01T01:00:00.000Z"], + ) + .unwrap(); + let worktree_id = conn.last_insert_rowid(); + + let columns: Vec = conn + .prepare("PRAGMA table_info(repository_observations)") + .unwrap() + .query_map([], |row| row.get(1)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!( + columns, + vec![ + "id", + "observation_key", + "repository_id", + "worktree_id", + "observed_at", + "observation_kind", + "old_head_sha", + "new_head_sha", + "summary", + "payload_json", + "created_at", + ] + ); + + let insert = |key: &str, observed_at: &str, kind: &str| { + conn.execute( + "INSERT INTO repository_observations + (observation_key, repository_id, worktree_id, observed_at, + observation_kind, summary, payload_json) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, '{}')", + rusqlite::params![key, repository_id, worktree_id, observed_at, kind, key], + ) + }; + insert("obs-1", "2026-08-01T01:00:00.000Z", "discovered").unwrap(); + insert("obs-2", "2026-08-01T01:00:01.000Z", "status").unwrap(); + insert("obs-3", "2026-08-01T01:00:01.000Z", "head").unwrap(); + + assert!( + insert("obs-1", "2026-08-01T01:00:02.000Z", "status").is_err(), + "observation_key must be globally unique" + ); + assert!( + conn.execute( + "INSERT INTO repository_observations + (observation_key, repository_id, observed_at, observation_kind, payload_json) + VALUES ('bad-json', ?1, ?2, 'error', '{')", + rusqlite::params![repository_id, "2026-08-01T01:00:03.000Z"], + ) + .is_err(), + "payload_json must be valid JSON" + ); + + let ordered: Vec = conn + .prepare( + "SELECT observation_key FROM repository_observations + WHERE repository_id = ?1 + ORDER BY observed_at DESC, id DESC", + ) + .unwrap() + .query_map([repository_id], |row| row.get(0)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!(ordered, vec!["obs-3", "obs-2", "obs-1"]); + + let repo_plan: Vec = conn + .prepare( + "EXPLAIN QUERY PLAN + SELECT id FROM repository_observations + WHERE repository_id = ?1 + ORDER BY observed_at DESC, id DESC LIMIT 10", + ) + .unwrap() + .query_map([repository_id], |row| row.get(3)) + .unwrap() + .collect::>() + .unwrap(); + assert!( + repo_plan + .iter() + .any(|detail| detail.contains("idx_repository_observations_repo_time")), + "repository timeline query must use its chronological index: {repo_plan:?}" + ); + + let indexes: Vec = conn + .prepare( + "SELECT name FROM sqlite_master + WHERE type = 'index' AND tbl_name = 'repository_observations' + ORDER BY name", + ) + .unwrap() + .query_map([], |row| row.get(0)) + .unwrap() + .collect::>() + .unwrap(); + assert!( + indexes + .iter() + .any(|name| name == "idx_repository_observations_repo_time") + ); + assert!( + indexes + .iter() + .any(|name| name == "idx_repository_observations_worktree_time") + ); +} + +#[test] +fn init_pool_creates_agent_observatory_git_commit_schema_scaffold() { + let dir = tempfile::tempdir().unwrap(); + let config = test_storage_config(dir.path().join("observatory-commits.db")); + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + + for (key, common_dir, path, name) in [ + ("repo-1", "/workspace/one/.git", "/workspace/one", "one"), + ("repo-2", "/workspace/two/.git", "/workspace/two", "two"), + ] { + conn.execute( + "INSERT INTO repositories + (repository_key, hostname, common_git_dir, primary_path, display_name, + first_seen_at, last_seen_at) + VALUES (?1, 'devhost', ?2, ?3, ?4, ?5, ?5)", + rusqlite::params![key, common_dir, path, name, "2026-08-01T01:00:00.000Z"], + ) + .unwrap(); + } + let repo_one: i64 = conn + .query_row( + "SELECT id FROM repositories WHERE repository_key = 'repo-1'", + [], + |row| row.get(0), + ) + .unwrap(); + let repo_two: i64 = conn + .query_row( + "SELECT id FROM repositories WHERE repository_key = 'repo-2'", + [], + |row| row.get(0), + ) + .unwrap(); + + let columns: Vec = conn + .prepare("PRAGMA table_info(git_commits)") + .unwrap() + .query_map([], |row| row.get(1)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!( + columns, + vec![ + "id", + "repository_id", + "sha", + "parent_shas_json", + "author_name", + "author_email_hash", + "authored_at", + "committed_at", + "subject", + "changed_files", + "insertions", + "deletions", + "changed_paths_json", + "first_observed_at", + "last_observed_at", + "reachable", + "metadata_json", + ] + ); + assert!( + !columns + .iter() + .any(|name| matches!(name.as_str(), "diff" | "patch" | "blob" | "author_email")) + ); + + let sha = "0123456789012345678901234567890123456789"; + let insert_commit = |repository_id: i64| { + conn.execute( + "INSERT INTO git_commits + (repository_id, sha, parent_shas_json, author_name, author_email_hash, + authored_at, committed_at, subject, changed_files, insertions, + deletions, changed_paths_json, first_observed_at, last_observed_at, + metadata_json) + VALUES (?1, ?2, '[]', 'Cortex Test', 'sha256:test', ?3, ?3, + 'test commit', 2, 10, 3, '[\"src/lib.rs\"]', ?3, ?3, '{}')", + rusqlite::params![repository_id, sha, "2026-08-01T01:00:00.000Z"], + ) + }; + insert_commit(repo_one).unwrap(); + assert!( + insert_commit(repo_one).is_err(), + "same SHA must dedupe within a repository" + ); + insert_commit(repo_two).unwrap(); + + assert!( + conn.execute( + "INSERT INTO git_commits + (repository_id, sha, parent_shas_json, changed_paths_json, + first_observed_at, last_observed_at) + VALUES (?1, 'bad-json', '{', '[]', ?2, ?2)", + rusqlite::params![repo_one, "2026-08-01T01:00:00.000Z"], + ) + .is_err(), + "commit JSON columns must reject invalid JSON" + ); + + conn.execute( + "UPDATE git_commits + SET reachable = 0, last_observed_at = ?1 + WHERE repository_id = ?2 AND sha = ?3", + rusqlite::params!["2026-08-01T02:00:00.000Z", repo_one, sha], + ) + .unwrap(); + let state: (i64, String, String) = conn + .query_row( + "SELECT reachable, subject, last_observed_at FROM git_commits + WHERE repository_id = ?1 AND sha = ?2", + rusqlite::params![repo_one, sha], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + ) + .unwrap(); + assert_eq!( + state, + ( + 0, + "test commit".to_string(), + "2026-08-01T02:00:00.000Z".to_string() + ) + ); + + let repo_one_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM git_commits WHERE repository_id = ?1", + [repo_one], + |row| row.get(0), + ) + .unwrap(); + let repo_two_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM git_commits WHERE repository_id = ?1", + [repo_two], + |row| row.get(0), + ) + .unwrap(); + assert_eq!((repo_one_count, repo_two_count), (1, 1)); + + let index_exists: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master + WHERE type = 'index' AND name = 'idx_git_commits_repo_time'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(index_exists, 1); +} + +#[test] +fn migration_44_applies_from_schema_43_and_is_idempotent() { + let dir = tempfile::tempdir().unwrap(); + let db_path = dir.path().join("observatory-migration-44.db"); + let config = test_storage_config(db_path.clone()); + + { + let pool = init_pool(&config).unwrap(); + drop(pool); + } + + { + let conn = rusqlite::Connection::open(&db_path).unwrap(); + conn.execute_batch( + "PRAGMA foreign_keys = OFF; + DROP TABLE IF EXISTS git_commits; + DROP TABLE IF EXISTS repository_observations; + DROP TABLE IF EXISTS repository_worktrees; + DROP TABLE IF EXISTS repositories; + DELETE FROM schema_migrations WHERE version = 44; + INSERT OR REPLACE INTO stream_last_seen + (hostname, source_kind, last_seen_at) + VALUES ('legacy-host', 'syslog-tcp', '2026-08-01T01:00:00.000Z'); + PRAGMA foreign_keys = ON;", + ) + .unwrap(); + } + + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + let max_version: i64 = conn + .query_row("SELECT MAX(version) FROM schema_migrations", [], |row| { + row.get(0) + }) + .unwrap(); + assert_eq!(max_version, 44); + let marker_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM schema_migrations WHERE version = 44", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(marker_count, 1); + + for table in [ + "repositories", + "repository_worktrees", + "repository_observations", + "git_commits", + ] { + let exists: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?1", + [table], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(exists, 1, "migration 44 must create {table}"); + } + + let legacy_rows: i64 = conn + .query_row( + "SELECT COUNT(*) FROM stream_last_seen + WHERE hostname = 'legacy-host' AND source_kind = 'syslog-tcp'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(legacy_rows, 1, "migration must preserve schema-43 data"); + let foreign_key_violation: Option = conn + .query_row("PRAGMA foreign_key_check", [], |row| row.get(0)) + .optional() + .unwrap(); + assert_eq!(foreign_key_violation, None); + let integrity: String = conn + .query_row("PRAGMA integrity_check", [], |row| row.get(0)) + .unwrap(); + assert_eq!(integrity, "ok"); + drop(conn); + drop(pool); + + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + let marker_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM schema_migrations WHERE version = 44", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(marker_count, 1, "reopening must not duplicate migration 44"); +} + +#[test] +fn init_pool_does_not_create_partial_agent_observatory_migration_45() { + let dir = tempfile::tempdir().unwrap(); + let config = test_storage_config(dir.path().join("observatory-runs.db")); + let pool = init_pool(&config).unwrap(); + let conn = pool.get().unwrap(); + + for table in ["agent_runs", "agent_run_actors", "agent_run_worktrees"] { + let exists: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?1", + [table], + |row| row.get(0), + ) + .unwrap(); + assert_eq!( + exists, 0, + "partial migration-45 table {table} must not exist" + ); + } +} + #[test] fn graph_schema_enforces_vocabulary_and_dedup_keys() { let dir = tempfile::tempdir().unwrap(); diff --git a/src/heartbeat_agent.rs b/src/heartbeat_agent.rs index bb40988c..93620aa4 100644 --- a/src/heartbeat_agent.rs +++ b/src/heartbeat_agent.rs @@ -23,9 +23,12 @@ pub const DEFAULT_COLLECTION_DEADLINE_MS: u64 = 5_000; pub const DEFAULT_RETRY_BUFFER_LIMIT: usize = 32; pub const DEFAULT_TARGET: &str = "http://127.0.0.1:3100"; pub const DEFAULT_DOCKER_URL: &str = "unix:///var/run/docker.sock"; +pub const AI_TRANSCRIPT_FORWARD_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD"; +pub const AI_TRANSCRIPT_FORWARD_LEGACY_ENV: &str = "CORTEX_AGENT_AI_TRANSCRIPTS"; pub const OPTIONAL_ENV_KEYS: &[&str] = &[ "CORTEX_AGENT_FILE_TAILS", - "CORTEX_AGENT_AI_TRANSCRIPTS", + AI_TRANSCRIPT_FORWARD_ENV, + AI_TRANSCRIPT_FORWARD_LEGACY_ENV, "CORTEX_AGENT_AI_TRANSCRIPT_CHECKPOINT", "CORTEX_AGENT_COMMAND_FORWARD", "CORTEX_AGENT_COMMAND_SPOOL", @@ -34,6 +37,63 @@ pub const OPTIONAL_ENV_KEYS: &[&str] = &[ "CORTEX_AGENT_AUTO_UPDATE", ]; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum TranscriptForwardEnvWarning { + LegacyAlias, + ConflictingValues, +} + +impl TranscriptForwardEnvWarning { + pub(crate) const fn code(self) -> &'static str { + match self { + Self::LegacyAlias => "agent_ai_transcript_forward_legacy_alias", + Self::ConflictingValues => "agent_ai_transcript_forward_conflict", + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) struct TranscriptForwardEnvResolution { + pub(crate) enabled: bool, + pub(crate) warning: Option, +} + +fn parse_forward_flag(value: &str) -> bool { + value.eq_ignore_ascii_case("true") || value == "1" +} + +pub(crate) fn resolve_ai_transcript_forward_env( + current: Option<&str>, + legacy: Option<&str>, +) -> TranscriptForwardEnvResolution { + match (current, legacy) { + (Some(current), Some(legacy)) => { + let enabled = parse_forward_flag(current); + let warning = if enabled == parse_forward_flag(legacy) { + TranscriptForwardEnvWarning::LegacyAlias + } else { + TranscriptForwardEnvWarning::ConflictingValues + }; + TranscriptForwardEnvResolution { + enabled, + warning: Some(warning), + } + } + (Some(current), None) => TranscriptForwardEnvResolution { + enabled: parse_forward_flag(current), + warning: None, + }, + (None, Some(legacy)) => TranscriptForwardEnvResolution { + enabled: parse_forward_flag(legacy), + warning: Some(TranscriptForwardEnvWarning::LegacyAlias), + }, + (None, None) => TranscriptForwardEnvResolution { + enabled: false, + warning: None, + }, + } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct HeartbeatAgentConfig { pub target: Option, @@ -102,10 +162,29 @@ impl HeartbeatAgentConfig { .map(|spec| crate::agent::syslog_file::parse_file_tails(&spec)) .unwrap_or_default(); let syslog_target = std::env::var("CORTEX_SYSLOG_TARGET").ok(); - let ai_transcripts = std::env::var("CORTEX_AGENT_AI_TRANSCRIPTS") - .ok() - .map(|v| v.eq_ignore_ascii_case("true") || v == "1") - .unwrap_or(false); + let current_transcript_forward = std::env::var(AI_TRANSCRIPT_FORWARD_ENV).ok(); + let legacy_transcript_forward = std::env::var(AI_TRANSCRIPT_FORWARD_LEGACY_ENV).ok(); + let transcript_forward = resolve_ai_transcript_forward_env( + current_transcript_forward.as_deref(), + legacy_transcript_forward.as_deref(), + ); + if let Some(warning) = transcript_forward.warning { + match warning { + TranscriptForwardEnvWarning::LegacyAlias => tracing::warn!( + warning_code = warning.code(), + legacy_env = AI_TRANSCRIPT_FORWARD_LEGACY_ENV, + replacement_env = AI_TRANSCRIPT_FORWARD_ENV, + "deprecated remote transcript-forwarding environment alias is configured" + ), + TranscriptForwardEnvWarning::ConflictingValues => tracing::warn!( + warning_code = warning.code(), + legacy_env = AI_TRANSCRIPT_FORWARD_LEGACY_ENV, + replacement_env = AI_TRANSCRIPT_FORWARD_ENV, + "conflicting remote transcript-forwarding environment values; replacement wins" + ), + } + } + let ai_transcripts = transcript_forward.enabled; let ai_transcript_checkpoint_path = std::env::var("CORTEX_AGENT_AI_TRANSCRIPT_CHECKPOINT") .ok() .map(PathBuf::from) diff --git a/src/heartbeat_agent_tests.rs b/src/heartbeat_agent_tests.rs index bbc239c1..f01acbd9 100644 --- a/src/heartbeat_agent_tests.rs +++ b/src/heartbeat_agent_tests.rs @@ -645,3 +645,72 @@ fn test_payload(sequence: i64) -> HeartbeatPayload { containers: None, } } + +// ENV-001: remote transcript-forwarding environment compatibility. + +#[test] +fn transcript_forward_env_resolution_precedence_is_stable() { + use TranscriptForwardEnvWarning::{ConflictingValues, LegacyAlias}; + + let cases = [ + (None, None, false, None), + (Some("true"), None, true, None), + (Some("1"), None, true, None), + (Some("false"), None, false, None), + (None, Some("true"), true, Some(LegacyAlias)), + (None, Some("false"), false, Some(LegacyAlias)), + (Some("true"), Some("true"), true, Some(LegacyAlias)), + (Some("false"), Some("false"), false, Some(LegacyAlias)), + (Some("true"), Some("false"), true, Some(ConflictingValues)), + (Some("false"), Some("true"), false, Some(ConflictingValues)), + ]; + + for (current, legacy, enabled, warning) in cases { + assert_eq!( + resolve_ai_transcript_forward_env(current, legacy), + TranscriptForwardEnvResolution { enabled, warning }, + "current={current:?}, legacy={legacy:?}" + ); + } + assert_eq!( + LegacyAlias.code(), + "agent_ai_transcript_forward_legacy_alias" + ); + assert_eq!( + ConflictingValues.code(), + "agent_ai_transcript_forward_conflict" + ); +} + +#[test] +#[serial] +fn transcript_forward_env_current_value_is_authoritative_in_config() { + let _new = EnvGuard::set(AI_TRANSCRIPT_FORWARD_ENV, "false"); + let _legacy = EnvGuard::set(AI_TRANSCRIPT_FORWARD_LEGACY_ENV, "true"); + let config = HeartbeatAgentConfig::from_env(PathBuf::from("/tmp/host-id")); + assert!(!config.ai_transcripts); +} + +#[test] +#[serial] +fn transcript_forward_env_legacy_value_is_honored_in_config() { + let _new = EnvGuard::unset(AI_TRANSCRIPT_FORWARD_ENV); + let _legacy = EnvGuard::set(AI_TRANSCRIPT_FORWARD_LEGACY_ENV, "true"); + let config = HeartbeatAgentConfig::from_env(PathBuf::from("/tmp/host-id")); + assert!(config.ai_transcripts); +} + +#[test] +fn transcript_forward_env_does_not_gate_local_sessions_watch_service() { + let unit = crate::setup::ai_watch_service_unit( + Path::new("/home/test/.local/bin/cortex"), + Path::new("/home/test/.config/cortex/sessions-watch.env"), + Path::new("/home/test/.cortex/data/cortex.db"), + Path::new("/home/test/.local/state/cortex"), + Path::new("/home/test"), + ); + + assert!(unit.contains("sessions watch --no-initial-scan --json")); + assert!(!unit.contains(AI_TRANSCRIPT_FORWARD_ENV)); + assert!(!unit.contains(AI_TRANSCRIPT_FORWARD_LEGACY_ENV)); +} diff --git a/src/lib.rs b/src/lib.rs index c6a6f975..67d8d749 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -8,6 +8,7 @@ pub mod agent; pub mod agent_command_ingest; pub mod agent_deploy; +pub mod agent_observatory; pub(crate) mod ai_project; pub mod ai_transcript_ingest; pub mod ai_watch; diff --git a/src/setup/doctor.rs b/src/setup/doctor.rs index dc037ed5..03f6045a 100644 --- a/src/setup/doctor.rs +++ b/src/setup/doctor.rs @@ -57,6 +57,7 @@ pub async fn run_setup_doctor(fix: bool, yes: bool) -> io::Result { .phases, ); phases.push(runtime_current_phase(&repo_path)); + phases.push(check_transcript_forward_env_migration(&env_path, fix, yes)); phases.push(stale_agent_command_units_phase(fix, yes).await); let elapsed_ms = started.elapsed().as_millis(); @@ -289,6 +290,116 @@ pub(crate) async fn stale_agent_command_units_phase(fix: bool, yes: bool) -> Set } } +/// ENV-003: Check for deprecated transcript forwarding environment variable +/// configurations and optionally migrate unambiguous files. This is called +/// from `run_setup_doctor` and integrates with the existing fix authorization +/// pattern (both `--fix` AND `--yes` required for mutation). +pub(crate) fn check_transcript_forward_env_migration( + env_path: &Path, + fix: bool, + yes: bool, +) -> SetupPhase { + check_transcript_forward_env_migration_with_hook(env_path, fix, yes, || {}) +} + +fn check_transcript_forward_env_migration_with_hook( + env_path: &Path, + fix: bool, + yes: bool, + after_read: F, +) -> SetupPhase +where + F: FnOnce(), +{ + use crate::heartbeat_agent::{AI_TRANSCRIPT_FORWARD_ENV, AI_TRANSCRIPT_FORWARD_LEGACY_ENV}; + + let timer = PhaseTimer::start("transcript-forward-env-migration"); + + // Replacing a symlink would detach it from its target. Refuse before any + // read or write so an operator can migrate the real file deliberately. + if std::fs::symlink_metadata(env_path).is_ok_and(|metadata| metadata.file_type().is_symlink()) { + return timer.finish( + SetupStatus::Error, + ".env is a symbolic link; migrate its target explicitly", + ); + } + + if !env_path.exists() { + return timer.finish(SetupStatus::Ok, ".env file not found; skipped"); + } + + let content = match std::fs::read_to_string(env_path) { + Ok(content) => content, + Err(error) => { + return timer.finish( + SetupStatus::Error, + format!("failed to read .env file: {error}"), + ); + } + }; + after_read(); + + let mut assignments = Vec::new(); + + for (line_num, line) in content.lines().enumerate() { + let trimmed = line.trim(); + for (key, legacy) in [ + (AI_TRANSCRIPT_FORWARD_ENV, false), + (AI_TRANSCRIPT_FORWARD_LEGACY_ENV, true), + ] { + if let Some(value) = trimmed + .strip_prefix(key) + .and_then(|rest| rest.strip_prefix('=')) + { + assignments.push((line_num, value.to_string(), legacy)); + break; + } + } + } + + if assignments.is_empty() { + return timer.finish(SetupStatus::Ok, "no transcript forwarding configuration"); + } + let value = &assignments[0].1; + if assignments + .iter() + .any(|(_, candidate, _)| candidate != value) + { + return timer.finish( + SetupStatus::Error, + "conflicting duplicate transcript forwarding values; manual resolution required", + ); + } + let has_legacy = assignments.iter().any(|(_, _, legacy)| *legacy); + if !has_legacy && assignments.len() == 1 { + return timer.finish( + SetupStatus::Ok, + "transcript forwarding uses current variable name", + ); + } + if !fix || !yes { + let authorization = if fix { "--yes" } else { "--fix --yes" }; + return timer.finish( + SetupStatus::Warn, + format!( + "deprecated {AI_TRANSCRIPT_FORWARD_LEGACY_ENV} or duplicate transcript forwarding configuration; run with {authorization} to normalize to {AI_TRANSCRIPT_FORWARD_ENV}" + ), + ); + } + + timer.finish( + SetupStatus::Error, + format!( + "automatic rewrite is disabled to prevent concurrent edit or symlink-swap data loss; edit {} manually, replace {AI_TRANSCRIPT_FORWARD_LEGACY_ENV} with {AI_TRANSCRIPT_FORWARD_ENV}, and collapse equal duplicates", + env_path.display() + ), + ) +} + #[cfg(test)] #[path = "doctor_tests.rs"] mod tests; + +#[cfg(test)] +#[path = "doctor_transcript_forward_tests.rs"] +mod transcript_forward_tests; diff --git a/src/setup/doctor_tests.rs b/src/setup/doctor_tests.rs index 58fbc460..d3741463 100644 --- a/src/setup/doctor_tests.rs +++ b/src/setup/doctor_tests.rs @@ -258,3 +258,235 @@ fn stale_agent_command_fix_requires_yes() { "yes alone without fix must not disable" ); } + +// ENV-003: transcript forwarding env migration tests + +#[test] +fn transcript_forward_migration_detects_legacy_only() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n").unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, false, false); + + assert!(matches!(result.status, SetupStatus::Warn)); + assert!(result.detail.contains("deprecated")); + assert!(result.detail.contains("CORTEX_AGENT_AI_TRANSCRIPTS")); + assert!(result.detail.contains("CORTEX_AGENT_AI_TRANSCRIPT_FORWARD")); +} + +#[test] +fn transcript_forward_migration_detects_both_equal() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write( + &env_path, + "CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true\n", + ) + .unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, false, false); + + assert!(matches!(result.status, SetupStatus::Warn)); + assert!(result.detail.contains("duplicate")); +} + +#[test] +fn transcript_forward_migration_detects_conflicting_values() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write( + &env_path, + "CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPT_FORWARD=false\n", + ) + .unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, false, false); + + assert!(matches!(result.status, SetupStatus::Error)); + assert!(result.detail.contains("conflicting")); +} + +#[test] +fn transcript_forward_migration_ok_when_new_only() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true\n").unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, false, false); + + assert!(matches!(result.status, SetupStatus::Ok)); +} + +#[test] +fn transcript_forward_migration_ok_when_neither_set() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_API_TOKEN=secret\n").unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, false, false); + + assert!(matches!(result.status, SetupStatus::Ok)); +} + +#[test] +fn transcript_forward_migration_ok_when_file_missing() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + + let result = check_transcript_forward_env_migration(&env_path, false, false); + + assert!(matches!(result.status, SetupStatus::Ok)); +} + +#[test] +fn transcript_forward_migration_fix_requires_yes() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n").unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, true, false); + + // Should still warn without --yes + assert!(matches!(result.status, SetupStatus::Warn)); + assert!(result.detail.contains("--yes")); +} + +#[cfg(unix)] +#[test] +fn transcript_forward_migration_fix_fails_closed_for_legacy_only() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n").unwrap(); + let mut perms = std::fs::metadata(&env_path).unwrap().permissions(); + perms.set_mode(0o600); + std::fs::set_permissions(&env_path, perms).unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, true, true); + + assert!(matches!(result.status, SetupStatus::Error)); + assert!(result.detail.contains("automatic rewrite is disabled")); + + let content = std::fs::read_to_string(&env_path).unwrap(); + assert_eq!(content, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n"); + + // Verify permissions remain private + let new_perms = std::fs::metadata(&env_path).unwrap().permissions(); + assert_eq!(new_perms.mode() & 0o777, 0o600); +} + +#[cfg(unix)] +#[test] +fn transcript_forward_migration_fix_preserves_comment_and_assignment() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + // A comment mentioning the legacy key as a substring, plus the real + // legacy assignment. A whole-file str::replace would rewrite both; + // the line-anchored migration must touch only the actual assignment. + std::fs::write( + &env_path, + "# migrated from CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPTS=true\n", + ) + .unwrap(); + let mut perms = std::fs::metadata(&env_path).unwrap().permissions(); + perms.set_mode(0o600); + std::fs::set_permissions(&env_path, perms).unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, true, true); + + assert!(matches!(result.status, SetupStatus::Error)); + + let content = std::fs::read_to_string(&env_path).unwrap(); + assert_eq!( + content, + "# migrated from CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPTS=true\n" + ); + + // Verify permissions remain private + let new_perms = std::fs::metadata(&env_path).unwrap().permissions(); + assert_eq!(new_perms.mode() & 0o777, 0o600); +} + +#[cfg(unix)] +#[test] +fn transcript_forward_migration_fix_preserves_equal_assignments() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write( + &env_path, + "CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true\n", + ) + .unwrap(); + let mut perms = std::fs::metadata(&env_path).unwrap().permissions(); + perms.set_mode(0o600); + std::fs::set_permissions(&env_path, perms).unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, true, true); + + assert!(matches!(result.status, SetupStatus::Error)); + + let content = std::fs::read_to_string(&env_path).unwrap(); + assert_eq!( + content, + "CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true\n" + ); + + // Verify permissions remain private + let new_perms = std::fs::metadata(&env_path).unwrap().permissions(); + assert_eq!(new_perms.mode() & 0o777, 0o600); +} + +#[cfg(unix)] +#[test] +fn transcript_forward_migration_fix_errors_on_conflict() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + let original = "CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPT_FORWARD=false\n"; + std::fs::write(&env_path, original).unwrap(); + let mut perms = std::fs::metadata(&env_path).unwrap().permissions(); + perms.set_mode(0o600); + std::fs::set_permissions(&env_path, perms).unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, true, true); + + assert!(matches!(result.status, SetupStatus::Error)); + assert!(result.detail.contains("conflicting")); + + // File should not be modified + let content = std::fs::read_to_string(&env_path).unwrap(); + assert_eq!(content, original); +} + +#[cfg(unix)] +#[test] +fn transcript_forward_migration_repeated_fix_attempts_preserve_file() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n").unwrap(); + let mut perms = std::fs::metadata(&env_path).unwrap().permissions(); + perms.set_mode(0o600); + std::fs::set_permissions(&env_path, perms).unwrap(); + + // Every run fails closed without mutating operator configuration. + let result1 = check_transcript_forward_env_migration(&env_path, true, true); + assert!(matches!(result1.status, SetupStatus::Error)); + + let content1 = std::fs::read_to_string(&env_path).unwrap(); + assert_eq!(content1, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n"); + + let result2 = check_transcript_forward_env_migration(&env_path, true, true); + assert!(matches!(result2.status, SetupStatus::Error)); + + let content2 = std::fs::read_to_string(&env_path).unwrap(); + assert_eq!(content1, content2); +} diff --git a/src/setup/doctor_transcript_forward_tests.rs b/src/setup/doctor_transcript_forward_tests.rs new file mode 100644 index 00000000..f60bd7d7 --- /dev/null +++ b/src/setup/doctor_transcript_forward_tests.rs @@ -0,0 +1,60 @@ +use super::*; + +#[test] +fn fix_fails_closed_and_preserves_equal_duplicates() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + let original = "CORTEX_AGENT_AI_TRANSCRIPTS=true\nCORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true\n"; + std::fs::write(&env_path, original).unwrap(); + + let result = check_transcript_forward_env_migration(&env_path, true, true); + assert!(matches!(result.status, SetupStatus::Error)); + assert!(result.detail.contains("automatic rewrite is disabled")); + assert_eq!(std::fs::read_to_string(&env_path).unwrap(), original); +} + +#[test] +fn noncooperative_edit_after_validated_read_is_never_lost() { + let dir = tempfile::tempdir().unwrap(); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n").unwrap(); + + let result = check_transcript_forward_env_migration_with_hook(&env_path, true, true, || { + std::fs::write(&env_path, "operator-change=true\n").unwrap(); + }); + + assert!(matches!(result.status, SetupStatus::Error)); + assert_eq!( + std::fs::read_to_string(&env_path).unwrap(), + "operator-change=true\n" + ); +} + +#[cfg(unix)] +#[test] +fn noncooperative_symlink_swap_after_validated_read_is_not_detached() { + use std::os::unix::fs::symlink; + + let dir = tempfile::tempdir().unwrap(); + let target = dir.path().join("shared.env"); + let env_path = dir.path().join(".env"); + std::fs::write(&env_path, "CORTEX_AGENT_AI_TRANSCRIPTS=true\n").unwrap(); + std::fs::write(&target, "operator-target=true\n").unwrap(); + + let result = check_transcript_forward_env_migration_with_hook(&env_path, true, true, || { + std::fs::remove_file(&env_path).unwrap(); + symlink(&target, &env_path).unwrap(); + }); + + assert!(matches!(result.status, SetupStatus::Error)); + assert!( + std::fs::symlink_metadata(&env_path) + .unwrap() + .file_type() + .is_symlink() + ); + assert_eq!( + std::fs::read_to_string(&target).unwrap(), + "operator-target=true\n" + ); +} diff --git a/src/setup/heartbeat_agent.rs b/src/setup/heartbeat_agent.rs index c015d6f1..282d2cf2 100644 --- a/src/setup/heartbeat_agent.rs +++ b/src/setup/heartbeat_agent.rs @@ -161,7 +161,25 @@ fn write_heartbeat_agent_env(env_path: &Path) -> io::Result { shell_safe_value(&token)? )); } + let transcript_forward = std::env::var(heartbeat_agent::AI_TRANSCRIPT_FORWARD_ENV) + .ok() + .or_else(|| std::env::var(heartbeat_agent::AI_TRANSCRIPT_FORWARD_LEGACY_ENV).ok()) + .filter(|value| !value.trim().is_empty()); + if let Some(value) = transcript_forward { + body.push_str(&format!( + "{}={}\n", + heartbeat_agent::AI_TRANSCRIPT_FORWARD_ENV, + shell_safe_value(&value)? + )); + } for key in heartbeat_agent::OPTIONAL_ENV_KEYS { + if matches!( + *key, + heartbeat_agent::AI_TRANSCRIPT_FORWARD_ENV + | heartbeat_agent::AI_TRANSCRIPT_FORWARD_LEGACY_ENV + ) { + continue; + } if let Ok(value) = std::env::var(key) && !value.trim().is_empty() { diff --git a/src/setup/heartbeat_agent_tests.rs b/src/setup/heartbeat_agent_tests.rs index 3b931dc9..0185eb7c 100644 --- a/src/setup/heartbeat_agent_tests.rs +++ b/src/setup/heartbeat_agent_tests.rs @@ -352,7 +352,8 @@ fn write_heartbeat_agent_env_reads_setup_env_fallbacks_and_optional_syslog() { assert!(raw.contains("CORTEX_AGENT_SYSLOG_FILE=/var/log/syslog\n")); assert!(raw.contains("CORTEX_SYSLOG_TARGET=127.0.0.1:1514\n")); assert!(raw.contains("CORTEX_AGENT_FILE_TAILS=/var/log/app.log:app\n")); - assert!(raw.contains("CORTEX_AGENT_AI_TRANSCRIPTS=true\n")); + assert!(raw.contains("CORTEX_AGENT_AI_TRANSCRIPT_FORWARD=true\n")); + assert!(!raw.contains("CORTEX_AGENT_AI_TRANSCRIPTS=")); assert!(raw.contains("CORTEX_AGENT_COMMAND_FORWARD=true\n")); assert!(raw.contains("CORTEX_AGENT_SHELL_HISTORY_FORWARD=true\n")); assert!(raw.contains("CORTEX_AGENT_AUTO_UPDATE=false\n"));