1use std::collections::BTreeMap;
8use std::fs::{File, OpenOptions};
9use std::io::{Read, Write};
10use std::path::{Component, Path, PathBuf};
11
12use bamboo_domain::{
13 LegacyMemoryFileDisposition, LegacyMemoryMigrationFile, LegacyMemoryMigrationPhase,
14 LegacyMemoryMigrationReport, LegacyMemoryReadAlias, ProjectId,
15};
16use chrono::Utc;
17use serde::{Deserialize, Serialize};
18use sha2::{Digest, Sha256};
19use uuid::Uuid;
20
21use super::{
22 assert_plain_directory, ensure_confined_directory, lock_exclusive, replace_path,
23 sync_directory, validate_existing_confined_directory, validate_legacy_project_key,
24 write_bytes_atomic, write_json_atomic, ProjectStore, ProjectStoreError, ProjectStoreResult,
25 TempCleanup,
26};
27
28const JOURNAL_SCHEMA_VERSION: u32 = 1;
29const LEGACY_MEMORY_STATE_DIR: &str = "legacy-memory-migration";
30
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct LegacyMemoryReadRoot {
33 pub legacy_project_key: String,
34 pub root: PathBuf,
35 pub read_only: bool,
36}
37
38#[derive(Debug, Clone, PartialEq, Eq)]
43pub struct ProjectMemoryReadRoots {
44 pub primary: PathBuf,
45 pub legacy_aliases: Vec<LegacyMemoryReadRoot>,
46}
47
48#[derive(Debug, Clone, Serialize, Deserialize)]
49struct LegacyMemoryMigrationJournal {
50 schema_version: u32,
51 report: LegacyMemoryMigrationReport,
52 #[serde(default)]
53 resource_revision_bumped: bool,
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57enum MigrationFault {
58 None,
59 AfterJournal,
60 AfterFirstStaged,
61 AfterFirstCommitted,
62}
63
64impl ProjectStore {
65 pub fn legacy_memory_source_root(
67 &self,
68 legacy_project_key: &str,
69 ) -> ProjectStoreResult<PathBuf> {
70 validate_legacy_project_key(legacy_project_key)?;
71 Ok(self
72 .paths()
73 .data_dir()
74 .join("memory")
75 .join("v1")
76 .join("scopes")
77 .join("projects")
78 .join(legacy_project_key))
79 }
80
81 pub fn migrate_legacy_memory(
84 &self,
85 project_id: &ProjectId,
86 legacy_project_key: &str,
87 ) -> ProjectStoreResult<LegacyMemoryMigrationReport> {
88 self.migrate_legacy_memory_inner(project_id, legacy_project_key, MigrationFault::None)
89 }
90
91 pub fn legacy_memory_migration_status(
93 &self,
94 project_id: &ProjectId,
95 legacy_project_key: &str,
96 ) -> ProjectStoreResult<Option<LegacyMemoryMigrationReport>> {
97 validate_legacy_project_key(legacy_project_key)?;
98 self.get(project_id)?;
99 let project_home = self.paths().project_home(project_id);
100 validate_existing_confined_directory(self.paths().data_dir(), &project_home)?;
101 let state_root = self.legacy_memory_state_root(project_id);
102 if std::fs::symlink_metadata(&state_root)
103 .is_err_and(|error| error.kind() == std::io::ErrorKind::NotFound)
104 {
105 return Ok(None);
106 }
107 validate_existing_confined_directory(&project_home, &state_root)?;
108 let _lock = lock_exclusive(state_root.join(".lock"))?;
109 Ok(self
110 .load_migration_journal(project_id, legacy_project_key)?
111 .map(|journal| journal.report))
112 }
113
114 pub fn legacy_memory_aliases(
120 &self,
121 project_id: &ProjectId,
122 ) -> ProjectStoreResult<Vec<LegacyMemoryReadAlias>> {
123 let manifest = self.get(project_id)?;
124 let mut aliases = Vec::with_capacity(manifest.legacy_project_keys.len());
125 for key in manifest.legacy_project_keys {
126 let source = self.legacy_memory_source_root(&key)?;
127 let source_available =
128 validate_existing_confined_directory(self.paths().data_dir(), &source).is_ok();
129 let migration_committed = self
130 .legacy_memory_migration_status(project_id, &key)?
131 .is_some_and(|report| report.phase == LegacyMemoryMigrationPhase::Committed);
132 aliases.push(LegacyMemoryReadAlias {
133 legacy_project_key: key,
134 read_only: true,
135 project_home_precedence: true,
136 source_available,
137 migration_committed,
138 });
139 }
140 Ok(aliases)
141 }
142
143 pub fn project_memory_read_roots(
144 &self,
145 project_id: &ProjectId,
146 ) -> ProjectStoreResult<ProjectMemoryReadRoots> {
147 let aliases = self
148 .legacy_memory_aliases(project_id)?
149 .into_iter()
150 .filter(|alias| alias.source_available)
151 .map(|alias| {
152 Ok(LegacyMemoryReadRoot {
153 root: self.legacy_memory_source_root(&alias.legacy_project_key)?,
154 legacy_project_key: alias.legacy_project_key,
155 read_only: true,
156 })
157 })
158 .collect::<ProjectStoreResult<Vec<_>>>()?;
159 Ok(ProjectMemoryReadRoots {
160 primary: self.paths().memory_v1_dir(project_id),
161 legacy_aliases: aliases,
162 })
163 }
164
165 fn migrate_legacy_memory_inner(
166 &self,
167 project_id: &ProjectId,
168 legacy_project_key: &str,
169 fault: MigrationFault,
170 ) -> ProjectStoreResult<LegacyMemoryMigrationReport> {
171 validate_legacy_project_key(legacy_project_key)?;
172 let manifest = self.get(project_id)?;
173 if !manifest
174 .legacy_project_keys
175 .iter()
176 .any(|key| key == legacy_project_key)
177 {
178 return Err(ProjectStoreError::Validation(format!(
179 "legacy memory key is not declared by project {}",
180 project_id
181 )));
182 }
183 let project_home = self.paths().project_home(project_id);
184 validate_existing_confined_directory(self.paths().data_dir(), &project_home)?;
185 let state_root = self.legacy_memory_state_root(project_id);
186 ensure_confined_directory(&project_home, &state_root)?;
187 let _lock = lock_exclusive(state_root.join(".lock"))?;
188
189 let manifest = self.get(project_id)?;
193 if !manifest
194 .legacy_project_keys
195 .iter()
196 .any(|key| key == legacy_project_key)
197 {
198 return Err(ProjectStoreError::Validation(format!(
199 "legacy memory key is not declared by project {}",
200 project_id
201 )));
202 }
203
204 let source_root = self.legacy_memory_source_root(legacy_project_key)?;
205 if validate_existing_confined_directory(self.paths().data_dir(), &source_root).is_err() {
206 return Err(ProjectStoreError::Validation(format!(
207 "legacy memory source is missing or not a plain directory: {}",
208 source_root.display()
209 )));
210 }
211 let target_root = self.paths().memory_v1_dir(project_id);
212 ensure_confined_directory(&project_home, &target_root)?;
213 let source_files = enumerate_source_files(&source_root)?;
214
215 let existing = self.load_migration_journal(project_id, legacy_project_key)?;
216 if let Some(mut journal) = existing.as_ref().cloned() {
217 if journal.report.phase == LegacyMemoryMigrationPhase::Committed {
218 self.ensure_resource_revision_bumped(&mut journal)?;
219 return Ok(journal.report.clone());
220 }
221 }
222
223 let mut journal = match existing {
224 Some(journal) if source_snapshot_matches(&journal.report.files, &source_files) => {
225 journal
226 }
227 _ => self.plan_migration(project_id, legacy_project_key, &target_root, source_files)?,
228 };
229 self.persist_migration_journal(&journal)?;
230 if fault == MigrationFault::AfterJournal {
231 return Err(injected_interruption("after journal"));
232 }
233
234 let stage_root = self
235 .legacy_memory_state_root(project_id)
236 .join("staging")
237 .join(&journal.report.transaction_id);
238 ensure_confined_directory(&project_home, &stage_root)?;
239
240 let mut staged = 0usize;
241 for index in 0..journal.report.files.len() {
242 let disposition = journal.report.files[index].disposition;
243 if !matches!(
244 disposition,
245 LegacyMemoryFileDisposition::Pending | LegacyMemoryFileDisposition::Staged
246 ) {
247 continue;
248 }
249 let file = journal.report.files[index].clone();
250 let relative = validated_relative_path(&file.relative_path)?;
251 let source = source_root.join(&relative);
252 let target = safe_target_path(&target_root, &relative)?;
253 match target_state(&target, &file.sha256, file.size)? {
254 TargetState::Identical => {
255 journal.report.files[index].disposition =
256 LegacyMemoryFileDisposition::ExistingIdentical;
257 }
258 TargetState::Missing | TargetState::Conflict => {
259 let staged_path = safe_target_path(&stage_root, &relative)?;
263 if !file_matches(&staged_path, &file.sha256, file.size)? {
264 copy_file_verified_atomic(&source, &staged_path, &file.sha256, file.size)?;
265 }
266 if let Some(diagnostic) =
267 canonical_topic_diagnostic(&file.relative_path, &staged_path)?
268 {
269 journal.report.files[index].disposition =
270 LegacyMemoryFileDisposition::SkippedInvalid;
271 journal.report.files[index].diagnostic = Some(diagnostic);
272 let _ = std::fs::remove_file(&staged_path);
273 touch_journal(&mut journal);
274 self.persist_migration_journal(&journal)?;
275 continue;
276 }
277 journal.report.files[index].disposition = LegacyMemoryFileDisposition::Staged;
278 staged += 1;
279 }
280 }
281 touch_journal(&mut journal);
282 self.persist_migration_journal(&journal)?;
283 if fault == MigrationFault::AfterFirstStaged && staged == 1 {
284 return Err(injected_interruption("after first staged file"));
285 }
286 }
287
288 for index in 0..journal.report.files.len() {
289 if journal.report.files[index].disposition == LegacyMemoryFileDisposition::Staged {
290 let relative_path = journal.report.files[index].relative_path.clone();
291 let relative = validated_relative_path(&relative_path)?;
292 let staged_path = safe_target_path(&stage_root, &relative)?;
293 if !file_matches(
294 &staged_path,
295 &journal.report.files[index].sha256,
296 journal.report.files[index].size,
297 )? {
298 return Err(ProjectStoreError::Validation(format!(
299 "staged legacy memory file failed verification: {}",
300 relative_path
301 )));
302 }
303 if let Some(diagnostic) = canonical_topic_diagnostic(&relative_path, &staged_path)?
304 {
305 journal.report.files[index].disposition =
306 LegacyMemoryFileDisposition::SkippedInvalid;
307 journal.report.files[index].diagnostic = Some(diagnostic);
308 let _ = std::fs::remove_file(&staged_path);
309 }
310 }
311 }
312 journal.report.phase = LegacyMemoryMigrationPhase::Verified;
313 touch_journal(&mut journal);
314 self.persist_migration_journal(&journal)?;
315
316 let mut committed = 0usize;
317 for index in 0..journal.report.files.len() {
318 let file = journal.report.files[index].clone();
319 let relative = validated_relative_path(&file.relative_path)?;
320 let target = safe_target_path(&target_root, &relative)?;
321 let disposition = match file.disposition {
322 LegacyMemoryFileDisposition::Staged => {
323 let staged_path = safe_target_path(&stage_root, &relative)?;
324 install_verified_no_clobber(&staged_path, &target, &file.sha256, file.size)?
325 }
326 LegacyMemoryFileDisposition::Copied => {
327 match target_state(&target, &file.sha256, file.size)? {
328 TargetState::Identical => LegacyMemoryFileDisposition::Copied,
329 TargetState::Conflict => LegacyMemoryFileDisposition::TargetConflict,
330 TargetState::Missing => {
331 let staged_path = safe_target_path(&stage_root, &relative)?;
332 install_verified_no_clobber(
333 &staged_path,
334 &target,
335 &file.sha256,
336 file.size,
337 )?
338 }
339 }
340 }
341 LegacyMemoryFileDisposition::ExistingIdentical => {
342 match target_state(&target, &file.sha256, file.size)? {
343 TargetState::Identical => LegacyMemoryFileDisposition::ExistingIdentical,
344 TargetState::Missing | TargetState::Conflict => {
345 LegacyMemoryFileDisposition::TargetConflict
346 }
347 }
348 }
349 other => other,
350 };
351 journal.report.files[index].disposition = disposition;
352 if disposition == LegacyMemoryFileDisposition::Copied {
353 committed += 1;
354 }
355 touch_journal(&mut journal);
356 self.persist_migration_journal(&journal)?;
357 if fault == MigrationFault::AfterFirstCommitted && committed == 1 {
358 return Err(injected_interruption("after first committed file"));
359 }
360 }
361
362 journal.report.phase = LegacyMemoryMigrationPhase::Committed;
363 let now = Utc::now();
364 journal.report.updated_at = now;
365 journal.report.committed_at = Some(now);
366 self.persist_migration_journal(&journal)?;
367 self.ensure_resource_revision_bumped(&mut journal)?;
368 Ok(journal.report)
369 }
370
371 fn plan_migration(
372 &self,
373 project_id: &ProjectId,
374 legacy_project_key: &str,
375 target_root: &Path,
376 source_files: Vec<SourceFile>,
377 ) -> ProjectStoreResult<LegacyMemoryMigrationJournal> {
378 let files = source_files
379 .into_iter()
380 .map(|source| {
381 let relative = validated_relative_path(&source.relative_path)?;
382 let target = safe_target_path(target_root, &relative)?;
383 let disposition = match source.diagnostic.as_ref() {
384 Some(_) => LegacyMemoryFileDisposition::SkippedInvalid,
385 None => match target_state(&target, &source.sha256, source.size)? {
386 TargetState::Missing => LegacyMemoryFileDisposition::Pending,
387 TargetState::Identical => LegacyMemoryFileDisposition::ExistingIdentical,
388 TargetState::Conflict => LegacyMemoryFileDisposition::Pending,
389 },
390 };
391 Ok(LegacyMemoryMigrationFile {
392 relative_path: source.relative_path,
393 size: source.size,
394 sha256: source.sha256,
395 disposition,
396 diagnostic: source.diagnostic,
397 })
398 })
399 .collect::<ProjectStoreResult<Vec<_>>>()?;
400 let now = Utc::now();
401 Ok(LegacyMemoryMigrationJournal {
402 schema_version: JOURNAL_SCHEMA_VERSION,
403 resource_revision_bumped: false,
404 report: LegacyMemoryMigrationReport {
405 project_id: project_id.clone(),
406 legacy_project_key: legacy_project_key.to_string(),
407 transaction_id: Uuid::new_v4().to_string(),
408 phase: LegacyMemoryMigrationPhase::Copying,
409 files,
410 started_at: now,
411 updated_at: now,
412 committed_at: None,
413 },
414 })
415 }
416
417 fn legacy_memory_state_root(&self, project_id: &ProjectId) -> PathBuf {
418 self.paths()
419 .state_dir(project_id)
420 .join(LEGACY_MEMORY_STATE_DIR)
421 }
422
423 fn migration_journal_path(&self, project_id: &ProjectId, legacy_project_key: &str) -> PathBuf {
424 let digest = hex::encode(Sha256::digest(legacy_project_key.as_bytes()));
425 self.legacy_memory_state_root(project_id)
426 .join("journals")
427 .join(format!("{digest}.json"))
428 }
429
430 fn load_migration_journal(
431 &self,
432 project_id: &ProjectId,
433 legacy_project_key: &str,
434 ) -> ProjectStoreResult<Option<LegacyMemoryMigrationJournal>> {
435 let path = self.migration_journal_path(project_id, legacy_project_key);
436 let journal_dir = path.parent().ok_or_else(|| {
437 ProjectStoreError::Validation(
438 "legacy memory migration journal has no parent".to_string(),
439 )
440 })?;
441 if std::fs::symlink_metadata(journal_dir)
442 .is_err_and(|error| error.kind() == std::io::ErrorKind::NotFound)
443 {
444 return Ok(None);
445 }
446 validate_existing_confined_directory(&self.paths().project_home(project_id), journal_dir)?;
447 let bytes = match std::fs::read(&path) {
448 Ok(bytes) => bytes,
449 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
450 Err(error) => return Err(error.into()),
451 };
452 let journal = match serde_json::from_slice::<LegacyMemoryMigrationJournal>(&bytes) {
453 Ok(journal) => journal,
454 Err(error) => {
455 let quarantine = path.with_file_name(format!(
456 "{}.corrupt.{}",
457 path.file_name()
458 .and_then(|name| name.to_str())
459 .unwrap_or("journal.json"),
460 Uuid::new_v4()
461 ));
462 write_bytes_atomic(&quarantine, &bytes)?;
463 tracing::warn!(
464 %error,
465 quarantine = %quarantine.display(),
466 "legacy memory migration journal was corrupt; replanning"
467 );
468 return Ok(None);
469 }
470 };
471 validate_journal(&journal, project_id, legacy_project_key)?;
472 Ok(Some(journal))
473 }
474
475 fn persist_migration_journal(
476 &self,
477 journal: &LegacyMemoryMigrationJournal,
478 ) -> ProjectStoreResult<()> {
479 validate_journal(
480 journal,
481 &journal.report.project_id,
482 &journal.report.legacy_project_key,
483 )?;
484 let path = self.migration_journal_path(
485 &journal.report.project_id,
486 &journal.report.legacy_project_key,
487 );
488 let journal_dir = path.parent().ok_or_else(|| {
489 ProjectStoreError::Validation(
490 "legacy memory migration journal has no parent".to_string(),
491 )
492 })?;
493 ensure_confined_directory(
494 &self.paths().project_home(&journal.report.project_id),
495 journal_dir,
496 )?;
497 write_json_atomic(&path, journal)
498 }
499
500 fn ensure_resource_revision_bumped(
501 &self,
502 journal: &mut LegacyMemoryMigrationJournal,
503 ) -> ProjectStoreResult<()> {
504 if journal.resource_revision_bumped {
505 return Ok(());
506 }
507 loop {
508 let project = self.get(&journal.report.project_id)?;
509 match self.bump_resource_revision(&project.id, project.revision) {
510 Ok(_) => break,
511 Err(ProjectStoreError::Conflict { .. }) => continue,
512 Err(error) => return Err(error),
513 }
514 }
515 journal.resource_revision_bumped = true;
516 touch_journal(journal);
517 self.persist_migration_journal(journal)
518 }
519
520 #[cfg(test)]
521 fn migrate_legacy_memory_with_fault(
522 &self,
523 project_id: &ProjectId,
524 legacy_project_key: &str,
525 fault: MigrationFault,
526 ) -> ProjectStoreResult<LegacyMemoryMigrationReport> {
527 self.migrate_legacy_memory_inner(project_id, legacy_project_key, fault)
528 }
529}
530
531fn validate_journal(
532 journal: &LegacyMemoryMigrationJournal,
533 project_id: &ProjectId,
534 legacy_project_key: &str,
535) -> ProjectStoreResult<()> {
536 if journal.schema_version != JOURNAL_SCHEMA_VERSION
537 || &journal.report.project_id != project_id
538 || journal.report.legacy_project_key != legacy_project_key
539 {
540 return Err(ProjectStoreError::Validation(
541 "legacy memory migration journal identity mismatch".to_string(),
542 ));
543 }
544 if Uuid::parse_str(&journal.report.transaction_id).is_err() {
545 return Err(ProjectStoreError::Validation(
546 "legacy memory migration transaction id is invalid".to_string(),
547 ));
548 }
549 let mut paths = BTreeMap::new();
550 for file in &journal.report.files {
551 validated_relative_path(&file.relative_path)?;
552 if file.sha256.len() != 64
553 || !file.sha256.bytes().all(|byte| byte.is_ascii_hexdigit())
554 || paths.insert(&file.relative_path, ()).is_some()
555 {
556 return Err(ProjectStoreError::Validation(
557 "legacy memory migration file manifest is invalid".to_string(),
558 ));
559 }
560 }
561 Ok(())
562}
563
564#[derive(Debug, Clone, PartialEq, Eq)]
565struct SourceFile {
566 relative_path: String,
567 size: u64,
568 sha256: String,
569 diagnostic: Option<String>,
570}
571
572fn enumerate_source_files(root: &Path) -> ProjectStoreResult<Vec<SourceFile>> {
573 assert_plain_directory(root)?;
574 let mut files = Vec::new();
575 enumerate_source_directory(root, root, &mut files)?;
576 files.sort_by(|left, right| left.relative_path.cmp(&right.relative_path));
577 Ok(files)
578}
579
580fn enumerate_source_directory(
581 root: &Path,
582 directory: &Path,
583 files: &mut Vec<SourceFile>,
584) -> ProjectStoreResult<()> {
585 for entry in std::fs::read_dir(directory)? {
586 let entry = entry?;
587 let metadata = std::fs::symlink_metadata(entry.path())?;
588 if metadata.file_type().is_symlink() {
589 return Err(ProjectStoreError::Validation(format!(
590 "legacy memory migration refuses symlink: {}",
591 entry.path().display()
592 )));
593 }
594 if metadata.is_dir() {
595 enumerate_source_directory(root, &entry.path(), files)?;
596 } else if metadata.is_file() {
597 let entry_path = entry.path();
598 let relative = entry_path.strip_prefix(root).map_err(|_| {
599 ProjectStoreError::Validation(
600 "legacy memory entry escaped its source root".to_string(),
601 )
602 })?;
603 let relative_path = relative_path_string(relative)?;
604 let (size, sha256) = hash_regular_file(&entry.path())?;
605 let diagnostic = canonical_topic_diagnostic(&relative_path, &entry.path())?;
606 files.push(SourceFile {
607 relative_path,
608 size,
609 sha256,
610 diagnostic,
611 });
612 } else {
613 return Err(ProjectStoreError::Validation(format!(
614 "legacy memory migration refuses special file: {}",
615 entry.path().display()
616 )));
617 }
618 }
619 Ok(())
620}
621
622fn source_snapshot_matches(planned: &[LegacyMemoryMigrationFile], source: &[SourceFile]) -> bool {
623 if planned.len() != source.len() {
624 return false;
625 }
626 planned.iter().zip(source).all(|(planned, source)| {
627 planned.relative_path == source.relative_path
628 && planned.size == source.size
629 && planned.sha256 == source.sha256
630 && (planned.disposition == LegacyMemoryFileDisposition::SkippedInvalid)
631 == source.diagnostic.is_some()
632 })
633}
634
635fn canonical_topic_diagnostic(
636 relative_path: &str,
637 file_path: &Path,
638) -> ProjectStoreResult<Option<String>> {
639 if !is_canonical_topic_path(relative_path) {
640 return Ok(None);
641 }
642 let bytes = std::fs::read(file_path)?;
643 let content = match std::str::from_utf8(&bytes) {
644 Ok(content) => content,
645 Err(_) => {
646 return Ok(Some(
647 "canonical memory topic is not valid UTF-8".to_string(),
648 ));
649 }
650 };
651 match parse_canonical_memory_topic(content) {
652 Ok(_) => Ok(None),
653 Err(error) => {
654 tracing::warn!(
655 path = %file_path.display(),
656 %error,
657 "isolating invalid legacy Project memory topic"
658 );
659 Ok(Some(
660 "canonical memory topic frontmatter is invalid".to_string(),
661 ))
662 }
663 }
664}
665
666fn parse_canonical_memory_topic(content: &str) -> Result<(), String> {
667 let trimmed = content.trim_start_matches('\u{feff}');
668 let rest = trimmed
669 .strip_prefix("---\n")
670 .ok_or_else(|| "missing frontmatter start marker".to_string())?;
671 let end_index = rest
672 .find("\n---\n")
673 .ok_or_else(|| "missing frontmatter end marker".to_string())?;
674 serde_yaml::from_str::<CanonicalMemoryTopicFrontmatter>(&rest[..end_index])
675 .map(|_| ())
676 .map_err(|error| format!("failed to parse memory frontmatter: {error}"))
677}
678
679#[allow(dead_code)]
680#[derive(Deserialize)]
681struct CanonicalMemoryTopicFrontmatter {
682 id: String,
683 title: String,
684 #[serde(rename = "type")]
685 memory_type: CanonicalMemoryType,
686 scope: CanonicalMemoryScope,
687 #[serde(default)]
688 project_key: Option<String>,
689 #[serde(default)]
690 granularity: Option<CanonicalTemporalGranularity>,
691 status: CanonicalMemoryStatus,
692 #[serde(default)]
693 freshness: Option<String>,
694 #[serde(default)]
695 confidence: Option<String>,
696 created_at: String,
697 updated_at: String,
698 created_by: CanonicalMemoryActor,
699 updated_by: CanonicalMemoryActor,
700 #[serde(default)]
701 sources: Vec<CanonicalMemorySource>,
702 #[serde(default)]
703 relations: CanonicalMemoryRelations,
704 #[serde(default)]
705 tags: Vec<String>,
706 #[serde(default)]
707 retrieval: CanonicalMemoryRetrieval,
708}
709
710#[derive(Deserialize)]
711#[serde(rename_all = "snake_case")]
712enum CanonicalMemoryType {
713 User,
714 Feedback,
715 Project,
716 Reference,
717}
718
719#[derive(Deserialize)]
720#[serde(rename_all = "snake_case")]
721enum CanonicalMemoryScope {
722 Session,
723 Project,
724 Global,
725}
726
727#[derive(Deserialize)]
728#[serde(rename_all = "snake_case")]
729enum CanonicalTemporalGranularity {
730 Day,
731 Week,
732 Month,
733 Quarter,
734 Year,
735}
736
737#[derive(Deserialize)]
738#[serde(rename_all = "snake_case")]
739enum CanonicalMemoryStatus {
740 Active,
741 Stale,
742 Superseded,
743 Contradicted,
744 Archived,
745}
746
747#[allow(dead_code)]
748#[derive(Deserialize)]
749struct CanonicalMemoryActor {
750 kind: String,
751 #[serde(default)]
752 id: Option<String>,
753 #[serde(default)]
754 actor: Option<String>,
755}
756
757#[allow(dead_code)]
758#[derive(Deserialize)]
759struct CanonicalMemorySource {
760 kind: String,
761 id: String,
762 #[serde(default)]
763 message_range: Vec<String>,
764}
765
766#[allow(dead_code)]
767#[derive(Default, Deserialize)]
768struct CanonicalMemoryRelations {
769 #[serde(default)]
770 supersedes: Vec<String>,
771 #[serde(default)]
772 contradicted_by: Vec<String>,
773 #[serde(default)]
774 related: Vec<String>,
775}
776
777#[allow(dead_code)]
778#[derive(Default, Deserialize)]
779struct CanonicalMemoryRetrieval {
780 #[serde(default)]
781 keywords: Vec<String>,
782 #[serde(default)]
783 entities: Vec<String>,
784 #[serde(default)]
785 embedding_ready: bool,
786 #[serde(default)]
787 last_accessed_at: Option<String>,
788}
789
790fn is_canonical_topic_path(relative_path: &str) -> bool {
791 let mut components = Path::new(relative_path).components();
792 matches!(
793 (
794 components.next(),
795 components.next(),
796 components.next(),
797 ),
798 (
799 Some(Component::Normal(directory)),
800 Some(Component::Normal(file)),
801 None
802 ) if directory == "topics" && Path::new(file).extension().is_some_and(|ext| ext == "md")
803 )
804}
805
806fn relative_path_string(path: &Path) -> ProjectStoreResult<String> {
807 let mut components = Vec::new();
808 for component in path.components() {
809 let Component::Normal(component) = component else {
810 return Err(ProjectStoreError::Validation(
811 "legacy memory relative path is invalid".to_string(),
812 ));
813 };
814 let component = component.to_str().ok_or_else(|| {
815 ProjectStoreError::Validation("legacy memory paths must be valid UTF-8".to_string())
816 })?;
817 if component.contains('\\') || component.is_empty() {
818 return Err(ProjectStoreError::Validation(
819 "legacy memory path component is not portable".to_string(),
820 ));
821 }
822 components.push(component);
823 }
824 if components.is_empty() {
825 return Err(ProjectStoreError::Validation(
826 "legacy memory relative path is empty".to_string(),
827 ));
828 }
829 Ok(components.join("/"))
830}
831
832fn validated_relative_path(path: &str) -> ProjectStoreResult<PathBuf> {
833 if path.is_empty() || path.contains('\\') {
834 return Err(ProjectStoreError::Validation(
835 "legacy memory relative path is invalid".to_string(),
836 ));
837 }
838 let path = Path::new(path);
839 if path.is_absolute()
840 || path
841 .components()
842 .any(|component| !matches!(component, Component::Normal(_)))
843 {
844 return Err(ProjectStoreError::Validation(
845 "legacy memory relative path is invalid".to_string(),
846 ));
847 }
848 Ok(path.to_path_buf())
849}
850
851fn safe_target_path(root: &Path, relative: &Path) -> ProjectStoreResult<PathBuf> {
852 let target = confined_path(root, relative)?;
853 let parent = target.parent().map(Path::to_path_buf).ok_or_else(|| {
854 ProjectStoreError::Validation("legacy memory target has no parent".to_string())
855 })?;
856 let canonical_root = ensure_confined_directory(root, &parent)?;
857 let root = std::fs::canonicalize(root)?;
858 if !canonical_root.starts_with(&root) {
859 return Err(ProjectStoreError::Validation(
860 "legacy memory target parent escapes destination root".to_string(),
861 ));
862 }
863 Ok(target)
864}
865
866fn confined_path(root: &Path, relative: &Path) -> ProjectStoreResult<PathBuf> {
867 let relative = validated_relative_path(relative.to_str().ok_or_else(|| {
868 ProjectStoreError::Validation("legacy memory target path must be valid UTF-8".to_string())
869 })?)?;
870 assert_plain_directory(root)?;
871 Ok(std::fs::canonicalize(root)?.join(relative))
872}
873
874fn hash_regular_file(path: &Path) -> ProjectStoreResult<(u64, String)> {
875 let metadata = std::fs::symlink_metadata(path)?;
876 if !metadata.is_file() || metadata.file_type().is_symlink() {
877 return Err(ProjectStoreError::Validation(format!(
878 "expected a regular file: {}",
879 path.display()
880 )));
881 }
882 let mut file = File::open(path)?;
883 let mut digest = Sha256::new();
884 let mut size = 0u64;
885 let mut buffer = [0u8; 64 * 1024];
886 loop {
887 let read = file.read(&mut buffer)?;
888 if read == 0 {
889 break;
890 }
891 digest.update(&buffer[..read]);
892 size = size
893 .checked_add(read as u64)
894 .ok_or_else(|| ProjectStoreError::Validation("file size overflow".to_string()))?;
895 }
896 Ok((size, hex::encode(digest.finalize())))
897}
898
899fn file_matches(
900 path: &Path,
901 expected_sha256: &str,
902 expected_size: u64,
903) -> ProjectStoreResult<bool> {
904 match std::fs::symlink_metadata(path) {
905 Ok(metadata) if metadata.is_file() && !metadata.file_type().is_symlink() => {
906 let (size, sha256) = hash_regular_file(path)?;
907 Ok(size == expected_size && sha256 == expected_sha256)
908 }
909 Ok(_) => Ok(false),
910 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
911 Err(error) => Err(error.into()),
912 }
913}
914
915enum TargetState {
916 Missing,
917 Identical,
918 Conflict,
919}
920
921fn target_state(
922 path: &Path,
923 expected_sha256: &str,
924 expected_size: u64,
925) -> ProjectStoreResult<TargetState> {
926 match std::fs::symlink_metadata(path) {
927 Ok(metadata) if metadata.is_file() && !metadata.file_type().is_symlink() => {
928 if file_matches(path, expected_sha256, expected_size)? {
929 Ok(TargetState::Identical)
930 } else {
931 Ok(TargetState::Conflict)
932 }
933 }
934 Ok(_) => Ok(TargetState::Conflict),
935 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(TargetState::Missing),
936 Err(error) => Err(error.into()),
937 }
938}
939
940fn copy_file_verified_atomic(
941 source: &Path,
942 target: &Path,
943 expected_sha256: &str,
944 expected_size: u64,
945) -> ProjectStoreResult<()> {
946 let parent = target.parent().unwrap_or_else(|| Path::new("."));
947 assert_plain_directory(parent)?;
948 let temp = parent.join(format!(".legacy-memory.tmp.{}", Uuid::new_v4()));
949 let mut cleanup = TempCleanup(Some(temp.clone()));
950 let mut input = File::open(source)?;
951 let mut output = OpenOptions::new()
952 .create_new(true)
953 .write(true)
954 .open(&temp)?;
955 let mut digest = Sha256::new();
956 let mut size = 0u64;
957 let mut buffer = [0u8; 64 * 1024];
958 loop {
959 let read = input.read(&mut buffer)?;
960 if read == 0 {
961 break;
962 }
963 output.write_all(&buffer[..read])?;
964 digest.update(&buffer[..read]);
965 size = size
966 .checked_add(read as u64)
967 .ok_or_else(|| ProjectStoreError::Validation("file size overflow".to_string()))?;
968 }
969 output.sync_all()?;
970 drop(output);
971 let actual = hex::encode(digest.finalize());
972 if size != expected_size || actual != expected_sha256 {
973 return Err(ProjectStoreError::Validation(format!(
974 "legacy memory source changed while copying: {}",
975 source.display()
976 )));
977 }
978 sync_directory(parent)?;
979 replace_path(&temp, target)?;
980 cleanup.0 = None;
981 sync_directory(parent)?;
982 Ok(())
983}
984
985fn install_verified_no_clobber(
986 staged: &Path,
987 target: &Path,
988 expected_sha256: &str,
989 expected_size: u64,
990) -> ProjectStoreResult<LegacyMemoryFileDisposition> {
991 match target_state(target, expected_sha256, expected_size)? {
992 TargetState::Identical => return Ok(LegacyMemoryFileDisposition::ExistingIdentical),
993 TargetState::Conflict => return Ok(LegacyMemoryFileDisposition::TargetConflict),
994 TargetState::Missing => {}
995 }
996 if !file_matches(staged, expected_sha256, expected_size)? {
997 return Err(ProjectStoreError::Validation(format!(
998 "staged legacy memory file failed verification: {}",
999 staged.display()
1000 )));
1001 }
1002 let parent = target.parent().unwrap_or_else(|| Path::new("."));
1003 assert_plain_directory(parent)?;
1004 let temp = parent.join(format!(".legacy-memory.commit.{}", Uuid::new_v4()));
1005 let mut cleanup = TempCleanup(Some(temp.clone()));
1006 copy_file_verified_atomic(staged, &temp, expected_sha256, expected_size)?;
1007 match std::fs::hard_link(&temp, target) {
1008 Ok(()) => {
1009 sync_directory(parent)?;
1010 std::fs::remove_file(&temp)?;
1011 cleanup.0 = None;
1012 sync_directory(parent)?;
1013 if !file_matches(target, expected_sha256, expected_size)? {
1014 return Err(ProjectStoreError::Validation(format!(
1015 "committed legacy memory file failed verification: {}",
1016 target.display()
1017 )));
1018 }
1019 Ok(LegacyMemoryFileDisposition::Copied)
1020 }
1021 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
1022 match target_state(target, expected_sha256, expected_size)? {
1023 TargetState::Identical => Ok(LegacyMemoryFileDisposition::ExistingIdentical),
1024 TargetState::Missing | TargetState::Conflict => {
1025 Ok(LegacyMemoryFileDisposition::TargetConflict)
1026 }
1027 }
1028 }
1029 Err(error) => Err(error.into()),
1030 }
1031}
1032
1033fn touch_journal(journal: &mut LegacyMemoryMigrationJournal) {
1034 journal.report.updated_at = Utc::now();
1035}
1036
1037fn injected_interruption(point: &str) -> ProjectStoreError {
1038 ProjectStoreError::Io(std::io::Error::other(format!(
1039 "injected legacy memory migration interruption {point}"
1040 )))
1041}
1042
1043#[cfg(test)]
1044mod tests {
1045 use super::*;
1046 use bamboo_domain::{LegacyMemoryFileDisposition, ProjectManifest};
1047 use bamboo_memory::memory_store::{MemoryQueryOptions, MemoryScope, MemoryStore};
1048 use tempfile::TempDir;
1049
1050 fn valid_topic(id: &str, title: &str, body: &str) -> String {
1051 format!(
1052 "---\n\
1053id: {id}\n\
1054title: {title}\n\
1055type: project\n\
1056scope: project\n\
1057project_key: zenith-deadbeef\n\
1058status: active\n\
1059created_at: 2026-01-01T00:00:00Z\n\
1060updated_at: 2026-01-01T00:00:00Z\n\
1061created_by:\n kind: session\n\
1062updated_by:\n kind: memory_write\n\
1063---\n\
1064{body}\n"
1065 )
1066 }
1067
1068 fn setup() -> (TempDir, ProjectStore, ProjectManifest, String, PathBuf) {
1069 let temp = tempfile::tempdir().unwrap();
1070 let store = ProjectStore::open(temp.path()).unwrap();
1071 let created = store.create("Memory", None).unwrap();
1072 let key = "zenith-deadbeef".to_string();
1073 let project = store
1074 .update(&created.id, created.revision, |manifest| {
1075 manifest.legacy_project_keys.push(key.clone());
1076 Ok(())
1077 })
1078 .unwrap();
1079 let source = store.legacy_memory_source_root(&key).unwrap();
1080 std::fs::create_dir_all(source.join("topics")).unwrap();
1081 std::fs::create_dir_all(source.join("views")).unwrap();
1082 std::fs::write(
1083 source.join("topics/one.md"),
1084 valid_topic(
1085 "one",
1086 "Project identity",
1087 "The Project identity remains stable.",
1088 ),
1089 )
1090 .unwrap();
1091 std::fs::write(source.join("views/dream.md"), b"dream").unwrap();
1092 (temp, store, project, key, source)
1093 }
1094
1095 #[test]
1096 fn migration_copies_to_project_home_and_never_deletes_source() {
1097 let (_temp, store, project, key, source) = setup();
1098 let report = store.migrate_legacy_memory(&project.id, &key).unwrap();
1099 assert_eq!(report.phase, LegacyMemoryMigrationPhase::Committed);
1100 assert!(report
1101 .files
1102 .iter()
1103 .all(|file| file.disposition == LegacyMemoryFileDisposition::Copied));
1104 let target = store.paths().memory_v1_dir(&project.id);
1105 assert_eq!(
1106 std::fs::read_to_string(target.join("topics/one.md")).unwrap(),
1107 valid_topic(
1108 "one",
1109 "Project identity",
1110 "The Project identity remains stable."
1111 )
1112 );
1113 assert_eq!(
1114 std::fs::read(target.join("views/dream.md")).unwrap(),
1115 b"dream"
1116 );
1117 assert!(source.join("topics/one.md").exists());
1118 assert!(source.join("views/dream.md").exists());
1119 assert!(
1120 target.starts_with(store.paths().project_home(&project.id)),
1121 "assigned memory must physically live under Project home"
1122 );
1123 }
1124
1125 #[test]
1126 fn target_content_wins_and_legacy_alias_remains_read_only() {
1127 let (_temp, store, project, key, source) = setup();
1128 let target = store.paths().memory_v1_dir(&project.id);
1129 std::fs::create_dir_all(target.join("topics")).unwrap();
1130 std::fs::write(target.join("topics/one.md"), b"new-data").unwrap();
1131
1132 let report = store.migrate_legacy_memory(&project.id, &key).unwrap();
1133 let conflict = report
1134 .files
1135 .iter()
1136 .find(|file| file.relative_path == "topics/one.md")
1137 .unwrap();
1138 assert_eq!(
1139 conflict.disposition,
1140 LegacyMemoryFileDisposition::TargetConflict
1141 );
1142 assert_eq!(
1143 std::fs::read(target.join("topics/one.md")).unwrap(),
1144 b"new-data"
1145 );
1146 assert_eq!(
1147 std::fs::read_to_string(source.join("topics/one.md")).unwrap(),
1148 valid_topic(
1149 "one",
1150 "Project identity",
1151 "The Project identity remains stable."
1152 )
1153 );
1154
1155 let roots = store.project_memory_read_roots(&project.id).unwrap();
1156 assert_eq!(roots.primary, target);
1157 assert_eq!(roots.legacy_aliases.len(), 1);
1158 assert!(roots.legacy_aliases[0].read_only);
1159 let aliases = store.legacy_memory_aliases(&project.id).unwrap();
1160 assert!(aliases[0].project_home_precedence);
1161 assert!(aliases[0].migration_committed);
1162 }
1163
1164 #[test]
1165 fn interrupted_staging_resumes_idempotently() {
1166 let (_temp, store, project, key, source) = setup();
1167 assert!(store
1168 .migrate_legacy_memory_with_fault(&project.id, &key, MigrationFault::AfterFirstStaged,)
1169 .is_err());
1170 let status = store
1171 .legacy_memory_migration_status(&project.id, &key)
1172 .unwrap()
1173 .unwrap();
1174 assert_eq!(status.phase, LegacyMemoryMigrationPhase::Copying);
1175
1176 let committed = store.migrate_legacy_memory(&project.id, &key).unwrap();
1177 assert_eq!(committed.phase, LegacyMemoryMigrationPhase::Committed);
1178 let resource_revision = store.get(&project.id).unwrap().resource_revision;
1179 let again = store.migrate_legacy_memory(&project.id, &key).unwrap();
1180 assert_eq!(again.transaction_id, committed.transaction_id);
1181 assert_eq!(again, committed);
1182 assert_eq!(
1183 store.get(&project.id).unwrap().resource_revision,
1184 resource_revision,
1185 "idempotent replay must not bump resource revision twice"
1186 );
1187 assert!(source.join("topics/one.md").exists());
1188 }
1189
1190 #[test]
1191 fn interrupted_commit_recovers_without_overwrite_or_duplicate() {
1192 let (_temp, store, project, key, _source) = setup();
1193 assert!(store
1194 .migrate_legacy_memory_with_fault(
1195 &project.id,
1196 &key,
1197 MigrationFault::AfterFirstCommitted,
1198 )
1199 .is_err());
1200 let committed = store.migrate_legacy_memory(&project.id, &key).unwrap();
1201 assert_eq!(committed.phase, LegacyMemoryMigrationPhase::Committed);
1202 assert_eq!(
1203 committed
1204 .files
1205 .iter()
1206 .filter(|file| matches!(
1207 file.disposition,
1208 LegacyMemoryFileDisposition::Copied
1209 | LegacyMemoryFileDisposition::ExistingIdentical
1210 ))
1211 .count(),
1212 2
1213 );
1214 }
1215
1216 #[test]
1217 fn changed_source_replans_uncommitted_transaction() {
1218 let (_temp, store, project, key, source) = setup();
1219 assert!(store
1220 .migrate_legacy_memory_with_fault(&project.id, &key, MigrationFault::AfterJournal,)
1221 .is_err());
1222 let old_transaction = store
1223 .legacy_memory_migration_status(&project.id, &key)
1224 .unwrap()
1225 .unwrap()
1226 .transaction_id;
1227 let updated = valid_topic(
1228 "one",
1229 "Project identity updated",
1230 "The Project identity remains stable after migration.",
1231 );
1232 std::fs::write(source.join("topics/one.md"), &updated).unwrap();
1233 let committed = store.migrate_legacy_memory(&project.id, &key).unwrap();
1234 assert_ne!(committed.transaction_id, old_transaction);
1235 assert_eq!(
1236 std::fs::read_to_string(
1237 store
1238 .paths()
1239 .memory_v1_dir(&project.id)
1240 .join("topics/one.md")
1241 )
1242 .unwrap(),
1243 updated
1244 );
1245 }
1246
1247 #[tokio::test]
1248 async fn invalid_topic_is_isolated_while_valid_topic_remains_recallable() {
1249 let (temp, store, project, key, source) = setup();
1250 std::fs::write(
1251 source.join("topics/corrupt.md"),
1252 "not a canonical memory document",
1253 )
1254 .unwrap();
1255
1256 let report = store.migrate_legacy_memory(&project.id, &key).unwrap();
1257 let corrupt = report
1258 .files
1259 .iter()
1260 .find(|file| file.relative_path == "topics/corrupt.md")
1261 .unwrap();
1262 assert_eq!(
1263 corrupt.disposition,
1264 LegacyMemoryFileDisposition::SkippedInvalid
1265 );
1266 assert_eq!(
1267 corrupt.diagnostic.as_deref(),
1268 Some("canonical memory topic frontmatter is invalid")
1269 );
1270 assert!(source.join("topics/corrupt.md").exists());
1271
1272 let primary = store.paths().memory_v1_dir(&project.id);
1273 assert!(!primary.join("topics/corrupt.md").exists());
1274 assert!(primary.join("topics/one.md").exists());
1275
1276 let memory = MemoryStore::new(temp.path()).for_project(&project.id);
1277 let recalled = memory
1278 .query_scope(
1279 MemoryScope::Project,
1280 Some(project.id.as_str()),
1281 Some("project identity stable"),
1282 None,
1283 None,
1284 None,
1285 &MemoryQueryOptions {
1286 limit: Some(5),
1287 max_chars: Some(3_000),
1288 cursor: None,
1289 include_related: false,
1290 },
1291 )
1292 .await
1293 .unwrap();
1294 assert_eq!(recalled.matched_count, 1);
1295 assert_eq!(recalled.items[0].id, "one");
1296 }
1297
1298 #[test]
1299 fn missing_project_queries_do_not_create_orphan_project_homes() {
1300 let temp = tempfile::tempdir().unwrap();
1301 let store = ProjectStore::open(temp.path()).unwrap();
1302 let missing: ProjectId = "01JMISSINGPROJECT00000000000".parse().unwrap();
1303 assert!(store.migrate_legacy_memory(&missing, "legacy-key").is_err());
1304 assert!(store
1305 .legacy_memory_migration_status(&missing, "legacy-key")
1306 .is_err());
1307 assert!(!store.paths().project_home(&missing).exists());
1308 }
1309
1310 #[cfg(unix)]
1311 #[test]
1312 fn source_symlinks_are_rejected_without_following() {
1313 use std::os::unix::fs::symlink;
1314
1315 let (temp, store, project, key, source) = setup();
1316 let outside = temp.path().join("outside-secret");
1317 std::fs::write(&outside, b"secret").unwrap();
1318 symlink(&outside, source.join("topics/link.md")).unwrap();
1319 assert!(store.migrate_legacy_memory(&project.id, &key).is_err());
1320 assert!(!store
1321 .paths()
1322 .memory_v1_dir(&project.id)
1323 .join("topics/link.md")
1324 .exists());
1325 }
1326
1327 #[cfg(unix)]
1328 #[test]
1329 fn source_ancestor_symlink_is_rejected_without_read_or_target_write() {
1330 use std::os::unix::fs::symlink;
1331
1332 let (temp, store, project, key, _source) = setup();
1333 let memory_root = temp.path().join("memory");
1334 let outside_memory = temp.path().join("outside-memory");
1335 std::fs::rename(&memory_root, &outside_memory).unwrap();
1336 symlink(&outside_memory, &memory_root).unwrap();
1337 let legacy_file = outside_memory
1338 .join("v1/scopes/projects")
1339 .join(&key)
1340 .join("topics/one.md");
1341 let before = std::fs::read(&legacy_file).unwrap();
1342
1343 assert!(store.migrate_legacy_memory(&project.id, &key).is_err());
1344 assert_eq!(std::fs::read(&legacy_file).unwrap(), before);
1345 assert!(!store.paths().memory_v1_dir(&project.id).exists());
1346 }
1347
1348 #[cfg(unix)]
1349 #[test]
1350 fn target_memory_symlink_has_zero_external_side_effects() {
1351 use std::os::unix::fs::symlink;
1352
1353 let (temp, store, project, key, _source) = setup();
1354 let outside = temp.path().join("outside-target");
1355 std::fs::create_dir(&outside).unwrap();
1356 let project_home = store.paths().project_home(&project.id);
1357 symlink(&outside, project_home.join("memory")).unwrap();
1358
1359 assert!(store.migrate_legacy_memory(&project.id, &key).is_err());
1360 assert_eq!(
1361 std::fs::read_dir(&outside).unwrap().count(),
1362 0,
1363 "migration must not create v1 or files through a target symlink"
1364 );
1365 }
1366
1367 #[cfg(unix)]
1368 #[test]
1369 fn staging_symlink_has_zero_external_side_effects() {
1370 use std::os::unix::fs::symlink;
1371
1372 let (temp, store, project, key, _source) = setup();
1373 let outside = temp.path().join("outside-stage");
1374 std::fs::create_dir(&outside).unwrap();
1375 let state = store.paths().state_dir(&project.id);
1376 std::fs::create_dir_all(&state).unwrap();
1377 symlink(&outside, state.join(LEGACY_MEMORY_STATE_DIR)).unwrap();
1378
1379 assert!(store.migrate_legacy_memory(&project.id, &key).is_err());
1380 assert_eq!(
1381 std::fs::read_dir(&outside).unwrap().count(),
1382 0,
1383 "migration must not create lock, journal, or staging files through a symlink"
1384 );
1385 }
1386}