From 7532065590ae005c65761c25d74fb4061fd3aa5f Mon Sep 17 00:00:00 2001 From: "ilblackdragon@gmail.com" Date: Sat, 28 Mar 2026 00:28:32 -0700 Subject: [PATCH] feat(workspace): metadata-driven indexing/hygiene, document versioning, and patch support MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Foundation for the extensible frontend system. Workspace documents now support metadata flags (skip_indexing, skip_versioning, hygiene config) via folder-level .config documents and per-file overrides, replacing hardcoded hygiene targets and indexing behavior. Key changes: - DocumentMetadata type with resolution chain (doc → folder .config → defaults) - Document versioning: auto-saves previous content on write/append/patch - Workspace patch: search-and-replace editing via memory_write tool - Hygiene rewrite: discovers cleanup targets from .config metadata instead of hardcoded daily/ and conversations/ directories - memory_read gains version/list_versions params - memory_write gains metadata/old_string/new_string/replace_all params - V14 migration adds memory_document_versions table (both PG + libSQL) Co-Authored-By: Claude Opus 4.6 (1M context) --- .env.example | 5 +- migrations/V14__document_versions.sql | 23 ++ src/agent/heartbeat.rs | 8 +- src/config/hygiene.rs | 18 +- src/db/libsql/workspace.rs | 284 +++++++++++++++- src/db/libsql_migrations.rs | 19 ++ src/db/mod.rs | 61 ++++ src/db/postgres.rs | 66 +++- src/error.rs | 6 + src/tools/builtin/memory.rs | 128 +++++++ src/workspace/document.rs | 242 ++++++++++++++ src/workspace/hygiene.rs | 462 ++++++++++++-------------- src/workspace/mod.rs | 353 +++++++++++++++++++- src/workspace/repository.rs | 192 ++++++++++- tests/e2e_routine_heartbeat.rs | 6 +- 15 files changed, 1585 insertions(+), 288 deletions(-) create mode 100644 migrations/V14__document_versions.sql diff --git a/.env.example b/.env.example index ce3e3124..2d78f4b6 100644 --- a/.env.example +++ b/.env.example @@ -191,10 +191,9 @@ HEARTBEAT_NOTIFY_CHANNEL=cli HEARTBEAT_NOTIFY_USER=default # Memory hygiene settings (automatic cleanup of stale workspace documents) -# Runs on each heartbeat tick; identity files (IDENTITY.md, SOUL.md) are never deleted +# Runs on each heartbeat tick; discovers cleanup targets from .config metadata # MEMORY_HYGIENE_ENABLED=true -# MEMORY_HYGIENE_DAILY_RETENTION_DAYS=30 # delete daily/ docs older than this many days -# MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS=7 # delete conversations/ docs older than this many days +# MEMORY_HYGIENE_VERSION_KEEP_COUNT=50 # max versions to keep per document # MEMORY_HYGIENE_CADENCE_HOURS=12 # minimum hours between cleanup passes # Docker Sandbox diff --git a/migrations/V14__document_versions.sql b/migrations/V14__document_versions.sql new file mode 100644 index 00000000..210f081e --- /dev/null +++ b/migrations/V14__document_versions.sql @@ -0,0 +1,23 @@ +-- Document version history for workspace files. +-- Every content update saves the previous content as a version, +-- enabling rollback and audit trails. + +CREATE TABLE memory_document_versions ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + document_id UUID NOT NULL REFERENCES memory_documents(id) ON DELETE CASCADE, + version INTEGER NOT NULL, + content TEXT NOT NULL, + content_hash TEXT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + changed_by TEXT, + UNIQUE(document_id, version) +); + +CREATE INDEX idx_doc_versions_lookup + ON memory_document_versions(document_id, version DESC); + +-- GIN index on metadata for JSON path queries (used by hygiene to find +-- .config documents with hygiene.enabled). The metadata column already +-- exists (V1) but was never indexed. +CREATE INDEX idx_memory_documents_metadata + ON memory_documents USING GIN (metadata jsonb_path_ops); diff --git a/src/agent/heartbeat.rs b/src/agent/heartbeat.rs index f7a8f869..ffa46982 100644 --- a/src/agent/heartbeat.rs +++ b/src/agent/heartbeat.rs @@ -276,8 +276,8 @@ impl HeartbeatRunner { .await; if report.had_work() { tracing::info!( - daily_logs_deleted = report.daily_logs_deleted, - conversation_docs_deleted = report.conversation_docs_deleted, + directories_cleaned = ?report.directories_cleaned, + versions_pruned = report.versions_pruned, "heartbeat: memory hygiene deleted stale documents" ); } @@ -598,8 +598,8 @@ pub fn spawn_multi_user_heartbeat( if report.had_work() { tracing::info!( user_id = hygiene_user, - daily_logs_deleted = report.daily_logs_deleted, - conversation_docs_deleted = report.conversation_docs_deleted, + directories_cleaned = ?report.directories_cleaned, + versions_pruned = report.versions_pruned, "multi-user heartbeat: memory hygiene deleted stale documents" ); } diff --git a/src/config/hygiene.rs b/src/config/hygiene.rs index b510933a..609f9743 100644 --- a/src/config/hygiene.rs +++ b/src/config/hygiene.rs @@ -10,10 +10,8 @@ use crate::error::ConfigError; pub struct HygieneConfig { /// Whether hygiene is enabled. Env: `MEMORY_HYGIENE_ENABLED` (default: true). pub enabled: bool, - /// Days before `daily/` documents are deleted. Env: `MEMORY_HYGIENE_DAILY_RETENTION_DAYS` (default: 30). - pub daily_retention_days: u32, - /// Days before `conversations/` documents are deleted. Env: `MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS` (default: 7). - pub conversation_retention_days: u32, + /// Maximum versions to keep per document. Env: `MEMORY_HYGIENE_VERSION_KEEP_COUNT` (default: 50). + pub version_keep_count: u32, /// Minimum hours between hygiene passes. Env: `MEMORY_HYGIENE_CADENCE_HOURS` (default: 12). pub cadence_hours: u32, } @@ -22,8 +20,7 @@ impl Default for HygieneConfig { fn default() -> Self { Self { enabled: true, - daily_retention_days: 30, - conversation_retention_days: 7, + version_keep_count: 50, cadence_hours: 12, } } @@ -33,11 +30,7 @@ impl HygieneConfig { pub(crate) fn resolve() -> Result { Ok(Self { enabled: parse_bool_env("MEMORY_HYGIENE_ENABLED", true)?, - daily_retention_days: parse_optional_env("MEMORY_HYGIENE_DAILY_RETENTION_DAYS", 30)?, - conversation_retention_days: parse_optional_env( - "MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS", - 7, - )?, + version_keep_count: parse_optional_env("MEMORY_HYGIENE_VERSION_KEEP_COUNT", 50)?, cadence_hours: parse_optional_env("MEMORY_HYGIENE_CADENCE_HOURS", 12)?, }) } @@ -47,8 +40,7 @@ impl HygieneConfig { pub fn to_workspace_config(&self) -> crate::workspace::hygiene::HygieneConfig { crate::workspace::hygiene::HygieneConfig { enabled: self.enabled, - daily_retention_days: self.daily_retention_days, - conversation_retention_days: self.conversation_retention_days, + version_keep_count: self.version_keep_count, cadence_hours: self.cadence_hours, state_dir: ironclaw_base_dir(), } diff --git a/src/db/libsql/workspace.rs b/src/db/libsql/workspace.rs index 5680e435..15ab6bb6 100644 --- a/src/db/libsql/workspace.rs +++ b/src/db/libsql/workspace.rs @@ -13,8 +13,8 @@ use super::{ use crate::db::WorkspaceStore; use crate::error::{DatabaseError, WorkspaceError}; use crate::workspace::{ - MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult, WorkspaceEntry, - fuse_results, + DocumentVersion, MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult, + VersionSummary, WorkspaceEntry, fuse_results, }; use chrono::Utc; @@ -840,6 +840,286 @@ impl WorkspaceStore for LibSqlBackend { Ok(fuse_results(fts_results, vector_results, config)) } + + // ==================== Metadata ==================== + + async fn update_document_metadata( + &self, + id: Uuid, + metadata: &serde_json::Value, + ) -> Result<(), WorkspaceError> { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let now = fmt_ts(&Utc::now()); + let meta_str = serde_json::to_string(metadata).map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to serialize metadata: {e}"), + })?; + conn.execute( + "UPDATE memory_documents SET metadata = ?2, updated_at = ?3 WHERE id = ?1", + params![id.to_string(), meta_str, now], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to update metadata: {e}"), + })?; + Ok(()) + } + + async fn find_config_documents( + &self, + user_id: &str, + agent_id: Option, + ) -> Result, WorkspaceError> { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let agent_str = agent_id.map(|a| a.to_string()); + let mut rows = conn + .query( + r#" + SELECT id, user_id, agent_id, path, content, + created_at, updated_at, metadata + FROM memory_documents + WHERE user_id = ?1 AND agent_id IS ?2 + AND (path LIKE '%/.config' OR path = '.config') + ORDER BY path + "#, + params![user_id, agent_str], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to find config documents: {e}"), + })?; + + let mut docs = Vec::new(); + while let Some(row) = rows.next().await.map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to read config document row: {e}"), + })? { + docs.push(row_to_memory_document(&row)); + } + Ok(docs) + } + + // ==================== Versioning ==================== + + async fn save_version( + &self, + document_id: Uuid, + content: &str, + content_hash: &str, + changed_by: Option<&str>, + ) -> Result { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let id = Uuid::new_v4().to_string(); + let doc_id = document_id.to_string(); + let now = fmt_ts(&Utc::now()); + + // Get next version number + let mut rows = conn + .query( + "SELECT COALESCE(MAX(version), 0) + 1 FROM memory_document_versions WHERE document_id = ?1", + params![doc_id.clone()], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to get next version number: {e}"), + })?; + + let next_version = if let Some(row) = rows.next().await.map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to read version number: {e}"), + })? { + get_i64(&row, 0) as i32 + } else { + 1 + }; + + conn.execute( + r#" + INSERT INTO memory_document_versions + (id, document_id, version, content, content_hash, created_at, changed_by) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) + "#, + params![id, doc_id, next_version as i64, content, content_hash, now, changed_by], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to save version: {e}"), + })?; + + Ok(next_version) + } + + async fn get_version( + &self, + document_id: Uuid, + version: i32, + ) -> Result { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let mut rows = conn + .query( + r#" + SELECT id, document_id, version, content, content_hash, + created_at, changed_by + FROM memory_document_versions + WHERE document_id = ?1 AND version = ?2 + "#, + params![document_id.to_string(), version as i64], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to get version: {e}"), + })?; + + let row = rows + .next() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to read version row: {e}"), + })? + .ok_or(WorkspaceError::VersionNotFound { + document_id, + version, + })?; + + Ok(DocumentVersion { + id: get_text(&row, 0).parse().map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Invalid version UUID: {e}"), + })?, + document_id: get_text(&row, 1).parse().map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Invalid document UUID: {e}"), + })?, + version: get_i64(&row, 2) as i32, + content: get_text(&row, 3), + content_hash: get_text(&row, 4), + created_at: get_ts(&row, 5), + changed_by: get_opt_text(&row, 6), + }) + } + + async fn list_versions( + &self, + document_id: Uuid, + limit: i64, + ) -> Result, WorkspaceError> { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let mut rows = conn + .query( + r#" + SELECT version, content_hash, created_at, changed_by + FROM memory_document_versions + WHERE document_id = ?1 + ORDER BY version DESC + LIMIT ?2 + "#, + params![document_id.to_string(), limit], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to list versions: {e}"), + })?; + + let mut versions = Vec::new(); + while let Some(row) = rows.next().await.map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to read version row: {e}"), + })? { + versions.push(VersionSummary { + version: get_i64(&row, 0) as i32, + content_hash: get_text(&row, 1), + created_at: get_ts(&row, 2), + changed_by: get_opt_text(&row, 3), + }); + } + Ok(versions) + } + + async fn get_latest_version_number( + &self, + document_id: Uuid, + ) -> Result, WorkspaceError> { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let mut rows = conn + .query( + "SELECT MAX(version) FROM memory_document_versions WHERE document_id = ?1", + params![document_id.to_string()], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to get latest version number: {e}"), + })?; + + if let Some(row) = rows.next().await.map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to read version number: {e}"), + })? { + // MAX returns NULL if no rows — libsql returns Null for the value + let val = row.get::(0).ok(); + match val { + Some(libsql::Value::Integer(v)) => Ok(Some(v as i32)), + _ => Ok(None), + } + } else { + Ok(None) + } + } + + async fn prune_versions( + &self, + document_id: Uuid, + keep_count: i32, + ) -> Result { + let conn = self + .connect() + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: e.to_string(), + })?; + let doc_id = document_id.to_string(); + let result = conn + .execute( + r#" + DELETE FROM memory_document_versions + WHERE document_id = ?1 + AND version NOT IN ( + SELECT version FROM memory_document_versions + WHERE document_id = ?1 + ORDER BY version DESC + LIMIT ?2 + ) + "#, + params![doc_id, keep_count as i64], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to prune versions: {e}"), + })?; + Ok(result) + } } #[cfg(test)] diff --git a/src/db/libsql_migrations.rs b/src/db/libsql_migrations.rs index d0ec20ef..2742e178 100644 --- a/src/db/libsql_migrations.rs +++ b/src/db/libsql_migrations.rs @@ -723,6 +723,25 @@ CREATE INDEX IF NOT EXISTS idx_routines_event_triggers WHERE enabled = 1 AND trigger_type IN ('event', 'system_event'); PRAGMA foreign_keys=ON; +"#, + ), + ( + 14, + "document_versions", + r#" +CREATE TABLE IF NOT EXISTS memory_document_versions ( + id TEXT PRIMARY KEY, + document_id TEXT NOT NULL REFERENCES memory_documents(id) ON DELETE CASCADE, + version INTEGER NOT NULL, + content TEXT NOT NULL, + content_hash TEXT NOT NULL, + created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')), + changed_by TEXT, + UNIQUE(document_id, version) +); + +CREATE INDEX IF NOT EXISTS idx_doc_versions_lookup + ON memory_document_versions(document_id, version DESC); "#, ), ]; diff --git a/src/db/mod.rs b/src/db/mod.rs index d89b976e..97e6942e 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -663,6 +663,67 @@ pub trait WorkspaceStore: Send + Sync { config: &SearchConfig, ) -> Result, WorkspaceError>; + // ==================== Metadata ==================== + + /// Update the metadata JSON field on a document (full replacement). + async fn update_document_metadata( + &self, + id: Uuid, + metadata: &serde_json::Value, + ) -> Result<(), WorkspaceError>; + + /// Find all `.config` documents in the workspace. + /// + /// Returns documents whose path ends with `/.config` or equals `.config`. + /// Used by the hygiene system to discover metadata-driven cleanup targets. + async fn find_config_documents( + &self, + user_id: &str, + agent_id: Option, + ) -> Result, WorkspaceError>; + + // ==================== Versioning ==================== + + /// Save the current content of a document as a new version. + /// + /// Returns the new version number (1-based, monotonically increasing). + async fn save_version( + &self, + document_id: Uuid, + content: &str, + content_hash: &str, + changed_by: Option<&str>, + ) -> Result; + + /// Get a specific version of a document. + async fn get_version( + &self, + document_id: Uuid, + version: i32, + ) -> Result; + + /// List versions of a document (newest first). + async fn list_versions( + &self, + document_id: Uuid, + limit: i64, + ) -> Result, WorkspaceError>; + + /// Get the latest version number for a document, or `None` if no versions exist. + async fn get_latest_version_number( + &self, + document_id: Uuid, + ) -> Result, WorkspaceError>; + + /// Delete old versions, keeping only the most recent `keep_count`. + /// + /// Returns the number of versions deleted. + async fn prune_versions( + &self, + document_id: Uuid, + keep_count: i32, + ) -> Result; + // ==================== Multi-scope read methods ==================== // // Default implementations loop over user_ids calling single-scope methods, diff --git a/src/db/postgres.rs b/src/db/postgres.rs index 9e5ea9ce..380afd2e 100644 --- a/src/db/postgres.rs +++ b/src/db/postgres.rs @@ -25,7 +25,8 @@ use crate::history::{ LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, Store, }; use crate::workspace::{ - MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, WorkspaceEntry, + DocumentVersion, MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, + VersionSummary, WorkspaceEntry, }; /// PostgreSQL database backend. @@ -785,4 +786,67 @@ impl WorkspaceStore for PgBackend { .list_directory_multi(user_ids, agent_id, directory) .await } + + // ==================== Metadata ==================== + + async fn update_document_metadata( + &self, + id: Uuid, + metadata: &serde_json::Value, + ) -> Result<(), WorkspaceError> { + self.repo.update_document_metadata(id, metadata).await + } + + async fn find_config_documents( + &self, + user_id: &str, + agent_id: Option, + ) -> Result, WorkspaceError> { + self.repo.find_config_documents(user_id, agent_id).await + } + + // ==================== Versioning ==================== + + async fn save_version( + &self, + document_id: Uuid, + content: &str, + content_hash: &str, + changed_by: Option<&str>, + ) -> Result { + self.repo + .save_version(document_id, content, content_hash, changed_by) + .await + } + + async fn get_version( + &self, + document_id: Uuid, + version: i32, + ) -> Result { + self.repo.get_version(document_id, version).await + } + + async fn list_versions( + &self, + document_id: Uuid, + limit: i64, + ) -> Result, WorkspaceError> { + self.repo.list_versions(document_id, limit).await + } + + async fn get_latest_version_number( + &self, + document_id: Uuid, + ) -> Result, WorkspaceError> { + self.repo.get_latest_version_number(document_id).await + } + + async fn prune_versions( + &self, + document_id: Uuid, + keep_count: i32, + ) -> Result { + self.repo.prune_versions(document_id, keep_count).await + } } diff --git a/src/error.rs b/src/error.rs index e4f1b957..07c4f6d2 100644 --- a/src/error.rs +++ b/src/error.rs @@ -315,6 +315,12 @@ pub enum WorkspaceError { #[error("Write rejected for '{path}': prompt injection detected ({reason})")] InjectionRejected { path: String, reason: String }, + + #[error("Version not found: document {document_id} version {version}")] + VersionNotFound { document_id: Uuid, version: i32 }, + + #[error("Patch failed for '{path}': {reason}")] + PatchFailed { path: String, reason: String }, } /// Orchestrator errors (internal API, container management). diff --git a/src/tools/builtin/memory.rs b/src/tools/builtin/memory.rs index 501ccf46..25ef8633 100644 --- a/src/tools/builtin/memory.rs +++ b/src/tools/builtin/memory.rs @@ -246,6 +246,23 @@ impl Tool for MemoryWriteTool { "type": "boolean", "description": "Skip privacy classification and write directly to the specified layer without redirect. Use when you're certain the content belongs in the target layer.", "default": false + }, + "metadata": { + "type": "object", + "description": "Optional metadata to set on the document (e.g., {\"skip_indexing\": true, \"hygiene\": {\"enabled\": true, \"retention_days\": 7}})" + }, + "old_string": { + "type": "string", + "description": "When present, switches to patch mode: finds and replaces this exact string in the document. Requires target to be a path (not 'memory' or 'daily_log')." + }, + "new_string": { + "type": "string", + "description": "Replacement string (required when old_string is present)." + }, + "replace_all": { + "type": "boolean", + "description": "If true, replace all occurrences of old_string. Default: false.", + "default": false } }, "required": ["content"] @@ -330,6 +347,46 @@ impl Tool for MemoryWriteTool { path => path.to_string(), }; + // Patch mode: if old_string is provided, do search-and-replace instead of write/append. + let old_string = params.get("old_string").and_then(|v| v.as_str()); + if let Some(old_str) = old_string { + let new_str = params + .get("new_string") + .and_then(|v| v.as_str()) + .ok_or_else(|| { + ToolError::InvalidParameters( + "new_string is required when old_string is provided".to_string(), + ) + })?; + let replace_all = params + .get("replace_all") + .and_then(|v| v.as_bool()) + .unwrap_or(false); + + let result = workspace + .patch(&resolved_path, old_str, new_str, replace_all) + .await + .map_err(map_write_err)?; + + // Apply metadata if provided + if let Some(meta) = params.get("metadata") + && meta.is_object() + { + workspace + .update_metadata(result.document.id, meta) + .await + .map_err(map_write_err)?; + } + + let output = serde_json::json!({ + "status": "patched", + "path": resolved_path, + "replacements": result.replacements, + "content_length": result.document.content.len(), + }); + return Ok(ToolOutput::success(output, start.elapsed())); + } + // When a layer is specified, route through layer-aware methods for ALL targets. // Otherwise, use default workspace methods (which include injection scanning). let layer_result = if let Some(layer_name) = layer { @@ -433,6 +490,18 @@ impl Tool for MemoryWriteTool { } } + // Apply metadata if provided (after write/append, works for all targets). + if let Some(meta) = params.get("metadata") + && meta.is_object() + { + // Read the document to get its ID + if let Ok(doc) = workspace.read(&resolved_path).await + && let Err(e) = workspace.update_metadata(doc.id, meta).await + { + tracing::warn!(path = %resolved_path, "failed to update metadata: {e}"); + } + } + let mut output = serde_json::json!({ "status": "written", "path": resolved_path, @@ -501,6 +570,15 @@ impl Tool for MemoryReadTool { "path": { "type": "string", "description": "Path to the file (e.g., 'MEMORY.md', 'daily/2024-01-15.md', 'projects/alpha/notes.md')" + }, + "version": { + "type": "integer", + "description": "Read a specific historical version of the document (omit for current content)" + }, + "list_versions": { + "type": "boolean", + "description": "If true, return version history instead of file content", + "default": false } }, "required": ["path"] @@ -525,11 +603,61 @@ impl Tool for MemoryReadTool { } let workspace = self.resolver.resolve(&ctx.user_id).await; + + let list_versions = params + .get("list_versions") + .and_then(|v| v.as_bool()) + .unwrap_or(false); + let version = params + .get("version") + .and_then(|v| v.as_i64()) + .map(|v| v as i32); + + // Read the document first (needed for document_id in all version operations) let doc = workspace .read(path) .await .map_err(|e| ToolError::ExecutionFailed(format!("Read failed: {}", e)))?; + // List versions mode + if list_versions { + let versions = workspace + .list_versions(doc.id, 50) + .await + .map_err(|e| ToolError::ExecutionFailed(format!("List versions failed: {}", e)))?; + + let output = serde_json::json!({ + "path": doc.path, + "versions": versions.iter().map(|v| serde_json::json!({ + "version": v.version, + "content_hash": v.content_hash, + "created_at": v.created_at.to_rfc3339(), + "changed_by": v.changed_by, + })).collect::>(), + "version_count": versions.len(), + }); + return Ok(ToolOutput::success(output, start.elapsed())); + } + + // Specific version mode + if let Some(ver) = version { + let version_doc = workspace + .get_version(doc.id, ver) + .await + .map_err(|e| ToolError::ExecutionFailed(format!("Get version failed: {}", e)))?; + + let output = serde_json::json!({ + "path": doc.path, + "version": version_doc.version, + "content": version_doc.content, + "content_hash": version_doc.content_hash, + "created_at": version_doc.created_at.to_rfc3339(), + "changed_by": version_doc.changed_by, + }); + return Ok(ToolOutput::success(output, start.elapsed())); + } + + // Normal read let output = serde_json::json!({ "path": doc.path, "content": doc.content, diff --git a/src/workspace/document.rs b/src/workspace/document.rs index b1fa176a..dafaf1e3 100644 --- a/src/workspace/document.rs +++ b/src/workspace/document.rs @@ -2,6 +2,7 @@ use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; use uuid::Uuid; /// Well-known document paths. @@ -37,6 +38,139 @@ pub mod paths { pub const ASSISTANT_DIRECTIVES: &str = "context/assistant-directives.md"; } +/// Name of the folder-level configuration document. +/// +/// A document at `{directory}/.config` carries metadata flags that apply +/// as defaults to all documents in that directory (e.g., `skip_indexing`, +/// `hygiene` settings). Individual document metadata overrides folder defaults. +pub const CONFIG_FILE_NAME: &str = ".config"; + +/// Typed overlay for the `metadata` JSON field on [`MemoryDocument`]. +/// +/// Fields use `Option` so that only explicitly set flags participate in +/// the merge chain (document metadata → folder `.config` → system defaults). +/// Unknown fields are preserved via `serde(flatten)`. +#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)] +pub struct DocumentMetadata { + /// When `true`, skip chunking and embedding for this document/folder. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub skip_indexing: Option, + + /// When `true`, skip automatic versioning for this document/folder. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub skip_versioning: Option, + + /// Hygiene (auto-cleanup) configuration for this folder. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub hygiene: Option, + + /// Preserve unknown fields for forward compatibility. + #[serde(flatten)] + pub extra: serde_json::Map, +} + +impl DocumentMetadata { + /// Parse from a raw JSON [`serde_json::Value`]. + /// + /// Returns [`Default`] if the value is not an object or cannot be parsed. + pub fn from_value(value: &serde_json::Value) -> Self { + serde_json::from_value(value.clone()).unwrap_or_default() + } + + /// Convert to a JSON [`serde_json::Value`]. + pub fn to_value(&self) -> serde_json::Value { + serde_json::to_value(self).unwrap_or(serde_json::json!({})) + } + + /// Merge two metadata values: `overlay` keys win over `base` keys. + /// + /// This is a shallow merge at the top-level keys — nested objects are + /// replaced wholesale, not recursively merged. This keeps the semantics + /// simple and predictable across both PostgreSQL and libSQL. + pub fn merge(base: &serde_json::Value, overlay: &serde_json::Value) -> serde_json::Value { + let mut merged = match base { + serde_json::Value::Object(map) => map.clone(), + _ => serde_json::Map::new(), + }; + if let serde_json::Value::Object(over) = overlay { + for (k, v) in over { + merged.insert(k.clone(), v.clone()); + } + } + serde_json::Value::Object(merged) + } +} + +/// Hygiene (auto-cleanup) settings for a folder. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct HygieneMetadata { + /// Whether this folder is a hygiene target. + pub enabled: bool, + + /// Delete documents older than this many days. + #[serde(default = "default_retention_days")] + pub retention_days: u32, +} + +fn default_retention_days() -> u32 { + 30 +} + +/// A historical version of a workspace document. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct DocumentVersion { + /// Version record ID. + pub id: Uuid, + /// Parent document ID. + pub document_id: Uuid, + /// Version number (1-based, monotonically increasing per document). + pub version: i32, + /// Full document content at this version. + pub content: String, + /// SHA-256 hash of `content` (hex-encoded, prefixed with `sha256:`). + pub content_hash: String, + /// When this version was created. + pub created_at: DateTime, + /// Who/what created this version (e.g. `"agent"`, `"user:alice"`). + pub changed_by: Option, +} + +/// Summary of a document version (without full content). +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct VersionSummary { + /// Version number. + pub version: i32, + /// SHA-256 hash of the version's content. + pub content_hash: String, + /// When this version was created. + pub created_at: DateTime, + /// Who/what created this version. + pub changed_by: Option, +} + +/// Result of a workspace patch operation. +#[derive(Debug, Clone)] +pub struct PatchResult { + /// The updated document. + pub document: MemoryDocument, + /// Number of replacements made. + pub replacements: usize, +} + +/// Compute a SHA-256 hash of content, returned as `"sha256:{hex}"`. +pub fn content_sha256(content: &str) -> String { + let mut hasher = Sha256::new(); + hasher.update(content.as_bytes()); + let result = hasher.finalize(); + format!("sha256:{:x}", result) +} + +/// Check if a path refers to a `.config` document. +pub fn is_config_path(path: &str) -> bool { + let file_name = path.rsplit('/').next().unwrap_or(path); + file_name == CONFIG_FILE_NAME +} + /// Paths treated as identity documents for multi-scope isolation. /// /// These files are always read from the primary scope only — never from @@ -360,6 +494,114 @@ mod tests { assert_eq!(result[0].updated_at, Some(ts)); } + #[test] + fn test_document_metadata_default_is_empty() { + let meta = DocumentMetadata::default(); + assert_eq!(meta.skip_indexing, None); + assert_eq!(meta.skip_versioning, None); + assert_eq!(meta.hygiene, None); + assert!(meta.extra.is_empty()); + } + + #[test] + fn test_document_metadata_from_value_full() { + let value = serde_json::json!({ + "skip_indexing": true, + "skip_versioning": false, + "hygiene": { "enabled": true, "retention_days": 7 } + }); + let meta = DocumentMetadata::from_value(&value); + assert_eq!(meta.skip_indexing, Some(true)); + assert_eq!(meta.skip_versioning, Some(false)); + let hygiene = meta.hygiene.unwrap(); + assert!(hygiene.enabled); + assert_eq!(hygiene.retention_days, 7); + } + + #[test] + fn test_document_metadata_from_value_partial() { + let value = serde_json::json!({"skip_indexing": true}); + let meta = DocumentMetadata::from_value(&value); + assert_eq!(meta.skip_indexing, Some(true)); + assert_eq!(meta.hygiene, None); + } + + #[test] + fn test_document_metadata_from_value_invalid() { + let meta = DocumentMetadata::from_value(&serde_json::json!("not an object")); + assert_eq!(meta, DocumentMetadata::default()); + } + + #[test] + fn test_document_metadata_preserves_unknown_fields() { + let value = serde_json::json!({ + "skip_indexing": true, + "custom_field": "hello" + }); + let meta = DocumentMetadata::from_value(&value); + assert_eq!(meta.skip_indexing, Some(true)); + assert_eq!(meta.extra.get("custom_field").and_then(|v| v.as_str()), Some("hello")); + + // Round-trip preserves the field + let back = meta.to_value(); + assert_eq!(back.get("custom_field").and_then(|v| v.as_str()), Some("hello")); + } + + #[test] + fn test_document_metadata_merge() { + let base = serde_json::json!({"skip_indexing": false, "hygiene": {"enabled": true, "retention_days": 30}}); + let overlay = serde_json::json!({"skip_indexing": true, "skip_versioning": true}); + let merged = DocumentMetadata::merge(&base, &overlay); + let meta = DocumentMetadata::from_value(&merged); + // Overlay wins + assert_eq!(meta.skip_indexing, Some(true)); + assert_eq!(meta.skip_versioning, Some(true)); + // Base preserved when not overridden + assert!(meta.hygiene.is_some()); + } + + #[test] + fn test_document_metadata_merge_empty_base() { + let base = serde_json::json!({}); + let overlay = serde_json::json!({"skip_indexing": true}); + let merged = DocumentMetadata::merge(&base, &overlay); + let meta = DocumentMetadata::from_value(&merged); + assert_eq!(meta.skip_indexing, Some(true)); + } + + #[test] + fn test_hygiene_metadata_default_retention() { + let value = serde_json::json!({"enabled": true}); + let hygiene: HygieneMetadata = serde_json::from_value(value).unwrap(); + assert!(hygiene.enabled); + assert_eq!(hygiene.retention_days, 30); + } + + #[test] + fn test_content_sha256_deterministic() { + let hash1 = content_sha256("hello world"); + let hash2 = content_sha256("hello world"); + assert_eq!(hash1, hash2); + assert!(hash1.starts_with("sha256:")); + } + + #[test] + fn test_content_sha256_different_content() { + let hash1 = content_sha256("hello"); + let hash2 = content_sha256("world"); + assert_ne!(hash1, hash2); + } + + #[test] + fn test_is_config_path() { + assert!(is_config_path(".config")); + assert!(is_config_path("daily/.config")); + assert!(is_config_path("frontend/widgets/.config")); + assert!(!is_config_path("daily/2024-01-15.md")); + assert!(!is_config_path("MEMORY.md")); + assert!(!is_config_path(".config.bak")); + } + #[test] fn test_merge_workspace_entries_sorted_by_path() { let entries = vec![ diff --git a/src/workspace/hygiene.rs b/src/workspace/hygiene.rs index d84d8f03..4a11d561 100644 --- a/src/workspace/hygiene.rs +++ b/src/workspace/hygiene.rs @@ -1,8 +1,10 @@ //! Memory hygiene: automatic cleanup of stale workspace documents. //! -//! Runs on a configurable cadence and deletes daily log entries and conversation -//! documents older than their respective retention periods. Identity files -//! (`IDENTITY.md`, `SOUL.md`, etc.) are never touched. +//! Runs on a configurable cadence and discovers which directories have hygiene +//! enabled by reading `.config` metadata documents. This is a **metadata-driven** +//! approach: instead of hardcoding `daily/` and `conversations/`, the system +//! respects `hygiene.enabled` and `hygiene.retention_days` set on each folder's +//! `.config` document. //! //! A global [`AtomicBool`] guard prevents concurrent hygiene passes, which //! avoids TOCTOU races on the state file and Windows file-locking errors @@ -10,18 +12,16 @@ //! pass completes. //! //! ```text -//! ┌─────────────────────────────────────────────┐ -//! │ Hygiene Pass │ -//! │ │ -//! │ 0. Acquire RUNNING guard (skip if held) │ -//! │ 1. Check cadence (skip if ran recently) │ -//! │ 2. Save state (claim the cadence window) │ -//! │ 3. List daily/ documents │ -//! │ 4. Delete those older than daily_retention │ -//! │ 5. List conversations/ documents │ -//! │ 6. Delete those older than conversation_ret │ -//! │ 7. Log summary │ -//! └─────────────────────────────────────────────┘ +//! ┌──────────────────────────────────────────────────┐ +//! │ Hygiene Pass │ +//! │ │ +//! │ 0. Acquire RUNNING guard (skip if held) │ +//! │ 1. Check cadence (skip if ran recently) │ +//! │ 2. Save state (claim the cadence window) │ +//! │ 3. Discover .config docs with hygiene.enabled │ +//! │ 4. For each: cleanup_directory(parent, retention)│ +//! │ 5. Log summary │ +//! └──────────────────────────────────────────────────┘ //! ``` use std::path::PathBuf; @@ -31,46 +31,22 @@ use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; use crate::bootstrap::ironclaw_base_dir; -use crate::workspace::Workspace; +use crate::workspace::{DocumentMetadata, Workspace, is_config_path}; /// Global guard preventing concurrent hygiene passes. static RUNNING: AtomicBool = AtomicBool::new(false); -/// Paths that must never be deleted by hygiene, regardless of age. -const IDENTITY_PATHS: &[&str] = &[ - crate::workspace::document::paths::MEMORY, - crate::workspace::document::paths::IDENTITY, - crate::workspace::document::paths::SOUL, - crate::workspace::document::paths::AGENTS, - crate::workspace::document::paths::USER, - crate::workspace::document::paths::HEARTBEAT, - crate::workspace::document::paths::README, - crate::workspace::document::paths::TOOLS, - crate::workspace::document::paths::BOOTSTRAP, -]; - -/// Check if a document path is an identity document that must never be deleted. -/// -/// Performs case-insensitive comparison to handle case-insensitive filesystems -/// (Windows, macOS) and prevent accidental deletion of identity docs with -/// different casing (e.g., memory.md, MEMORY.MD, Memory.md). -fn is_identity_path(path: &str) -> bool { - let file_name = path.rsplit('/').next().unwrap_or(path); - let file_name_lower = file_name.to_lowercase(); - IDENTITY_PATHS - .iter() - .any(|&p| p.to_lowercase() == file_name_lower) -} - /// Configuration for workspace hygiene. #[derive(Debug, Clone)] pub struct HygieneConfig { /// Whether hygiene is enabled at all. pub enabled: bool, - /// Documents in `daily/` older than this many days are deleted. - pub daily_retention_days: u32, - /// Documents in `conversations/` older than this many days are deleted. - pub conversation_retention_days: u32, + /// Maximum number of versions to keep per document. + /// + /// TODO: Wire up global version pruning once per-document iteration + /// is efficient (e.g., via a dedicated DB query). For now this field + /// is stored in config but not actively enforced during hygiene passes. + pub version_keep_count: u32, /// Minimum hours between hygiene passes. pub cadence_hours: u32, /// Directory to store state file (default: `~/.ironclaw`). @@ -81,8 +57,7 @@ impl Default for HygieneConfig { fn default() -> Self { Self { enabled: true, - daily_retention_days: 30, - conversation_retention_days: 7, + version_keep_count: 50, cadence_hours: 12, state_dir: ironclaw_base_dir(), } @@ -98,10 +73,10 @@ struct HygieneState { /// Summary of what a hygiene pass cleaned up. #[derive(Debug, Default)] pub struct HygieneReport { - /// Number of daily log documents deleted. - pub daily_logs_deleted: u32, - /// Number of conversation documents deleted. - pub conversation_docs_deleted: u32, + /// Per-directory cleanup results: `(directory_path, deleted_count)`. + pub directories_cleaned: Vec<(String, u32)>, + /// Number of document versions pruned across all documents. + pub versions_pruned: u64, /// Whether the run was skipped (cadence not yet elapsed). pub skipped: bool, } @@ -109,7 +84,7 @@ pub struct HygieneReport { impl HygieneReport { /// True if any cleanup work was done. pub fn had_work(&self) -> bool { - self.daily_logs_deleted > 0 || self.conversation_docs_deleted > 0 + self.directories_cleaned.iter().any(|(_, n)| *n > 0) || self.versions_pruned > 0 } } @@ -168,30 +143,55 @@ pub async fn run_if_due(workspace: &Workspace, config: &HygieneConfig) -> Hygien // TOCTOU races where another task reads stale state. save_state(&state_file); - tracing::info!( - daily_retention_days = config.daily_retention_days, - conversation_retention_days = config.conversation_retention_days, - "memory hygiene: starting cleanup pass" - ); + tracing::info!("memory hygiene: starting cleanup pass"); let mut report = HygieneReport::default(); - // Delete old daily logs - match cleanup_daily_logs(workspace, config.daily_retention_days).await { - Ok(count) => report.daily_logs_deleted = count, - Err(e) => tracing::warn!("memory hygiene: failed to clean daily logs: {e}"), - } + // Discover directories that have hygiene enabled via .config metadata. + let config_docs = match workspace.find_config_documents().await { + Ok(docs) => docs, + Err(e) => { + tracing::warn!("memory hygiene: failed to discover .config documents: {e}"); + return report; + } + }; - // Delete old conversation documents - match cleanup_conversation_docs(workspace, config.conversation_retention_days).await { - Ok(count) => report.conversation_docs_deleted = count, - Err(e) => tracing::warn!("memory hygiene: failed to clean conversation docs: {e}"), + for doc in &config_docs { + let meta = DocumentMetadata::from_value(&doc.metadata); + let Some(hygiene) = meta.hygiene else { + continue; + }; + if !hygiene.enabled { + continue; + } + + // Derive the parent directory from the .config path. + let directory = match doc.path.rsplit_once('/') { + Some((dir, _)) => format!("{dir}/"), + None => continue, // root-level .config — skip + }; + + match cleanup_directory(workspace, &directory, hygiene.retention_days).await { + Ok(deleted) => { + if deleted > 0 { + tracing::info!( + directory, + deleted, + "memory hygiene: cleaned directory" + ); + } + report.directories_cleaned.push((directory, deleted)); + } + Err(e) => { + tracing::warn!(directory, "memory hygiene: failed to clean directory: {e}"); + } + } } if report.had_work() { tracing::info!( - daily_logs_deleted = report.daily_logs_deleted, - conversation_docs_deleted = report.conversation_docs_deleted, + directories_cleaned = ?report.directories_cleaned, + versions_pruned = report.versions_pruned, "memory hygiene: cleanup complete" ); } else { @@ -210,88 +210,41 @@ impl Drop for RunningGuard { } } -/// Delete daily log documents older than `retention_days`. -async fn cleanup_daily_logs( +/// Delete documents in `directory` that are older than `retention_days`. +/// +/// Skips directories and `.config` files (which must never be deleted by +/// hygiene). Returns the number of documents deleted. +async fn cleanup_directory( workspace: &Workspace, + directory: &str, retention_days: u32, ) -> Result { let cutoff = Utc::now() - chrono::Duration::days(i64::from(retention_days)); - let entries = workspace.list("daily/").await?; - + let entries = workspace.list(directory).await?; let mut deleted = 0u32; for entry in entries { if entry.is_directory { continue; } - - // Never delete identity documents - if is_identity_path(&entry.path) { + if is_config_path(&entry.path) { continue; } - - // Check if the document is old enough to delete if let Some(updated_at) = entry.updated_at && updated_at < cutoff { - let path = if entry.path.starts_with("daily/") { + let path = if entry.path.starts_with(directory) { entry.path.clone() } else { - format!("daily/{}", entry.path) + format!("{}{}", directory, entry.path) }; - if let Err(e) = workspace.delete(&path).await { tracing::warn!(path, "memory hygiene: failed to delete: {e}"); } else { - tracing::debug!(path, "memory hygiene: deleted old daily log"); + tracing::debug!(path, "memory hygiene: deleted stale document"); deleted += 1; } } } - - Ok(deleted) -} - -/// Delete conversation documents older than `retention_days`. -async fn cleanup_conversation_docs( - workspace: &Workspace, - retention_days: u32, -) -> Result { - let cutoff = Utc::now() - chrono::Duration::days(i64::from(retention_days)); - let entries = workspace.list("conversations/").await?; - - let mut deleted = 0u32; - for entry in entries { - if entry.is_directory { - continue; - } - - // Never delete identity documents - if is_identity_path(&entry.path) { - continue; - } - - // Check if the document is old enough to delete - if let Some(updated_at) = entry.updated_at - && updated_at < cutoff - { - let path = if entry.path.starts_with("conversations/") { - entry.path.clone() - } else { - format!("conversations/{}", entry.path) - }; - - if let Err(e) = workspace.delete(&path).await { - tracing::warn!( - path, - "memory hygiene: failed to delete conversation doc: {e}" - ); - } else { - tracing::debug!(path, "memory hygiene: deleted old conversation doc"); - deleted += 1; - } - } - } - Ok(deleted) } @@ -349,8 +302,7 @@ mod tests { fn default_config_is_reasonable() { let cfg = HygieneConfig::default(); assert!(cfg.enabled); - assert_eq!(cfg.daily_retention_days, 30); - assert_eq!(cfg.conversation_retention_days, 7); + assert_eq!(cfg.version_keep_count, 50); assert_eq!(cfg.cadence_hours, 12); } @@ -362,84 +314,33 @@ mod tests { } #[test] - fn report_had_work_when_deleted() { + fn report_had_work_when_directories_cleaned() { let report = HygieneReport { - daily_logs_deleted: 3, - conversation_docs_deleted: 0, + directories_cleaned: vec![("daily/".to_string(), 3)], + versions_pruned: 0, skipped: false, }; assert!(report.had_work()); } #[test] - fn report_had_work_when_conversation_deleted() { + fn report_had_work_when_versions_pruned() { let report = HygieneReport { - daily_logs_deleted: 0, - conversation_docs_deleted: 2, + directories_cleaned: vec![], + versions_pruned: 5, skipped: false, }; assert!(report.had_work()); } #[test] - fn is_identity_path_excludes_sacred_docs() { - for name in [ - "MEMORY.md", - "IDENTITY.md", - "SOUL.md", - "AGENTS.md", - "USER.md", - "HEARTBEAT.md", - "README.md", - "TOOLS.md", - "BOOTSTRAP.md", - ] { - assert!(is_identity_path(name), "{name} should be excluded"); - assert!( - is_identity_path(&format!("conversations/{name}")), - "conversations/{name} should be excluded via path" - ); - } - } - - #[test] - fn is_identity_path_case_insensitive() { - // Verify case-insensitive matching for case-insensitive filesystems - assert!( - is_identity_path("memory.md"), - "lowercase memory.md should be excluded" - ); - assert!( - is_identity_path("Memory.md"), - "mixed case Memory.md should be excluded" - ); - assert!( - is_identity_path("MEMORY.MD"), - "uppercase MEMORY.MD should be excluded" - ); - assert!( - is_identity_path("identity.md"), - "lowercase identity.md should be excluded" - ); - assert!( - is_identity_path("conversations/soul.md"), - "conversations/soul.md should be excluded" - ); - assert!( - is_identity_path("conversations/SOUL.MD"), - "conversations/SOUL.MD should be excluded" - ); - } - - #[test] - fn is_identity_path_allows_normal_docs() { - for path in [ - "daily/2024-01-01.md", - "conversations/chat-abc.md", - "notes.md", - ] { - assert!(!is_identity_path(path), "{path} should not be excluded"); - } + fn report_no_work_when_zero_deletions() { + let report = HygieneReport { + directories_cleaned: vec![("daily/".to_string(), 0)], + versions_pruned: 0, + skipped: false, + }; + assert!(!report.had_work()); } #[test] @@ -552,61 +453,111 @@ mod tests { Arc::new(Workspace::new_with_db("default", db.clone())) } - #[tokio::test] - async fn cleanup_daily_logs_preserves_identity_documents() { - let (db, _tmp) = create_test_db().await; - let ws = create_workspace(&db); - - // Write several regular documents (non-identity) - ws.write("daily/2024-01-15.md", "Old log") + /// Helper to seed a .config document with hygiene metadata on a directory. + async fn seed_hygiene_config( + workspace: &Workspace, + directory: &str, + retention_days: u32, + ) { + let config_path = format!("{}.config", directory); + // Create the .config document with empty content + workspace + .write(&config_path, "") .await - .expect("write log 1"); - ws.write("daily/2024-01-20.md", "Another log") + .expect("write .config"); + // Read back to get the document ID + let doc = workspace + .read(&config_path) .await - .expect("write log 2"); - - // Write an identity document - ws.write("MEMORY.md", "Long-term curated memory") + .expect("read .config doc"); + // Set hygiene metadata + workspace + .update_metadata( + doc.id, + &serde_json::json!({ + "hygiene": {"enabled": true, "retention_days": retention_days}, + "skip_versioning": true + }), + ) .await - .expect("write identity"); - - // List before cleanup - let before = ws.list("daily/").await.expect("list before"); - let daily_count_before = before.iter().filter(|e| !e.is_directory).count(); - assert!(daily_count_before >= 2, "should have at least 2 daily logs"); - - // Run cleanup with 0-day retention (deletes everything old) - // This tests that even with aggressive cleanup, identity docs survive - let deleted = cleanup_daily_logs(&ws, 0) - .await - .expect("cleanup_daily_logs"); - - // Should have deleted some documents (the daily logs) - assert!(deleted > 0, "should have deleted old daily documents"); - - // Verify identity doc still exists - let identity = db - .get_document_by_path("default", None, "MEMORY.md") - .await - .expect("get identity doc"); - assert_eq!(identity.path, "MEMORY.md"); - assert_eq!(identity.content, "Long-term curated memory"); + .expect("set metadata"); } #[tokio::test] - async fn cleanup_conversation_docs_handles_empty_directory() { + async fn cleanup_directory_skips_config_files() { let (db, _tmp) = create_test_db().await; let ws = create_workspace(&db); - // Run cleanup on an empty directory (conversations/ doesn't exist) - let deleted = cleanup_conversation_docs(&ws, 7) + // Write documents including a .config + ws.write("daily/2024-01-15.md", "Old log") .await - .expect("cleanup_conversation_docs"); + .expect("write log"); + ws.write("daily/.config", "") + .await + .expect("write config"); + + // Run cleanup with 0-day retention (deletes everything old) + let deleted = cleanup_directory(&ws, "daily/", 0) + .await + .expect("cleanup_directory"); + + // Should have deleted the log but not the .config + assert!(deleted > 0, "should have deleted old daily documents"); + + // Verify .config still exists + let config_doc = db + .get_document_by_path("default", None, "daily/.config") + .await + .expect("get .config doc"); + assert_eq!(config_doc.path, "daily/.config"); + } + + #[tokio::test] + async fn cleanup_directory_handles_empty_directory() { + let (db, _tmp) = create_test_db().await; + let ws = create_workspace(&db); + + // Run cleanup on an empty directory + let deleted = cleanup_directory(&ws, "conversations/", 7) + .await + .expect("cleanup_directory"); - // Should delete 0 (nothing to delete) assert_eq!(deleted, 0, "should delete 0 from empty directory"); } + #[tokio::test] + async fn metadata_driven_cleanup_discovers_directories() { + let (db, _tmp) = create_test_db().await; + let ws = create_workspace(&db); + + // Seed .config with hygiene enabled on daily/ + seed_hygiene_config(&ws, "daily/", 0).await; + + // Write some documents + ws.write("daily/log1.md", "content 1") + .await + .expect("write doc 1"); + ws.write("daily/log2.md", "content 2") + .await + .expect("write doc 2"); + + let config = HygieneConfig { + enabled: true, + version_keep_count: 50, + cadence_hours: 12, + state_dir: _tmp.path().to_path_buf(), + }; + + // First run should discover daily/ and clean it + let report = run_if_due(&ws, &config).await; + assert!(!report.skipped, "first run should not be skipped"); + assert!(report.had_work(), "should have cleaned documents"); + assert!( + !report.directories_cleaned.is_empty(), + "should have at least one directory cleaned" + ); + } + #[tokio::test] async fn cleanup_respects_cadence_prevents_concurrent_runs() { let (db, _tmp) = create_test_db().await; @@ -614,8 +565,7 @@ mod tests { let config = HygieneConfig { enabled: true, - daily_retention_days: 30, - conversation_retention_days: 7, + version_keep_count: 50, cadence_hours: 12, state_dir: _tmp.path().to_path_buf(), }; @@ -627,13 +577,6 @@ mod tests { // Second run immediately should be skipped (cadence not elapsed) let report2 = run_if_due(&ws, &config).await; assert!(report2.skipped, "second run should be skipped by cadence"); - - // Report structure should be correct - assert_eq!( - report1.daily_logs_deleted + report1.conversation_docs_deleted, - 0, - "first run should have clean counts" - ); } #[tokio::test] @@ -641,6 +584,10 @@ mod tests { let (db, _tmp) = create_test_db().await; let ws = create_workspace(&db); + // Seed hygiene on both directories + seed_hygiene_config(&ws, "daily/", 0).await; + seed_hygiene_config(&ws, "conversations/", 0).await; + // Write some documents ws.write("daily/log1.md", "content 1") .await @@ -652,35 +599,34 @@ mod tests { .await .expect("write doc 3"); - // Run with 0-day retention to delete everything non-identity - let deleted_daily = cleanup_daily_logs(&ws, 0).await.expect("cleanup daily"); - let deleted_conv = cleanup_conversation_docs(&ws, 0) + // Run with 0-day retention via direct cleanup_directory calls + let deleted_daily = cleanup_directory(&ws, "daily/", 0) + .await + .expect("cleanup daily"); + let deleted_conv = cleanup_directory(&ws, "conversations/", 0) .await .expect("cleanup conversations"); - // Both should report deletions assert!(deleted_daily > 0, "should report deleted daily logs"); assert_eq!(deleted_conv, 1, "should report 1 deleted conversation doc"); - // Create a HygieneReport and verify aggregation works + // Verify HygieneReport aggregation let report = HygieneReport { - daily_logs_deleted: deleted_daily, - conversation_docs_deleted: deleted_conv, + directories_cleaned: vec![ + ("daily/".to_string(), deleted_daily), + ("conversations/".to_string(), deleted_conv), + ], + versions_pruned: 0, skipped: false, }; - // Verify HygieneReport structure assert!(!report.skipped, "should not be skipped"); assert!(report.had_work(), "report should indicate work was done"); - assert!( - report.daily_logs_deleted > 0 || report.conversation_docs_deleted > 0, - "report should have at least one deletion count > 0" - ); - // Verify had_work() correctly combines both counts + // Verify had_work() correctly checks directory counts let no_work = HygieneReport { - daily_logs_deleted: 0, - conversation_docs_deleted: 0, + directories_cleaned: vec![], + versions_pruned: 0, skipped: false, }; assert!(!no_work.had_work(), "empty report should indicate no work"); diff --git a/src/workspace/mod.rs b/src/workspace/mod.rs index 51d7d2fc..8ffcd369 100644 --- a/src/workspace/mod.rs +++ b/src/workspace/mod.rs @@ -53,8 +53,9 @@ mod search; pub use chunker::{ChunkConfig, chunk_document}; pub use document::{ - IDENTITY_PATHS, MemoryChunk, MemoryDocument, WorkspaceEntry, is_identity_path, - merge_workspace_entries, paths, + CONFIG_FILE_NAME, DocumentMetadata, DocumentVersion, HygieneMetadata, IDENTITY_PATHS, + MemoryChunk, MemoryDocument, PatchResult, VersionSummary, WorkspaceEntry, content_sha256, + is_config_path, is_identity_path, merge_workspace_entries, paths, }; pub use embedding_cache::{CachedEmbeddingProvider, EmbeddingCacheConfig}; pub use embeddings::{ @@ -366,6 +367,101 @@ impl WorkspaceStorage { } } } + + // ==================== Metadata ==================== + + async fn update_document_metadata( + &self, + id: Uuid, + metadata: &serde_json::Value, + ) -> Result<(), WorkspaceError> { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => repo.update_document_metadata(id, metadata).await, + Self::Db(db) => db.update_document_metadata(id, metadata).await, + } + } + + async fn find_config_documents( + &self, + user_id: &str, + agent_id: Option, + ) -> Result, WorkspaceError> { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => repo.find_config_documents(user_id, agent_id).await, + Self::Db(db) => db.find_config_documents(user_id, agent_id).await, + } + } + + // ==================== Versioning ==================== + + async fn save_version( + &self, + document_id: Uuid, + content: &str, + content_hash: &str, + changed_by: Option<&str>, + ) -> Result { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => { + repo.save_version(document_id, content, content_hash, changed_by) + .await + } + Self::Db(db) => { + db.save_version(document_id, content, content_hash, changed_by) + .await + } + } + } + + async fn get_version( + &self, + document_id: Uuid, + version: i32, + ) -> Result { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => repo.get_version(document_id, version).await, + Self::Db(db) => db.get_version(document_id, version).await, + } + } + + async fn list_versions( + &self, + document_id: Uuid, + limit: i64, + ) -> Result, WorkspaceError> { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => repo.list_versions(document_id, limit).await, + Self::Db(db) => db.list_versions(document_id, limit).await, + } + } + + async fn get_latest_version_number( + &self, + document_id: Uuid, + ) -> Result, WorkspaceError> { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => repo.get_latest_version_number(document_id).await, + Self::Db(db) => db.get_latest_version_number(document_id).await, + } + } + + async fn prune_versions( + &self, + document_id: Uuid, + keep_count: i32, + ) -> Result { + match self { + #[cfg(feature = "postgres")] + Self::Repo(repo) => repo.prune_versions(document_id, keep_count).await, + Self::Db(db) => db.prune_versions(document_id, keep_count).await, + } + } } /// Default template seeded into HEARTBEAT.md on first access. @@ -696,10 +792,199 @@ impl Workspace { .await } + // ==================== Metadata ==================== + + /// Update the metadata JSON on a document by ID (full replacement). + pub async fn update_metadata( + &self, + id: Uuid, + metadata: &serde_json::Value, + ) -> Result<(), WorkspaceError> { + self.storage.update_document_metadata(id, metadata).await + } + + /// Prune old versions for a document, keeping only the most recent `keep_count`. + /// + /// Returns the number of versions deleted. + pub async fn prune_versions( + &self, + document_id: Uuid, + keep_count: i32, + ) -> Result { + self.storage.prune_versions(document_id, keep_count).await + } + + /// Find all `.config` documents in this workspace scope. + pub async fn find_config_documents(&self) -> Result, WorkspaceError> { + self.storage + .find_config_documents(&self.user_id, self.agent_id) + .await + } + + /// Resolve effective metadata for a document path. + /// + /// Resolution chain: document's own metadata → nearest ancestor `.config` → defaults. + pub async fn resolve_metadata(&self, path: &str) -> DocumentMetadata { + // 1. Document's own metadata + let doc_meta = self + .storage + .get_document_by_path(&self.user_id, self.agent_id, path) + .await + .ok() + .map(|d| d.metadata); + + // 2. Walk up parent directories looking for .config + let mut config_meta = None; + let normalized = normalize_path(path); + let mut current = normalized.as_str(); + while let Some(slash_pos) = current.rfind('/') { + let parent = ¤t[..slash_pos]; + let config_path = format!("{}/{CONFIG_FILE_NAME}", parent); + if let Ok(doc) = self + .storage + .get_document_by_path(&self.user_id, self.agent_id, &config_path) + .await + { + config_meta = Some(doc.metadata); + break; + } + current = parent; + } + // Also check root-level .config + if config_meta.is_none() + && let Ok(doc) = self + .storage + .get_document_by_path(&self.user_id, self.agent_id, CONFIG_FILE_NAME) + .await + { + config_meta = Some(doc.metadata); + } + + // 3. Merge: config as base, document metadata as overlay + let base = config_meta.unwrap_or(serde_json::json!({})); + let overlay = doc_meta.unwrap_or(serde_json::json!({})); + let merged = DocumentMetadata::merge(&base, &overlay); + DocumentMetadata::from_value(&merged) + } + + // ==================== Versioning ==================== + + /// List versions of a document (newest first). + pub async fn list_versions( + &self, + document_id: Uuid, + limit: i64, + ) -> Result, WorkspaceError> { + self.storage.list_versions(document_id, limit).await + } + + /// Get a specific version of a document. + pub async fn get_version( + &self, + document_id: Uuid, + version: i32, + ) -> Result { + self.storage.get_version(document_id, version).await + } + + /// Save the current content as a version if it differs from the latest. + /// + /// Returns the new version number, or `None` if skipped (empty content, + /// identical hash, or versioning disabled via metadata). + async fn maybe_save_version( + &self, + document_id: Uuid, + current_content: &str, + path: &str, + changed_by: Option<&str>, + ) -> Result, WorkspaceError> { + // Don't version empty documents + if current_content.is_empty() { + return Ok(None); + } + + // Check metadata for skip_versioning flag + let metadata = self.resolve_metadata(path).await; + if metadata.skip_versioning == Some(true) { + return Ok(None); + } + + let hash = content_sha256(current_content); + + // Check if latest version already has this hash (skip duplicate saves) + if let Ok(Some(latest)) = self.storage.get_latest_version_number(document_id).await + && let Ok(ver) = self.storage.get_version(document_id, latest).await + && ver.content_hash == hash + { + return Ok(None); + } + + let version = self + .storage + .save_version(document_id, current_content, &hash, changed_by) + .await?; + Ok(Some(version)) + } + + // ==================== Patch ==================== + + /// Apply a search-and-replace patch to a workspace document. + /// + /// Finds `old_string` in the document and replaces it with `new_string`. + /// If `replace_all` is true, replaces all occurrences; otherwise only the first. + /// Auto-versions before applying the patch. + pub async fn patch( + &self, + path: &str, + old_string: &str, + new_string: &str, + replace_all: bool, + ) -> Result { + let path = normalize_path(path); + let doc = self + .storage + .get_document_by_path(&self.user_id, self.agent_id, &path) + .await?; + + if !doc.content.contains(old_string) { + return Err(WorkspaceError::PatchFailed { + path, + reason: "old_string not found in document".to_string(), + }); + } + + let (new_content, count) = if replace_all { + let count = doc.content.matches(old_string).count(); + (doc.content.replace(old_string, new_string), count) + } else { + (doc.content.replacen(old_string, new_string, 1), 1) + }; + + // Injection scan for system prompt files + if is_system_prompt_file(&path) && !new_content.is_empty() { + reject_if_injected(&path, &new_content)?; + } + + // Auto-version before updating + let _ = self + .maybe_save_version(doc.id, &doc.content, &path, None) + .await; + + self.storage.update_document(doc.id, &new_content).await?; + self.reindex_document(doc.id).await?; + + let updated = self.storage.get_document_by_id(doc.id).await?; + Ok(PatchResult { + document: updated, + replacements: count, + }) + } + /// Write (create or update) a file. /// /// Creates parent directories implicitly (they're virtual in the DB). /// Re-indexes the document for search after writing. + /// Auto-versions the previous content before overwriting. /// /// # Example /// ```ignore @@ -715,6 +1000,12 @@ impl Workspace { .storage .get_or_create_document_by_path(&self.user_id, self.agent_id, &path) .await?; + + // Auto-version previous content before overwriting + let _ = self + .maybe_save_version(doc.id, &doc.content, &path, None) + .await; + self.storage.update_document(doc.id, content).await?; self.reindex_document(doc.id).await?; @@ -754,6 +1045,11 @@ impl Workspace { reject_if_injected(&path, &new_content)?; } + // Auto-version previous content before appending + let _ = self + .maybe_save_version(doc.id, &doc.content, &path, None) + .await; + self.storage.update_document(doc.id, &new_content).await?; self.reindex_document(doc.id).await?; Ok(()) @@ -1585,6 +1881,14 @@ impl Workspace { // Get the document let doc = self.storage.get_document_by_id(document_id).await?; + // Check metadata for skip_indexing flag + let metadata = self.resolve_metadata(&doc.path).await; + if metadata.skip_indexing == Some(true) { + // Delete any existing chunks and skip indexing + self.storage.delete_chunks(document_id).await?; + return Ok(()); + } + // Chunk the content let chunks = chunk_document(&doc.content, ChunkConfig::default()); @@ -1670,6 +1974,51 @@ impl Workspace { } } + // Seed folder-level .config documents for hygiene defaults. + let config_seeds: &[(&str, serde_json::Value)] = &[ + ( + "daily/.config", + serde_json::json!({ + "hygiene": {"enabled": true, "retention_days": 30}, + "skip_versioning": true + }), + ), + ( + "conversations/.config", + serde_json::json!({ + "hygiene": {"enabled": true, "retention_days": 7}, + "skip_versioning": true + }), + ), + ]; + + for (config_path, metadata_value) in config_seeds { + match self.read_primary(config_path).await { + Ok(_) => continue, // Already exists, don't overwrite + Err(WorkspaceError::DocumentNotFound { .. }) => {} + Err(e) => { + tracing::debug!("Failed to check {}: {}", config_path, e); + continue; + } + } + // Create empty document with metadata + if let Ok(doc) = self + .storage + .get_or_create_document_by_path(&self.user_id, self.agent_id, config_path) + .await + { + if let Err(e) = self + .storage + .update_document_metadata(doc.id, metadata_value) + .await + { + tracing::debug!("Failed to set metadata on {}: {}", config_path, e); + } else { + count += 1; + } + } + } + // BOOTSTRAP.md is only seeded on truly fresh workspaces (no identity // files existed before seeding) AND when no profile exists yet (the user // may already have a profile from a previous install and doesn't need diff --git a/src/workspace/repository.rs b/src/workspace/repository.rs index 13f6816b..0622d876 100644 --- a/src/workspace/repository.rs +++ b/src/workspace/repository.rs @@ -11,7 +11,9 @@ use uuid::Uuid; use crate::error::WorkspaceError; -use crate::workspace::document::{MemoryChunk, MemoryDocument, WorkspaceEntry}; +use crate::workspace::document::{ + DocumentVersion, MemoryChunk, MemoryDocument, VersionSummary, WorkspaceEntry, +}; use crate::workspace::search::{RankedResult, SearchConfig, SearchResult, fuse_results}; /// Database repository for workspace operations. @@ -702,4 +704,192 @@ impl Repository { } Ok(crate::workspace::merge_workspace_entries(all_entries)) } + + // ==================== Metadata ==================== + + pub async fn update_document_metadata( + &self, + id: Uuid, + metadata: &serde_json::Value, + ) -> Result<(), WorkspaceError> { + let conn = self.conn().await?; + conn.execute( + "UPDATE memory_documents SET metadata = $2, updated_at = NOW() WHERE id = $1", + &[&id, &metadata], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to update metadata: {e}"), + })?; + Ok(()) + } + + pub async fn find_config_documents( + &self, + user_id: &str, + agent_id: Option, + ) -> Result, WorkspaceError> { + let conn = self.conn().await?; + let rows = conn + .query( + r#" + SELECT id, user_id, agent_id, path, content, + created_at, updated_at, metadata + FROM memory_documents + WHERE user_id = $1 AND agent_id IS NOT DISTINCT FROM $2 + AND (path LIKE '%/.config' OR path = '.config') + ORDER BY path + "#, + &[&user_id, &agent_id], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to find config documents: {e}"), + })?; + Ok(rows.iter().map(|r| self.row_to_document(r)).collect()) + } + + // ==================== Versioning ==================== + + pub async fn save_version( + &self, + document_id: Uuid, + content: &str, + content_hash: &str, + changed_by: Option<&str>, + ) -> Result { + let conn = self.conn().await?; + let row = conn + .query_one( + r#" + INSERT INTO memory_document_versions + (id, document_id, version, content, content_hash, changed_by) + VALUES ( + gen_random_uuid(), + $1, + (SELECT COALESCE(MAX(version), 0) + 1 + FROM memory_document_versions WHERE document_id = $1), + $2, $3, $4 + ) + RETURNING version + "#, + &[&document_id, &content, &content_hash, &changed_by], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to save version: {e}"), + })?; + Ok(row.get(0)) + } + + pub async fn get_version( + &self, + document_id: Uuid, + version: i32, + ) -> Result { + let conn = self.conn().await?; + let row = conn + .query_opt( + r#" + SELECT id, document_id, version, content, content_hash, + created_at, changed_by + FROM memory_document_versions + WHERE document_id = $1 AND version = $2 + "#, + &[&document_id, &version], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to get version: {e}"), + })? + .ok_or(WorkspaceError::VersionNotFound { + document_id, + version, + })?; + Ok(DocumentVersion { + id: row.get(0), + document_id: row.get(1), + version: row.get(2), + content: row.get(3), + content_hash: row.get(4), + created_at: row.get(5), + changed_by: row.get(6), + }) + } + + pub async fn list_versions( + &self, + document_id: Uuid, + limit: i64, + ) -> Result, WorkspaceError> { + let conn = self.conn().await?; + let rows = conn + .query( + r#" + SELECT version, content_hash, created_at, changed_by + FROM memory_document_versions + WHERE document_id = $1 + ORDER BY version DESC + LIMIT $2 + "#, + &[&document_id, &limit], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to list versions: {e}"), + })?; + Ok(rows + .iter() + .map(|row| VersionSummary { + version: row.get(0), + content_hash: row.get(1), + created_at: row.get(2), + changed_by: row.get(3), + }) + .collect()) + } + + pub async fn get_latest_version_number( + &self, + document_id: Uuid, + ) -> Result, WorkspaceError> { + let conn = self.conn().await?; + let row = conn + .query_one( + "SELECT MAX(version) FROM memory_document_versions WHERE document_id = $1", + &[&document_id], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to get latest version number: {e}"), + })?; + Ok(row.get(0)) + } + + pub async fn prune_versions( + &self, + document_id: Uuid, + keep_count: i32, + ) -> Result { + let conn = self.conn().await?; + let result = conn + .execute( + r#" + DELETE FROM memory_document_versions + WHERE document_id = $1 + AND version NOT IN ( + SELECT version FROM memory_document_versions + WHERE document_id = $1 + ORDER BY version DESC + LIMIT $2 + ) + "#, + &[&document_id, &(keep_count as i64)], + ) + .await + .map_err(|e| WorkspaceError::SearchFailed { + reason: format!("Failed to prune versions: {e}"), + })?; + Ok(result) + } } diff --git a/tests/e2e_routine_heartbeat.rs b/tests/e2e_routine_heartbeat.rs index 6849ee05..daf6185f 100644 --- a/tests/e2e_routine_heartbeat.rs +++ b/tests/e2e_routine_heartbeat.rs @@ -955,8 +955,7 @@ mod tests { let hygiene_config = HygieneConfig { enabled: false, - daily_retention_days: 30, - conversation_retention_days: 7, + version_keep_count: 50, cadence_hours: 24, state_dir: _tmp.path().to_path_buf(), }; @@ -1002,8 +1001,7 @@ mod tests { let hygiene_config = HygieneConfig { enabled: false, - daily_retention_days: 30, - conversation_retention_days: 7, + version_keep_count: 50, cadence_hours: 24, state_dir: _tmp.path().to_path_buf(), };