diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index fe000eb..5c20e7e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -113,3 +113,15 @@ jobs: npm run build npx tsc --noEmit git diff --exit-code -- main.js + + windows-storage: + name: Windows storage + recovery + runs-on: windows-2025 + steps: + - uses: actions/checkout@v7 + - name: Select Rust toolchain + run: | + rustup toolchain install 1.88.0 --profile minimal + rustup default 1.88.0 + - name: Test storage and maintenance + run: cargo test --locked --test nosync_layout --test database_maintenance diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index b0e5f67..4250cf1 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -51,6 +51,9 @@ jobs: package_version="$(sed -n 's/^version = "\([^"]*\)"/\1/p' Cargo.toml | head -n 1)" test "$tag_version" = "$package_version" + - name: Test storage and maintenance + run: cargo test --locked --test nosync_layout --test database_maintenance + - name: Build release binary env: RUSTFLAGS: -D warnings @@ -141,15 +144,19 @@ jobs: case "$GITHUB_REF_NAME" in *-*) release_flags+=(--prerelease) ;; esac + notes_flags=(--generate-notes) + if [ -f "docs/releases/$GITHUB_REF_NAME.md" ]; then + notes_flags=(--notes-file "docs/releases/$GITHUB_REF_NAME.md") + fi if gh release view "$GITHUB_REF_NAME" >/dev/null 2>&1; then gh release upload "$GITHUB_REF_NAME" dist/* --clobber else gh release create "$GITHUB_REF_NAME" dist/* \ - --title "Noema $GITHUB_REF_NAME" --generate-notes "${release_flags[@]}" + --title "Noema $GITHUB_REF_NAME" "${notes_flags[@]}" "${release_flags[@]}" fi - name: Update Homebrew tap - if: env.HOMEBREW_TAP_TOKEN != '' + if: env.HOMEBREW_TAP_TOKEN != '' && !contains(github.ref_name, '-') env: GH_TOKEN: ${{ secrets.HOMEBREW_TAP_TOKEN }} HOMEBREW_TAP_TOKEN: ${{ secrets.HOMEBREW_TAP_TOKEN }} diff --git a/Cargo.lock b/Cargo.lock index 1b1634e..ab22a10 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1859,7 +1859,7 @@ dependencies = [ [[package]] name = "noema" -version = "0.21.9" +version = "0.22.0-rc.1" dependencies = [ "anyhow", "assert_cmd", diff --git a/Cargo.toml b/Cargo.toml index 4a1aa0e..57e0cfb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "noema" -version = "0.21.9" +version = "0.22.0-rc.1" edition = "2024" rust-version = "1.88" description = "The intentional memory layer for your AI agents" diff --git a/README.md b/README.md index 17e17d9..43fba69 100644 --- a/README.md +++ b/README.md @@ -200,6 +200,12 @@ noema get 20260329-we-chose-local-sqlite noema init --name [--path ] Create a new Cortex noema use Set the default Cortex noema cortex list List all known Cortexes +noema cortex storage [--json] Inspect database, WAL, and reusable-space statistics +noema cortex storage --database --backup + Safely change database storage with a backup +noema cortex storage --resume Finish an interrupted storage migration +noema cortex compact --backup [--json] + Compact a stopped database with backup and integrity checks noema cortex remove [--purge] [--force] Unregister a Cortex (--purge also deletes its directory) noema cortex backup [-o ] [--force] @@ -346,6 +352,22 @@ noema version Print version, commit, and build date 2. `NOEMA_CORTEX` environment variable 3. Default set via `noema use ` +### Database storage in iCloud Drive + +Use opt-in `nosync` storage to keep the database and its SQLite sidecar files in +`db.nosync/` while Markdown traces continue syncing. The storage command creates +a backup and safely moves an existing database. See +[database storage](docs/database-storage.md) for enabling, reversing, recovering, +and backing up this layout, including iCloud's parent-folder eviction limits. + +### Database maintenance + +`noema cortex storage ` reports database size, reusable free space, WAL size, +and compaction headroom. Add `--json` for structured output. For deliberate +compaction, stop clients and use `noema cortex compact --backup `. +Noema requires a complete external backup and verifies database integrity before +and after compaction. See [database maintenance](docs/database-maintenance.md). + ### Durability profiles Noema defaults to the `standard` durability profile. It matches the mutation @@ -373,6 +395,16 @@ for measurements and the full trade-off. ## Agent Integrations +MCP read-only hints describe the purpose of a tool. Retrieval tools such as +`get_trace`, `search_traces`, `recall_context`, and `find_similar_traces` advertise +`readOnlyHint: true` while retaining Noema's built-in access and search-hit +tracking. These signals help Noema assess memory importance and are part of +retrieval, not a separate operation agents must remember to perform. The existing +`record_usage` argument controls explicit read counts; it does not change the +read hint or disable automatic search-hit tracking. Read hints do not promise +zero database writes: internal metrics and metrics retention also continue. +Tools that can explicitly change trace content, tags, or lifecycle remain writes. + `noema integrate` connects a Cortex to supported coding agents and installs the small startup bootstrap that tells each agent to call `get_instructions`. An MCP entry without that bootstrap is callable but not reliably memory-aware after a @@ -555,10 +587,29 @@ Noema can run as an [MCP](https://modelcontextprotocol.io) server, giving any MC Call `get_instructions` first in any new agent session for concise Markdown guidance. Use `cortex_usage` when a client needs structured JSON context. MCP tool discovery and each tool's schema remain the authoritative callable tool reference. -Set `NOEMA_MCP_TOOL_PROFILE=continuity-read` on an MCP server process to expose only -`recall_context`, or `NOEMA_MCP_TOOL_PROFILE=continuity-capture` to expose only -`get_instructions` and `create_traces`. These opt-in profiles reduce model-facing tool schemas for -bounded continuity tasks; the default remains the full tool set. +HTTP servers expose the full tool set at `/mcp` and optional role endpoints: + +| Endpoint | Tool set | +| --- | --- | +| `/mcp` | All tools, including future additions; existing clients remain compatible | +| `/mcp/agent` | Everyday memory retrieval, capture, editing, tagging, voting, and lifecycle tools | +| `/mcp/maintainer` | Agent tools plus tag cleanup and operational diagnostics | +| `/mcp/curator` | Agent tools plus consolidation candidates and distilled-memory creation | +| `/mcp/federation` | Cortex identity, event exchange, and usage-signal exchange | + +All endpoints share the same cortex and usage tracking. Role endpoints enforce +explicit tool allowlists for both discovery and execution; they share the existing +authentication and are not separate authorization roles. Existing clients and +federation peers can keep using `/mcp`. See [MCP tool endpoints](docs/mcp-endpoints.md) +for exact tool sets and configuration. + +For **stdio** servers, `NOEMA_MCP_TOOL_PROFILE` accepts `full` (the default), +`agent`, `maintainer`, `curator`, or `federation`. The existing bounded profiles +remain available: `continuity-read` exposes only `recall_context`, and +`continuity-capture` exposes only `get_instructions` and `create_traces`. +HTTP servers reject a restricted `NOEMA_MCP_TOOL_PROFILE` at startup; unset it +and choose a role URL instead. This prevents an existing restricted configuration +from silently widening access when `/mcp` becomes the full interface. ### stdio (Claude Desktop, Claude Code, any MCP client) diff --git a/docs/database-maintenance.md b/docs/database-maintenance.md new file mode 100644 index 0000000..b1e4efb --- /dev/null +++ b/docs/database-maintenance.md @@ -0,0 +1,86 @@ +# Database maintenance + +Noema exposes storage diagnostics and explicit compaction for both `db/` and +`db.nosync/`. It does not schedule full vacuuming or change SQLite's auto-vacuum +setting automatically. + +## Inspect storage + +```sh +noema cortex storage mycortex +noema cortex storage mycortex --json +``` + +Diagnostics can run while Noema clients are active. They read database statistics +without running cortex startup, rebuilding search, pruning records, or requesting +a WAL checkpoint. A missing database is reported rather than created. + +The report includes: + +- The selected layout and database path. +- Physical database file size and logical size, including pages represented in WAL. +- Completely free pages, their byte size, and their percentage of the database. +- WAL file size and the current auto-vacuum mode. +- Available space on the database filesystem and conservative compaction headroom. + +Free pages remain available for future database writes. Their total is not an exact +prediction of compaction savings: a full vacuum can also repack partially filled +pages. A large WAL is a separate issue from free space in the main database. +Statistics are a snapshot; file sizes can change while other clients write. +JSON size fields use bytes and `reusable_percent` uses the range 0–100. + +## Compact explicitly + +Stop all servers, watchers, and agent connections using the cortex, then choose +a new backup archive outside its directory: + +```sh +noema cortex compact mycortex --backup /path/to/backups/before-compact.tar.gz +``` + +Add `--json` for a structured result containing `before`, `after`, `backup_path`, +and `reclaimed_bytes`. The ordinary output reports before/after file sizes, free +space, the final WAL size, and the backup location. `reclaimed_bytes` measures +reduction of the main database file, not changes in the WAL. + +The command: + +1. Acquires the exclusive Noema storage lock and opens the existing database. +2. Checks database integrity and available space. +3. Checkpoints and closes SQLite while retaining the Noema maintenance lock, then + creates a complete cortex backup. Existing backup files + are never overwritten, and backup paths inside the cortex are refused. +4. Reopens SQLite, rechecks space after creating the backup, then runs `VACUUM`. +5. Checks integrity again, checkpoints/truncates the WAL, and reports the result. + +Supported Noema clients prevent compaction while their database connection is +open. Older clients and direct SQLite tools must also be stopped; they do not +participate in Noema's process lock. Compaction does not stop or restart services +for you. Choose a maintenance window and restart clients after success. + +The space check conservatively requires twice the larger of the logical database +size and the main file size to be available on the database filesystem, in addition +to the space used by the backup. SQLite may also use temporary storage; ensure +that filesystem has room if configured separately. Space checks cannot reserve +capacity against other applications writing concurrently. + +Compaction preserves retained traces, event history, embeddings, federation state, +and the selected storage layout. It does not implement retention or remove old +events, and it does not rewrite Markdown files or configuration. + +## Interrupted or unsuccessful compaction + +Noema uses SQLite's transactional `VACUUM`; it does not replace the database file +with a separately generated copy. SQLite's normal journal/WAL recovery applies +after a process interruption. Leave the database and its sidecars together. + +The complete backup is written before vacuuming begins. If a later step fails, +the error identifies the retained backup path. After stopping competing clients +or resolving disk-space errors, verify the cortex and retry using a new backup +filename. There is no compaction-specific `--resume` journal. The existing +`noema cortex restore` workflow can restore the backup if recovery is needed. + +A full vacuum rewrites the database and can generate substantial I/O. Consider it +after large deletions or when diagnostics show enough unused space to justify a +maintenance window. A fixed calendar interval is not required for normal use. +See [SQLite's VACUUM documentation](https://www.sqlite.org/lang_vacuum.html). diff --git a/docs/database-storage.md b/docs/database-storage.md new file mode 100644 index 0000000..4a09441 --- /dev/null +++ b/docs/database-storage.md @@ -0,0 +1,113 @@ +# Keeping the database local in iCloud Drive + +For size diagnostics and explicit compaction, see [database maintenance](database-maintenance.md). + +Noema normally stores its SQLite database in `db/` inside each cortex. Database +writes, including background federation bookkeeping, can cause frequent iCloud +uploads even when no Markdown traces change. Opt-in `nosync` storage moves the +whole database directory to `db.nosync/`. Traces remain in their usual folders. + +## Enable or disable + +Stop every Noema server, watcher, and agent connection using this cortex before +changing storage. Upgrade all clients that will open it. New clients enforce a +local process lock outside the cortex, under `noema/storage-locks/` in +`XDG_RUNTIME_DIR` (or the operating system's temporary directory). Keep that +runtime directory local and consistent across clients. The lock uses the +canonical cortex path, so path aliases share the same lock. Replacing the +compatibility lock inside a synced cortex cannot bypass this runtime lock. + +Clients also retain the previous cortex-local lock for upgrade compatibility. +Restart existing clients to gain the runtime-lock protection; versions predating +storage locking must be stopped manually. Locks do not coordinate different +computers over iCloud. + +This is an offline maintenance operation. Stop direct SQLite tools and older +clients too, and prevent automatic restarts until migration finishes. A successful +initial checkpoint does not reserve the database against a new, nonparticipating +SQLite writer. Noema's process lock coordinates compatible Noema clients; it +cannot guarantee safe directory migration while another program ignores that +lock. The required backup captures the state before migration, not subsequent +writes by an unsupported concurrent client. + +Choose a new backup filename outside the cortex, preferably outside iCloud Drive: + +```sh +noema cortex storage mycortex +noema cortex storage mycortex --database nosync --backup /path/to/backups/before-nosync.tar.gz +``` + +The command checks database integrity, checkpoints the WAL, writes a complete +backup, and moves the database directory. It preserves event history, federation +state, embeddings, pending recovery records, and files in the database directory. +It refuses conflicting database directories and never overwrites an existing +backup. Restart your Noema clients after it succeeds. + +To return to the original layout, stop clients and run: + +```sh +noema cortex storage mycortex --database default --backup /path/to/backups/before-default.tar.gz +``` + +Reversing the setting makes the database eligible for iCloud syncing again. +Repeating an already completed setting is a no-op. A previously renamed +`db.nosync/` directory can be adopted using the same enable command and backup. + +## Configuration and compatibility + +The command maintains `storage.yaml` in the cortex root: + +```yaml +database: nosync # default or nosync +``` + +Use the command to change an existing cortex; editing this field alone does not +move the database. Normal cortexes without this file continue using `db/`. +The command preserves other YAML keys, comments, and their order; the `database` +entry must use a top-level `database:` line for command-based editing. +On Unix, existing `storage.yaml` permissions are preserved, including when +resuming an interrupted migration; a new file defaults to `0640`. +It leaves `cortex.md` and the global registry untouched. + +With `nosync` enabled, `db/noema.db` is a small, static compatibility guard, not a +second database. Older Noema versions fail to open it instead of creating an +empty database. Do not delete or replace this guard. Use a compatible version +to reverse the migration before downgrading. Do not rename the WAL or SHM files +individually, or replace the database directory with a symlink. + +## Interrupted migrations and backups + +An interrupted migration blocks normal database access. After stopping clients, +finish it with: + +```sh +noema cortex storage mycortex --resume +``` + +Resume completes the original direction. To undo it, finish the interrupted +migration and then run the command for the opposite storage mode. The original +backup is also available through Noema's ordinary cortex restore workflow. +Do not delete the migration journal to bypass recovery. + +`noema cortex backup` includes `db.nosync`, its recovery files, the configuration, +and the compatibility guard. Restore preserves this layout. Keep complete Noema +backups: Markdown alone does not preserve the database's event and federation +state. Historical event records may name a recovery artifact under the previous +directory prefix; the artifact moves with the database directory. + +## iCloud scope + +Apple documents `.nosync` as an exclusion from iCloud transfer. It also documents +that deleting or evicting a parent directory removes its `.nosync` children. +Keep the cortex's parent folder downloaded and maintain backups outside iCloud. +See [Apple's iCloud storage documentation](https://developer.apple.com/library/archive/documentation/General/Conceptual/iCloudDesignGuide/Chapters/iCloudFundametals.html). + +This setting does not disable Markdown syncing, change SQLite's WAL behavior, +or promise exclusion from Dropbox, OneDrive, or other providers. The database +can continue changing locally while iCloud ignores its directory. + +On another device, iCloud can deliver the setting and guard without the excluded +database. Noema reports the missing local database rather than silently creating +one. Restore a complete Noema backup to establish that device's local copy; +do not treat iCloud as replication for the database or run copied cortex +identities as independent federation peers. diff --git a/docs/mcp-endpoints.md b/docs/mcp-endpoints.md new file mode 100644 index 0000000..8f01911 --- /dev/null +++ b/docs/mcp-endpoints.md @@ -0,0 +1,84 @@ +# MCP tool endpoints + +A Noema HTTP server serves multiple MCP endpoints on the same host and port. +Existing clients can keep using `https://noema.example.com:3000/mcp`: this endpoint +always exposes all tools, including tools added in future versions. Narrower +endpoints are optional views of the same cortex, not separate memory stores. + +| Path | Current tools | Intended use | +| --- | ---: | --- | +| `/mcp` | 35 | Complete interface and compatibility with existing clients | +| `/mcp/agent` | 24 | Everyday agent memory work | +| `/mcp/maintainer` | 30 | Everyday work plus diagnosis and maintenance | +| `/mcp/curator` | 26 | Everyday work plus deliberate memory consolidation | +| `/mcp/federation` | 3 | Exchange between Noema instances | + +Counts describe the current tool inventory. `tools/list` is authoritative. +All five paths also accept a trailing slash. Unknown paths return 404 rather than +falling back to the full interface. + +## Tool membership + +The agent endpoint exposes: + +- Guidance and identity: `get_instructions`, `cortex_usage`, `cortex_identity`. +- Retrieval: `list_traces`, `get_trace`, `search_traces`, `recall_context`, + `find_similar_traces`. +- Capture and editing: `create_trace`, `create_traces`, `update_trace`, `append_trace`. +- Tags and engagement: `set_trace_tags`, `append_trace_tags`, `tag_stats`, + `vote_trace`, `search_activity`. +- Lifecycle and provenance: `archive_trace`, `unarchive_trace`, `delete_trace`, + `recover_trace`, `trace_history`, `trace_lineage`, `resolve_divergence`. + +The maintainer endpoint includes all agent tools plus `tag_doctor`, `rename_tag`, +`delete_tag`, `metrics_summary`, `consolidation_health`, and `federation_status`. + +The curator endpoint includes all agent tools plus `list_consolidation_candidates` +and `record_consolidation_result`. + +The federation endpoint exposes only `cortex_identity`, `sync_events`, and +`sync_read_signal`. Noema's existing federation client continues to construct +`/mcp` from the configured peer base URL; existing peer configurations need no +changes. The dedicated endpoint is available to clients that select an MCP URL +directly. + +`announce_peer` remains available through the full interface only. It acknowledges +an announcement and reports whether a peer is configured; it does not configure +or connect that peer. + +## Behavior and access + +Each role uses an explicit allowlist. New tools automatically appear at `/mcp`, +but must be deliberately added to narrower roles. Omitted tools are unavailable +through both `tools/list` and `tools/call`. MCP sessions are scoped to the endpoint +that created them: reconnect when changing endpoint URLs. + +All endpoints share the cortex connection and background workers. Retrieval keeps +its existing access counters, search-hit tracking, and read hints regardless of +endpoint. Role selection does not make usage tracking a separate operation. +`get_instructions` includes a concise, role-aware directory of the other endpoints +and explains when a different MCP connection is needed. It does not include the +other tools' schemas or imply that the client can automatically switch endpoints. +`cortex_usage` reports `runtime.mcp_tool_profile`, `runtime.mcp_tool_count`, and +`runtime.mcp_endpoints`. The directory provides relative paths, profile purposes, +the current endpoint, and the count of additional tools beyond the current profile. +Stdio guidance makes HTTP availability conditional rather than assuming an HTTP +server is running. + +Existing authentication, Host/Origin validation, TLS configuration, and federation +write restrictions apply across the endpoints. A role URL is not an authorization +boundary: a credential accepted by the server can also access `/mcp`. Use these +endpoints to select an appropriate tool set, not to grant different users different +permissions. + +## Stdio and existing process profiles + +Stdio clients can set `NOEMA_MCP_TOOL_PROFILE` to `agent`, `maintainer`, `curator`, +`federation`, or `full` (the default). The existing `continuity-read` and +`continuity-capture` profiles remain supported for bounded continuity integrations. + +For HTTP, choose the endpoint URL instead of a process-wide profile. An unset or +`full` profile is accepted. A restricted or unknown `NOEMA_MCP_TOOL_PROFILE` causes +HTTP startup to fail with migration guidance, rather than silently turning a +previously restricted process into a full-access interface. Remove that environment +setting and configure clients to use the appropriate endpoint. diff --git a/docs/releases/v0.22.0-rc.1.md b/docs/releases/v0.22.0-rc.1.md new file mode 100644 index 0000000..5df869b --- /dev/null +++ b/docs/releases/v0.22.0-rc.1.md @@ -0,0 +1,36 @@ +# Noema v0.22.0-rc.1 + +This release candidate adds opt-in local database storage for iCloud-backed cortexes, explicit database diagnostics and compaction, and focused MCP endpoints for different agent workflows. It also fixes repeated search-index rebuilds and makes retrieval tools appear as reads in MCP clients while preserving memory usage tracking. + +## Added + +- `noema cortex storage`: inspect database size, WAL size, reusable pages, and available space; safely migrate between `db/` and `db.nosync/` with a required full backup and interrupted-migration recovery. +- `noema cortex compact`: explicitly reclaim unused database space with a required backup, integrity checks, and space checks. No automatic periodic VACUUM is enabled. +- Focused MCP endpoints: `/mcp/agent` (24 tools), `/mcp/maintainer` (30), `/mcp/curator` (26), and `/mcp/federation` (3). `/mcp` continues to expose all 35 tools. +- Profile-aware instructions describe sibling endpoints and their additional capabilities. + +## Changed + +- Retrieval and diagnostic tools advertise read-only intent to MCP clients. Normal retrieval still records usage signals used by Noema's memory model. +- Storage-aware clients hold a local runtime lock outside the cortex, plus a compatibility lock for earlier storage-aware builds. Replacing a synced compatibility lock does not bypass the runtime lock. +- Prereleases do not update the stable Homebrew formula or cask. + +## Fixed + +- Reopening a cortex with unavailable trace files no longer repeatedly rebuilds an otherwise current search index. +- Storage migration preserves existing Unix configuration permissions, including recovery after interruption, along with comments and unknown YAML fields. + +- Compaction closes SQLite before its required archive is copied, and backups omit the transient storage lock file, avoiding Windows locked-file read failures. + +## Upgrade and testing notes + +- This is a prerelease for qualification, not the stable v0.22.0 release. +- Restart all Noema clients after upgrading so they acquire the new runtime lock. Use a consistent local runtime directory across clients. +- Migration and compaction require an offline maintenance window: stop servers, watchers, older clients, and direct SQLite tools, and prevent restarts until the command finishes. Locks coordinate participating Noema clients; the initial SQLite checkpoint cannot exclude a new writer that ignores those locks. +- `.nosync` keeps the database local while Markdown traces can continue syncing. Another device needs a complete backup restore; iCloud does not replicate the excluded database. Keep the parent folder downloaded and retain backups outside iCloud. +- Keep the compatibility guard at `db/noema.db` when using `db.nosync`; do not rename individual WAL/SHM files. +- Qualification has exercised macOS and Linux storage round trips, backup failures from real disk exhaustion, and recovery after interruption during VACUUM. Disk-full cases during the later maintenance stages remain follow-up coverage; no claim is made that every concurrent external writer or fault timing is supported. + +See [database storage](https://github.com/Fail-Safe/Noema/blob/v0.22.0-rc.1/docs/database-storage.md), [database maintenance](https://github.com/Fail-Safe/Noema/blob/v0.22.0-rc.1/docs/database-maintenance.md), and [MCP endpoints](https://github.com/Fail-Safe/Noema/blob/v0.22.0-rc.1/docs/mcp-endpoints.md) for details. + +Full changelog: https://github.com/Fail-Safe/Noema/compare/v0.21.9...v0.22.0-rc.1 diff --git a/src/cli.rs b/src/cli.rs index 3fb4bbc..9dca034 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -506,6 +506,30 @@ enum EmbeddingCommand { #[derive(Debug, Subcommand)] enum CortexCommand { List, + /// Inspect or safely change database storage (stop all clients before changing) + Storage { + name: String, + #[arg(long, value_enum, conflicts_with = "resume")] + database: Option, + /// Required migration backup, outside the cortex + #[arg(long, requires = "database", conflicts_with = "resume")] + backup: Option, + /// Finish an interrupted storage migration + #[arg(long)] + resume: bool, + /// Emit storage statistics as JSON + #[arg(long, conflicts_with_all = ["database", "backup", "resume"])] + json: bool, + }, + /// Compact a database after stopping all clients; requires an external backup + Compact { + name: String, + /// New backup archive outside the cortex (never overwritten) + #[arg(long)] + backup: PathBuf, + #[arg(long)] + json: bool, + }, /// Write a gzipped tarball of a cortex Backup { name: String, @@ -1354,8 +1378,13 @@ fn migrate_command(selected: Option<&str>, command: MigrateCommand) -> Result<() ); } println!( - " backups: cortex.md.{}.bak, db/noema.db.{}.bak", - result.stamp, result.stamp + " backups: cortex.md.{}.bak, {}/noema.db.{}.bak", + result.stamp, + crate::db::directory(&entry.path)? + .file_name() + .unwrap() + .to_string_lossy(), + result.stamp ); Ok(()) } @@ -1495,6 +1524,88 @@ fn use_cortex(name: &str) -> Result<()> { fn cortex_command(command: CortexCommand) -> Result<()> { match command { + CortexCommand::Storage { + name, + database, + backup, + resume, + json, + } => { + let cfg = Config::load()?; + let entry = cfg + .cortexes + .get(&name) + .ok_or_else(|| anyhow::anyhow!("unknown cortex"))?; + if database.is_some() || resume { + let mode = + crate::storage::migrate(&entry.path, database, backup.as_deref(), resume)?; + println!( + "Database storage: {} ({})", + mode.as_str(), + mode.directory_name() + ); + if mode == crate::storage::DatabaseStorage::Nosync { + println!( + "The database stays local in iCloud Drive. Keep separate backups and keep the parent folder downloaded." + ); + } + } else { + let stats = crate::maintenance::storage_stats(&entry.path)?; + if json { + println!("{}", serde_json::to_string_pretty(&stats)?); + } else { + println!("Database storage: {}", stats.database.as_str()); + println!("Database path: {}", stats.database_path.display()); + println!("Database file: {}", human_bytes(stats.database_bytes)); + println!("Logical database: {}", human_bytes(stats.logical_bytes)); + println!( + "Reusable free space: {} ({:.1}%; {} pages)", + human_bytes(stats.reusable_bytes), + stats.reusable_percent, + stats.free_pages + ); + println!("WAL file: {}", human_bytes(stats.wal_bytes)); + println!("Auto-vacuum: {}", stats.auto_vacuum); + println!( + "Available disk space: {}", + human_bytes(stats.available_disk_bytes) + ); + println!( + "Compaction headroom required: {} (plus backup space)", + human_bytes(stats.compact_required_free_bytes) + ); + println!( + "Reusable space can serve future writes; it is not an exact estimate of compaction savings." + ); + } + } + } + CortexCommand::Compact { name, backup, json } => { + let cfg = Config::load()?; + let entry = cfg + .cortexes + .get(&name) + .ok_or_else(|| anyhow::anyhow!("unknown cortex"))?; + let result = crate::maintenance::compact(&entry.path, &backup)?; + if json { + println!("{}", serde_json::to_string_pretty(&result)?); + } else { + println!("Database compacted and integrity verified."); + println!( + "Database file: {} -> {} ({} reclaimed)", + human_bytes(result.before.database_bytes), + human_bytes(result.after.database_bytes), + human_bytes(result.reclaimed_bytes) + ); + println!( + "Reusable free space: {} -> {}", + human_bytes(result.before.reusable_bytes), + human_bytes(result.after.reusable_bytes) + ); + println!("WAL file: {}", human_bytes(result.after.wal_bytes)); + println!("Backup: {}", result.backup_path.display()); + } + } CortexCommand::List => { let cfg = Config::load()?; for (name, entry) in cfg.cortexes { @@ -3250,7 +3361,7 @@ fn check_cortex_layout(cx: &Cortex) -> CheckResult { ("traces/", cx.traces_dir()), ("archive/traces/", cx.archive_dir()), ("trash/traces/", cx.trash_dir()), - ("db/", cx.dir.join("db")), + ("database directory", cx.db_dir.clone()), ]; let missing = required .into_iter() diff --git a/src/cortex.rs b/src/cortex.rs index 0660e63..b2744ec 100644 --- a/src/cortex.rs +++ b/src/cortex.rs @@ -791,8 +791,9 @@ pub struct Cortex { pub id: String, pub name: String, pub dir: PathBuf, + pub db_dir: PathBuf, pub manifest: Manifest, - connection: Connection, + connection: db::Database, force_source_lock: bool, signing_key: Option, durability: DurabilityProfile, @@ -1014,7 +1015,8 @@ pub fn inspect_recovery_status(dir: &Path) -> RecoveryStatus { } fn inspect_recovery_status_inner(dir: &Path) -> Result { - let database_path = dir.join("db/noema.db"); + let _storage_lock = crate::storage::StorageLock::acquire(dir, false)?; + let database_path = db::directory(dir)?.join("noema.db"); if !fs::metadata(&database_path)?.is_file() { bail!("cortex database is not a regular file") } @@ -1132,11 +1134,13 @@ impl Cortex { ); } let connection = db::open(&dir)?; + let db_dir = db::directory(&dir)?; let signing_key = load_signing_key(&dir, &manifest)?; let mut cortex = Self { id: manifest.id.clone(), name, dir, + db_dir, manifest, connection, force_source_lock: false, @@ -1742,7 +1746,7 @@ impl Cortex { } fn pending_mutation_lock_directory(&self) -> PathBuf { - self.dir.join("db/pending-mutations") + self.db_dir.join("pending-mutations") } pub fn resolve(name_override: Option<&str>) -> Result { @@ -4480,7 +4484,7 @@ impl Cortex { let path = self.file_path(&row); let original_bytes = fs::read(&path) .with_context(|| format!("reading drifted trace {id:?} for recovery artifact"))?; - let artifact_directory = self.dir.join("db/reconciliations"); + let artifact_directory = self.db_dir.join("reconciliations"); fs::create_dir_all(&artifact_directory)?; #[cfg(unix)] { @@ -4490,7 +4494,10 @@ impl Cortex { let artifact_name = format!("{id}-{}.md", ulid::Ulid::new()); let artifact_path = artifact_directory.join(&artifact_name); trace::write_bytes_atomic_with_mode(&artifact_path, &original_bytes, 0o600)?; - let recovery_artifact = format!("db/reconciliations/{artifact_name}"); + let recovery_artifact = artifact_path + .strip_prefix(&self.dir)? + .to_string_lossy() + .replace('\\', "/"); let (_, mut canonical) = self.canonical_long_term_trace(&row)?; canonical.frontmatter.extra = Trace::parse(&original_bytes)?.frontmatter.extra; let now = trace::now_rfc3339(); @@ -4542,7 +4549,7 @@ impl Cortex { current.row.id ) })?; - let artifact_directory = self.dir.join("db/markdown-normalizations"); + let artifact_directory = self.db_dir.join("markdown-normalizations"); fs::create_dir_all(&artifact_directory)?; #[cfg(unix)] { @@ -4552,7 +4559,10 @@ impl Cortex { let artifact_name = format!("{}-{}.md", current.row.id, ulid::Ulid::new()); let artifact_path = artifact_directory.join(&artifact_name); trace::write_bytes_atomic_with_mode(&artifact_path, &original_bytes, 0o600)?; - let recovery_artifact = format!("db/markdown-normalizations/{artifact_name}"); + let recovery_artifact = artifact_path + .strip_prefix(&self.dir)? + .to_string_lossy() + .replace('\\', "/"); let mut normalized = current.trace.clone(); normalized.body = current.normalized_body; @@ -5326,18 +5336,8 @@ impl Cortex { } fn rebuild_fts_if_stale(&mut self) -> Result<()> { - let traces: i64 = self - .connection - .query_row("SELECT COUNT(*) FROM traces", [], |row| row.get(0))?; - let fts: i64 = self - .connection - .query_row("SELECT COUNT(*) FROM traces_fts", [], |row| row.get(0))?; - if traces == fts { - return Ok(()); - } let tx = self.connection.unchecked_transaction()?; - tx.execute("DELETE FROM traces_fts", [])?; - let ids: Vec = { + let paths: BTreeMap = { let mut statement = tx.prepare("SELECT id,archived_at,trashed_at FROM traces")?; statement .query_map([], |row| { @@ -5349,16 +5349,44 @@ impl Cortex { })? .collect::>>()? .into_iter() - .filter_map(|(id, archived, trashed)| { + .map(|(id, archived, trashed)| { let path = if trashed.is_some() { self.trash_dir().join(format!("{id}.md")) } else { self.trace_file(&id, archived.is_some()) }; - path.exists().then_some(id) + (id, path) }) .collect() }; + let indexed: Vec = { + let mut statement = tx.prepare("SELECT id FROM traces_fts")?; + statement + .query_map([], |row| row.get(0))? + .collect::>()? + }; + let indexed_ids: BTreeSet<&String> = indexed.iter().collect(); + let mut stale = indexed_ids.len() != indexed.len() + || indexed_ids.iter().any(|id| !paths.contains_key(*id)); + // Missing files cannot be indexed. Only a missing index entry with an + // available file warrants a rebuild; otherwise every open repeats it. + for (id, path) in &paths { + if !indexed_ids.contains(id) && path.try_exists()? { + stale = true; + break; + } + } + if !stale { + tx.commit()?; + return Ok(()); + } + let mut ids = Vec::new(); + for (id, path) in paths { + if path.try_exists()? { + ids.push(id); + } + } + tx.execute("DELETE FROM traces_fts", [])?; for id in ids { let row = self.get_from_tx(&tx, &id)?; let trace = Trace::parse_file(&self.file_path(&row))?; @@ -6165,6 +6193,111 @@ mod tests { (temp, cx) } + #[test] + fn fts_reopen_skips_missing_files_and_indexes_them_when_restored() { + for location in ["active", "archive", "trash"] { + let (_temp, cx) = cortex(); + let mut present = Trace::new("Present", "fact", "", vec![], "searchable quartz"); + cx.add(&mut present).unwrap(); + let mut absent = Trace::new("Absent", "fact", "", vec![], "restored zircon"); + cx.add(&mut absent).unwrap(); + let id = absent.frontmatter.id.clone(); + match location { + "archive" => cx.archive(&id).unwrap(), + "trash" => cx.trash(&id).unwrap(), + _ => (), + } + let path = match location { + "archive" => cx.archive_dir(), + "trash" => cx.trash_dir(), + _ => cx.traces_dir(), + } + .join(format!("{id}.md")); + let saved = fs::read(&path).unwrap(); + fs::remove_file(&path).unwrap(); + cx.connection + .execute("DELETE FROM traces_fts WHERE id=?1", [&id]) + .unwrap(); + let root = cx.dir.clone(); + drop(cx); + for _ in 0..3 { + let cx = Cortex::open("test", &root).unwrap(); + assert_eq!( + cx.connection.total_changes(), + 0, + "unnecessary writes for {location}" + ); + let hits: i64 = cx + .connection + .query_row( + "SELECT count(*) FROM traces_fts WHERE traces_fts MATCH 'quartz'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(hits, 1); + } + fs::write(&path, saved).unwrap(); + let cx = Cortex::open("test", &root).unwrap(); + let hits: i64 = cx + .connection + .query_row( + "SELECT count(*) FROM traces_fts WHERE traces_fts MATCH 'zircon'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(hits, 1); + drop(cx); + let cx = Cortex::open("test", &root).unwrap(); + assert_eq!(cx.connection.total_changes(), 0); + } + } + + #[test] + fn fts_reopen_repairs_wrong_ids_and_duplicates_even_when_counts_match() { + for duplicate in [false, true] { + let (_temp, cx) = cortex(); + let mut first = Trace::new("First", "fact", "", vec![], "quartz"); + let mut second = Trace::new("Second", "fact", "", vec![], "zircon"); + cx.add(&mut first).unwrap(); + cx.add(&mut second).unwrap(); + cx.connection + .execute( + "DELETE FROM traces_fts WHERE id=?1", + [&second.frontmatter.id], + ) + .unwrap(); + cx.connection + .execute( + "INSERT INTO traces_fts(id,title,body,tags) VALUES (?1,'','stale','')", + [if duplicate { + first.frontmatter.id.as_str() + } else { + "orphan" + }], + ) + .unwrap(); + let root = cx.dir.clone(); + drop(cx); + let cx = Cortex::open("test", &root).unwrap(); + let ids: Vec = cx + .connection + .prepare("SELECT id FROM traces_fts ORDER BY id") + .unwrap() + .query_map([], |row| row.get(0)) + .unwrap() + .collect::>() + .unwrap(); + let mut expected = vec![first.frontmatter.id, second.frontmatter.id]; + expected.sort(); + assert_eq!(ids, expected); + drop(cx); + let cx = Cortex::open("test", &root).unwrap(); + assert_eq!(cx.connection.total_changes(), 0); + } + } + #[test] fn durability_profile_parser_is_explicit_and_fail_closed() { assert_eq!( diff --git a/src/db.rs b/src/db.rs index fc0bcb8..b964e0e 100644 --- a/src/db.rs +++ b/src/db.rs @@ -1,4 +1,8 @@ -use std::{fs, path::Path, time::Duration}; +use std::{ + fs, + path::{Path, PathBuf}, + time::Duration, +}; use anyhow::{Context, Result}; use include_dir::{Dir, include_dir}; @@ -6,17 +10,39 @@ use rusqlite::{Connection, OpenFlags, OptionalExtension, params}; static MIGRATIONS: Dir<'_> = include_dir!("$CARGO_MANIFEST_DIR/migrations"); -pub fn open(cortex_dir: &Path) -> Result { - let db_dir = cortex_dir.join("db"); +pub fn directory(cortex_dir: &Path) -> Result { + crate::storage::directory(cortex_dir) +} + +#[derive(Debug)] +pub struct Database { + connection: Connection, + _storage_lock: crate::storage::StorageLock, +} + +impl std::ops::Deref for Database { + type Target = Connection; + fn deref(&self) -> &Connection { + &self.connection + } +} + +pub fn open(cortex_dir: &Path) -> Result { + let storage_lock = crate::storage::StorageLock::acquire(cortex_dir, false)?; + let db_dir = directory(cortex_dir)?; fs::create_dir_all(&db_dir)?; let connection = Connection::open(db_dir.join("noema.db"))?; configure(&connection)?; migrate(&connection)?; - Ok(connection) + Ok(Database { + connection, + _storage_lock: storage_lock, + }) } -pub fn open_existing_without_migrations(cortex_dir: &Path) -> Result> { - let path = cortex_dir.join("db/noema.db"); +pub fn open_existing_without_migrations(cortex_dir: &Path) -> Result> { + let storage_lock = crate::storage::StorageLock::acquire(cortex_dir, false)?; + let path = directory(cortex_dir)?.join("noema.db"); if !path.exists() { return Ok(None); } @@ -25,7 +51,10 @@ pub fn open_existing_without_migrations(cortex_dir: &Path) -> Result Result<()> { diff --git a/src/lib.rs b/src/lib.rs index cbf11a0..1cfc151 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -9,12 +9,14 @@ pub mod eventsig; pub mod federation; pub mod integrate; pub mod lock; +pub mod maintenance; pub mod markdown_normalization; pub mod mcp; pub mod migration; pub mod plugin; pub mod prefetch; pub mod restore; +pub mod storage; pub mod tag; pub mod tlsutil; pub mod trace; diff --git a/src/maintenance.rs b/src/maintenance.rs new file mode 100644 index 0000000..c8f1fef --- /dev/null +++ b/src/maintenance.rs @@ -0,0 +1,223 @@ +use std::{ + fs, + path::{Path, PathBuf}, + time::Duration, +}; + +use anyhow::{Context, Result, bail}; +use rusqlite::{Connection, OpenFlags}; +use serde::Serialize; + +use crate::storage::{self, DatabaseStorage, StorageLock}; + +#[derive(Debug, Serialize)] +pub struct StorageStats { + pub database: DatabaseStorage, + pub database_path: PathBuf, + pub database_bytes: u64, + pub logical_bytes: u64, + pub wal_bytes: u64, + pub page_size: u64, + pub page_count: u64, + pub free_pages: u64, + pub reusable_bytes: u64, + pub reusable_percent: f64, + pub auto_vacuum: String, + pub available_disk_bytes: u64, + pub compact_required_free_bytes: u64, +} + +#[derive(Debug, Serialize)] +pub struct CompactResult { + pub before: StorageStats, + pub after: StorageStats, + pub backup_path: PathBuf, + pub reclaimed_bytes: u64, +} + +fn file_size(path: &Path, required: bool) -> Result { + match fs::symlink_metadata(path) { + Ok(metadata) if metadata.is_file() => Ok(metadata.len()), + Ok(_) => bail!( + "database and sidecars must be regular files: {}", + path.display() + ), + Err(error) if !required && error.kind() == std::io::ErrorKind::NotFound => Ok(0), + Err(error) => Err(error).with_context(|| format!("reading {}", path.display())), + } +} + +fn database_path(root: &Path) -> Result { + let directory = storage::directory(root)?; + let database = directory.join("noema.db"); + file_size(&database, true)?; + for sidecar in ["noema.db-wal", "noema.db-shm", "noema.db-journal"] { + file_size(&directory.join(sidecar), false)?; + } + Ok(database) +} + +fn read_stats(connection: &Connection, database_path: &Path) -> Result { + let tx = connection.unchecked_transaction()?; + let page_size = u64::from(tx.query_row("PRAGMA page_size", [], |row| row.get::<_, u32>(0))?); + let page_count = u64::from(tx.query_row("PRAGMA page_count", [], |row| row.get::<_, u32>(0))?); + let free_pages = + u64::from(tx.query_row("PRAGMA freelist_count", [], |row| row.get::<_, u32>(0))?); + let auto_vacuum: u32 = tx.query_row("PRAGMA auto_vacuum", [], |row| row.get(0))?; + tx.commit()?; + let directory = database_path + .parent() + .context("database has no parent directory")?; + let logical_bytes = page_size * page_count; + let database_bytes = file_size(database_path, true)?; + let compact_required_free_bytes = logical_bytes + .max(database_bytes) + .checked_mul(2) + .context("database is too large to calculate compaction headroom")?; + Ok(StorageStats { + database: if directory + .file_name() + .is_some_and(|name| name == "db.nosync") + { + DatabaseStorage::Nosync + } else { + DatabaseStorage::Default + }, + database_path: database_path.to_owned(), + database_bytes, + logical_bytes, + wal_bytes: file_size(&directory.join("noema.db-wal"), false)?, + page_size, + page_count, + free_pages, + reusable_bytes: page_size * free_pages, + reusable_percent: if page_count == 0 { + 0.0 + } else { + free_pages as f64 * 100.0 / page_count as f64 + }, + auto_vacuum: match auto_vacuum { + 0 => "none".into(), + 1 => "full".into(), + 2 => "incremental".into(), + other => format!("unknown ({other})"), + }, + available_disk_bytes: fs2::available_space(directory)?, + compact_required_free_bytes, + }) +} + +pub fn storage_stats(root: &Path) -> Result { + let _lock = StorageLock::acquire(root, false)?; + let path = database_path(root)?; + let connection = Connection::open_with_flags(&path, OpenFlags::SQLITE_OPEN_READ_ONLY)?; + connection.busy_timeout(Duration::from_secs(5))?; + connection.execute_batch("PRAGMA query_only=ON")?; + read_stats(&connection, &path) +} + +fn integrity_check(connection: &Connection) -> Result<()> { + let rows = connection + .prepare("PRAGMA integrity_check")? + .query_map([], |row| row.get::<_, String>(0))? + .collect::>>()?; + if rows != ["ok"] { + bail!("database integrity check failed"); + } + Ok(()) +} + +fn checkpoint(connection: &Connection) -> Result<()> { + let busy: i64 = + connection.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| row.get(0))?; + if busy != 0 { + bail!("database checkpoint is busy; stop all Noema clients before compacting"); + } + Ok(()) +} + +fn require_headroom(available: u64, required: u64) -> Result<()> { + if available < required { + bail!( + "insufficient disk space for compaction: {required} bytes required, {available} bytes available" + ); + } + Ok(()) +} + +pub fn compact(root: &Path, backup: &Path) -> Result { + let _lock = StorageLock::acquire(root, true)?; + crate::cortex::read_manifest(root)?; + let path = database_path(root)?; + let connection = Connection::open_with_flags(&path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; + connection.busy_timeout(Duration::from_secs(5))?; + connection.execute_batch("PRAGMA synchronous=FULL")?; + integrity_check(&connection).context("checking database before compaction")?; + let before = read_stats(&connection, &path)?; + require_headroom( + before.available_disk_bytes, + before.compact_required_free_bytes, + )?; + checkpoint(&connection)?; + connection + .close() + .map_err(|(_, error)| error) + .context("closing database before compaction backup")?; + crate::restore::backup_without_storage_lock(root, backup, false) + .context("creating required compaction backup")?; + let result = (|| -> Result { + let connection = Connection::open_with_flags(&path, OpenFlags::SQLITE_OPEN_READ_WRITE)?; + connection.busy_timeout(Duration::from_secs(5))?; + connection.execute_batch("PRAGMA synchronous=FULL")?; + // The backup may share this volume, so check again after it has consumed space. + require_headroom( + fs2::available_space(path.parent().unwrap())?, + before.compact_required_free_bytes, + )?; + pause_for_test("backed-up")?; + connection + .execute_batch("VACUUM") + .context("vacuuming database")?; + pause_for_test("vacuumed")?; + integrity_check(&connection).context("checking database after compaction")?; + checkpoint(&connection)?; + read_stats(&connection, &path) + })(); + let after = result.with_context(|| { + format!( + "compaction did not finish verification; backup retained at {}", + backup.display() + ) + })?; + Ok(CompactResult { + reclaimed_bytes: before.database_bytes.saturating_sub(after.database_bytes), + before, + after, + backup_path: backup.to_owned(), + }) +} + +fn pause_for_test(phase: &str) -> Result<()> { + #[cfg(debug_assertions)] + if std::env::var("NOEMA_TEST_COMPACT_PHASE").ok().as_deref() == Some(phase) + && let Some(marker) = std::env::var_os("NOEMA_TEST_COMPACT_PAUSE") + { + fs::write(marker, phase)?; + loop { + std::thread::sleep(Duration::from_millis(50)); + } + } + let _ = phase; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::require_headroom; + + #[test] + fn insufficient_disk_space_is_rejected() { + assert!(require_headroom(99, 100).is_err()); + assert!(require_headroom(100, 100).is_ok()); + } +} diff --git a/src/mcp.rs b/src/mcp.rs index f8e9f75..fbe8158 100644 --- a/src/mcp.rs +++ b/src/mcp.rs @@ -49,11 +49,60 @@ const MAX_RECALL_BODY_CHARS: usize = 12_000; const MAX_RECALL_PREFERENCES: usize = 24; const MAX_CREATE_TRACES: usize = 16; +const AGENT_TOOLS: &[&str] = &[ + "get_instructions", + "cortex_usage", + "list_traces", + "get_trace", + "create_trace", + "create_traces", + "search_traces", + "recall_context", + "find_similar_traces", + "delete_trace", + "recover_trace", + "archive_trace", + "unarchive_trace", + "update_trace", + "set_trace_tags", + "append_trace_tags", + "tag_stats", + "vote_trace", + "search_activity", + "append_trace", + "trace_history", + "trace_lineage", + "resolve_divergence", + "cortex_identity", +]; +const MAINTAINER_TOOLS: &[&str] = &[ + "tag_doctor", + "rename_tag", + "delete_tag", + "metrics_summary", + "consolidation_health", + "federation_status", +]; +const CURATOR_TOOLS: &[&str] = &[ + "list_consolidation_candidates", + "record_consolidation_result", +]; +const FEDERATION_TOOLS: &[&str] = &["cortex_identity", "sync_events", "sync_read_signal"]; +const HTTP_TOOL_ENDPOINTS: &[(&str, &str)] = &[ + ("/mcp", "full"), + ("/mcp/agent", "agent"), + ("/mcp/maintainer", "maintainer"), + ("/mcp/curator", "curator"), + ("/mcp/federation", "federation"), +]; + #[derive(Clone)] pub struct NoemaServer { cortex: Arc>, federation_mode: String, tool_router: ToolRouter, + tool_profile: String, + http_transport: bool, } impl NoemaServer { @@ -62,6 +111,8 @@ impl NoemaServer { path: impl Into, remote_transport: bool, ) -> Result { + let tool_profile = std::env::var("NOEMA_MCP_TOOL_PROFILE").unwrap_or_default(); + validate_http_tool_profile(remote_transport, &tool_profile)?; let name = name.into(); let cortex = Cortex::open(&name, path.into())?; let federation_mode = if remote_transport { @@ -76,11 +127,16 @@ impl NoemaServer { } else { String::new() }; - let tool_profile = std::env::var("NOEMA_MCP_TOOL_PROFILE").unwrap_or_default(); Ok(Self { cortex: Arc::new(Mutex::new(cortex)), + http_transport: remote_transport, federation_mode, tool_router: Self::tool_router_for_profile(&tool_profile)?, + tool_profile: if tool_profile.is_empty() { + "full".into() + } else { + tool_profile + }, }) } @@ -88,6 +144,18 @@ impl NoemaServer { let mut router = Self::tool_router(); match profile { "" | "full" => {} + "agent" | "maintainer" | "curator" | "federation" => { + router.map.retain(|name, _| { + if profile == "federation" { + FEDERATION_TOOLS.contains(&name.as_ref()) + } else { + AGENT_TOOLS.contains(&name.as_ref()) + || (profile == "maintainer" + && MAINTAINER_TOOLS.contains(&name.as_ref())) + || (profile == "curator" && CURATOR_TOOLS.contains(&name.as_ref())) + } + }); + } "continuity-read" => { router.map.retain(|name, _| name == "recall_context"); } @@ -101,6 +169,53 @@ impl NoemaServer { Ok(router) } + fn endpoint_guidance(&self) -> (String, serde_json::Value) { + let http = self.http_transport; + let mut guidance = format!( + "\n\n## Additional capabilities\n\nCurrent tool profile: `{}`. ", + self.tool_profile + ); + guidance.push_str(if http { + "These MCP endpoints are on the same server:\n" + } else { + "This is a stdio connection. If this cortex is also served over HTTP, these paths are available on that HTTP server:\n" + }); + let endpoints = HTTP_TOOL_ENDPOINTS.iter().map(|&(path, profile)| { + let purpose = match profile { + "full" => "Complete tool set", + "agent" => "Everyday memory retrieval, capture, editing, and lifecycle operations", + "maintainer" => "Everyday tools plus cortex-wide tag cleanup and operational diagnostics", + "curator" => "Everyday tools plus consolidation candidates and distilled-memory creation", + "federation" => "Cortex identity, event exchange, and usage-signal exchange", + _ => unreachable!("built-in endpoint profile"), + }; + let router = Self::tool_router_for_profile(profile).expect("built-in endpoint profile"); + let additional_tools = router.map.keys().filter(|name| !self.tool_router.map.contains_key(*name)).count(); + let current = http && profile == self.tool_profile; + let availability = if current { + "current endpoint" + } else if additional_tools == 0 { + "capabilities already available in this profile" + } else { + "additional capabilities require another connection" + }; + let _ = writeln!(guidance, "- `{path}`: {purpose} ({availability})."); + json!({"path":path,"profile":profile,"purpose":purpose,"current_endpoint":current,"additional_tool_count":additional_tools}) + }).collect::>(); + let policy = "Only tools advertised by this connection are callable here. If a task needs additional capabilities, use an appropriately configured MCP connection if your client supports it; otherwise explain which endpoint is needed. Endpoint awareness does not automatically establish a connection or grant access."; + guidance.push_str(policy); + ( + guidance, + json!({ + "current_profile":self.tool_profile, + "transport":if http { "http" } else { "stdio" }, + "path_base":if http { "same HTTP server" } else { "HTTP server for this cortex, if configured" }, + "endpoints":endpoints, + "connection_policy":policy, + }), + ) + } + async fn open(&self) -> Result, ErrorData> { Ok(self.cortex.clone().lock_owned().await) } @@ -544,25 +659,40 @@ struct CortexContractOutput { #[tool_router(router = tool_router)] impl NoemaServer { #[tool( - description = "Returns concise Markdown guidance for agent use of this Cortex. Call this first if you are unfamiliar with Noema; use cortex_usage for structured MCP/client context." + description = "Returns concise Markdown guidance for agent use of this Cortex. Call this first if you are unfamiliar with Noema; use cortex_usage for structured MCP/client context.", + annotations(read_only_hint = true) )] async fn get_instructions(&self, _: Parameters) -> Result { let cx = self.open().await?; - Ok(render_instructions(&cx.manifest)) + Ok(render_instructions(&cx.manifest) + &self.endpoint_guidance().0) } #[tool( - description = "Returns structured JSON context for MCP clients: active Cortex identity, trace semantics, startup preference pattern, runtime posture, and operational constraints. Tool discovery remains authoritative for callable tools." + description = "Returns structured JSON context for MCP clients: active Cortex identity, trace semantics, startup preference pattern, runtime posture, and operational constraints. Tool discovery remains authoritative for callable tools.", + annotations(read_only_hint = true) )] async fn cortex_usage( &self, _: Parameters, ) -> Result, ErrorData> { let cx = self.open().await?; - build_cortex_usage(&cx).map(Json).map_err(mcp_error) + let mut output = build_cortex_usage(&cx).map_err(mcp_error)?; + output + .runtime + .insert("mcp_tool_profile".into(), json!(self.tool_profile)); + output + .runtime + .insert("mcp_tool_count".into(), json!(self.tool_router.map.len())); + output + .runtime + .insert("mcp_endpoints".into(), self.endpoint_guidance().1); + Ok(Json(output)) } - #[tool(description = "List traces in the cortex")] + #[tool( + description = "List traces in the cortex", + annotations(read_only_hint = true) + )] async fn list_traces( &self, Parameters(p): Parameters, @@ -582,7 +712,10 @@ impl NoemaServer { Ok(format_rows(&rows)) } - #[tool(description = "Get a trace by ID, including its full body")] + #[tool( + description = "Get a trace by ID, including its full body", + annotations(read_only_hint = true) + )] async fn get_trace(&self, Parameters(p): Parameters) -> Result { let cx = self.open().await?; let started = Instant::now(); @@ -705,7 +838,10 @@ impl NoemaServer { })) } - #[tool(description = "Full-text search across traces")] + #[tool( + description = "Full-text search across traces", + annotations(read_only_hint = true) + )] async fn search_traces( &self, Parameters(p): Parameters, @@ -735,7 +871,8 @@ impl NoemaServer { } #[tool( - description = "Retrieve startup preferences and bounded full trace bodies for several task queries in one read-only call. Use this continuity fast path instead of separate list/search/get calls when available." + description = "Retrieve startup preferences and bounded full trace bodies for several task queries in one read-only call. Use this continuity fast path instead of separate list/search/get calls when available.", + annotations(read_only_hint = true) )] async fn recall_context( &self, @@ -843,7 +980,10 @@ impl NoemaServer { outcome.map(Json).map_err(mcp_error) } - #[tool(description = "Find traces related to a given trace")] + #[tool( + description = "Find traces related to a given trace", + annotations(read_only_hint = true) + )] async fn find_similar_traces( &self, Parameters(p): Parameters, @@ -1016,7 +1156,8 @@ impl NoemaServer { })) } #[tool( - description = "Tag taxonomy statistics across active and archived traces, including assignment, tier, visibility, and engagement counts." + description = "Tag taxonomy statistics across active and archived traces, including assignment, tier, visibility, and engagement counts.", + annotations(read_only_hint = true) )] async fn tag_stats( &self, @@ -1159,7 +1300,8 @@ impl NoemaServer { } #[tool( - description = "Internal tool. Returns short-term traces within the rolling consolidation window along with their usage signals (read_count, modify_count, tier_votes, derived_from_count). Consumer scores these and submits distilled mid-tier traces via record_consolidation_result." + description = "Internal tool. Returns short-term traces within the rolling consolidation window along with their usage signals (read_count, modify_count, tier_votes, derived_from_count). Consumer scores these and submits distilled mid-tier traces via record_consolidation_result.", + annotations(read_only_hint = true) )] async fn list_consolidation_candidates( &self, @@ -1198,7 +1340,8 @@ impl NoemaServer { }))) } #[tool( - description = "Recent consolidation pipeline health: daily success/fail/promote/distill counts within the lookback window, short→mid and mid→long promotion-latency percentiles, and the 1-source mid leak detector. Lets an agent or operator answer 'is consolidation actually happening, and is anything leaking?' without raw SQL against the events table." + description = "Recent consolidation pipeline health: daily success/fail/promote/distill counts within the lookback window, short→mid and mid→long promotion-latency percentiles, and the 1-source mid leak detector. Lets an agent or operator answer 'is consolidation actually happening, and is anything leaking?' without raw SQL against the events table.", + annotations(read_only_hint = true) )] async fn consolidation_health( &self, @@ -1214,7 +1357,8 @@ impl NoemaServer { }))) } #[tool( - description = "Local MCP/CLI operation metrics for the lookback window: per-op counts and p50/p95 latency, daily volume (opens/searches/creates), source mix, and get_trace usage-open rate. Local-only and pruneable — not federated telemetry. Lets an agent answer 'is search getting slow?' or 'are agents actually opening traces?' without SQL." + description = "Local MCP/CLI operation metrics for the lookback window: per-op counts and p50/p95 latency, daily volume (opens/searches/creates), source mix, and get_trace usage-open rate. Local-only and pruneable — not federated telemetry. Lets an agent answer 'is search getting slow?' or 'are agents actually opening traces?' without SQL.", + annotations(read_only_hint = true) )] async fn metrics_summary( &self, @@ -1228,7 +1372,8 @@ impl NoemaServer { }))) } #[tool( - description = "Top-N traces by federation-wide search popularity (search_hit_count then read_count) plus top-N tags by aggregate engagement. Lets an agent answer 'what's worth reading?' or 'which topics are hot?' without scanning every trace. Active traces only; archived/trashed are excluded." + description = "Top-N traces by federation-wide search popularity (search_hit_count then read_count) plus top-N tags by aggregate engagement. Lets an agent answer 'what's worth reading?' or 'which topics are hot?' without scanning every trace. Active traces only; archived/trashed are excluded.", + annotations(read_only_hint = true) )] async fn search_activity( &self, @@ -1290,7 +1435,8 @@ impl NoemaServer { Ok(format!("Content appended to trace {}.", p.id)) } #[tool( - description = "Show the event log (audit trail) for a trace: all mutations in chronological order." + description = "Show the event log (audit trail) for a trace: all mutations in chronological order.", + annotations(read_only_hint = true) )] async fn trace_history(&self, Parameters(p): Parameters) -> Result { let events = self.open().await?.history(&p.id).map_err(mcp_error)?; @@ -1307,7 +1453,8 @@ impl NoemaServer { Ok(output) } #[tool( - description = "Show the derivation graph for a trace: what it was derived from and what was derived from it." + description = "Show the derivation graph for a trace: what it was derived from and what was derived from it.", + annotations(read_only_hint = true) )] async fn trace_lineage(&self, Parameters(p): Parameters) -> Result { let (from, by) = self.open().await?.lineage(&p.id).map_err(mcp_error)?; @@ -1350,7 +1497,8 @@ impl NoemaServer { } } #[tool( - description = "Returns this cortex's stable identity (ULID, name, manifest version). Federation peers call this on every sync to verify the remote endpoint still belongs to the cortex they originally paired with." + description = "Returns this cortex's stable identity (ULID, name, manifest version). Federation peers call this on every sync to verify the remote endpoint still belongs to the cortex they originally paired with.", + annotations(read_only_hint = true) )] async fn cortex_identity(&self, _: Parameters) -> Result { let cx = self.open().await?; @@ -1379,7 +1527,8 @@ impl NoemaServer { Ok(json_text(payload)) } #[tool( - description = "Returns events from this cortex for federation sync. Remote peers call this to pull new events. Returns a JSON array of event objects." + description = "Returns events from this cortex for federation sync. Remote peers call this to pull new events. Returns a JSON array of event objects.", + annotations(read_only_hint = true) )] async fn sync_events( &self, @@ -1413,7 +1562,8 @@ impl NoemaServer { serde_json::to_string(&cx.events_since(since, limit).map_err(mcp_error)?).map_err(mcp_error) } #[tool( - description = "Returns per-peer tier-usage deltas (read_count, modify_count, search_hit_count, last_read_at) for federation sync. Each peer publishes only its own rows — the ring aggregates by SUMing over every peer's contribution, so consolidation decisions operate on a federation-wide signal rather than the local slice. Returns a JSON array of trace_usage rows owned by this cortex with updated_at > since. search_hit_count is omitted when zero for wire compatibility with pre-migration-015 peers." + description = "Returns per-peer tier-usage deltas (read_count, modify_count, search_hit_count, last_read_at) for federation sync. Each peer publishes only its own rows — the ring aggregates by SUMing over every peer's contribution, so consolidation decisions operate on a federation-wide signal rather than the local slice. Returns a JSON array of trace_usage rows owned by this cortex with updated_at > since. search_hit_count is omitted when zero for wire compatibility with pre-migration-015 peers.", + annotations(read_only_hint = true) )] async fn sync_read_signal( &self, @@ -1443,7 +1593,10 @@ impl NoemaServer { ) .map_err(mcp_error) } - #[tool(description = "Show federation configuration, peer sync state, and local vector clock.")] + #[tool( + description = "Show federation configuration, peer sync state, and local vector clock.", + annotations(read_only_hint = true) + )] async fn federation_status(&self, _: Parameters) -> Result { let cx = self.open().await?; render_federation_status(&cx).map_err(mcp_error) @@ -1681,6 +1834,38 @@ fn local_interface_addresses() -> Result> { bail!("dynamic interface discovery is not implemented on this platform") } +fn validate_http_tool_profile(remote_transport: bool, profile: &str) -> Result<()> { + if remote_transport && !matches!(profile, "" | "full") { + bail!( + "NOEMA_MCP_TOOL_PROFILE is only supported for stdio; HTTP /mcp always exposes all tools. Unset it and select a role endpoint such as /mcp/agent instead." + ); + } + Ok(()) +} + +fn build_http_tool_router( + server: &NoemaServer, + allowed_hosts: Vec, +) -> Result { + let mut router = axum::Router::new(); + for &(endpoint, profile) in HTTP_TOOL_ENDPOINTS { + let mut role_server = server.clone(); + role_server.http_transport = true; + role_server.tool_router = NoemaServer::tool_router_for_profile(profile)?; + role_server.tool_profile = profile.into(); + let service: StreamableHttpService = + StreamableHttpService::new( + move || Ok(role_server.clone()), + Default::default(), + StreamableHttpServerConfig::default().with_allowed_hosts(allowed_hosts.clone()), + ); + router = router + .route_service(endpoint, service.clone()) + .route_service(&format!("{endpoint}/"), service); + } + Ok(router) +} + pub async fn serve_http( name: String, path: PathBuf, @@ -1705,16 +1890,10 @@ pub async fn serve_http( ) }; let background_lock = CortexLock::try_acquire_background(&cortex_id)?; - let service: StreamableHttpService = - StreamableHttpService::new( - move || Ok(server.clone()), - Default::default(), - StreamableHttpServerConfig::default().with_allowed_hosts(allowed_http_hosts( - &hosts, - &dynamic_hosts, - &allowed_hosts, - )), - ); + let http_router = build_http_tool_router( + &server, + allowed_http_hosts(&hosts, &dynamic_hosts, &allowed_hosts), + )?; let certificate_path = tls.as_ref().map(|(certificate, _)| certificate.clone()); let tls_config = match tls { Some((certificate, private_key)) => Some( @@ -1795,10 +1974,7 @@ pub async fn serve_http( listener.local_addr()? ); } - let router = apply_http_middleware( - axum::Router::new().nest_service("/mcp", service), - &access_key, - ); + let router = apply_http_middleware(http_router, &access_key); let server_handles = listeners .iter() .map(|_| axum_server::Handle::new()) @@ -2807,6 +2983,44 @@ mod tests { assert!(!schema.contains("\"format\":\"uint\"")); } + #[test] + fn tool_discovery_classifies_reads_by_purpose_including_usage_tracking() { + let tools = NoemaServer::tool_router_for_profile("").unwrap().list_all(); + let expected = [ + "get_trace", + "search_traces", + "recall_context", + "find_similar_traces", + "metrics_summary", + "get_instructions", + "cortex_usage", + "list_traces", + "tag_stats", + "list_consolidation_candidates", + "consolidation_health", + "search_activity", + "trace_history", + "trace_lineage", + "cortex_identity", + "sync_events", + "sync_read_signal", + "federation_status", + ]; + assert_eq!(tools.len(), 35); + for tool in tools { + let wire = serde_json::to_value(&tool).unwrap(); + let read_only = wire["annotations"]["readOnlyHint"] + .as_bool() + .unwrap_or(false); + assert_eq!( + read_only, + expected.contains(&tool.name.as_ref()), + "{}", + tool.name + ); + } + } + #[test] fn continuity_profiles_expose_only_their_bounded_tools() { let full = NoemaServer::tool_router_for_profile("").unwrap(); @@ -3075,6 +3289,38 @@ mod tests { "The deployment color is ULTRAVIOLET" ); assert!(!output.usage_recorded); + let cx = server.open().await.unwrap(); + let usage = cx.local_usage_since("", 100).unwrap(); + let row = usage.iter().find(|row| row.trace_id == fact_id).unwrap(); + assert_eq!(row.search_hit_count, 1); + assert_eq!(row.read_count, 0); + } + + #[tokio::test] + async fn read_hint_preserves_get_trace_usage_tracking() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("test", temp.path()).unwrap(); + let root = temp.path().join("test"); + let cx = Cortex::open("test", &root).unwrap(); + let mut trace = Trace::new("Example", "fact", "test", vec![], "body"); + let id = trace.frontmatter.id.clone(); + cx.add(&mut trace).unwrap(); + drop(cx); + let server = NoemaServer::new("test", &root, false).unwrap(); + for record_usage in [false, true, true] { + server + .get_trace(Parameters(GetParams { + id: id.clone(), + record_usage, + })) + .await + .unwrap(); + } + let cx = server.open().await.unwrap(); + let usage = cx.local_usage_since("", 100).unwrap(); + let row = usage.iter().find(|row| row.trace_id == id).unwrap(); + assert_eq!(row.read_count, 2); + assert_eq!(cx.get_trace(&id).unwrap().1.body, "body"); } #[tokio::test] @@ -3287,3 +3533,6 @@ mod tests { assert!(addresses.contains(&std::net::IpAddr::V6(std::net::Ipv6Addr::LOCALHOST))); } } + +#[cfg(test)] +mod http_role_tests; diff --git a/src/mcp/http_role_tests.rs b/src/mcp/http_role_tests.rs new file mode 100644 index 0000000..9c00dc4 --- /dev/null +++ b/src/mcp/http_role_tests.rs @@ -0,0 +1,331 @@ +use super::*; +use rmcp::{ + model::{CallToolRequestParams, ClientInfo}, + transport::{ + StreamableHttpClientTransport, streamable_http_client::StreamableHttpClientTransportConfig, + }, +}; + +#[test] +fn role_profiles_are_explicit_and_preserve_continuity_profiles() { + let full = NoemaServer::tool_router_for_profile("full").unwrap(); + let agent = NoemaServer::tool_router_for_profile("agent").unwrap(); + assert_eq!(full.map.len(), 35); + assert_eq!(agent.map.len(), 24); + for (role, count) in [("maintainer", 30), ("curator", 26), ("federation", 3)] { + let router = NoemaServer::tool_router_for_profile(role).unwrap(); + assert_eq!(router.map.len(), count); + assert!(router.map.keys().all(|name| full.map.contains_key(name))); + if role != "federation" { + assert!(agent.map.keys().all(|name| router.map.contains_key(name))); + } + } + for name in [ + "sync_events", + "tag_doctor", + "record_consolidation_result", + "announce_peer", + ] { + assert!(!agent.map.contains_key(name)); + } + assert!(validate_http_tool_profile(true, "full").is_ok()); + assert!(validate_http_tool_profile(true, "continuity-read").is_err()); + assert!(validate_http_tool_profile(false, "continuity-read").is_ok()); + assert_eq!( + NoemaServer::tool_router_for_profile("continuity-read") + .unwrap() + .map + .len(), + 1 + ); + assert_eq!( + NoemaServer::tool_router_for_profile("continuity-capture") + .unwrap() + .map + .len(), + 2 + ); +} + +#[tokio::test] +async fn role_endpoints_enforce_discovery_calls_and_share_usage() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("test", temp.path()).unwrap(); + let root = temp.path().join("test"); + let server = NoemaServer::new("test", &root, false).unwrap(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let router = build_http_tool_router(&server, vec![address.to_string()]).unwrap(); + let router = apply_http_middleware( + router, + &AccessKey { + value: "test-secret".into(), + ..Default::default() + }, + ); + let task = tokio::spawn(async move { + axum::serve(listener, router).await.unwrap(); + }); + let http = reqwest::Client::new(); + let mut clients = Vec::new(); + for &(path, profile) in HTTP_TOOL_ENDPOINTS { + let endpoint = format!("http://{address}{path}"); + assert_eq!( + http.post(&endpoint).send().await.unwrap().status(), + StatusCode::UNAUTHORIZED + ); + let preflight = http + .request(reqwest::Method::OPTIONS, &endpoint) + .header(header::ORIGIN, "app://obsidian.md") + .header(header::ACCESS_CONTROL_REQUEST_METHOD, "POST") + .send() + .await + .unwrap(); + assert_eq!(preflight.status(), StatusCode::NO_CONTENT); + for (name, value) in [ + (header::HOST, "untrusted.example"), + (header::ORIGIN, "https://untrusted.example"), + ] { + let response = http + .post(&endpoint) + .bearer_auth("test-secret") + .header(name, value) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::FORBIDDEN); + } + + let transport = StreamableHttpClientTransport::with_client( + http.clone(), + StreamableHttpClientTransportConfig::with_uri(endpoint).auth_header("test-secret"), + ); + let client = ClientInfo::default().serve(transport).await.unwrap(); + let listed = client.list_all_tools().await.unwrap(); + let mut actual = listed + .iter() + .map(|t| t.name.to_string()) + .collect::>(); + actual.sort(); + let mut expected = NoemaServer::tool_router_for_profile(profile) + .unwrap() + .map + .keys() + .map(ToString::to_string) + .collect::>(); + expected.sort(); + assert_eq!(actual, expected, "{profile}"); + assert!( + listed + .iter() + .find(|t| t.name == "cortex_identity") + .unwrap() + .annotations + .as_ref() + .unwrap() + .read_only_hint + .unwrap() + ); + let identity = client + .call_tool(CallToolRequestParams::new("cortex_identity")) + .await + .unwrap(); + assert_ne!(identity.is_error, Some(true)); + if profile != "full" { + // Even valid arguments cannot invoke an omitted tool by name. + let result = client + .call_tool( + CallToolRequestParams::new("announce_peer").with_arguments( + json!({"name":"peer-b","endpoint":"https://example.com"}) + .as_object() + .unwrap() + .clone(), + ), + ) + .await; + assert!( + result.unwrap_err().to_string().contains("tool not found"), + "{profile}" + ); + } + if profile == "federation" { + let result = client + .call_tool( + CallToolRequestParams::new("create_trace").with_arguments( + json!({"title":"Forbidden", "type":"fact", "body":"must not persist"}) + .as_object() + .unwrap() + .clone(), + ), + ) + .await; + assert!(result.unwrap_err().to_string().contains("tool not found")); + } + let allowed = match profile { + "maintainer" => "metrics_summary", + "curator" => "list_consolidation_candidates", + "federation" => "sync_events", + _ => "tag_stats", + }; + assert_ne!( + client + .call_tool(CallToolRequestParams::new(allowed)) + .await + .unwrap() + .is_error, + Some(true) + ); + clients.push((profile, client)); + } + let agent = &clients.iter().find(|(p, _)| *p == "agent").unwrap().1; + let created = agent + .call_tool( + CallToolRequestParams::new("create_trace").with_arguments( + json!({"title":"Shared role memory","type":"fact","body":"quartz"}) + .as_object() + .unwrap() + .clone(), + ), + ) + .await + .unwrap(); + assert_ne!(created.is_error, Some(true)); + let id = { + let cx = server.open().await.unwrap(); + cx.list(&ListOptions::default()).unwrap()[0].id.clone() + }; + for (profile, client) in &clients { + if *profile == "federation" { + continue; + } + let read = client + .call_tool( + CallToolRequestParams::new("get_trace").with_arguments( + json!({"id":id,"record_usage":true}) + .as_object() + .unwrap() + .clone(), + ), + ) + .await + .unwrap(); + assert_ne!(read.is_error, Some(true)); + let usage = client + .call_tool(CallToolRequestParams::new("cortex_usage")) + .await + .unwrap(); + let output = usage.structured_content.unwrap(); + assert_eq!(output["runtime"]["mcp_tool_profile"], *profile); + let directory = &output["runtime"]["mcp_endpoints"]; + assert_eq!(directory["transport"], "http"); + let endpoints = directory["endpoints"].as_array().unwrap(); + assert_eq!(endpoints.len(), 5); + assert_eq!( + endpoints + .iter() + .filter(|e| e["current_endpoint"] == true) + .count(), + 1 + ); + let full = endpoints.iter().find(|e| e["profile"] == "full").unwrap(); + assert_eq!( + full["additional_tool_count"], + 35 - NoemaServer::tool_router_for_profile(profile) + .unwrap() + .map + .len() + ); + let instructions = client + .call_tool(CallToolRequestParams::new("get_instructions")) + .await + .unwrap(); + let text = serde_json::to_string(&instructions.content).unwrap(); + assert!(text.contains("/mcp/maintainer")); + assert!(text.contains("/mcp/curator")); + assert!(text.contains("Only tools advertised by this connection are callable here")); + } + { + let cx = server.open().await.unwrap(); + let usage = cx.local_usage_since("", 100).unwrap(); + assert_eq!( + usage.iter().find(|u| u.trace_id == id).unwrap().read_count, + 4 + ); + } + for path in ["/mcp/agnt", "/mcp/agent/unknown", "/mcp/full"] { + assert_eq!( + http.post(format!("http://{address}{path}")) + .bearer_auth("test-secret") + .send() + .await + .unwrap() + .status(), + StatusCode::NOT_FOUND + ); + } + for (_, client) in clients { + client.cancel().await.unwrap(); + } + task.abort(); +} + +#[tokio::test] +async fn sessions_cannot_cross_role_endpoints() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("test", temp.path()).unwrap(); + let server = NoemaServer::new("test", temp.path().join("test"), false).unwrap(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let router = build_http_tool_router(&server, vec![address.to_string()]).unwrap(); + let task = tokio::spawn(async move { + axum::serve(listener, router).await.unwrap(); + }); + let http = reqwest::Client::new(); + for (source, destination) in [("/mcp", "/mcp/agent"), ("/mcp/agent/", "/mcp")] { + let response = http.post(format!("http://{address}{source}")) + .header(header::ACCEPT, "application/json, text/event-stream") + .json(&json!({"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-03-26","capabilities":{},"clientInfo":{"name":"test","version":"1"}}})) + .send().await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let session = response.headers()["mcp-session-id"].clone(); + response.text().await.unwrap(); + let response = http + .post(format!("http://{address}{destination}")) + .header(header::ACCEPT, "application/json, text/event-stream") + .header("mcp-session-id", session) + .json(&json!({"jsonrpc":"2.0","id":2,"method":"tools/list","params":{}})) + .send() + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } + task.abort(); +} + +#[test] +fn stdio_endpoint_guidance_does_not_claim_an_http_connection() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("test", temp.path()).unwrap(); + let mut server = NoemaServer::new("test", temp.path().join("test"), false).unwrap(); + server.tool_profile = "agent".into(); + server.tool_router = NoemaServer::tool_router_for_profile("agent").unwrap(); + let (text, directory) = server.endpoint_guidance(); + assert!(text.contains("If this cortex is also served over HTTP")); + assert_eq!(directory["transport"], "stdio"); + let endpoints = directory["endpoints"].as_array().unwrap(); + assert!(endpoints.iter().all(|e| e["current_endpoint"] == false)); + assert_eq!( + endpoints + .iter() + .find(|e| e["profile"] == "maintainer") + .unwrap()["additional_tool_count"], + 6 + ); + assert_eq!( + endpoints + .iter() + .find(|e| e["profile"] == "curator") + .unwrap()["additional_tool_count"], + 2 + ); +} diff --git a/src/migration.rs b/src/migration.rs index 14e3200..dfc2d25 100644 --- a/src/migration.rs +++ b/src/migration.rs @@ -100,6 +100,7 @@ pub fn migrate_cortex_id( config_path: Option<&Path>, ) -> Result { let dir = &entry.path; + let _storage_lock = crate::storage::StorageLock::acquire(dir, false)?; let _lock = MigrationLock::acquire(dir)?; let journal_path = dir.join(JOURNAL_NAME); let existing_journal = read_journal(&journal_path)?; @@ -143,7 +144,7 @@ pub fn migrate_cortex_id( db::checkpoint_wal(dir).context("checkpointing database before backup")?; let manifest_path = dir.join("cortex.md"); - let database_path = dir.join("db/noema.db"); + let database_path = db::directory(dir)?.join("noema.db"); copy_backup_atomic( &manifest_path, &PathBuf::from(format!("{}.{}.bak", manifest_path.display(), journal.stamp)), diff --git a/src/restore.rs b/src/restore.rs index c67141c..3545ff8 100644 --- a/src/restore.rs +++ b/src/restore.rs @@ -91,6 +91,16 @@ struct RestoreTransaction { } pub fn backup(source: &Path, output: &Path, force: bool) -> Result { + let _storage_lock = crate::storage::StorageLock::acquire(source, false)?; + crate::storage::directory(source)?; + backup_without_storage_lock(source, output, force) +} + +pub(crate) fn backup_without_storage_lock( + source: &Path, + output: &Path, + force: bool, +) -> Result { let source_metadata = fs::metadata(source) .with_context(|| format!("reading cortex directory {}", source.display()))?; if !source_metadata.is_dir() { @@ -173,6 +183,13 @@ fn append_backup_entry( source: &Path, archive_path: &Path, ) -> Result<()> { + if archive_path.components().count() == 2 + && archive_path + .file_name() + .is_some_and(|name| name == ".noema-storage.lock") + { + return Ok(()); + } let metadata = fs::symlink_metadata(source)?; let mut header = tar::Header::new_gnu(); header.set_uid(0); @@ -205,7 +222,9 @@ fn append_backup_entry( header.set_size(metadata.len()); header.set_mode(metadata_mode(&metadata, 0o640)); let mut file = File::open(source)?; - builder.append_data(&mut header, archive_path, &mut file)?; + builder + .append_data(&mut header, archive_path, &mut file) + .with_context(|| format!("archiving {}", archive_path.display()))?; } else if metadata.file_type().is_symlink() && is_legacy_trash_alias(archive_path, &fs::read_link(source)?) { @@ -774,6 +793,7 @@ where let scratch = ScratchDirectory::create()?; let staged_cortex = extract_archive(tarball, &scratch.path) .with_context(|| format!("extracting restore archive {}", tarball.display()))?; + crate::storage::directory(&staged_cortex).context("validating restored database storage")?; let mut manifest = read_manifest(&staged_cortex).context("reading cortex.md from archive")?; validate_name(&manifest.name).context("invalid cortex name in archive")?; if !manifest.id.is_empty() { @@ -836,6 +856,8 @@ where Err(error) if error.kind() == io::ErrorKind::NotFound => false, Err(error) => return Err(error).context("inspecting restore destination"), }; + // Recreate transient coordination state before recording the recovery hash. + drop(crate::storage::StorageLock::acquire(&staged_cortex, false)?); let restored_hash = tree_hash(&staged_cortex).context("hashing staged restored cortex")?; let previous_hash = existing_destination .then(|| tree_hash(&final_path).context("hashing existing restore destination")) diff --git a/src/storage.rs b/src/storage.rs new file mode 100644 index 0000000..8f4b4f7 --- /dev/null +++ b/src/storage.rs @@ -0,0 +1,379 @@ +use std::{ + fs::{self, File, OpenOptions}, + path::{Path, PathBuf}, +}; + +use anyhow::{Context, Result, bail}; +use fs2::FileExt; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; + +use crate::trace::{sync_directory, write_bytes_atomic, write_bytes_atomic_with_mode}; + +const CONFIG: &str = "storage.yaml"; +const JOURNAL: &str = ".noema-storage-migration.json"; +const GUARD: &[u8] = b"Noema database storage is nosync. Upgrade Noema to access db.nosync; do not replace this guard.\n"; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, clap::ValueEnum)] +#[serde(rename_all = "lowercase")] +pub enum DatabaseStorage { + Default, + Nosync, +} + +impl DatabaseStorage { + pub fn directory_name(self) -> &'static str { + match self { + Self::Default => "db", + Self::Nosync => "db.nosync", + } + } + pub fn as_str(self) -> &'static str { + match self { + Self::Default => "default", + Self::Nosync => "nosync", + } + } +} + +#[derive(Debug)] +pub struct StorageLock(Vec); + +impl StorageLock { + pub fn acquire(root: &Path, exclusive: bool) -> Result { + let primary_path = runtime_lock_path(root)?; + let legacy_path = root.join(".noema-storage.lock"); + let mut files = Vec::with_capacity(2); + // Keep the legacy lock during upgrades; the runtime lock remains authoritative + // even if a sync provider replaces the file inside the cortex. + for path in [&primary_path, &legacy_path] { + exists(path)?; + let file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .truncate(false) + .open(path)?; + let result = if exclusive { + FileExt::try_lock_exclusive(&file) + } else { + FileExt::try_lock_shared(&file) + }; + result.context("database storage is in use; stop Noema servers, watchers, and other clients before changing storage")?; + files.push(file); + } + Ok(Self(files)) + } +} + +fn runtime_lock_path(root: &Path) -> Result { + let canonical_root = + fs::canonicalize(root).context("resolving cortex storage lock identity")?; + let key = format!( + "{:x}", + Sha256::digest(canonical_root.as_os_str().as_encoded_bytes()) + ); + let base = std::env::var_os("XDG_RUNTIME_DIR") + .map(PathBuf::from) + .unwrap_or_else(std::env::temp_dir); + let directory = base.join("noema").join("storage-locks").join(key); + exists(&directory)?; + fs::create_dir_all(&directory)?; + if fs::canonicalize(&directory)?.starts_with(&canonical_root) { + bail!("storage runtime locks must be outside the cortex; choose a local XDG_RUNTIME_DIR"); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + fs::set_permissions(&directory, fs::Permissions::from_mode(0o700))?; + } + Ok(directory.join("storage.lock")) +} + +impl Drop for StorageLock { + fn drop(&mut self) { + for file in &self.0 { + let _ = FileExt::unlock(file); + } + } +} + +fn exists(path: &Path) -> Result { + match fs::symlink_metadata(path) { + Ok(metadata) if metadata.file_type().is_symlink() => { + bail!("storage paths must not be symlinks: {}", path.display()) + } + Ok(_) => Ok(true), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false), + Err(error) => Err(error.into()), + } +} + +fn configured(root: &Path) -> Result> { + let path = root.join(CONFIG); + if !exists(&path)? { + return Ok(None); + } + let value: serde_yaml::Value = serde_yaml::from_slice(&fs::read(path)?)?; + let mapping = value + .as_mapping() + .context("storage.yaml must be a YAML mapping")?; + mapping + .get(serde_yaml::Value::String("database".into())) + .map(|value| { + serde_yaml::from_value(value.clone()).context("database must be default or nosync") + }) + .transpose() +} + +fn guard_directory(root: &Path, allow_empty: bool) -> Result { + let path = root.join("db"); + if !exists(&path)? { + return Ok(false); + } + let entries = fs::read_dir(&path)?.collect::>>()?; + if entries.is_empty() { + return Ok(allow_empty); + } + if entries.len() != 1 || entries[0].file_name() != "noema.db" { + return Ok(false); + } + let file = path.join("noema.db"); + Ok(exists(&file)? + && fs::metadata(&file)?.len() == GUARD.len() as u64 + && fs::read(file)? == GUARD) +} + +pub fn directory(root: &Path) -> Result { + if exists(&root.join(JOURNAL))? { + bail!( + "interrupted storage migration; run noema cortex storage --resume before opening this cortex" + ); + } + let regular = exists(&root.join("db"))?; + let local = exists(&root.join("db.nosync"))?; + let mode = configured(root)?; + if local { + if regular && !guard_directory(root, false)? { + bail!("both db and db.nosync exist; refusing ambiguous database storage"); + } + if mode == Some(DatabaseStorage::Default) { + bail!( + "storage.yaml selects default but db.nosync exists; use the storage migration command" + ); + } + let path = root.join("db.nosync"); + if !path.is_dir() || !exists(&path.join("noema.db"))? { + bail!("db.nosync database is missing; restore a full Noema backup on this device"); + } + return Ok(path); + } + if mode == Some(DatabaseStorage::Nosync) || guard_directory(root, false)? { + bail!("db.nosync database is missing; restore a full Noema backup on this device"); + } + Ok(root.join("db")) +} + +// Preserve comments, unknown keys, and ordering rather than serializing the configuration. +fn configuration_bytes(root: &Path, mode: DatabaseStorage) -> Result> { + let path = root.join(CONFIG); + let previous = if exists(&path)? { + fs::read_to_string(path)? + } else { + String::new() + }; + let parsed = configured(root)?; + let mut found = false; + let mut next = String::new(); + for line in previous.split_inclusive('\n') { + if let Some(rest) = line.strip_prefix("database:") { + if found { + bail!("duplicate database setting"); + } + found = true; + next.push_str(&format!("database: {}", mode.as_str())); + if let Some((_, comment)) = rest.split_once('#') { + next.push_str(" #"); + next.push_str(comment.trim_end_matches(['\r', '\n'])); + } + if line.ends_with('\n') { + next.push('\n'); + } + } else { + next.push_str(line); + } + } + if found && parsed.is_none() { + bail!("cannot safely identify the database setting in storage.yaml"); + } + if !found { + if parsed.is_some() { + bail!("use a top-level database: setting in storage.yaml before migrating"); + } + if !next.is_empty() && !next.ends_with('\n') { + next.push('\n'); + } + next.push_str(&format!("database: {}\n", mode.as_str())); + } + let value: serde_yaml::Value = serde_yaml::from_str(&next)?; + if value.get("database").and_then(|v| v.as_str()) != Some(mode.as_str()) { + bail!("cannot safely update storage.yaml"); + } + Ok(next.into_bytes()) +} + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct Migration { + version: u32, + target: DatabaseStorage, + config: String, +} + +pub fn migrate( + root: &Path, + target: Option, + backup: Option<&Path>, + resume: bool, +) -> Result { + let _lock = StorageLock::acquire(root, true)?; + crate::cortex::read_manifest(root)?; + let journal_path = root.join(JOURNAL); + let migration = if exists(&journal_path)? { + if !resume || target.is_some() || backup.is_some() { + bail!("interrupted storage migration; use --resume alone"); + } + let journal: Migration = serde_json::from_slice(&fs::read(&journal_path)?)?; + if journal.version != 1 { + bail!("unsupported storage migration version"); + } + let config: serde_yaml::Value = serde_yaml::from_str(&journal.config)?; + if config.get("database").and_then(|value| value.as_str()) != Some(journal.target.as_str()) + { + bail!("storage migration configuration does not match its target"); + } + journal + } else { + if resume { + bail!("no interrupted storage migration"); + } + let target = target.context("select --database default or --database nosync")?; + let source = directory(root)?; + let config = configuration_bytes(root, target)?; + if source.file_name().and_then(|v| v.to_str()) == Some(target.directory_name()) + && configured(root)? == Some(target) + && (target == DatabaseStorage::Default || guard_directory(root, false)?) + { + return Ok(target); + } + let backup = backup.context("--backup outside the cortex is required before changing database storage")?; + let database = source.join("noema.db"); + if !exists(&database)? { + bail!("source database is missing"); + } + let connection = rusqlite::Connection::open_with_flags( + &database, + rusqlite::OpenFlags::SQLITE_OPEN_READ_WRITE, + )?; + let integrity: String = connection.query_row("PRAGMA quick_check", [], |row| row.get(0))?; + if integrity != "ok" { + bail!("database integrity check failed"); + } + let busy: i64 = + connection.query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| row.get(0))?; + if busy != 0 { + bail!("database is busy; stop all Noema clients before migrating"); + } + drop(connection); + crate::restore::backup_without_storage_lock(root, backup, false)?; + let journal = Migration { + version: 1, + target, + config: String::from_utf8(config)?, + }; + write_bytes_atomic(&journal_path, &serde_json::to_vec(&journal)?)?; + pause_for_test("journal")?; + journal + }; + match migration.target { + DatabaseStorage::Nosync => { + if !exists(&root.join("db.nosync"))? { + if guard_directory(root, true)? { + bail!("source database is missing; restore the migration backup"); + } + fs::rename(root.join("db"), root.join("db.nosync"))?; + sync_directory(root)?; + } + if !exists(&root.join("db.nosync/noema.db"))? { + bail!("nosync database is missing"); + } + pause_for_test("moved")?; + if exists(&root.join("db"))? && !guard_directory(root, true)? { + bail!("unexpected files at db; refusing to overwrite them"); + } + fs::create_dir_all(root.join("db"))?; + let guard = root.join(".noema-storage-guard"); + exists(&guard)?; + write_bytes_atomic(&guard, GUARD)?; + fs::rename(guard, root.join("db/noema.db"))?; + sync_directory(&root.join("db"))?; + sync_directory(root)?; + } + DatabaseStorage::Default => { + if exists(&root.join("db.nosync"))? { + if exists(&root.join("db"))? { + if !guard_directory(root, true)? { + bail!("unexpected files at db; refusing to overwrite them"); + } + if exists(&root.join("db/noema.db"))? { + fs::remove_file(root.join("db/noema.db"))?; + sync_directory(&root.join("db"))?; + } + fs::remove_dir(root.join("db"))?; + sync_directory(root)?; + } + fs::rename(root.join("db.nosync"), root.join("db"))?; + sync_directory(root)?; + } + if !exists(&root.join("db/noema.db"))? { + bail!("default database is missing"); + } + pause_for_test("moved")?; + } + } + let config_path = root.join(CONFIG); + let existed = exists(&config_path)?; + #[cfg(unix)] + let mode = { + use std::os::unix::fs::PermissionsExt; + if existed { + fs::metadata(&config_path)?.permissions().mode() & 0o7777 + } else { + 0o640 + } + }; + #[cfg(not(unix))] + let mode = { + let _ = existed; + 0o640 + }; + write_bytes_atomic_with_mode(&config_path, migration.config.as_bytes(), mode)?; + pause_for_test("configured")?; + fs::remove_file(journal_path)?; + sync_directory(root)?; + Ok(migration.target) +} + +fn pause_for_test(phase: &str) -> Result<()> { + #[cfg(debug_assertions)] + if std::env::var("NOEMA_TEST_STORAGE_PHASE").ok().as_deref() == Some(phase) + && let Some(marker) = std::env::var_os("NOEMA_TEST_STORAGE_PAUSE") + { + fs::write(marker, phase)?; + loop { + std::thread::sleep(std::time::Duration::from_millis(50)); + } + } + let _ = phase; + Ok(()) +} diff --git a/tests/database_maintenance.rs b/tests/database_maintenance.rs new file mode 100644 index 0000000..cf1035d --- /dev/null +++ b/tests/database_maintenance.rs @@ -0,0 +1,326 @@ +use std::{ + fs, + path::Path, + process::{Command, Stdio}, + thread, + time::{Duration, Instant}, +}; + +use noema::{ + cortex::Cortex, + db, embedding, maintenance, + storage::{self, DatabaseStorage}, + trace::Trace, +}; +use rusqlite::{Connection, params}; + +fn cli(config: &Path, args: &[&str]) -> std::process::Output { + Command::new(env!("CARGO_BIN_EXE_noema")) + .env("XDG_CONFIG_HOME", config) + .args(args) + .output() + .unwrap() +} + +fn success(config: &Path, args: &[&str]) -> std::process::Output { + let output = cli(config, args); + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + output +} + +fn populate(root: &Path) -> String { + let cx = Cortex::open("sample", root).unwrap(); + let mut trace = Trace::new( + "Preserved history", + "fact", + "", + vec!["storage".into()], + "searchable quartz", + ); + cx.add(&mut trace).unwrap(); + let id = trace.frontmatter.id.clone(); + cx.set_federation_state("maintenance-test", "durable cursor") + .unwrap(); + drop(cx); + let database = db::open(root).unwrap(); + database.execute("INSERT INTO trace_embeddings(trace_id,embedding_model,dim,embedding,source_hash,updated_at) VALUES (?1,'test-model',3,?2,'test-hash','2026-01-01T00:00:00Z')", params![id, embedding::encode(&[1.0, 0.0, 0.0])]).unwrap(); + database + .execute_batch( + "CREATE TABLE space_fixture(value BLOB); + WITH RECURSIVE n(x) AS (SELECT 1 UNION ALL SELECT x+1 FROM n WHERE x<32) + INSERT INTO space_fixture SELECT zeroblob(65536) FROM n; + DROP TABLE space_fixture;", + ) + .unwrap(); + id +} + +fn assert_preserved(root: &Path, id: &str, history: usize) { + let cx = Cortex::open("sample", root).unwrap(); + assert_eq!(cx.get_trace(id).unwrap().1.body, "searchable quartz"); + assert_eq!(cx.history(id).unwrap().len(), history); + assert_eq!( + cx.federation_state("maintenance-test").unwrap(), + "durable cursor" + ); + let connection = db::open(root).unwrap(); + let vector: Vec = connection + .query_row( + "SELECT embedding FROM trace_embeddings WHERE trace_id=?1", + [id], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(vector, embedding::encode(&[1.0, 0.0, 0.0])); + let hits: i64 = connection + .query_row( + "SELECT count(*) FROM traces_fts WHERE traces_fts MATCH 'quartz'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(hits, 1); +} + +#[test] +fn compact_preserves_data_and_restorable_backup_in_both_layouts() { + for mode in [DatabaseStorage::Default, DatabaseStorage::Nosync] { + let temp = tempfile::tempdir().unwrap(); + let config = temp.path().join("config"); + success( + &config, + &[ + "init", + "--name", + "sample", + "--path", + temp.path().to_str().unwrap(), + ], + ); + let root = temp.path().join("sample"); + let id = populate(&root); + let history = Cortex::open("sample", &root) + .unwrap() + .history(&id) + .unwrap() + .len(); + if mode == DatabaseStorage::Nosync { + storage::migrate( + &root, + Some(mode), + Some(&temp.path().join("migration.tar.gz")), + false, + ) + .unwrap(); + } + let manifest = fs::read(root.join("cortex.md")).unwrap(); + let trace = fs::read(root.join("traces").join(format!("{id}.md"))).unwrap(); + let before = maintenance::storage_stats(&root).unwrap(); + assert!(before.reusable_bytes > 1024 * 1024); + let backup = temp.path().join("before-compact.tar.gz"); + let output = success( + &config, + &[ + "cortex", + "compact", + "sample", + "--backup", + backup.to_str().unwrap(), + "--json", + ], + ); + let report: serde_json::Value = serde_json::from_slice(&output.stdout).unwrap(); + assert!(report["reclaimed_bytes"].as_u64().unwrap() > 1024 * 1024); + assert_eq!(report["after"]["reusable_bytes"], 0); + assert_eq!(report["after"]["wal_bytes"], 0); + assert_eq!(fs::read(root.join("cortex.md")).unwrap(), manifest); + assert_eq!( + fs::read(root.join("traces").join(format!("{id}.md"))).unwrap(), + trace + ); + assert_preserved(&root, &id, history); + let restored_parent = temp.path().join("restored"); + success( + &temp.path().join("restore-config"), + &[ + "cortex", + "restore", + backup.to_str().unwrap(), + "--path", + restored_parent.to_str().unwrap(), + "--name", + "restored", + ], + ); + assert_preserved(&restored_parent.join("restored"), &id, history); + assert!( + restored_parent + .join("restored") + .join(mode.directory_name()) + .join("noema.db") + .is_file() + ); + } +} + +#[test] +fn status_is_read_only_while_compact_refuses_live_clients() { + let temp = tempfile::tempdir().unwrap(); + let config = temp.path().join("config"); + success( + &config, + &[ + "init", + "--name", + "sample", + "--path", + temp.path().to_str().unwrap(), + ], + ); + let root = temp.path().join("sample"); + populate(&root); + let cx = Cortex::open("sample", &root).unwrap(); + let observer = db::open(&root).unwrap(); + let version: i64 = observer + .query_row("PRAGMA data_version", [], |row| row.get(0)) + .unwrap(); + let stats = success(&config, &["cortex", "storage", "sample", "--json"]); + let stats: serde_json::Value = serde_json::from_slice(&stats.stdout).unwrap(); + assert_eq!(stats["database"], "default"); + assert!(stats["reusable_bytes"].as_u64().unwrap() > 0); + let after: i64 = observer + .query_row("PRAGMA data_version", [], |row| row.get(0)) + .unwrap(); + assert_eq!(version, after); + let backup = temp.path().join("busy.tar.gz"); + let output = cli( + &config, + &[ + "cortex", + "compact", + "sample", + "--backup", + backup.to_str().unwrap(), + ], + ); + assert!(!output.status.success()); + assert!(String::from_utf8_lossy(&output.stderr).contains("storage is in use")); + assert!(!backup.exists()); + drop(cx); +} + +#[test] +fn compact_refuses_unsafe_backup_and_missing_database() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + populate(&root); + let before = maintenance::storage_stats(&root).unwrap(); + assert!(maintenance::compact(&root, &root.join("backup.tar.gz")).is_err()); + let backup = temp.path().join("existing.tar.gz"); + fs::write(&backup, b"preserve existing archive").unwrap(); + assert!(maintenance::compact(&root, &backup).is_err()); + assert_eq!(fs::read(&backup).unwrap(), b"preserve existing archive"); + assert_eq!( + maintenance::storage_stats(&root).unwrap().reusable_bytes, + before.reusable_bytes + ); + fs::rename(root.join("db/noema.db"), temp.path().join("saved.db")).unwrap(); + assert!(maintenance::storage_stats(&root).is_err()); + assert!(maintenance::compact(&root, &temp.path().join("missing.tar.gz")).is_err()); + assert!(!root.join("db/noema.db").exists()); +} + +#[test] +#[cfg(debug_assertions)] +fn killed_compaction_keeps_database_and_backup_readable() { + for phase in ["backed-up", "vacuumed"] { + let temp = tempfile::tempdir().unwrap(); + let config = temp.path().join("config"); + success( + &config, + &[ + "init", + "--name", + "sample", + "--path", + temp.path().to_str().unwrap(), + ], + ); + let root = temp.path().join("sample"); + let id = populate(&root); + let history = Cortex::open("sample", &root) + .unwrap() + .history(&id) + .unwrap() + .len(); + let marker = temp.path().join("paused"); + let backup = temp.path().join("before.tar.gz"); + let mut child = Command::new(env!("CARGO_BIN_EXE_noema")) + .env("XDG_CONFIG_HOME", &config) + .env("NOEMA_TEST_COMPACT_PHASE", phase) + .env("NOEMA_TEST_COMPACT_PAUSE", &marker) + .args([ + "cortex", + "compact", + "sample", + "--backup", + backup.to_str().unwrap(), + ]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .unwrap(); + let deadline = Instant::now() + Duration::from_secs(15); + while !marker.exists() && Instant::now() < deadline { + assert!( + child.try_wait().unwrap().is_none(), + "compaction exited before {phase}" + ); + thread::sleep(Duration::from_millis(20)); + } + assert!(marker.exists()); + assert!(maintenance::storage_stats(&root).is_err()); + child.kill().unwrap(); + child.wait().unwrap(); + assert_preserved(&root, &id, history); + let restored_parent = temp.path().join("restored"); + success( + &temp.path().join("restore-config"), + &[ + "cortex", + "restore", + backup.to_str().unwrap(), + "--path", + restored_parent.to_str().unwrap(), + "--name", + "restored", + ], + ); + assert_preserved(&restored_parent.join("restored"), &id, history); + let connection = Connection::open(root.join("db/noema.db")).unwrap(); + let integrity: String = connection + .query_row("PRAGMA integrity_check", [], |row| row.get(0)) + .unwrap(); + assert_eq!(integrity, "ok"); + } +} + +#[test] +fn compact_refuses_sqlite_writer_outside_noema_lock() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + populate(&root); + let writer = Connection::open(root.join("db/noema.db")).unwrap(); + writer.execute_batch("BEGIN IMMEDIATE").unwrap(); + let backup = temp.path().join("busy.tar.gz"); + let error = maintenance::compact(&root, &backup).unwrap_err(); + assert!(format!("{error:#}").contains("busy")); + assert!(!backup.exists()); + writer.execute_batch("ROLLBACK").unwrap(); +} diff --git a/tests/nosync_layout.rs b/tests/nosync_layout.rs new file mode 100644 index 0000000..53a07f1 --- /dev/null +++ b/tests/nosync_layout.rs @@ -0,0 +1,509 @@ +use std::{ + fs, + path::Path, + process::{Command, Stdio}, + thread, + time::{Duration, Instant}, +}; + +use noema::{ + cortex::Cortex, + db, + storage::{self, DatabaseStorage}, + trace::Trace, +}; + +fn cli(config: &Path, args: &[&str]) { + let output = Command::new(env!("CARGO_BIN_EXE_noema")) + .env("XDG_CONFIG_HOME", config) + .args(args) + .output() + .unwrap(); + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); +} + +#[test] +fn nosync_preserves_history_and_round_trips_through_backup_restore() { + let temp = tempfile::tempdir().unwrap(); + let config = temp.path().join("config"); + cli( + &config, + &[ + "init", + "--name", + "sample", + "--path", + temp.path().to_str().unwrap(), + ], + ); + let root = temp.path().join("sample"); + let cx = Cortex::open("sample", &root).unwrap(); + let mut trace = Trace::new("Storage trial", "fact", "", vec![], "preserved body"); + cx.add(&mut trace).unwrap(); + let id = trace.frontmatter.id.clone(); + let history = cx.history(&id).unwrap().len(); + cx.set_federation_state("trial-cursor", "preserved cursor") + .unwrap(); + drop(cx); + cli( + &config, + &[ + "cortex", + "storage", + "sample", + "--database", + "nosync", + "--backup", + temp.path().join("before.tar.gz").to_str().unwrap(), + ], + ); + cli(&config, &["verify", "cortex"]); + let archive = temp.path().join("backup.tar.gz"); + cli( + &config, + &[ + "cortex", + "backup", + "sample", + "--output", + archive.to_str().unwrap(), + ], + ); + let decoder = flate2::read::GzDecoder::new(fs::File::open(&archive).unwrap()); + let mut archived = tar::Archive::new(decoder); + assert!(archived.entries().unwrap().all(|entry| { + entry.unwrap().path().unwrap() != Path::new("sample/.noema-storage.lock") + })); + let destination = temp.path().join("restored"); + cli( + &temp.path().join("restore-config"), + &[ + "cortex", + "restore", + archive.to_str().unwrap(), + "--path", + destination.to_str().unwrap(), + "--name", + "restored", + ], + ); + let restored = destination.join("restored"); + let cx = Cortex::open("restored", &restored).unwrap(); + assert_eq!(cx.get_trace(&id).unwrap().1.body, "preserved body"); + assert_eq!(cx.history(&id).unwrap().len(), history); + assert_eq!( + cx.federation_state("trial-cursor").unwrap(), + "preserved cursor" + ); + assert!(restored.join("db.nosync/noema.db").is_file()); + assert!(restored.join("db/noema.db").is_file()); + assert!(root.join("db/noema.db").is_file()); + drop(cx); + storage::migrate( + &restored, + Some(DatabaseStorage::Default), + Some(&temp.path().join("reverse.tar.gz")), + false, + ) + .unwrap(); + let cx = Cortex::open("restored", &restored).unwrap(); + assert_eq!(cx.history(&id).unwrap().len(), history); + assert_eq!( + cx.federation_state("trial-cursor").unwrap(), + "preserved cursor" + ); + assert!(!restored.join("db.nosync").exists()); +} + +#[test] +fn migration_refuses_live_clients_and_preserves_configuration_bytes() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + let manifest = fs::read(root.join("cortex.md")).unwrap(); + fs::write( + root.join("storage.yaml"), + "# Storage preference\nfuture: {keep: true}\ndatabase: default # local choice\n", + ) + .unwrap(); + let backup = temp.path().join("before.tar.gz"); + let cx = Cortex::open("sample", &root).unwrap(); + let error = + storage::migrate(&root, Some(DatabaseStorage::Nosync), Some(&backup), false).unwrap_err(); + assert!(format!("{error:#}").contains("storage is in use")); + assert!(!backup.exists()); + drop(cx); + storage::migrate(&root, Some(DatabaseStorage::Nosync), Some(&backup), false).unwrap(); + assert_eq!(fs::read(root.join("cortex.md")).unwrap(), manifest); + assert_eq!( + fs::read_to_string(root.join("storage.yaml")).unwrap(), + "# Storage preference\nfuture: {keep: true}\ndatabase: nosync # local choice\n" + ); + // A pre-feature client opening the legacy SQLite path cannot initialize a second database. + let old = rusqlite::Connection::open(root.join("db/noema.db")).unwrap(); + assert!( + old.execute_batch("PRAGMA journal_mode=WAL; CREATE TABLE accidental(id);") + .is_err() + ); + drop(old); + storage::migrate(&root, Some(DatabaseStorage::Nosync), None, false).unwrap(); +} + +#[test] +fn trial_layout_is_adopted_and_missing_local_database_fails_closed() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + fs::rename(root.join("db"), root.join("db.nosync")).unwrap(); + let before = fs::read(root.join("db.nosync/noema.db")).unwrap(); + storage::migrate( + &root, + Some(DatabaseStorage::Nosync), + Some(&temp.path().join("before.tar.gz")), + false, + ) + .unwrap(); + assert_eq!(fs::read(root.join("db.nosync/noema.db")).unwrap(), before); + fs::rename(root.join("db.nosync"), temp.path().join("saved-db")).unwrap(); + assert!( + format!("{:#}", Cortex::open("sample", &root).err().unwrap()) + .contains("restore a full Noema backup") + ); + assert!(!root.join("db.nosync").exists()); +} + +#[test] +fn migration_requires_external_backup_and_rejects_conflicting_layouts() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + assert!(storage::migrate(&root, Some(DatabaseStorage::Nosync), None, false).is_err()); + assert!( + storage::migrate( + &root, + Some(DatabaseStorage::Nosync), + Some(&root.join("backup.tar.gz")), + false + ) + .is_err() + ); + assert!(!root.join(".noema-storage-migration.json").exists()); + fs::create_dir(root.join("db.nosync")).unwrap(); + assert!( + storage::migrate( + &root, + Some(DatabaseStorage::Nosync), + Some(&temp.path().join("backup.tar.gz")), + false + ) + .is_err() + ); + assert!(root.join("db/noema.db").is_file()); +} + +#[test] +fn unsupported_yaml_edit_preserves_unknown_keys_and_database() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + let config = "database:future: retained\n"; + fs::write(root.join("storage.yaml"), config).unwrap(); + let backup = temp.path().join("before.tar.gz"); + assert!(storage::migrate(&root, Some(DatabaseStorage::Nosync), Some(&backup), false).is_err()); + assert_eq!( + fs::read_to_string(root.join("storage.yaml")).unwrap(), + config + ); + assert!(root.join("db/noema.db").is_file()); + assert!(!backup.exists()); + assert!(!root.join(".noema-storage-migration.json").exists()); +} + +#[test] +#[cfg(debug_assertions)] +fn killed_storage_migrations_resume_in_both_directions() { + for direction in ["nosync", "default"] { + for phase in ["journal", "moved", "configured"] { + let temp = tempfile::tempdir().unwrap(); + let config = temp.path().join("config"); + cli( + &config, + &[ + "init", + "--name", + "sample", + "--path", + temp.path().to_str().unwrap(), + ], + ); + let root = temp.path().join("sample"); + let cx = Cortex::open("sample", &root).unwrap(); + cx.set_federation_state("storage-test", "durable cursor") + .unwrap(); + drop(cx); + if direction == "default" { + storage::migrate( + &root, + Some(DatabaseStorage::Nosync), + Some(&temp.path().join("initial.tar.gz")), + false, + ) + .unwrap(); + } + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let path = root.join("storage.yaml"); + if !path.exists() { + fs::write(&path, "database: default\n").unwrap(); + } + fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap(); + } + let marker = temp.path().join("paused"); + let mut child = Command::new(env!("CARGO_BIN_EXE_noema")) + .env("XDG_CONFIG_HOME", &config) + .env("NOEMA_TEST_STORAGE_PHASE", phase) + .env("NOEMA_TEST_STORAGE_PAUSE", &marker) + .args([ + "cortex", + "storage", + "sample", + "--database", + direction, + "--backup", + temp.path().join("before.tar.gz").to_str().unwrap(), + ]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .unwrap(); + let deadline = Instant::now() + Duration::from_secs(15); + while !marker.exists() && Instant::now() < deadline { + assert!( + child.try_wait().unwrap().is_none(), + "migration exited before {phase}" + ); + thread::sleep(Duration::from_millis(20)); + } + assert!(marker.exists(), "migration never reached {phase}"); + assert!(Cortex::open("sample", &root).is_err()); + child.kill().unwrap(); + child.wait().unwrap(); + assert!( + format!("{:#}", Cortex::open("sample", &root).err().unwrap()).contains("--resume") + ); + cli(&config, &["cortex", "storage", "sample", "--resume"]); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + assert_eq!( + fs::metadata(root.join("storage.yaml")) + .unwrap() + .permissions() + .mode() + & 0o777, + 0o600 + ); + } + let cx = Cortex::open("sample", &root).unwrap(); + assert_eq!( + cx.federation_state("storage-test").unwrap(), + "durable cursor" + ); + assert_eq!( + cx.db_dir.file_name().unwrap(), + if direction == "nosync" { + "db.nosync" + } else { + "db" + } + ); + assert!(!root.join(".noema-storage-migration.json").exists()); + } + } +} + +#[test] +fn ambiguous_directories_fail_without_creating_a_database() { + let temp = tempfile::tempdir().unwrap(); + fs::create_dir(temp.path().join("db")).unwrap(); + fs::create_dir(temp.path().join("db.nosync")).unwrap(); + let error = db::open(temp.path()).unwrap_err(); + assert!(error.to_string().contains("both db and db.nosync")); + assert!(!temp.path().join("db/noema.db").exists()); + assert!(!temp.path().join("db.nosync/noema.db").exists()); +} + +#[test] +fn recovery_artifact_uses_the_selected_directory() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + fs::rename(root.join("db"), root.join("db.nosync")).unwrap(); + let cx = Cortex::open("sample", &root).unwrap(); + let mut trace = Trace::new("Canonical policy", "fact", "", vec![], "canonical body"); + cx.add(&mut trace).unwrap(); + let id = trace.frontmatter.id.clone(); + cx.promote(&id, "mid").unwrap(); + cx.promote(&id, "long").unwrap(); + let path = cx.trace_file(&id, false); + let mut drifted = Trace::parse_file(&path).unwrap(); + drifted.body = "external edit".into(); + drifted.write_preserving_updated(&path).unwrap(); + let original = fs::read(&path).unwrap(); + let result = cx.reconcile_long_term(&id).unwrap(); + assert!( + result + .recovery_artifact + .starts_with("db.nosync/reconciliations/") + ); + assert_eq!( + fs::read(root.join(result.recovery_artifact)).unwrap(), + original + ); + assert!(!root.join("db").exists()); +} + +#[test] +#[cfg(debug_assertions)] +fn interrupted_mutation_recovers_from_nosync_database() { + let temp = tempfile::tempdir().unwrap(); + let config = temp.path().join("config"); + cli( + &config, + &[ + "init", + "--name", + "sample", + "--path", + temp.path().to_str().unwrap(), + ], + ); + let root = temp.path().join("sample"); + fs::rename(root.join("db"), root.join("db.nosync")).unwrap(); + let cx = Cortex::open("sample", &root).unwrap(); + let mut trace = Trace::new("Recovery trial", "fact", "", vec![], "original body"); + cx.add(&mut trace).unwrap(); + let id = trace.frontmatter.id.clone(); + let path = cx.trace_file(&id, false); + let original = fs::read(&path).unwrap(); + drop(cx); + let marker = temp.path().join("mutation-complete"); + let mut child = Command::new(env!("CARGO_BIN_EXE_noema")) + .env("XDG_CONFIG_HOME", &config) + .env("NOEMA_DURABILITY", "strong") + .env("NOEMA_RUST_TEST_PAUSE_AFTER_FILESYSTEM_MUTATION", &marker) + .args(["append", &id, "--content", "interrupted change"]) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .unwrap(); + let deadline = Instant::now() + Duration::from_secs(10); + while !marker.exists() && Instant::now() < deadline { + assert!( + child.try_wait().unwrap().is_none(), + "child exited before mutation" + ); + thread::sleep(Duration::from_millis(10)); + } + if !marker.exists() { + let _ = child.kill(); + let _ = child.wait(); + panic!("mutation marker timed out"); + } + assert_ne!(fs::read(&path).unwrap(), original); + child.kill().unwrap(); + child.wait().unwrap(); + let recovered = Cortex::open("sample", &root).unwrap(); + assert_eq!(recovered.get_trace(&id).unwrap().1.body, "original body"); + assert_eq!(fs::read(path).unwrap(), original); + assert!(!root.join("db").exists()); +} + +#[test] +#[cfg(unix)] +fn migration_preserves_private_configuration_permissions_in_both_directions() { + use std::os::unix::fs::PermissionsExt; + for mode in [0o600, 0o640] { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + let path = root.join("storage.yaml"); + fs::write(&path, "database: default\n").unwrap(); + fs::set_permissions(&path, fs::Permissions::from_mode(mode)).unwrap(); + for target in [DatabaseStorage::Nosync, DatabaseStorage::Default] { + storage::migrate( + &root, + Some(target), + Some(&temp.path().join(format!("{}.tar.gz", target.as_str()))), + false, + ) + .unwrap(); + assert_eq!( + fs::metadata(&path).unwrap().permissions().mode() & 0o777, + mode + ); + } + } +} + +#[test] +#[cfg(unix)] +fn replaced_legacy_lock_cannot_bypass_live_clients_even_through_path_alias() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + let alias = temp.path().join("alias"); + std::os::unix::fs::symlink(&root, &alias).unwrap(); + let cx = Cortex::open("sample", &root).unwrap(); + fs::rename( + root.join(".noema-storage.lock"), + temp.path().join("old-lock"), + ) + .unwrap(); + fs::write(root.join(".noema-storage.lock"), "").unwrap(); + for path in [&root, &alias] { + let backup = temp.path().join("blocked.tar.gz"); + let error = storage::migrate(path, Some(DatabaseStorage::Nosync), Some(&backup), false) + .unwrap_err(); + assert!(format!("{error:#}").contains("storage is in use")); + let error = noema::maintenance::compact(path, &backup).unwrap_err(); + assert!(format!("{error:#}").contains("storage is in use")); + assert!(!backup.exists()); + assert!(!root.join(".noema-storage-migration.json").exists()); + assert!(root.join("db/noema.db").is_file()); + } + drop(cx); + storage::migrate( + &root, + Some(DatabaseStorage::Nosync), + Some(&temp.path().join("allowed.tar.gz")), + false, + ) + .unwrap(); + noema::maintenance::compact(&root, &temp.path().join("compact.tar.gz")).unwrap(); +} + +#[test] +fn legacy_clients_still_block_maintenance() { + let temp = tempfile::tempdir().unwrap(); + Cortex::create("sample", temp.path()).unwrap(); + let root = temp.path().join("sample"); + let file = fs::OpenOptions::new() + .read(true) + .write(true) + .open(root.join(".noema-storage.lock")) + .unwrap(); + fs2::FileExt::try_lock_shared(&file).unwrap(); + let backup = temp.path().join("allowed.tar.gz"); + let error = + storage::migrate(&root, Some(DatabaseStorage::Nosync), Some(&backup), false).unwrap_err(); + assert!(format!("{error:#}").contains("storage is in use")); + assert!(!backup.exists()); + drop(file); + storage::migrate(&root, Some(DatabaseStorage::Nosync), Some(&backup), false).unwrap(); +}