已开启
feat(datalog): add retention policy (max_age_days / max_total_size_mb) #1090
TrueFurina创建于 21 天前
feat(datalog): add retention policy (max_age_days / max_total_size_mb) #1090
已开启
共 3 个文件变更+375-8
| @@ -88,14 +88,32 @@ impl DatalogHook { | |||
| 88 | model: impl Into<String>, | 88 | model: impl Into<String>, |
| 89 | context_window: u32, | 89 | context_window: u32, |
| 90 | ) -> Option<Self> { | 90 | ) -> Option<Self> { |
| 91 | - config.enabled.then(|| Self { | 91 | + config.enabled.then(|| { |
| 92 | - working_dir: working_dir.into(), | 92 | + let hook = Self { |
| 93 | - configured_dir: config.dir.clone(), | 93 | + working_dir: working_dir.into(), |
| 94 | - model: model.into(), | 94 | + configured_dir: config.dir.clone(), |
| 95 | - context_window, | 95 | + model: model.into(), |
| 96 | - state: Mutex::new(TurnLog::default()), | 96 | + context_window, |
| 97 | - writer: DatalogWriter::start(), | 97 | + state: Mutex::new(TurnLog::default()), |
| 98 | - instance_id: HOOK_SEQUENCE.fetch_add(1, Ordering::Relaxed), | 98 | + writer: DatalogWriter::start(), |
| 99 | + instance_id: HOOK_SEQUENCE.fetch_add(1, Ordering::Relaxed), | ||
| 100 | + }; | ||
| 101 | + // One-shot retention pass, executed on the writer thread so it | ||
| 102 | + // never delays turn handling. Enqueued first, so it completes | ||
| 103 | + // before any log writes from this hook. | ||
| 104 | + if config.max_age_days.is_some() || config.max_total_size_mb.is_some() { | ||
| 105 | + let root = Self::resolve_log_dir(&hook.working_dir, config.dir.as_deref()); | ||
| 106 | + // Walk the datalog ROOT (all project buckets), not just this | ||
| 107 | + // project's slug directory. | ||
| 108 | + let root = root | ||
| 109 | + .ancestors() | ||
| 110 | + .nth(1) | ||
| 111 | + .map(Path::to_path_buf) | ||
| 112 | + .unwrap_or(root); | ||
| 113 | + hook.writer | ||
| 114 | + .retain(root, config.max_age_days, config.max_total_size_mb); | ||
| 115 | + } | ||
| 116 | + hook | ||
| 99 | }) | 117 | }) |
| 100 | } | 118 | } |
| 101 | 119 | ||
| @@ -420,6 +438,19 @@ impl DatalogWriter { | |||
| 420 | let _ = self.tx.send(WriteOp::Append { path, content }); | 438 | let _ = self.tx.send(WriteOp::Append { path, content }); |
| 421 | } | 439 | } |
| 422 | 440 | ||
| 441 | + /// Run the retention pass on its own detached thread (PR #1090 review, | ||
| 442 | + /// P2). NOT routed through the writer queue: a slow scan over a huge | ||
| 443 | + /// datalog root would sit ahead of the first Initialize in the queue and | ||
| 444 | + /// burn through its 500ms `IO_WAIT_TIMEOUT`, silently losing the whole | ||
| 445 | + /// first turn's log. A dedicated thread keeps the queue write-only; the | ||
| 446 | + /// scan only ever touches files older than the cutoff, so it cannot race | ||
| 447 | + /// with the new files this hook is about to create. | ||
| 448 | + fn retain(&self, root: PathBuf, max_age_days: Option<u64>, max_total_size_mb: Option<u64>) { | ||
| 449 | + let _ = std::thread::Builder::new() | ||
| 450 | + .name("atomcode-datalog-retain".into()) | ||
| 451 | + .spawn(move || apply_retention(&root, max_age_days, max_total_size_mb)); | ||
| 452 | + } | ||
| 453 | + | ||
| 423 | async fn barrier(&self) { | 454 | async fn barrier(&self) { |
| 424 | let (reply, receive) = tokio::sync::oneshot::channel(); | 455 | let (reply, receive) = tokio::sync::oneshot::channel(); |
| 425 | let _ = self.tx.send(WriteOp::Barrier { reply }); | 456 | let _ = self.tx.send(WriteOp::Barrier { reply }); |
| @@ -454,6 +485,104 @@ fn writer_loop(rx: mpsc::Receiver<WriteOp>) { | |||
| 454 | } | 485 | } |
| 455 | } | 486 | } |
| 456 | 487 | ||
| 488 | +/// One-shot retention pass over the whole datalog root. All failures are | ||
| 489 | +/// ignored: retention is best-effort cleanup and must never affect a turn. | ||
| 490 | +/// | ||
| 491 | +/// Order matters for the size cap: age deletions run first so expired files | ||
| 492 | +/// are not counted against the budget, then oldest-first deletion brings the | ||
| 493 | +/// total back under `max_total_size_mb` if it still exceeds it. | ||
| 494 | +fn apply_retention(root: &Path, max_age_days: Option<u64>, max_total_size_mb: Option<u64>) { | ||
| 495 | + if max_age_days.is_none() && max_total_size_mb.is_none() { | ||
| 496 | + return; | ||
| 497 | + } | ||
| 498 | + let Ok(entries) = collect_datalog_files(root) else { | ||
| 499 | + return; | ||
| 500 | + }; | ||
| 501 | + // Newest last, so popping from the end in the size pass skips survivors. | ||
| 502 | + let mut entries: Vec<(PathBuf, std::time::SystemTime, u64)> = entries; | ||
| 503 | + entries.sort_by_key(|(_, modified, _)| *modified); | ||
| 504 | + | ||
| 505 | + if let Some(days) = max_age_days { | ||
| 506 | + let cutoff = std::time::SystemTime::now() | ||
| 507 | + .checked_sub(std::time::Duration::from_secs(days.saturating_mul(86_400))) | ||
| 508 | + .unwrap_or(std::time::SystemTime::UNIX_EPOCH); | ||
| 509 | + entries.retain(|(path, modified, _)| { | ||
| 510 | + if *modified >= cutoff { | ||
| 511 | + return true; | ||
| 512 | + } | ||
| 513 | + let _ = fs::remove_file(path); | ||
| 514 | + false | ||
| 515 | + }); | ||
| 516 | + } | ||
| 517 | + | ||
| 518 | + if let Some(max_bytes) = max_total_size_mb.and_then(|mb| mb.checked_mul(1024 * 1024)) { | ||
| 519 | + let mut total: u64 = entries.iter().map(|(_, _, size)| size).sum(); | ||
| 520 | + if total > max_bytes { | ||
| 521 | + // Drain oldest-first until back under the cap. | ||
| 522 | + let mut index = 0; | ||
| 523 | + while total > max_bytes && index < entries.len() { | ||
| 524 | + let (path, _, size) = &entries[index]; | ||
| 525 | + if fs::remove_file(path).is_ok() { | ||
| 526 | + total = total.saturating_sub(*size); | ||
| 527 | + } | ||
| 528 | + index += 1; | ||
| 529 | + } | ||
| 530 | + } | ||
| 531 | + } | ||
| 532 | +} | ||
| 533 | + | ||
| 534 | +/// Collect every `.md`/`.jsonl` datalog file under `root` as | ||
| 535 | +/// `(path, last_modified, size)`. | ||
| 536 | +/// | ||
| 537 | +/// Scope is deliberately narrow so retention can never expand into arbitrary | ||
| 538 | +/// user data (PR #1090 review, P1): only files sitting DIRECTLY inside a | ||
| 539 | +/// first-level subdirectory of `root` whose name matches the project-bucket | ||
| 540 | +/// shape `<slug>-<8 hex chars>` are considered. No recursion — datalog files | ||
| 541 | +/// are always written one level deep (`initialize_files` joins the bucket), | ||
| 542 | +/// and anything else (loose files at the root, nested user trees, buckets | ||
| 543 | +/// with unexpected names) is left untouched. This matters because the root | ||
| 544 | +/// is derived from the user-configured `dir`, which may legitimately be | ||
| 545 | +/// `~`, `.`, or a shared directory full of unrelated `.md`/`.jsonl` files. | ||
| 546 | +fn collect_datalog_files(root: &Path) -> std::io::Result<Vec<(PathBuf, std::time::SystemTime, u64)>> { | ||
| 547 | + let mut files = Vec::new(); | ||
| 548 | + for entry in fs::read_dir(root)?.flatten() { | ||
| 549 | + // Only first-level directories count as project buckets. | ||
| 550 | + if !entry.file_type().map(|t| t.is_dir()).unwrap_or(false) { | ||
| 551 | + continue; | ||
| 552 | + } | ||
| 553 | + if !is_project_bucket_name(&entry.file_name().to_string_lossy()) { | ||
| 554 | + continue; | ||
| 555 | + } | ||
| 556 | + for file in fs::read_dir(entry.path())?.flatten() { | ||
| 557 | + if file.file_type().map(|t| t.is_dir()).unwrap_or(true) { | ||
| 558 | + continue; | ||
| 559 | + } | ||
| 560 | + let path = file.path(); | ||
| 561 | + let extension = path.extension().and_then(|ext| ext.to_str()); | ||
| 562 | + if extension != Some("md") && extension != Some("jsonl") { | ||
| 563 | + continue; | ||
| 564 | + } | ||
| 565 | + if let Ok(metadata) = file.metadata() { | ||
| 566 | + if let Ok(modified) = metadata.modified() { | ||
| 567 | + files.push((path, modified, metadata.len())); | ||
| 568 | + } | ||
| 569 | + } | ||
| 570 | + } | ||
| 571 | + } | ||
| 572 | + Ok(files) | ||
| 573 | +} | ||
| 574 | + | ||
| 575 | +/// A bucket directory is `<sanitized-slug>-<8 lowercase hex>` (see | ||
| 576 | +/// [`project_slug`]/[`atomcode_config::util::stable_project_hash`]). Anything | ||
| 577 | +/// else — `Documents`, `notes`, `my.logs` — is not a datalog bucket and its | ||
| 578 | +/// contents must never be touched, even when they end in `.md`/`.jsonl`. | ||
| 579 | +fn is_project_bucket_name(name: &str) -> bool { | ||
| 580 | + let Some((_, hash)) = name.rsplit_once('-') else { | ||
| 581 | + return false; | ||
| 582 | + }; | ||
| 583 | + hash.len() == 8 && hash.bytes().all(|b| b.is_ascii_hexdigit()) | ||
| 584 | +} | ||
| 585 | + | ||
| 457 | fn initialize_files( | 586 | fn initialize_files( |
| 458 | directory: &Path, | 587 | directory: &Path, |
| 459 | filename_stem: &str, | 588 | filename_stem: &str, |
| @@ -579,6 +708,8 @@ mod tests { | |||
| 579 | let config = DatalogConfig { | 708 | let config = DatalogConfig { |
| 580 | enabled: false, | 709 | enabled: false, |
| 581 | dir: None, | 710 | dir: None, |
| 711 | + max_age_days: None, | ||
| 712 | + max_total_size_mb: None, | ||
| 582 | }; | 713 | }; |
| 583 | assert!(DatalogHook::new("/repo", &config, "model", 128_000).is_none()); | 714 | assert!(DatalogHook::new("/repo", &config, "model", 128_000).is_none()); |
| 584 | } | 715 | } |
| @@ -683,6 +814,8 @@ mod tests { | |||
| 683 | let config = DatalogConfig { | 814 | let config = DatalogConfig { |
| 684 | enabled: true, | 815 | enabled: true, |
| 685 | dir: Some(output.display().to_string()), | 816 | dir: Some(output.display().to_string()), |
| 817 | + max_age_days: None, | ||
| 818 | + max_total_size_mb: None, | ||
| 686 | }; | 819 | }; |
| 687 | let hook = DatalogHook::new(&project, &config, "test-model", 128_000).unwrap(); | 820 | let hook = DatalogHook::new(&project, &config, "test-model", 128_000).unwrap(); |
| 688 | 821 | ||
| @@ -763,6 +896,8 @@ mod tests { | |||
| 763 | let config = DatalogConfig { | 896 | let config = DatalogConfig { |
| 764 | enabled: true, | 897 | enabled: true, |
| 765 | dir: Some(output.display().to_string()), | 898 | dir: Some(output.display().to_string()), |
| 899 | + max_age_days: None, | ||
| 900 | + max_total_size_mb: None, | ||
| 766 | }; | 901 | }; |
| 767 | let first = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); | 902 | let first = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); |
| 768 | let second = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); | 903 | let second = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); |
| @@ -819,6 +954,8 @@ mod tests { | |||
| 819 | let config = DatalogConfig { | 954 | let config = DatalogConfig { |
| 820 | enabled: true, | 955 | enabled: true, |
| 821 | dir: Some(output.display().to_string()), | 956 | dir: Some(output.display().to_string()), |
| 957 | + max_age_days: None, | ||
| 958 | + max_total_size_mb: None, | ||
| 822 | }; | 959 | }; |
| 823 | let hook = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); | 960 | let hook = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); |
| 824 | hook.user_prompt_submit(&mut "prompt".to_string()) | 961 | hook.user_prompt_submit(&mut "prompt".to_string()) |
| @@ -853,4 +990,186 @@ mod tests { | |||
| 853 | assert_eq!(mode, 0o600); | 990 | assert_eq!(mode, 0o600); |
| 854 | } | 991 | } |
| 855 | } | 992 | } |
| 993 | + | ||
| 994 | + /// Backdate a file's mtime so age-based retention can be tested without | ||
| 995 | + /// sleeping. | ||
| 996 | + fn backdate(path: &Path, days_ago: u64) { | ||
| 997 | + let old = std::time::SystemTime::now() | ||
| 998 | + .checked_sub(std::time::Duration::from_secs(days_ago * 86_400)) | ||
| 999 | + .unwrap(); | ||
| 1000 | + let file = File::options().write(true).open(path).unwrap(); | ||
| 1001 | + file.set_modified(old).unwrap(); | ||
| 1002 | + } | ||
| 1003 | + | ||
| 1004 | + fn seed(root: &Path, name: &str, contents: &[u8]) -> PathBuf { | ||
| 1005 | + // Bucket name must match the `<slug>-<8 hex>` shape that | ||
| 1006 | + // `is_project_bucket_name` accepts (see the P1 fix). | ||
| 1007 | + let dir = root.join("project-1a2b3c4d"); | ||
| 1008 | + fs::create_dir_all(&dir).unwrap(); | ||
| 1009 | + let path = dir.join(name); | ||
| 1010 | + fs::write(&path, contents).unwrap(); | ||
| 1011 | + path | ||
| 1012 | + } | ||
| 1013 | + | ||
| 1014 | + | ||
| 1015 | + fn retention_deletes_files_older_than_max_age_days() { | ||
| 1016 | + let root = tempdir().unwrap(); | ||
| 1017 | + let old = seed(root.path(), "old.jsonl", b"x"); | ||
| 1018 | + let fresh = seed(root.path(), "fresh.md", b"y"); | ||
| 1019 | + backdate(&old, 10); | ||
| 1020 | + | ||
| 1021 | + apply_retention(root.path(), Some(7), None); | ||
| 1022 | + | ||
| 1023 | + assert!(!old.exists(), "10-day-old file must be deleted at 7-day cutoff"); | ||
| 1024 | + assert!(fresh.exists(), "fresh file must survive"); | ||
| 1025 | + } | ||
| 1026 | + | ||
| 1027 | + | ||
| 1028 | + fn retention_size_cap_deletes_oldest_files_first() { | ||
| 1029 | + let root = tempdir().unwrap(); | ||
| 1030 | + let a = seed(root.path(), "a.jsonl", &vec![b'a'; 600_000]); | ||
| 1031 | + let b = seed(root.path(), "b.jsonl", &vec![b'b'; 600_000]); | ||
| 1032 | + let c = seed(root.path(), "c.md", &vec![b'c'; 600_000]); | ||
| 1033 | + // Force a deterministic age order: a oldest, c newest. | ||
| 1034 | + backdate(&a, 3); | ||
| 1035 | + backdate(&b, 2); | ||
| 1036 | + // c keeps its (newest) current mtime. | ||
| 1037 | + | ||
| 1038 | + // 1.8 MB total; 1 MB cap must evict a and b, keep c. | ||
| 1039 | + apply_retention(root.path(), None, Some(1)); | ||
| 1040 | + | ||
| 1041 | + assert!(!a.exists(), "oldest file must be evicted first"); | ||
| 1042 | + assert!(!b.exists(), "second-oldest file must be evicted while over cap"); | ||
| 1043 | + assert!(c.exists(), "newest file must survive under the cap"); | ||
| 1044 | + } | ||
| 1045 | + | ||
| 1046 | + | ||
| 1047 | + fn retention_without_limits_is_a_no_op() { | ||
| 1048 | + let root = tempdir().unwrap(); | ||
| 1049 | + let file = seed(root.path(), "keep.jsonl", b"data"); | ||
| 1050 | + | ||
| 1051 | + apply_retention(root.path(), None, None); | ||
| 1052 | + | ||
| 1053 | + assert!(file.exists()); | ||
| 1054 | + } | ||
| 1055 | + | ||
| 1056 | + | ||
| 1057 | + fn retention_leaves_non_log_files_alone() { | ||
| 1058 | + let root = tempdir().unwrap(); | ||
| 1059 | + let stray = seed(root.path(), "notes.txt", b"keep me"); | ||
| 1060 | + backdate(&stray, 30); | ||
| 1061 | + | ||
| 1062 | + apply_retention(root.path(), Some(7), None); | ||
| 1063 | + | ||
| 1064 | + assert!(stray.exists(), "only .md/.jsonl datalog files may be deleted"); | ||
| 1065 | + } | ||
| 1066 | + | ||
| 1067 | + /// Integration test for the P1 fix (PR #1090 review): the retention root | ||
| 1068 | + /// is derived from the user-configured `dir`, which may be a directory | ||
| 1069 | + /// full of unrelated user data. The pass must only delete `.md`/`.jsonl` | ||
| 1070 | + /// files sitting inside a `<slug>-<8hex>` bucket at the FIRST level — | ||
| 1071 | + /// never loose files at the root, nested user trees, or look-alike | ||
| 1072 | + /// directories whose names don't match the bucket shape. | ||
| 1073 | + | ||
| 1074 | + async fn retention_via_hook_never_touches_non_bucket_user_files() { | ||
| 1075 | + let root = tempdir().unwrap(); | ||
| 1076 | + let project = root.path().join("project"); | ||
| 1077 | + fs::create_dir_all(&project).unwrap(); | ||
| 1078 | + let output = root.path().join("datalog-root"); | ||
| 1079 | + let config = DatalogConfig { | ||
| 1080 | + enabled: true, | ||
| 1081 | + dir: Some(output.display().to_string()), | ||
| 1082 | + max_age_days: Some(0), | ||
| 1083 | + max_total_size_mb: None, | ||
| 1084 | + }; | ||
| 1085 | + | ||
| 1086 | + // Real bucket: `<slug>-<8 hex>`, containing an old datalog file. | ||
| 1087 | + let bucket = output.join("project-1a2b3c4d"); | ||
| 1088 | + fs::create_dir_all(&bucket).unwrap(); | ||
| 1089 | + fs::write(bucket.join("old.jsonl"), b"datalog").unwrap(); | ||
| 1090 | + | ||
| 1091 | + // User data that must survive, in every dangerous shape: | ||
| 1092 | + let user_note = output.join("my-notes.md"); // loose .md at the root | ||
| 1093 | + fs::write(&user_note, b"user note").unwrap(); | ||
| 1094 | + let lookalike = output.join("documents-e1a2b3c4"); // 8-hex name BUT a file, not a dir | ||
| 1095 | + fs::write(&lookalike, b"actually a file").unwrap(); | ||
| 1096 | + let nested_dir = output.join("Documents"); // non-bucket dir… | ||
| 1097 | + fs::create_dir_all(nested_dir.join("deep")).unwrap(); | ||
| 1098 | + let nested = nested_dir.join("deep").join("report.md"); // …with nested .md | ||
| 1099 | + fs::write(&nested, b"nested doc").unwrap(); | ||
| 1100 | + let bad_hash = output.join("project-nothex"); // slug-ish name, hash not 8 hex | ||
| 1101 | + fs::create_dir_all(&bad_hash).unwrap(); | ||
| 1102 | + let bad_hash_file = bad_hash.join("file.jsonl"); | ||
| 1103 | + fs::write(&bad_hash_file, b"not a bucket").unwrap(); | ||
| 1104 | + | ||
| 1105 | + let hook = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); | ||
| 1106 | + hook.turn_complete(&Conversation::new(), &StopReason::Stopped, &TurnCtx::default()) | ||
| 1107 | + .await; | ||
| 1108 | + hook.writer.barrier().await; | ||
| 1109 | + // Retention runs on its own thread; give it a moment to finish. | ||
| 1110 | + for _ in 0..50 { | ||
| 1111 | + if !bucket.join("old.jsonl").exists() { | ||
| 1112 | + break; | ||
| 1113 | + } | ||
| 1114 | + std::thread::sleep(std::time::Duration::from_millis(20)); | ||
| 1115 | + } | ||
| 1116 | + | ||
| 1117 | + assert!(!bucket.join("old.jsonl").exists(), | ||
| 1118 | + "datlog file inside a real bucket must be deleted"); | ||
| 1119 | + assert!(user_note.exists(), "loose .md at the root is user data, not datalog"); | ||
| 1120 | + assert!(lookalike.exists(), "8-hex-named FILE is not a bucket"); | ||
| 1121 | + assert!(nested.exists(), "recursion into non-bucket dirs is forbidden"); | ||
| 1122 | + assert!(bad_hash_file.exists(), "dir whose hash is not 8 hex is not a bucket"); | ||
| 1123 | + } | ||
| 1124 | + | ||
| 1125 | + /// Integration test for the P2 fix (PR #1090 review): retention must NOT | ||
| 1126 | + /// occupy the writer queue ahead of the first Initialize — a turn's log | ||
| 1127 | + /// files must be created and written even while a retention pass is | ||
| 1128 | + /// pending. Before the fix, Retain sat at the head of the queue and a | ||
| 1129 | + /// slow scan could burn Initialize's 500ms IO_WAIT_TIMEOUT. | ||
| 1130 | + | ||
| 1131 | + async fn retention_does_not_block_first_turn_log_creation() { | ||
| 1132 | + let root = tempdir().unwrap(); | ||
| 1133 | + let project = root.path().join("project"); | ||
| 1134 | + fs::create_dir_all(&project).unwrap(); | ||
| 1135 | + let output = root.path().join("logs"); | ||
| 1136 | + let config = DatalogConfig { | ||
| 1137 | + enabled: true, | ||
| 1138 | + dir: Some(output.display().to_string()), | ||
| 1139 | + max_age_days: None, | ||
| 1140 | + max_total_size_mb: None, | ||
| 1141 | + }; | ||
| 1142 | + let hook = DatalogHook::new(&project, &config, "model", 128_000).unwrap(); | ||
| 1143 | + | ||
| 1144 | + hook.user_prompt_submit(&mut "prompt".to_string()).await.unwrap(); | ||
| 1145 | + let ctx = TurnCtx { | ||
| 1146 | + turn_id: 1, | ||
| 1147 | + request_id: 1, | ||
| 1148 | + round: 1, | ||
| 1149 | + ..TurnCtx::default() | ||
| 1150 | + }; | ||
| 1151 | + hook.on_request(&[Message::user("prompt")], &[], &ChatOptions::default(), &ctx) | ||
| 1152 | + .await; | ||
| 1153 | + hook.turn_complete(&Conversation::new(), &StopReason::Stopped, &ctx) | ||
| 1154 | + .await; | ||
| 1155 | + | ||
| 1156 | + // Initialize's reply is bounded by IO_WAIT_TIMEOUT (500ms); the | ||
| 1157 | + // barrier below would time out and the markdown would be missing if | ||
| 1158 | + // the retention pass had blocked the queue. | ||
| 1159 | + hook.writer.barrier().await; | ||
| 1160 | + let project_dir = fs::read_dir(output) | ||
| 1161 | + .unwrap() | ||
| 1162 | + .next() | ||
| 1163 | + .unwrap() | ||
| 1164 | + .unwrap() | ||
| 1165 | + .path(); | ||
| 1166 | + let files: Vec<PathBuf> = fs::read_dir(project_dir) | ||
| 1167 | + .unwrap() | ||
| 1168 | + .map(|entry| entry.unwrap().path()) | ||
| 1169 | + .collect(); | ||
| 1170 | + assert!(files.iter().any(|p| p.extension().and_then(|e| e.to_str()) == Some("md")), | ||
| 1171 | + "first turn's markdown must exist even with retention enabled"); | ||
| 1172 | + assert!(files.iter().any(|p| p.extension().and_then(|e| e.to_str()) == Some("jsonl")), | ||
| 1173 | + "first turn's jsonl must exist even with retention enabled"); | ||
| 1174 | + } | ||
| 856 | } | 1175 | } |
| @@ -956,6 +956,8 @@ mod tests { | |||
| 956 | source.datalog = atomcode_config::config::DatalogConfig { | 956 | source.datalog = atomcode_config::config::DatalogConfig { |
| 957 | enabled: false, | 957 | enabled: false, |
| 958 | dir: Some("/var/tmp/atomcode-datalog".into()), | 958 | dir: Some("/var/tmp/atomcode-datalog".into()), |
| 959 | + max_age_days: None, | ||
| 960 | + max_total_size_mb: None, | ||
| 959 | }; | 961 | }; |
| 960 | let runtime = CodingRuntimeConfig::from_config( | 962 | let runtime = CodingRuntimeConfig::from_config( |
| 961 | &source, | 963 | &source, |
| @@ -1452,6 +1452,24 @@ pub struct DatalogConfig { | |||
| 1452 | /// - Relative path → resolved against working_dir, follows /cd | 1452 | /// - Relative path → resolved against working_dir, follows /cd |
| 1453 | 1453 | ||
| 1454 | pub dir: Option<String>, | 1454 | pub dir: Option<String>, |
| 1455 | + /// Maximum age of datalog files, in days. Applied to `.md`/`.jsonl` | ||
| 1456 | + /// files under the whole datalog root (all project buckets): files whose | ||
| 1457 | + /// last modification is older than this are deleted once, when a datalog | ||
| 1458 | + /// hook is created. `None` (default) keeps files forever. | ||
| 1459 | + /// | ||
| 1460 | + /// This exists because datalog JSONL appends a full request snapshot per | ||
| 1461 | + /// round, so long sessions accumulate tens of GB within weeks (issue | ||
| 1462 | + /// #1551). A value of `0` treats every existing file as expired and | ||
| 1463 | + /// clears them on startup. | ||
| 1464 | + | ||
| 1465 | + pub max_age_days: Option<u64>, | ||
| 1466 | + /// Soft cap on the total size of `.md`/`.jsonl` datalog files under the | ||
| 1467 | + /// whole datalog root, in megabytes. When the total exceeds the cap, the | ||
| 1468 | + /// oldest files are deleted first until it is back under. `None` | ||
| 1469 | + /// (default) applies no cap. Evaluated once when a datalog hook is | ||
| 1470 | + /// created, not on every write. | ||
| 1471 | + | ||
| 1472 | + pub max_total_size_mb: Option<u64>, | ||
| 1455 | } | 1473 | } |
| 1456 | 1474 | ||
| 1457 | /// Controls long-running task completion notifications. | 1475 | /// Controls long-running task completion notifications. |
| @@ -1661,6 +1679,10 @@ impl Default for DatalogConfig { | |||
| 1661 | // discover that "unset == ~/.atomcode/datalog". Resolver still | 1679 | // discover that "unset == ~/.atomcode/datalog". Resolver still |
| 1662 | // treats this string the same as `None` (project slug appended). | 1680 | // treats this string the same as `None` (project slug appended). |
| 1663 | dir: Some("~/.atomcode/datalog".to_string()), | 1681 | dir: Some("~/.atomcode/datalog".to_string()), |
| 1682 | + // Retention is off by default: existing installs must not have | ||
| 1683 | + // files deleted from under them by an upgrade. | ||
| 1684 | + max_age_days: None, | ||
| 1685 | + max_total_size_mb: None, | ||
| 1664 | } | 1686 | } |
| 1665 | } | 1687 | } |
| 1666 | } | 1688 | } |
| @@ -1701,11 +1723,29 @@ fn render_datalog_section(cfg: &DatalogConfig) -> String { | |||
| 1701 | ); | 1723 | ); |
| 1702 | out.push_str("# - dir = \"/abs/path\" -> absolute, fixed (unaffected by /cd)\n"); | 1724 | out.push_str("# - dir = \"/abs/path\" -> absolute, fixed (unaffected by /cd)\n"); |
| 1703 | out.push_str("# - dir = \"rel/path\" -> joined with current working_dir, follows /cd\n"); | 1725 | out.push_str("# - dir = \"rel/path\" -> joined with current working_dir, follows /cd\n"); |
| 1726 | + out.push_str( | ||
| 1727 | + "# Retention (applied once per startup, to .md/.jsonl files under the whole\n", | ||
| 1728 | + ); | ||
| 1729 | + out.push_str( | ||
| 1730 | + "# datalog root, oldest first; 0 clears everything on startup):\n", | ||
| 1731 | + ); | ||
| 1732 | + out.push_str( | ||
| 1733 | + "# - max_age_days = 30 -> delete files not modified for 30+ days (unset = keep forever)\n", | ||
| 1734 | + ); | ||
| 1735 | + out.push_str( | ||
| 1736 | + "# - max_total_size_mb = 2048 -> keep total size under 2 GB, deleting oldest files first\n", | ||
| 1737 | + ); | ||
| 1704 | out.push_str("[datalog]\n"); | 1738 | out.push_str("[datalog]\n"); |
| 1705 | out.push_str(&format!("enabled = {}\n", cfg.enabled)); | 1739 | out.push_str(&format!("enabled = {}\n", cfg.enabled)); |
| 1706 | let dir_value = cfg.dir.as_deref().unwrap_or("~/.atomcode/datalog"); | 1740 | let dir_value = cfg.dir.as_deref().unwrap_or("~/.atomcode/datalog"); |
| 1707 | let escaped = dir_value.replace('\\', "\\\\").replace('"', "\\\""); | 1741 | let escaped = dir_value.replace('\\', "\\\\").replace('"', "\\\""); |
| 1708 | out.push_str(&format!("dir = \"{}\"\n", escaped)); | 1742 | out.push_str(&format!("dir = \"{}\"\n", escaped)); |
| 1743 | + if let Some(days) = cfg.max_age_days { | ||
| 1744 | + out.push_str(&format!("max_age_days = {days}\n")); | ||
| 1745 | + } | ||
| 1746 | + if let Some(size) = cfg.max_total_size_mb { | ||
| 1747 | + out.push_str(&format!("max_total_size_mb = {size}\n")); | ||
| 1748 | + } | ||
| 1709 | out | 1749 | out |
| 1710 | } | 1750 | } |
| 1711 | 1751 | ||
| @@ -3114,6 +3154,8 @@ model = "missing-type" | |||
| 3114 | let cfg = DatalogConfig { | 3154 | let cfg = DatalogConfig { |
| 3115 | enabled: true, | 3155 | enabled: true, |
| 3116 | dir: None, | 3156 | dir: None, |
| 3157 | + max_age_days: None, | ||
| 3158 | + max_total_size_mb: None, | ||
| 3117 | }; | 3159 | }; |
| 3118 | let rendered = render_datalog_section(&cfg); | 3160 | let rendered = render_datalog_section(&cfg); |
| 3119 | assert!(rendered.contains("\ndir = \"~/.atomcode/datalog\"\n")); | 3161 | assert!(rendered.contains("\ndir = \"~/.atomcode/datalog\"\n")); |
| @@ -3124,6 +3166,8 @@ model = "missing-type" | |||
| 3124 | let cfg = DatalogConfig { | 3166 | let cfg = DatalogConfig { |
| 3125 | enabled: false, | 3167 | enabled: false, |
| 3126 | dir: Some("~/.atomcode/logs".to_string()), | 3168 | dir: Some("~/.atomcode/logs".to_string()), |
| 3169 | + max_age_days: None, | ||
| 3170 | + max_total_size_mb: None, | ||
| 3127 | }; | 3171 | }; |
| 3128 | let rendered = render_datalog_section(&cfg); | 3172 | let rendered = render_datalog_section(&cfg); |
| 3129 | assert!(rendered.contains("enabled = false")); | 3173 | assert!(rendered.contains("enabled = false")); |
| @@ -3144,6 +3188,8 @@ model = "missing-type" | |||
| 3144 | datalog: DatalogConfig { | 3188 | datalog: DatalogConfig { |
| 3145 | enabled: false, | 3189 | enabled: false, |
| 3146 | dir: Some("/var/log/ac".to_string()), | 3190 | dir: Some("/var/log/ac".to_string()), |
| 3191 | + max_age_days: None, | ||
| 3192 | + max_total_size_mb: None, | ||
| 3147 | }, | 3193 | }, |
| 3148 | notifications: NotificationConfig::default(), | 3194 | notifications: NotificationConfig::default(), |
| 3149 | network: NetworkConfig::default(), | 3195 | network: NetworkConfig::default(), |
🟠 High Priority
变更行(datalog.rs:110-117):
DatalogHook::new用resolve_log_dir(...)得到<base>/<slug>后,取ancestors().nth(1)(即父目录)作为保留策略的遍历根。resolve_log_dir(datalog.rs:126-147)对任何非默认dir都直接返回base.join(project_slug(working_dir)),因此父目录就是用户配置的dir本身。受影响行为/契约:注释声称"遍历 datalog 根目录下所有项目 bucket",但
collect_datalog_files(datalog.rs:545-574)是全递归 + 仅按扩展名.md/.jsonl过滤,对父目录下任意深度的任意文件都纳入删除范围,且不校验文件是否位于 bucket 目录内。而resolve_log_dir显式支持的配置值dir = "~"(datalog.rs:131)会让父目录变成整个家目录;dir = "."(相对路径、文档明确支持的写法)会让父目录变成整个项目目录。失败模式:当用户配置非专用目录(
~、.、~/Documents、共享绝对路径等)并启用max_age_days或max_total_size_mb时,保留策略会递归遍历该目录并删除其中所有超过年龄阈值或按大小逐出的.md/.jsonl文件——包括用户笔记、项目文档、其他应用的 jsonl 日志等与 datalog 无关的数据。.md是用户文档的常见扩展名,这是真实的数据丢失(P1),且 4 个新增测试全部直接调用apply_retention传入显式 root,绕过了new里这段根目录推导,未覆盖此路径。建议:不要用 slug 目录的父目录盲推保留根。改为:仅在根目录第一层子目录(bucket)内删除 .md/.jsonl,且不无限递归到任意用户目录;对非默认 dir 额外校验目录形态(只含 slug 桶)后再遍历;补一个经 DatalogHook::new 的集成测试覆盖根目录推导。