Compare commits

..
Author SHA1 Message Date
Claude f61a0b747f fix(clippy): use any() instead of find().is_none() in user stats test
https://claude.ai/code/session_01Nm95eCjdrwxDjwkHTRieZs
2026-03-28 19:21:15 +00:00
ZakiandClaude 65c374f53c fix(routines): broaden strip_html_tags to cover all HTML forms
The whitelist-based regex missed self-closing tags without whitespace
(<br/>, <img/>), HTML comments (<!--...-->), SVG/MathML tags, and
custom elements (<custom-element>). This weakened the HTML stripping
guarantee for untrusted routine/job summaries in notifications.

Changes:
- Add separate regex for HTML comments (<!--...-->)
- Add SVG tags (svg, path, circle, etc.) and MathML tags (math, mrow,
  etc.) to the known tag list
- Add regex for custom elements (tags containing hyphens per web
  components spec)
- Fix self-closing tag matching to handle <br/> without whitespace
  by making the whitespace before /> optional
- Add regression tests for all four cases plus generics preservation
- Fix pre-existing compilation error in tunnel/mod.rs test helpers
  (missing GatewayConfig fields)

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 19:16:32 +00:00
Claude faf3012a36 fix: remove .expect() from strip_html_tags to pass no-panics check
Use Option<Regex> with graceful fallback instead of panicking on
regex compilation failure.

https://claude.ai/code/session_01CsP5wMZ2evEMHgGghjAfR1
2026-03-28 19:16:03 +00:00
ZakiandClaude ce193ff2d2 fix(routines): set run.job_id in-memory and revert Cargo.lock drift
Address PR #1470 review feedback:

- Set run.job_id = Some(job_id) after link_routine_run_to_job succeeds
  so send_notification reads the correct value instead of always None.
- Revert Cargo.lock to staging baseline: the PR had accumulated
  unrelated dependency changes (openssl, native-tls, crossterm 0.28.1
  downgrade, foreign-types, vcpkg) from a dirty lockfile resolution.

[skip-regression-check]

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 19:16:03 +00:00
ZakiandClaude 7881998d93 fix: preserve non-HTML angle brackets in sanitize_summary
The previous strip_html_tags() implementation blindly removed all content
between angle brackets, mangling legitimate text like Vec<String>, shell
redirects (cat < input.txt), and comparison operators in LLM/error output.

Replace the naive char-by-char scanner with a regex that only matches
known HTML tag names (div, script, img, a, b, etc.), preserving generic
angle-bracket content that appears in code snippets and error messages.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 19:16:03 +00:00
Claude d967e620cc fix: resolve mutability mismatch in execute_full_job call after rebase
[skip-regression-check]

https://claude.ai/code/session_01ABGWibdKVQ3b6pEKtxPPkM
2026-03-28 19:16:03 +00:00
ZakiandClaude 7b884c4a24 fix(routines): propagate job_id to notification metadata (#1321)
Update run.job_id after execute_full_job() links the routine run to the
job, ensuring send_notification receives the actual job ID instead of None
on the normal completion path.

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 19:16:03 +00:00
Claude d45d5977a0 fix(deps): resolve RUSTSEC-2026-0049 rustls-webpki CRL advisory
Update rustls-webpki 0.103.9 -> 0.103.10 and ignore the advisory for
0.102.8 which is pinned by libsql's rustls 0.22.4 dependency chain.

https://claude.ai/code/session_01MmxvBgAMn4m45pZguFKBEX
2026-03-28 19:16:03 +00:00
Claude d29811db6b style: run cargo fmt to fix formatting
https://claude.ai/code/session_01Nv2TJ3so5WQqpRUcirAhT3
2026-03-28 19:16:03 +00:00
ZakiandClaude 25dcbf78b8 fix(routines): normalize notification summaries with truncation and metadata (#1321)
- Capitalize status labels in notifications (ok -> Completed, attention -> Needs attention)
- Sanitize and truncate long summaries to 500 chars with UTF-8-safe ellipsis
- Include job_id in notification metadata for full-job routines
- Move sanitize_summary/strip_html_tags out of #[cfg(test)] for production use
- Add regression tests for truncation, status labels, job_id metadata

Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
2026-03-28 19:16:03 +00:00
18 changed files with 734 additions and 1676 deletions
+3 -2
View File
@@ -191,9 +191,10 @@ HEARTBEAT_NOTIFY_CHANNEL=cli
HEARTBEAT_NOTIFY_USER=default HEARTBEAT_NOTIFY_USER=default
# Memory hygiene settings (automatic cleanup of stale workspace documents) # Memory hygiene settings (automatic cleanup of stale workspace documents)
# Runs on each heartbeat tick; discovers cleanup targets from .config metadata # Runs on each heartbeat tick; identity files (IDENTITY.md, SOUL.md) are never deleted
# MEMORY_HYGIENE_ENABLED=true # MEMORY_HYGIENE_ENABLED=true
# MEMORY_HYGIENE_VERSION_KEEP_COUNT=50 # max versions to keep per document # 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_CADENCE_HOURS=12 # minimum hours between cleanup passes # MEMORY_HYGIENE_CADENCE_HOURS=12 # minimum hours between cleanup passes
# Docker Sandbox # Docker Sandbox
Generated
+125 -6
View File
@@ -1510,7 +1510,7 @@ version = "1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "980c2afde4af43d6a05c5be738f9eae595cff86dce1f38f88b95058a98c027f3" checksum = "980c2afde4af43d6a05c5be738f9eae595cff86dce1f38f88b95058a98c027f3"
dependencies = [ dependencies = [
"crossterm", "crossterm 0.29.0",
] ]
[[package]] [[package]]
@@ -1731,7 +1731,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "04a63daf06a168535c74ab97cdba3ed4fa5d4f32cb36e437dcceb83d66854b7c" checksum = "04a63daf06a168535c74ab97cdba3ed4fa5d4f32cb36e437dcceb83d66854b7c"
dependencies = [ dependencies = [
"crokey-proc_macros", "crokey-proc_macros",
"crossterm", "crossterm 0.29.0",
"once_cell", "once_cell",
"serde", "serde",
"strict", "strict",
@@ -1743,7 +1743,7 @@ version = "1.4.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "847f11a14855fc490bd5d059821895c53e77eeb3c2b73ee3dded7ce77c93b231" checksum = "847f11a14855fc490bd5d059821895c53e77eeb3c2b73ee3dded7ce77c93b231"
dependencies = [ dependencies = [
"crossterm", "crossterm 0.29.0",
"proc-macro2", "proc-macro2",
"quote", "quote",
"strict", "strict",
@@ -1817,6 +1817,22 @@ version = "0.8.21"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28" checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28"
[[package]]
name = "crossterm"
version = "0.28.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "829d955a0bb380ef178a640b91779e3987da38c9aea133b20614cfed8cdea9c6"
dependencies = [
"bitflags 2.11.0",
"crossterm_winapi",
"mio",
"parking_lot",
"rustix 0.38.44",
"signal-hook",
"signal-hook-mio",
"winapi",
]
[[package]] [[package]]
name = "crossterm" name = "crossterm"
version = "0.29.0" version = "0.29.0"
@@ -2476,6 +2492,21 @@ version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "77ce24cb58228fbb8aa041425bb1050850ac19177686ea6e0f41a70416f56fdb" checksum = "77ce24cb58228fbb8aa041425bb1050850ac19177686ea6e0f41a70416f56fdb"
[[package]]
name = "foreign-types"
version = "0.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f6f339eb8adc052cd2ca78910fda869aefa38d22d5cb648e6485e4d3fc06f3b1"
dependencies = [
"foreign-types-shared",
]
[[package]]
name = "foreign-types-shared"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "00b0228411908ca8685dba7fc2cdd70ec9990a6e753e89b6ac91a84c40fbaf4b"
[[package]] [[package]]
name = "form_urlencoded" name = "form_urlencoded"
version = "1.2.2" version = "1.2.2"
@@ -3118,7 +3149,6 @@ dependencies = [
"tokio", "tokio",
"tokio-rustls 0.26.4", "tokio-rustls 0.26.4",
"tower-service", "tower-service",
"webpki-roots 1.0.6",
] ]
[[package]] [[package]]
@@ -3133,6 +3163,22 @@ dependencies = [
"tokio-io-timeout", "tokio-io-timeout",
] ]
[[package]]
name = "hyper-tls"
version = "0.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "70206fc6890eaca9fde8a0bf71caa2ddfc9fe045ac9e5c70df101a7dbde866e0"
dependencies = [
"bytes",
"http-body-util",
"hyper 1.8.1",
"hyper-util",
"native-tls",
"tokio",
"tokio-native-tls",
"tower-service",
]
[[package]] [[package]]
name = "hyper-util" name = "hyper-util"
version = "0.1.20" version = "0.1.20"
@@ -3410,7 +3456,7 @@ dependencies = [
"clap_complete", "clap_complete",
"criterion", "criterion",
"cron", "cron",
"crossterm", "crossterm 0.28.1",
"deadpool-postgres", "deadpool-postgres",
"dirs 6.0.0", "dirs 6.0.0",
"dotenvy", "dotenvy",
@@ -4089,6 +4135,23 @@ dependencies = [
"rand 0.8.5", "rand 0.8.5",
] ]
[[package]]
name = "native-tls"
version = "0.2.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "465500e14ea162429d264d44189adc38b199b62b1c21eea9f69e4b73cb03bbf2"
dependencies = [
"libc",
"log",
"openssl",
"openssl-probe 0.2.1",
"openssl-sys",
"schannel",
"security-framework 3.7.0",
"security-framework-sys",
"tempfile",
]
[[package]] [[package]]
name = "new_debug_unreachable" name = "new_debug_unreachable"
version = "1.0.6" version = "1.0.6"
@@ -4311,6 +4374,32 @@ dependencies = [
"pathdiff", "pathdiff",
] ]
[[package]]
name = "openssl"
version = "0.10.76"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "951c002c75e16ea2c65b8c7e4d3d51d5530d8dfa7d060b4776828c88cfb18ecf"
dependencies = [
"bitflags 2.11.0",
"cfg-if",
"foreign-types",
"libc",
"once_cell",
"openssl-macros",
"openssl-sys",
]
[[package]]
name = "openssl-macros"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a948666b637a0f465e8564c73e89d4dde00d72d4d473cc972f390fc3dcee7d9c"
dependencies = [
"proc-macro2",
"quote",
"syn 2.0.117",
]
[[package]] [[package]]
name = "openssl-probe" name = "openssl-probe"
version = "0.1.6" version = "0.1.6"
@@ -4323,6 +4412,18 @@ version = "0.2.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe"
[[package]]
name = "openssl-sys"
version = "0.9.112"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "57d55af3b3e226502be1526dfdba67ab0e9c96fc293004e79576b2b9edb0dbdb"
dependencies = [
"cc",
"libc",
"pkg-config",
"vcpkg",
]
[[package]] [[package]]
name = "option-ext" name = "option-ext"
version = "0.2.0" version = "0.2.0"
@@ -5312,11 +5413,13 @@ dependencies = [
"http-body-util", "http-body-util",
"hyper 1.8.1", "hyper 1.8.1",
"hyper-rustls 0.27.7", "hyper-rustls 0.27.7",
"hyper-tls",
"hyper-util", "hyper-util",
"js-sys", "js-sys",
"log", "log",
"mime", "mime",
"mime_guess", "mime_guess",
"native-tls",
"percent-encoding", "percent-encoding",
"pin-project-lite", "pin-project-lite",
"quinn", "quinn",
@@ -5328,6 +5431,7 @@ dependencies = [
"serde_urlencoded", "serde_urlencoded",
"sync_wrapper 1.0.2", "sync_wrapper 1.0.2",
"tokio", "tokio",
"tokio-native-tls",
"tokio-rustls 0.26.4", "tokio-rustls 0.26.4",
"tokio-util", "tokio-util",
"tower 0.5.3", "tower 0.5.3",
@@ -5338,7 +5442,6 @@ dependencies = [
"wasm-bindgen-futures", "wasm-bindgen-futures",
"wasm-streams", "wasm-streams",
"web-sys", "web-sys",
"webpki-roots 1.0.6",
] ]
[[package]] [[package]]
@@ -6671,6 +6774,16 @@ dependencies = [
"syn 2.0.117", "syn 2.0.117",
] ]
[[package]]
name = "tokio-native-tls"
version = "0.3.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bbae76ab933c85776efabc971569dd6119c580d8f5d448769dec1764bf796ef2"
dependencies = [
"native-tls",
"tokio",
]
[[package]] [[package]]
name = "tokio-postgres" name = "tokio-postgres"
version = "0.7.16" version = "0.7.16"
@@ -7354,6 +7467,12 @@ version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65"
[[package]]
name = "vcpkg"
version = "0.2.15"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426"
[[package]] [[package]]
name = "version_check" name = "version_check"
version = "0.9.5" version = "0.9.5"
-23
View File
@@ -1,23 +0,0 @@
-- 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);
+4 -21
View File
@@ -276,8 +276,8 @@ impl HeartbeatRunner {
.await; .await;
if report.had_work() { if report.had_work() {
tracing::info!( tracing::info!(
directories_cleaned = ?report.directories_cleaned, daily_logs_deleted = report.daily_logs_deleted,
versions_pruned = report.versions_pruned, conversation_docs_deleted = report.conversation_docs_deleted,
"heartbeat: memory hygiene deleted stale documents" "heartbeat: memory hygiene deleted stale documents"
); );
} }
@@ -590,23 +590,6 @@ pub fn spawn_multi_user_heartbeat(
let workspace = Arc::new(Workspace::new_with_db(user_id, Arc::clone(store.db()))); let workspace = Arc::new(Workspace::new_with_db(user_id, Arc::clone(store.db())));
// Run memory hygiene per user (same as single-user heartbeat).
let hygiene_ws = Arc::clone(&workspace);
let hygiene_cfg = hygiene_config.clone();
let hygiene_user = user_id.clone();
tokio::spawn(async move {
let report =
crate::workspace::hygiene::run_if_due(&hygiene_ws, &hygiene_cfg).await;
if report.had_work() {
tracing::info!(
user_id = hygiene_user,
directories_cleaned = ?report.directories_cleaned,
versions_pruned = report.versions_pruned,
"multi-user heartbeat: memory hygiene deleted stale documents"
);
}
});
// Drain completed tasks to stay within the concurrency cap. // Drain completed tasks to stay within the concurrency cap.
while join_set.len() >= MAX_CONCURRENT_HEARTBEATS { while join_set.len() >= MAX_CONCURRENT_HEARTBEATS {
if let Some(join_result) = join_set.join_next().await { if let Some(join_result) = join_set.join_next().await {
@@ -633,8 +616,8 @@ pub fn spawn_multi_user_heartbeat(
if report.had_work() { if report.had_work() {
tracing::info!( tracing::info!(
user_id = uid, user_id = uid,
directories_cleaned = ?report.directories_cleaned, daily_logs_deleted = report.daily_logs_deleted,
versions_pruned = report.versions_pruned, conversation_docs_deleted = report.conversation_docs_deleted,
"multi-user heartbeat: memory hygiene deleted stale documents" "multi-user heartbeat: memory hygiene deleted stale documents"
); );
} }
+326 -25
View File
@@ -715,6 +715,7 @@ impl RoutineEngine {
status, status,
Some(summary), Some(summary),
thread_id.as_deref(), thread_id.as_deref(),
run.job_id,
) )
.await; .await;
@@ -1085,7 +1086,7 @@ struct EngineContext {
} }
/// Execute a routine run. Handles both lightweight and full_job modes. /// Execute a routine run. Handles both lightweight and full_job modes.
async fn execute_routine(ctx: EngineContext, routine: Routine, run: RoutineRun) { async fn execute_routine(ctx: EngineContext, routine: Routine, mut run: RoutineRun) {
// Increment running count (atomic: survives panics in the execution below) // Increment running count (atomic: survives panics in the execution below)
ctx.running_count.fetch_add(1, Ordering::Relaxed); ctx.running_count.fetch_add(1, Ordering::Relaxed);
@@ -1118,7 +1119,7 @@ async fn execute_routine(ctx: EngineContext, routine: Routine, run: RoutineRun)
description, description,
max_iterations: *max_iterations, max_iterations: *max_iterations,
}; };
execute_full_job(&ctx, &routine, &run, &execution).await execute_full_job(&ctx, &routine, &mut run, &execution).await
} }
}; };
@@ -1218,6 +1219,7 @@ async fn execute_routine(ctx: EngineContext, routine: Routine, run: RoutineRun)
status, status,
summary.as_deref(), summary.as_deref(),
thread_id.as_deref(), thread_id.as_deref(),
run.job_id,
) )
.await; .await;
} }
@@ -1253,7 +1255,7 @@ struct FullJobExecutionConfig<'a> {
async fn execute_full_job( async fn execute_full_job(
ctx: &EngineContext, ctx: &EngineContext,
routine: &Routine, routine: &Routine,
run: &RoutineRun, run: &mut RoutineRun,
execution: &FullJobExecutionConfig<'_>, execution: &FullJobExecutionConfig<'_>,
) -> Result<(RunStatus, Option<String>, Option<i32>), RoutineError> { ) -> Result<(RunStatus, Option<String>, Option<i32>), RoutineError> {
match ctx.sandbox_readiness { match ctx.sandbox_readiness {
@@ -1315,6 +1317,9 @@ async fn execute_full_job(
reason: format!("failed to link run to job: {e}"), reason: format!("failed to link run to job: {e}"),
})?; })?;
// Keep the in-memory struct in sync so send_notification can read run.job_id.
run.job_id = Some(job_id);
tracing::info!( tracing::info!(
routine = %routine.name, routine = %routine.name,
job_id = %job_id, job_id = %job_id,
@@ -1827,7 +1832,18 @@ async fn execute_routine_tool(
Ok(result_str) Ok(result_str)
} }
/// Human-readable label for a run status, suitable for user-facing notifications.
fn status_display_label(status: RunStatus) -> &'static str {
match status {
RunStatus::Ok => "Completed",
RunStatus::Attention => "Needs attention",
RunStatus::Failed => "Failed",
RunStatus::Running => "Running",
}
}
/// Send a notification based on the routine's notify config and run status. /// Send a notification based on the routine's notify config and run status.
#[allow(clippy::too_many_arguments)]
async fn send_notification( async fn send_notification(
tx: &mpsc::Sender<OutgoingResponse>, tx: &mpsc::Sender<OutgoingResponse>,
notify: &NotifyConfig, notify: &NotifyConfig,
@@ -1836,6 +1852,7 @@ async fn send_notification(
status: RunStatus, status: RunStatus,
summary: Option<&str>, summary: Option<&str>,
thread_id: Option<&str>, thread_id: Option<&str>,
job_id: Option<Uuid>,
) { ) {
let should_notify = match status { let should_notify = match status {
RunStatus::Ok => notify.on_success, RunStatus::Ok => notify.on_success,
@@ -1855,23 +1872,37 @@ async fn send_notification(
RunStatus::Running => "", RunStatus::Running => "",
}; };
let label = status_display_label(status);
let message = match summary { let message = match summary {
Some(s) => format!("{} *Routine '{}'*: {}\n\n{}", icon, routine_name, status, s), Some(s) => {
None => format!("{} *Routine '{}'*: {}", icon, routine_name, status), let sanitized = sanitize_summary(s);
format!(
"{} *Routine '{}'*: {}\n\n{}",
icon, routine_name, label, sanitized
)
}
None => format!("{} *Routine '{}'*: {}", icon, routine_name, label),
}; };
let mut metadata = serde_json::json!({
"source": "routine",
"routine_name": routine_name,
"status": status.to_string(),
"owner_id": owner_id,
"notify_user": notify.user,
"notify_channel": notify.channel,
});
if let Some(jid) = job_id {
metadata["job_id"] = serde_json::json!(jid.to_string());
}
let response = OutgoingResponse { let response = OutgoingResponse {
content: message, content: message,
thread_id: thread_id.map(String::from), thread_id: thread_id.map(String::from),
attachments: Vec::new(), attachments: Vec::new(),
metadata: serde_json::json!({ metadata,
"source": "routine",
"routine_name": routine_name,
"status": status.to_string(),
"owner_id": owner_id,
"notify_user": notify.user,
"notify_channel": notify.channel,
}),
}; };
if let Err(e) = tx.send(response).await { if let Err(e) = tx.send(response).await {
@@ -1934,7 +1965,6 @@ fn truncate(s: &str, max: usize) -> String {
/// 2. Strip HTML tags to prevent injection in web-rendered notifications /// 2. Strip HTML tags to prevent injection in web-rendered notifications
/// 3. Collapse multiple whitespace/newlines to single spaces for cleaner output /// 3. Collapse multiple whitespace/newlines to single spaces for cleaner output
/// 4. Truncate to 500 chars to prevent oversized notifications /// 4. Truncate to 500 chars to prevent oversized notifications
#[cfg(test)]
fn sanitize_summary(s: &str) -> String { fn sanitize_summary(s: &str) -> String {
// Strip control characters (keep newline for now, collapse later) // Strip control characters (keep newline for now, collapse later)
let no_control: String = s let no_control: String = s
@@ -1961,19 +1991,59 @@ fn sanitize_summary(s: &str) -> String {
} }
} }
/// Remove HTML/XML tags from a string. /// Remove actual HTML tags from a string while preserving non-HTML angle brackets.
#[cfg(test)] ///
/// Only strips patterns that look like real HTML/XML tags (e.g. `<div>`, `</p>`,
/// `<img src=...>`), not generic angle-bracket content like `Vec<String>`,
/// `cat < input.txt`, or comparison operators.
///
/// Also strips HTML comments (`<!--...-->`), SVG/MathML tags, and custom elements
/// (tags containing hyphens like `<custom-element>`).
fn strip_html_tags(s: &str) -> String { fn strip_html_tags(s: &str) -> String {
let mut result = String::with_capacity(s.len()); use std::sync::LazyLock;
let mut in_tag = false;
for c in s.chars() { // HTML comment pattern: <!--...-->
match c { static COMMENT_RE: LazyLock<Option<Regex>> =
'<' => in_tag = true, LazyLock::new(|| Regex::new(r"<!--[\s\S]*?-->").ok());
'>' if in_tag => in_tag = false,
_ if !in_tag => result.push(c), // Known HTML/SVG/MathML tag names. Includes SVG tags (svg, path, circle, etc.)
_ => {} // and MathML tags (math, mrow, etc.) that can carry event handlers.
} static HTML_TAG_RE: LazyLock<Option<Regex>> = LazyLock::new(|| {
let tags = "a|abbr|address|area|article|aside|audio|b|base|bdi|bdo|blockquote|\
body|br|button|canvas|caption|cite|code|col|colgroup|data|datalist|dd|del|\
details|dfn|dialog|div|dl|dt|em|embed|fieldset|figcaption|figure|footer|\
form|h[1-6]|head|header|hgroup|hr|html|i|iframe|img|input|ins|kbd|label|\
legend|li|link|main|map|mark|meta|meter|nav|noscript|object|ol|optgroup|\
option|output|p|param|picture|pre|progress|q|rp|rt|ruby|s|samp|script|\
section|select|slot|small|source|span|strong|style|sub|summary|sup|table|\
tbody|td|template|textarea|tfoot|th|thead|time|title|tr|track|u|ul|var|\
video|wbr|\
svg|g|path|circle|ellipse|line|polyline|polygon|rect|text|tspan|defs|\
clippath|mask|pattern|image|use|symbol|marker|lineargradient|\
radialgradient|stop|filter|foreignobject|animate|animatetransform|\
math|mrow|mi|mo|mn|ms|mtext|mfrac|msqrt|mroot|msub|msup|msubsup|\
munder|mover|munderover|mtable|mtr|mtd|mspace|mpadded|mfenced|menclose";
// Handles: <tag>, </tag>, <tag/>, <tag />, <tag attr="val">, <tag attr="val"/>
Regex::new(&format!(r"(?i)</?(?:{})(?:\s[^>]*)?\s*/?>", tags)).ok()
});
// Custom elements: tags containing a hyphen (web components spec requires it).
// E.g. <custom-element>, <my-widget foo="bar">, </x-foo>
static CUSTOM_ELEMENT_RE: LazyLock<Option<Regex>> =
LazyLock::new(|| Regex::new(r"(?i)</?\w+-[\w-]*(?:\s[^>]*)?\s*/?>").ok());
let mut result = s.to_string();
if let Some(re) = COMMENT_RE.as_ref() {
result = re.replace_all(&result, "").into_owned();
} }
if let Some(re) = HTML_TAG_RE.as_ref() {
result = re.replace_all(&result, "").into_owned();
}
if let Some(re) = CUSTOM_ELEMENT_RE.as_ref() {
result = re.replace_all(&result, "").into_owned();
}
result result
} }
@@ -2548,6 +2618,33 @@ mod tests {
assert_eq!(sanitize_summary("<img src=x onerror=alert(1)>"), ""); assert_eq!(sanitize_summary("<img src=x onerror=alert(1)>"), "");
} }
#[test]
fn test_sanitize_summary_preserves_non_html_angle_brackets() {
use super::sanitize_summary;
// Rust/Java generics must pass through unchanged
assert_eq!(
sanitize_summary("expected Vec<String>"),
"expected Vec<String>"
);
assert_eq!(
sanitize_summary("HashMap<String, Vec<u8>>"),
"HashMap<String, Vec<u8>>"
);
// Shell redirects must pass through unchanged
assert_eq!(sanitize_summary("cat < input.txt"), "cat < input.txt");
// Comparison operators must pass through unchanged
assert_eq!(sanitize_summary("x < 10 && y > 20"), "x < 10 && y > 20");
// Mixed: real HTML stripped but generics preserved
assert_eq!(
sanitize_summary("Error in Vec<String>: <b>failed</b>"),
"Error in Vec<String>: failed"
);
}
#[test] #[test]
fn test_sanitize_summary_multibyte_truncation() { fn test_sanitize_summary_multibyte_truncation() {
use super::sanitize_summary; use super::sanitize_summary;
@@ -2558,4 +2655,208 @@ mod tests {
assert!(result.len() <= 503); assert!(result.len() <= 503);
assert!(result.ends_with("...")); assert!(result.ends_with("..."));
} }
#[test]
fn test_sanitize_summary_truncates_long_text() {
use super::sanitize_summary;
let short = "This is a short summary.";
assert_eq!(sanitize_summary(short), short);
let long = "x".repeat(600);
let result = sanitize_summary(&long);
assert!(
result.len() <= 503,
"Truncated summary should be at most 503 bytes (500 + '...')"
);
assert!(
result.ends_with("..."),
"Truncated summary should end with ellipsis"
);
}
#[test]
fn test_sanitize_summary_strips_all_html_forms() {
use super::sanitize_summary;
// Self-closing tags without whitespace: <br/>, <img/>
assert_eq!(sanitize_summary("line1<br/>line2"), "line1line2");
assert_eq!(sanitize_summary("text<img/>more"), "textmore");
assert_eq!(sanitize_summary("text<br />more"), "textmore");
// HTML comments
assert_eq!(sanitize_summary("before<!--x-->after"), "beforeafter");
assert_eq!(sanitize_summary("a<!-- multi\nline -->b"), "ab");
// SVG tags (can carry event handlers)
assert_eq!(
sanitize_summary("<svg onload=alert(1)>payload</svg>"),
"payload"
);
assert_eq!(sanitize_summary("<svg><circle r=10/></svg>"), "");
// MathML tags
assert_eq!(sanitize_summary("<math><mrow>x</mrow></math>"), "x");
// Custom elements (web components with hyphens)
assert_eq!(
sanitize_summary("before<custom-element>inner</custom-element>after"),
"beforeinnerafter"
);
assert_eq!(
sanitize_summary("<my-widget foo=\"bar\">content</my-widget>"),
"content"
);
// Generics must still be preserved
assert_eq!(
sanitize_summary("expected Vec<String>"),
"expected Vec<String>"
);
}
#[test]
fn test_status_display_label_readable() {
use super::status_display_label;
assert_eq!(status_display_label(RunStatus::Ok), "Completed");
assert_eq!(status_display_label(RunStatus::Failed), "Failed");
assert_eq!(
status_display_label(RunStatus::Attention),
"Needs attention"
);
assert_eq!(status_display_label(RunStatus::Running), "Running");
}
#[tokio::test]
async fn test_notification_message_uses_readable_status() {
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
let notify = NotifyConfig {
on_success: true,
on_failure: true,
on_attention: true,
..Default::default()
};
super::send_notification(
&tx,
&notify,
"user-1",
"my-routine",
RunStatus::Ok,
Some("All good"),
None,
None,
)
.await;
let msg = rx.recv().await.expect("should receive notification");
assert!(
msg.content.contains("Completed"),
"Notification should use readable label 'Completed', got: {}",
msg.content
);
assert!(
!msg.content.contains(": ok"),
"Notification should not contain raw lowercase status"
);
}
#[tokio::test]
async fn test_notification_includes_job_id_in_metadata() {
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
let notify = NotifyConfig {
on_failure: true,
..Default::default()
};
let job_id = uuid::Uuid::new_v4();
super::send_notification(
&tx,
&notify,
"user-1",
"my-routine",
RunStatus::Failed,
Some("something broke"),
None,
Some(job_id),
)
.await;
let msg = rx.recv().await.expect("should receive notification");
let meta_job_id = msg.metadata["job_id"]
.as_str()
.expect("metadata should contain job_id");
assert_eq!(meta_job_id, job_id.to_string());
}
#[tokio::test]
async fn test_notification_omits_job_id_when_none() {
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
let notify = NotifyConfig {
on_success: true,
..Default::default()
};
super::send_notification(
&tx,
&notify,
"user-1",
"my-routine",
RunStatus::Ok,
Some("done"),
None,
None,
)
.await;
let msg = rx.recv().await.expect("should receive notification");
assert!(
msg.metadata.get("job_id").is_none(),
"metadata should not contain job_id when None"
);
}
#[tokio::test]
async fn test_notification_truncates_long_summary() {
use tokio::sync::mpsc;
let (tx, mut rx) = mpsc::channel(1);
let notify = NotifyConfig {
on_failure: true,
..Default::default()
};
let long_summary = "z".repeat(1000);
super::send_notification(
&tx,
&notify,
"user-1",
"my-routine",
RunStatus::Failed,
Some(&long_summary),
None,
None,
)
.await;
let msg = rx.recv().await.expect("should receive notification");
// The sanitized summary should be truncated to ~500 chars + "..."
// The full message includes icon + routine name + label, so just check
// it doesn't contain the full 1000-char string.
assert!(
!msg.content.contains(&long_summary),
"Notification should truncate long summaries"
);
assert!(
msg.content.contains("..."),
"Truncated notification should contain ellipsis"
);
}
} }
+13 -5
View File
@@ -10,8 +10,10 @@ use crate::error::ConfigError;
pub struct HygieneConfig { pub struct HygieneConfig {
/// Whether hygiene is enabled. Env: `MEMORY_HYGIENE_ENABLED` (default: true). /// Whether hygiene is enabled. Env: `MEMORY_HYGIENE_ENABLED` (default: true).
pub enabled: bool, pub enabled: bool,
/// Maximum versions to keep per document. Env: `MEMORY_HYGIENE_VERSION_KEEP_COUNT` (default: 50). /// Days before `daily/` documents are deleted. Env: `MEMORY_HYGIENE_DAILY_RETENTION_DAYS` (default: 30).
pub version_keep_count: u32, pub daily_retention_days: u32,
/// Days before `conversations/` documents are deleted. Env: `MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS` (default: 7).
pub conversation_retention_days: u32,
/// Minimum hours between hygiene passes. Env: `MEMORY_HYGIENE_CADENCE_HOURS` (default: 12). /// Minimum hours between hygiene passes. Env: `MEMORY_HYGIENE_CADENCE_HOURS` (default: 12).
pub cadence_hours: u32, pub cadence_hours: u32,
} }
@@ -20,7 +22,8 @@ impl Default for HygieneConfig {
fn default() -> Self { fn default() -> Self {
Self { Self {
enabled: true, enabled: true,
version_keep_count: 50, daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 12, cadence_hours: 12,
} }
} }
@@ -30,7 +33,11 @@ impl HygieneConfig {
pub(crate) fn resolve() -> Result<Self, ConfigError> { pub(crate) fn resolve() -> Result<Self, ConfigError> {
Ok(Self { Ok(Self {
enabled: parse_bool_env("MEMORY_HYGIENE_ENABLED", true)?, enabled: parse_bool_env("MEMORY_HYGIENE_ENABLED", true)?,
version_keep_count: parse_optional_env("MEMORY_HYGIENE_VERSION_KEEP_COUNT", 50)?, daily_retention_days: parse_optional_env("MEMORY_HYGIENE_DAILY_RETENTION_DAYS", 30)?,
conversation_retention_days: parse_optional_env(
"MEMORY_HYGIENE_CONVERSATION_RETENTION_DAYS",
7,
)?,
cadence_hours: parse_optional_env("MEMORY_HYGIENE_CADENCE_HOURS", 12)?, cadence_hours: parse_optional_env("MEMORY_HYGIENE_CADENCE_HOURS", 12)?,
}) })
} }
@@ -40,7 +47,8 @@ impl HygieneConfig {
pub fn to_workspace_config(&self) -> crate::workspace::hygiene::HygieneConfig { pub fn to_workspace_config(&self) -> crate::workspace::hygiene::HygieneConfig {
crate::workspace::hygiene::HygieneConfig { crate::workspace::hygiene::HygieneConfig {
enabled: self.enabled, enabled: self.enabled,
version_keep_count: self.version_keep_count, daily_retention_days: self.daily_retention_days,
conversation_retention_days: self.conversation_retention_days,
cadence_hours: self.cadence_hours, cadence_hours: self.cadence_hours,
state_dir: ironclaw_base_dir(), state_dir: ironclaw_base_dir(),
} }
+1 -1
View File
@@ -981,7 +981,7 @@ mod tests {
assert!(alice_stats.last_active_at.is_some()); assert!(alice_stats.last_active_at.is_some());
// Bob has no LLM calls so doesn't appear in summary stats // Bob has no LLM calls so doesn't appear in summary stats
assert!(stats.iter().find(|s| s.user_id == "bob").is_none()); assert!(!stats.iter().any(|s| s.user_id == "bob"));
// Filter to single user // Filter to single user
let alice_only = db.user_summary_stats(Some("alice")).await.unwrap(); let alice_only = db.user_summary_stats(Some("alice")).await.unwrap();
+2 -326
View File
@@ -13,8 +13,8 @@ use super::{
use crate::db::WorkspaceStore; use crate::db::WorkspaceStore;
use crate::error::{DatabaseError, WorkspaceError}; use crate::error::{DatabaseError, WorkspaceError};
use crate::workspace::{ use crate::workspace::{
DocumentVersion, MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult, MemoryChunk, MemoryDocument, RankedResult, SearchConfig, SearchResult, WorkspaceEntry,
VersionSummary, WorkspaceEntry, fuse_results, fuse_results,
}; };
use chrono::Utc; use chrono::Utc;
@@ -840,330 +840,6 @@ impl WorkspaceStore for LibSqlBackend {
Ok(fuse_results(fts_results, vector_results, config)) 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<Uuid>,
) -> Result<Vec<MemoryDocument>, 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<i32, WorkspaceError> {
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());
// Use a transaction to prevent race conditions: the SELECT and INSERT
// must be atomic so concurrent writers don't allocate the same version.
let tx = conn
.transaction()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to start transaction: {e}"),
})?;
// Get next version number (inside transaction — serializes writers)
let mut rows = tx
.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
};
drop(rows);
tx.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}"),
})?;
tx.commit()
.await
.map_err(|e| WorkspaceError::SearchFailed {
reason: format!("Failed to commit version: {e}"),
})?;
Ok(next_version)
}
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
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<Vec<VersionSummary>, 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<Option<i32>, 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::<libsql::Value>(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<u64, WorkspaceError> {
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)] #[cfg(test)]
-19
View File
@@ -785,25 +785,6 @@ CREATE TABLE IF NOT EXISTS api_tokens (
); );
CREATE INDEX IF NOT EXISTS idx_api_tokens_user ON api_tokens(user_id); CREATE INDEX IF NOT EXISTS idx_api_tokens_user ON api_tokens(user_id);
CREATE INDEX IF NOT EXISTS idx_api_tokens_hash ON api_tokens(token_hash); CREATE INDEX IF NOT EXISTS idx_api_tokens_hash ON api_tokens(token_hash);
"#,
),
(
15,
"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);
"#, "#,
), ),
]; ];
-61
View File
@@ -700,67 +700,6 @@ pub trait WorkspaceStore: Send + Sync {
config: &SearchConfig, config: &SearchConfig,
) -> Result<Vec<SearchResult>, WorkspaceError>; ) -> Result<Vec<SearchResult>, 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<Uuid>,
) -> Result<Vec<MemoryDocument>, 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<i32, WorkspaceError>;
/// Get a specific version of a document.
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<crate::workspace::DocumentVersion, WorkspaceError>;
/// List versions of a document (newest first).
async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<crate::workspace::VersionSummary>, 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<Option<i32>, 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<u64, WorkspaceError>;
// ==================== Multi-scope read methods ==================== // ==================== Multi-scope read methods ====================
// //
// Default implementations loop over user_ids calling single-scope methods, // Default implementations loop over user_ids calling single-scope methods,
+1 -65
View File
@@ -25,8 +25,7 @@ use crate::history::{
LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, Store, LlmCallRecord, SandboxJobRecord, SandboxJobSummary, SettingRow, Store,
}; };
use crate::workspace::{ use crate::workspace::{
DocumentVersion, MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, MemoryChunk, MemoryDocument, Repository, SearchConfig, SearchResult, WorkspaceEntry,
VersionSummary, WorkspaceEntry,
}; };
/// PostgreSQL database backend. /// PostgreSQL database backend.
@@ -786,69 +785,6 @@ impl WorkspaceStore for PgBackend {
.list_directory_multi(user_ids, agent_id, directory) .list_directory_multi(user_ids, agent_id, directory)
.await .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<Uuid>,
) -> Result<Vec<MemoryDocument>, 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<i32, WorkspaceError> {
self.repo
.save_version(document_id, content, content_hash, changed_by)
.await
}
async fn get_version(
&self,
document_id: Uuid,
version: i32,
) -> Result<DocumentVersion, WorkspaceError> {
self.repo.get_version(document_id, version).await
}
async fn list_versions(
&self,
document_id: Uuid,
limit: i64,
) -> Result<Vec<VersionSummary>, WorkspaceError> {
self.repo.list_versions(document_id, limit).await
}
async fn get_latest_version_number(
&self,
document_id: Uuid,
) -> Result<Option<i32>, WorkspaceError> {
self.repo.get_latest_version_number(document_id).await
}
async fn prune_versions(
&self,
document_id: Uuid,
keep_count: i32,
) -> Result<u64, WorkspaceError> {
self.repo.prune_versions(document_id, keep_count).await
}
} }
// ==================== UserStore ==================== // ==================== UserStore ====================
-6
View File
@@ -315,12 +315,6 @@ pub enum WorkspaceError {
#[error("Write rejected for '{path}': prompt injection detected ({reason})")] #[error("Write rejected for '{path}': prompt injection detected ({reason})")]
InjectionRejected { path: String, reason: String }, 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). /// Orchestrator errors (internal API, container management).
+4 -140
View File
@@ -246,26 +246,9 @@ impl Tool for MemoryWriteTool {
"type": "boolean", "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.", "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 "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": [] "required": ["content"]
}) })
} }
@@ -276,9 +259,7 @@ impl Tool for MemoryWriteTool {
) -> Result<ToolOutput, ToolError> { ) -> Result<ToolOutput, ToolError> {
let start = std::time::Instant::now(); let start = std::time::Instant::now();
// In patch mode (old_string present), content is not required. let content = require_str(&params, "content")?;
let is_patch_mode = params.get("old_string").and_then(|v| v.as_str()).is_some();
let content = params.get("content").and_then(|v| v.as_str()).unwrap_or("");
let target = params let target = params
.get("target") .get("target")
@@ -318,9 +299,9 @@ impl Tool for MemoryWriteTool {
return Ok(ToolOutput::success(output, start.elapsed())); return Ok(ToolOutput::success(output, start.elapsed()));
} }
if !is_patch_mode && content.trim().is_empty() { if content.trim().is_empty() {
return Err(ToolError::InvalidParameters( return Err(ToolError::InvalidParameters(
"content cannot be empty (use old_string/new_string for patch mode)".to_string(), "content cannot be empty".to_string(),
)); ));
} }
@@ -349,46 +330,6 @@ impl Tool for MemoryWriteTool {
path => path.to_string(), 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. // When a layer is specified, route through layer-aware methods for ALL targets.
// Otherwise, use default workspace methods (which include injection scanning). // Otherwise, use default workspace methods (which include injection scanning).
let layer_result = if let Some(layer_name) = layer { let layer_result = if let Some(layer_name) = layer {
@@ -492,24 +433,6 @@ impl Tool for MemoryWriteTool {
} }
} }
// Apply metadata if provided (after write/append, works for all targets).
// We read the document once to get its ID — this is a hot read right
// after the write, so it's effectively free (same DB connection/cache).
if let Some(meta) = params.get("metadata")
&& meta.is_object()
{
match workspace.read(&resolved_path).await {
Ok(doc) => {
if let Err(e) = workspace.update_metadata(doc.id, meta).await {
tracing::warn!(path = %resolved_path, "failed to update metadata: {e}");
}
}
Err(e) => {
tracing::warn!(path = %resolved_path, "failed to read doc for metadata update: {e}");
}
}
}
let mut output = serde_json::json!({ let mut output = serde_json::json!({
"status": "written", "status": "written",
"path": resolved_path, "path": resolved_path,
@@ -578,15 +501,6 @@ impl Tool for MemoryReadTool {
"path": { "path": {
"type": "string", "type": "string",
"description": "Path to the file (e.g., 'MEMORY.md', 'daily/2024-01-15.md', 'projects/alpha/notes.md')" "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"] "required": ["path"]
@@ -611,61 +525,11 @@ impl Tool for MemoryReadTool {
} }
let workspace = self.resolver.resolve(&ctx.user_id).await; 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 let doc = workspace
.read(path) .read(path)
.await .await
.map_err(|e| ToolError::ExecutionFailed(format!("Read failed: {}", e)))?; .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::<Vec<_>>(),
"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!({ let output = serde_json::json!({
"path": doc.path, "path": doc.path,
"content": doc.content, "content": doc.content,
-248
View File
@@ -2,7 +2,6 @@
use chrono::{DateTime, Utc}; use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use uuid::Uuid; use uuid::Uuid;
/// Well-known document paths. /// Well-known document paths.
@@ -38,139 +37,6 @@ pub mod paths {
pub const ASSISTANT_DIRECTIVES: &str = "context/assistant-directives.md"; 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<bool>,
/// When `true`, skip automatic versioning for this document/folder.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub skip_versioning: Option<bool>,
/// Hygiene (auto-cleanup) configuration for this folder.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hygiene: Option<HygieneMetadata>,
/// Preserve unknown fields for forward compatibility.
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
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<Utc>,
/// Who/what created this version (e.g. `"agent"`, `"user:alice"`).
pub changed_by: Option<String>,
}
/// 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<Utc>,
/// Who/what created this version.
pub changed_by: Option<String>,
}
/// 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. /// Paths treated as identity documents for multi-scope isolation.
/// ///
/// These files are always read from the primary scope only — never from /// These files are always read from the primary scope only — never from
@@ -494,120 +360,6 @@ mod tests {
assert_eq!(result[0].updated_at, Some(ts)); 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] #[test]
fn test_merge_workspace_entries_sorted_by_path() { fn test_merge_workspace_entries_sorted_by_path() {
let entries = vec![ let entries = vec![
+248 -184
View File
@@ -1,10 +1,8 @@
//! Memory hygiene: automatic cleanup of stale workspace documents. //! Memory hygiene: automatic cleanup of stale workspace documents.
//! //!
//! Runs on a configurable cadence and discovers which directories have hygiene //! Runs on a configurable cadence and deletes daily log entries and conversation
//! enabled by reading `.config` metadata documents. This is a **metadata-driven** //! documents older than their respective retention periods. Identity files
//! approach: instead of hardcoding `daily/` and `conversations/`, the system //! (`IDENTITY.md`, `SOUL.md`, etc.) are never touched.
//! respects `hygiene.enabled` and `hygiene.retention_days` set on each folder's
//! `.config` document.
//! //!
//! A global [`AtomicBool`] guard prevents concurrent hygiene passes, which //! A global [`AtomicBool`] guard prevents concurrent hygiene passes, which
//! avoids TOCTOU races on the state file and Windows file-locking errors //! avoids TOCTOU races on the state file and Windows file-locking errors
@@ -12,16 +10,18 @@
//! pass completes. //! pass completes.
//! //!
//! ```text //! ```text
//! ┌────────────────────────────────────────────────── //! ┌─────────────────────────────────────────────┐
//! │ Hygiene Pass //! │ Hygiene Pass │
//! │ //! │ │
//! │ 0. Acquire RUNNING guard (skip if held) //! │ 0. Acquire RUNNING guard (skip if held) │
//! │ 1. Check cadence (skip if ran recently) //! │ 1. Check cadence (skip if ran recently) │
//! │ 2. Save state (claim the cadence window) //! │ 2. Save state (claim the cadence window) │
//! │ 3. Discover .config docs with hygiene.enabled //! │ 3. List daily/ documents
//! │ 4. For each: cleanup_directory(parent, retention) //! │ 4. Delete those older than daily_retention
//! │ 5. Log summary //! │ 5. List conversations/ documents
//! └──────────────────────────────────────────────────┘ //! │ 6. Delete those older than conversation_ret │
//! │ 7. Log summary │
//! └─────────────────────────────────────────────┘
//! ``` //! ```
use std::path::PathBuf; use std::path::PathBuf;
@@ -31,22 +31,46 @@ use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use crate::bootstrap::ironclaw_base_dir; use crate::bootstrap::ironclaw_base_dir;
use crate::workspace::{DocumentMetadata, Workspace, is_config_path}; use crate::workspace::Workspace;
/// Global guard preventing concurrent hygiene passes. /// Global guard preventing concurrent hygiene passes.
static RUNNING: AtomicBool = AtomicBool::new(false); 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. /// Configuration for workspace hygiene.
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct HygieneConfig { pub struct HygieneConfig {
/// Whether hygiene is enabled at all. /// Whether hygiene is enabled at all.
pub enabled: bool, pub enabled: bool,
/// Maximum number of versions to keep per document. /// Documents in `daily/` older than this many days are deleted.
/// pub daily_retention_days: u32,
/// TODO: Wire up global version pruning once per-document iteration /// Documents in `conversations/` older than this many days are deleted.
/// is efficient (e.g., via a dedicated DB query). For now this field pub conversation_retention_days: u32,
/// is stored in config but not actively enforced during hygiene passes.
pub version_keep_count: u32,
/// Minimum hours between hygiene passes. /// Minimum hours between hygiene passes.
pub cadence_hours: u32, pub cadence_hours: u32,
/// Directory to store state file (default: `~/.ironclaw`). /// Directory to store state file (default: `~/.ironclaw`).
@@ -57,7 +81,8 @@ impl Default for HygieneConfig {
fn default() -> Self { fn default() -> Self {
Self { Self {
enabled: true, enabled: true,
version_keep_count: 50, daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 12, cadence_hours: 12,
state_dir: ironclaw_base_dir(), state_dir: ironclaw_base_dir(),
} }
@@ -73,10 +98,10 @@ struct HygieneState {
/// Summary of what a hygiene pass cleaned up. /// Summary of what a hygiene pass cleaned up.
#[derive(Debug, Default)] #[derive(Debug, Default)]
pub struct HygieneReport { pub struct HygieneReport {
/// Per-directory cleanup results: `(directory_path, deleted_count)`. /// Number of daily log documents deleted.
pub directories_cleaned: Vec<(String, u32)>, pub daily_logs_deleted: u32,
/// Number of document versions pruned across all documents. /// Number of conversation documents deleted.
pub versions_pruned: u64, pub conversation_docs_deleted: u32,
/// Whether the run was skipped (cadence not yet elapsed). /// Whether the run was skipped (cadence not yet elapsed).
pub skipped: bool, pub skipped: bool,
} }
@@ -84,7 +109,7 @@ pub struct HygieneReport {
impl HygieneReport { impl HygieneReport {
/// True if any cleanup work was done. /// True if any cleanup work was done.
pub fn had_work(&self) -> bool { pub fn had_work(&self) -> bool {
self.directories_cleaned.iter().any(|(_, n)| *n > 0) || self.versions_pruned > 0 self.daily_logs_deleted > 0 || self.conversation_docs_deleted > 0
} }
} }
@@ -143,51 +168,30 @@ pub async fn run_if_due(workspace: &Workspace, config: &HygieneConfig) -> Hygien
// TOCTOU races where another task reads stale state. // TOCTOU races where another task reads stale state.
save_state(&state_file); save_state(&state_file);
tracing::info!("memory hygiene: starting cleanup pass"); tracing::info!(
daily_retention_days = config.daily_retention_days,
conversation_retention_days = config.conversation_retention_days,
"memory hygiene: starting cleanup pass"
);
let mut report = HygieneReport::default(); let mut report = HygieneReport::default();
// Discover directories that have hygiene enabled via .config metadata. // Delete old daily logs
let config_docs = match workspace.find_config_documents().await { match cleanup_daily_logs(workspace, config.daily_retention_days).await {
Ok(docs) => docs, Ok(count) => report.daily_logs_deleted = count,
Err(e) => { Err(e) => tracing::warn!("memory hygiene: failed to clean daily logs: {e}"),
tracing::warn!("memory hygiene: failed to discover .config documents: {e}"); }
return report;
}
};
for doc in &config_docs { // Delete old conversation documents
let meta = DocumentMetadata::from_value(&doc.metadata); match cleanup_conversation_docs(workspace, config.conversation_retention_days).await {
let Some(hygiene) = meta.hygiene else { Ok(count) => report.conversation_docs_deleted = count,
continue; Err(e) => tracing::warn!("memory hygiene: failed to clean conversation docs: {e}"),
};
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() { if report.had_work() {
tracing::info!( tracing::info!(
directories_cleaned = ?report.directories_cleaned, daily_logs_deleted = report.daily_logs_deleted,
versions_pruned = report.versions_pruned, conversation_docs_deleted = report.conversation_docs_deleted,
"memory hygiene: cleanup complete" "memory hygiene: cleanup complete"
); );
} else { } else {
@@ -206,41 +210,88 @@ impl Drop for RunningGuard {
} }
} }
/// Delete documents in `directory` that are older than `retention_days`. /// Delete daily log documents older than `retention_days`.
/// async fn cleanup_daily_logs(
/// Skips directories and `.config` files (which must never be deleted by
/// hygiene). Returns the number of documents deleted.
async fn cleanup_directory(
workspace: &Workspace, workspace: &Workspace,
directory: &str,
retention_days: u32, retention_days: u32,
) -> Result<u32, anyhow::Error> { ) -> Result<u32, anyhow::Error> {
let cutoff = Utc::now() - chrono::Duration::days(i64::from(retention_days)); let cutoff = Utc::now() - chrono::Duration::days(i64::from(retention_days));
let entries = workspace.list(directory).await?; let entries = workspace.list("daily/").await?;
let mut deleted = 0u32; let mut deleted = 0u32;
for entry in entries { for entry in entries {
if entry.is_directory { if entry.is_directory {
continue; continue;
} }
if is_config_path(&entry.path) {
// Never delete identity documents
if is_identity_path(&entry.path) {
continue; continue;
} }
// Check if the document is old enough to delete
if let Some(updated_at) = entry.updated_at if let Some(updated_at) = entry.updated_at
&& updated_at < cutoff && updated_at < cutoff
{ {
let path = if entry.path.starts_with(directory) { let path = if entry.path.starts_with("daily/") {
entry.path.clone() entry.path.clone()
} else { } else {
format!("{}{}", directory, entry.path) format!("daily/{}", entry.path)
}; };
if let Err(e) = workspace.delete(&path).await { if let Err(e) = workspace.delete(&path).await {
tracing::warn!(path, "memory hygiene: failed to delete: {e}"); tracing::warn!(path, "memory hygiene: failed to delete: {e}");
} else { } else {
tracing::debug!(path, "memory hygiene: deleted stale document"); tracing::debug!(path, "memory hygiene: deleted old daily log");
deleted += 1; deleted += 1;
} }
} }
} }
Ok(deleted)
}
/// Delete conversation documents older than `retention_days`.
async fn cleanup_conversation_docs(
workspace: &Workspace,
retention_days: u32,
) -> Result<u32, anyhow::Error> {
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) Ok(deleted)
} }
@@ -298,7 +349,8 @@ mod tests {
fn default_config_is_reasonable() { fn default_config_is_reasonable() {
let cfg = HygieneConfig::default(); let cfg = HygieneConfig::default();
assert!(cfg.enabled); assert!(cfg.enabled);
assert_eq!(cfg.version_keep_count, 50); assert_eq!(cfg.daily_retention_days, 30);
assert_eq!(cfg.conversation_retention_days, 7);
assert_eq!(cfg.cadence_hours, 12); assert_eq!(cfg.cadence_hours, 12);
} }
@@ -310,33 +362,84 @@ mod tests {
} }
#[test] #[test]
fn report_had_work_when_directories_cleaned() { fn report_had_work_when_deleted() {
let report = HygieneReport { let report = HygieneReport {
directories_cleaned: vec![("daily/".to_string(), 3)], daily_logs_deleted: 3,
versions_pruned: 0, conversation_docs_deleted: 0,
skipped: false, skipped: false,
}; };
assert!(report.had_work()); assert!(report.had_work());
} }
#[test] #[test]
fn report_had_work_when_versions_pruned() { fn report_had_work_when_conversation_deleted() {
let report = HygieneReport { let report = HygieneReport {
directories_cleaned: vec![], daily_logs_deleted: 0,
versions_pruned: 5, conversation_docs_deleted: 2,
skipped: false, skipped: false,
}; };
assert!(report.had_work()); assert!(report.had_work());
} }
#[test] #[test]
fn report_no_work_when_zero_deletions() { fn is_identity_path_excludes_sacred_docs() {
let report = HygieneReport { for name in [
directories_cleaned: vec![("daily/".to_string(), 0)], "MEMORY.md",
versions_pruned: 0, "IDENTITY.md",
skipped: false, "SOUL.md",
}; "AGENTS.md",
assert!(!report.had_work()); "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");
}
} }
#[test] #[test]
@@ -449,105 +552,61 @@ mod tests {
Arc::new(Workspace::new_with_db("default", db.clone())) Arc::new(Workspace::new_with_db("default", db.clone()))
} }
/// 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 .config");
// Read back to get the document ID
let doc = workspace
.read(&config_path)
.await
.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("set metadata");
}
#[tokio::test] #[tokio::test]
async fn cleanup_directory_skips_config_files() { async fn cleanup_daily_logs_preserves_identity_documents() {
let (db, _tmp) = create_test_db().await; let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db); let ws = create_workspace(&db);
// Write documents including a .config // Write several regular documents (non-identity)
ws.write("daily/2024-01-15.md", "Old log") ws.write("daily/2024-01-15.md", "Old log")
.await .await
.expect("write log"); .expect("write log 1");
ws.write("daily/.config", "").await.expect("write config"); ws.write("daily/2024-01-20.md", "Another log")
.await
.expect("write log 2");
// Write an identity document
ws.write("MEMORY.md", "Long-term curated memory")
.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) // Run cleanup with 0-day retention (deletes everything old)
let deleted = cleanup_directory(&ws, "daily/", 0) // This tests that even with aggressive cleanup, identity docs survive
let deleted = cleanup_daily_logs(&ws, 0)
.await .await
.expect("cleanup_directory"); .expect("cleanup_daily_logs");
// Should have deleted the log but not the .config // Should have deleted some documents (the daily logs)
assert!(deleted > 0, "should have deleted old daily documents"); assert!(deleted > 0, "should have deleted old daily documents");
// Verify .config still exists // Verify identity doc still exists
let config_doc = db let identity = db
.get_document_by_path("default", None, "daily/.config") .get_document_by_path("default", None, "MEMORY.md")
.await .await
.expect("get .config doc"); .expect("get identity doc");
assert_eq!(config_doc.path, "daily/.config"); assert_eq!(identity.path, "MEMORY.md");
assert_eq!(identity.content, "Long-term curated memory");
} }
#[tokio::test] #[tokio::test]
async fn cleanup_directory_handles_empty_directory() { async fn cleanup_conversation_docs_handles_empty_directory() {
let (db, _tmp) = create_test_db().await; let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db); let ws = create_workspace(&db);
// Run cleanup on an empty directory // Run cleanup on an empty directory (conversations/ doesn't exist)
let deleted = cleanup_directory(&ws, "conversations/", 7) let deleted = cleanup_conversation_docs(&ws, 7)
.await .await
.expect("cleanup_directory"); .expect("cleanup_conversation_docs");
// Should delete 0 (nothing to delete)
assert_eq!(deleted, 0, "should delete 0 from empty directory"); 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] #[tokio::test]
async fn cleanup_respects_cadence_prevents_concurrent_runs() { async fn cleanup_respects_cadence_prevents_concurrent_runs() {
let (db, _tmp) = create_test_db().await; let (db, _tmp) = create_test_db().await;
@@ -555,7 +614,8 @@ mod tests {
let config = HygieneConfig { let config = HygieneConfig {
enabled: true, enabled: true,
version_keep_count: 50, daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 12, cadence_hours: 12,
state_dir: _tmp.path().to_path_buf(), state_dir: _tmp.path().to_path_buf(),
}; };
@@ -567,6 +627,13 @@ mod tests {
// Second run immediately should be skipped (cadence not elapsed) // Second run immediately should be skipped (cadence not elapsed)
let report2 = run_if_due(&ws, &config).await; let report2 = run_if_due(&ws, &config).await;
assert!(report2.skipped, "second run should be skipped by cadence"); 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] #[tokio::test]
@@ -574,10 +641,6 @@ mod tests {
let (db, _tmp) = create_test_db().await; let (db, _tmp) = create_test_db().await;
let ws = create_workspace(&db); 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 // Write some documents
ws.write("daily/log1.md", "content 1") ws.write("daily/log1.md", "content 1")
.await .await
@@ -589,34 +652,35 @@ mod tests {
.await .await
.expect("write doc 3"); .expect("write doc 3");
// Run with 0-day retention via direct cleanup_directory calls // Run with 0-day retention to delete everything non-identity
let deleted_daily = cleanup_directory(&ws, "daily/", 0) let deleted_daily = cleanup_daily_logs(&ws, 0).await.expect("cleanup daily");
.await let deleted_conv = cleanup_conversation_docs(&ws, 0)
.expect("cleanup daily");
let deleted_conv = cleanup_directory(&ws, "conversations/", 0)
.await .await
.expect("cleanup conversations"); .expect("cleanup conversations");
// Both should report deletions
assert!(deleted_daily > 0, "should report deleted daily logs"); assert!(deleted_daily > 0, "should report deleted daily logs");
assert_eq!(deleted_conv, 1, "should report 1 deleted conversation doc"); assert_eq!(deleted_conv, 1, "should report 1 deleted conversation doc");
// Verify HygieneReport aggregation // Create a HygieneReport and verify aggregation works
let report = HygieneReport { let report = HygieneReport {
directories_cleaned: vec![ daily_logs_deleted: deleted_daily,
("daily/".to_string(), deleted_daily), conversation_docs_deleted: deleted_conv,
("conversations/".to_string(), deleted_conv),
],
versions_pruned: 0,
skipped: false, skipped: false,
}; };
// Verify HygieneReport structure
assert!(!report.skipped, "should not be skipped"); assert!(!report.skipped, "should not be skipped");
assert!(report.had_work(), "report should indicate work was done"); 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 checks directory counts // Verify had_work() correctly combines both counts
let no_work = HygieneReport { let no_work = HygieneReport {
directories_cleaned: vec![], daily_logs_deleted: 0,
versions_pruned: 0, conversation_docs_deleted: 0,
skipped: false, skipped: false,
}; };
assert!(!no_work.had_work(), "empty report should indicate no work"); assert!(!no_work.had_work(), "empty report should indicate no work");
+2 -351
View File
@@ -53,9 +53,8 @@ mod search;
pub use chunker::{ChunkConfig, chunk_document}; pub use chunker::{ChunkConfig, chunk_document};
pub use document::{ pub use document::{
CONFIG_FILE_NAME, DocumentMetadata, DocumentVersion, HygieneMetadata, IDENTITY_PATHS, IDENTITY_PATHS, MemoryChunk, MemoryDocument, WorkspaceEntry, is_identity_path,
MemoryChunk, MemoryDocument, PatchResult, VersionSummary, WorkspaceEntry, content_sha256, merge_workspace_entries, paths,
is_config_path, is_identity_path, merge_workspace_entries, paths,
}; };
pub use embedding_cache::{CachedEmbeddingProvider, EmbeddingCacheConfig}; pub use embedding_cache::{CachedEmbeddingProvider, EmbeddingCacheConfig};
pub use embeddings::{ pub use embeddings::{
@@ -367,101 +366,6 @@ 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<Uuid>,
) -> Result<Vec<MemoryDocument>, 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<i32, WorkspaceError> {
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<DocumentVersion, WorkspaceError> {
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<Vec<VersionSummary>, 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<Option<i32>, 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<u64, WorkspaceError> {
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. /// Default template seeded into HEARTBEAT.md on first access.
@@ -792,199 +696,10 @@ impl Workspace {
.await .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<u64, WorkspaceError> {
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<Vec<MemoryDocument>, 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 = &current[..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<Vec<VersionSummary>, 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<DocumentVersion, WorkspaceError> {
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<Option<i32>, 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<PatchResult, WorkspaceError> {
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. /// Write (create or update) a file.
/// ///
/// Creates parent directories implicitly (they're virtual in the DB). /// Creates parent directories implicitly (they're virtual in the DB).
/// Re-indexes the document for search after writing. /// Re-indexes the document for search after writing.
/// Auto-versions the previous content before overwriting.
/// ///
/// # Example /// # Example
/// ```ignore /// ```ignore
@@ -1000,12 +715,6 @@ impl Workspace {
.storage .storage
.get_or_create_document_by_path(&self.user_id, self.agent_id, &path) .get_or_create_document_by_path(&self.user_id, self.agent_id, &path)
.await?; .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.storage.update_document(doc.id, content).await?;
self.reindex_document(doc.id).await?; self.reindex_document(doc.id).await?;
@@ -1045,11 +754,6 @@ impl Workspace {
reject_if_injected(&path, &new_content)?; 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.storage.update_document(doc.id, &new_content).await?;
self.reindex_document(doc.id).await?; self.reindex_document(doc.id).await?;
Ok(()) Ok(())
@@ -1881,14 +1585,6 @@ impl Workspace {
// Get the document // Get the document
let doc = self.storage.get_document_by_id(document_id).await?; 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 // Chunk the content
let chunks = chunk_document(&doc.content, ChunkConfig::default()); let chunks = chunk_document(&doc.content, ChunkConfig::default());
@@ -1974,51 +1670,6 @@ 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 // BOOTSTRAP.md is only seeded on truly fresh workspaces (no identity
// files existed before seeding) AND when no profile exists yet (the user // 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 // may already have a profile from a previous install and doesn't need
+1 -191
View File
@@ -11,9 +11,7 @@ use uuid::Uuid;
use crate::error::WorkspaceError; use crate::error::WorkspaceError;
use crate::workspace::document::{ use crate::workspace::document::{MemoryChunk, MemoryDocument, WorkspaceEntry};
DocumentVersion, MemoryChunk, MemoryDocument, VersionSummary, WorkspaceEntry,
};
use crate::workspace::search::{RankedResult, SearchConfig, SearchResult, fuse_results}; use crate::workspace::search::{RankedResult, SearchConfig, SearchResult, fuse_results};
/// Database repository for workspace operations. /// Database repository for workspace operations.
@@ -704,192 +702,4 @@ impl Repository {
} }
Ok(crate::workspace::merge_workspace_entries(all_entries)) 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<Uuid>,
) -> Result<Vec<MemoryDocument>, 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<i32, WorkspaceError> {
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<DocumentVersion, WorkspaceError> {
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<Vec<VersionSummary>, 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<Option<i32>, 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<u64, WorkspaceError> {
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)
}
} }
+4 -2
View File
@@ -955,7 +955,8 @@ mod tests {
let hygiene_config = HygieneConfig { let hygiene_config = HygieneConfig {
enabled: false, enabled: false,
version_keep_count: 50, daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 24, cadence_hours: 24,
state_dir: _tmp.path().to_path_buf(), state_dir: _tmp.path().to_path_buf(),
}; };
@@ -1001,7 +1002,8 @@ mod tests {
let hygiene_config = HygieneConfig { let hygiene_config = HygieneConfig {
enabled: false, enabled: false,
version_keep_count: 50, daily_retention_days: 30,
conversation_retention_days: 7,
cadence_hours: 24, cadence_hours: 24,
state_dir: _tmp.path().to_path_buf(), state_dir: _tmp.path().to_path_buf(),
}; };