From a3913bfc7f8f42efa95453d5af40241de44b421d Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Thu, 1 Oct 2026 14:16:22 +0300 Subject: [PATCH 1/3] test: move inline test modules into *_tests.rs files Extract each inline `#[cfg(test)] mod` into a sibling `_tests.rs` declared with `#[path]`, and record the rule in the repo guidance. Co-authored-by: Medulla --- AGENTS.md | 22 ++ src/memory/graph/bfs.rs | 37 +-- src/memory/graph/bfs_tests.rs | 33 +++ src/memory/graph/edge_store.rs | 24 +- src/memory/graph/edge_store_tests.rs | 20 ++ src/memory/people/address_book.rs | 91 +----- src/memory/people/address_book_tests.rs | 87 ++++++ src/memory/people/migrations.rs | 47 +-- src/memory/people/migrations_tests.rs | 43 +++ src/memory/people/resolver.rs | 353 +---------------------- src/memory/people/resolver_tests.rs | 349 ++++++++++++++++++++++ src/memory/people/scorer.rs | 116 +------- src/memory/people/scorer_tests.rs | 112 +++++++ src/memory/people/store.rs | 74 +---- src/memory/people/store_tests.rs | 70 +++++ src/memory/people/types.rs | 39 +-- src/memory/people/types_tests.rs | 35 +++ src/memory/sync/audit.rs | 62 +--- src/memory/sync/audit_tests.rs | 58 ++++ src/memory/sync/composio/client.rs | 33 +-- src/memory/sync/composio/client_tests.rs | 29 ++ src/memory/sync/github.rs | 26 +- src/memory/sync/github_tests.rs | 22 ++ src/memory/sync/periodic.rs | 78 +---- src/memory/sync/periodic_tests.rs | 74 +++++ src/memory/sync/rebuild.rs | 82 +----- src/memory/sync/rebuild_tests.rs | 78 +++++ src/memory/sync/state.rs | 68 +---- src/memory/sync/state_tests.rs | 60 ++++ src/memory/sync/status.rs | 72 +---- src/memory/sync/status_tests.rs | 68 +++++ src/memory/sync/workspace.rs | 303 +------------------ src/memory/sync/workspace_tests.rs | 295 +++++++++++++++++++ src/memory/tree/direct_ingest.rs | 41 +-- src/memory/tree/direct_ingest_tests.rs | 37 +++ 35 files changed, 1526 insertions(+), 1512 deletions(-) create mode 100644 src/memory/graph/bfs_tests.rs create mode 100644 src/memory/graph/edge_store_tests.rs create mode 100644 src/memory/people/address_book_tests.rs create mode 100644 src/memory/people/migrations_tests.rs create mode 100644 src/memory/people/resolver_tests.rs create mode 100644 src/memory/people/scorer_tests.rs create mode 100644 src/memory/people/store_tests.rs create mode 100644 src/memory/people/types_tests.rs create mode 100644 src/memory/sync/audit_tests.rs create mode 100644 src/memory/sync/composio/client_tests.rs create mode 100644 src/memory/sync/github_tests.rs create mode 100644 src/memory/sync/periodic_tests.rs create mode 100644 src/memory/sync/rebuild_tests.rs create mode 100644 src/memory/sync/state_tests.rs create mode 100644 src/memory/sync/status_tests.rs create mode 100644 src/memory/sync/workspace_tests.rs create mode 100644 src/memory/tree/direct_ingest_tests.rs diff --git a/AGENTS.md b/AGENTS.md index 62f9d445..d2084f94 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -29,3 +29,25 @@ Make small commits as much as possible, keeping each commit focused on one coher ## Parallel Work Use git worktrees when running tasks in parallel or when launching sub-agents in parallel, so each concurrent effort has an isolated checkout and does not disturb another worktree's files, build output, or branch state. + +## Tests live in `*_tests.rs` files + +- Unit tests are never inline. Do not write a `#[cfg(test)] mod tests { ... }` + block in a source file. Put the tests in a sibling `_tests.rs` + (`mod_tests.rs` beside a `mod.rs`, `lib_tests.rs` beside `lib.rs`) and declare + it at the bottom of the module: + + ```rust + #[cfg(test)] + #[path = "foo_tests.rs"] + mod tests; + ``` + +- The test file starts with `use super::*;` and carries no `#[cfg(test)]` of its + own. It is still a child module, so it reaches private items exactly as an + inline module did. +- Name test files `_tests.rs`; a second group for the same module is + `__tests.rs`. Never `test.rs`, `tests.rs` or `_test.rs`. +- Integration tests stay in the crate's `tests/` directory. +- OpenHuman's `scripts/externalize-inline-tests.mjs --write` moves + inline test modules out mechanically; without `--write` it only reports. diff --git a/src/memory/graph/bfs.rs b/src/memory/graph/bfs.rs index 2bc0f662..0c1be2ec 100644 --- a/src/memory/graph/bfs.rs +++ b/src/memory/graph/bfs.rs @@ -81,38 +81,5 @@ pub fn pair_distances( } #[cfg(test)] -mod tests { - use super::*; - use crate::memory::graph::{pairs_from_entities, upsert_edges}; - - #[test] - fn bounded_bfs_finds_two_hop_pair() { - let temp = tempfile::tempdir().unwrap(); - let config = MemoryConfig::new(temp.path()); - upsert_edges( - &config, - &pairs_from_entities(&["alice".into(), "bob".into()]), - 1, - ) - .unwrap(); - upsert_edges( - &config, - &pairs_from_entities(&["bob".into(), "carol".into()]), - 1, - ) - .unwrap(); - assert!( - pair_distances(&config, &["alice".into(), "carol".into()], 1) - .unwrap() - .is_empty() - ); - assert_eq!( - pair_distances(&config, &["alice".into(), "carol".into()], 2).unwrap(), - vec![PairDistance { - a: "alice".into(), - b: "carol".into(), - dist: 2, - }] - ); - } -} +#[path = "bfs_tests.rs"] +mod tests; diff --git a/src/memory/graph/bfs_tests.rs b/src/memory/graph/bfs_tests.rs new file mode 100644 index 00000000..e67c91f3 --- /dev/null +++ b/src/memory/graph/bfs_tests.rs @@ -0,0 +1,33 @@ +use super::*; +use crate::memory::graph::{pairs_from_entities, upsert_edges}; + +#[test] +fn bounded_bfs_finds_two_hop_pair() { + let temp = tempfile::tempdir().unwrap(); + let config = MemoryConfig::new(temp.path()); + upsert_edges( + &config, + &pairs_from_entities(&["alice".into(), "bob".into()]), + 1, + ) + .unwrap(); + upsert_edges( + &config, + &pairs_from_entities(&["bob".into(), "carol".into()]), + 1, + ) + .unwrap(); + assert!( + pair_distances(&config, &["alice".into(), "carol".into()], 1) + .unwrap() + .is_empty() + ); + assert_eq!( + pair_distances(&config, &["alice".into(), "carol".into()], 2).unwrap(), + vec![PairDistance { + a: "alice".into(), + b: "carol".into(), + dist: 2, + }] + ); +} diff --git a/src/memory/graph/edge_store.rs b/src/memory/graph/edge_store.rs index 4c7e9277..e4f3f8b2 100644 --- a/src/memory/graph/edge_store.rs +++ b/src/memory/graph/edge_store.rs @@ -106,25 +106,5 @@ pub fn count_edges(config: &MemoryConfig) -> Result { } #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn persisted_edges_are_canonical_weighted_and_symmetric() { - let temp = tempfile::tempdir().unwrap(); - let config = MemoryConfig::new(temp.path()); - let pairs = pairs_from_entities(&[ - "person:bob".into(), - "person:alice".into(), - "person:alice".into(), - ]); - assert_eq!(pairs, vec![("person:alice".into(), "person:bob".into())]); - upsert_edges(&config, &pairs, 1).unwrap(); - upsert_edges(&config, &pairs, 2).unwrap(); - assert_eq!( - edge_neighbors(&config, "person:bob").unwrap(), - vec![("person:alice".into(), 2)] - ); - assert_eq!(count_edges(&config).unwrap(), 1); - } -} +#[path = "edge_store_tests.rs"] +mod tests; diff --git a/src/memory/graph/edge_store_tests.rs b/src/memory/graph/edge_store_tests.rs new file mode 100644 index 00000000..6f9c67d0 --- /dev/null +++ b/src/memory/graph/edge_store_tests.rs @@ -0,0 +1,20 @@ +use super::*; + +#[test] +fn persisted_edges_are_canonical_weighted_and_symmetric() { + let temp = tempfile::tempdir().unwrap(); + let config = MemoryConfig::new(temp.path()); + let pairs = pairs_from_entities(&[ + "person:bob".into(), + "person:alice".into(), + "person:alice".into(), + ]); + assert_eq!(pairs, vec![("person:alice".into(), "person:bob".into())]); + upsert_edges(&config, &pairs, 1).unwrap(); + upsert_edges(&config, &pairs, 2).unwrap(); + assert_eq!( + edge_neighbors(&config, "person:bob").unwrap(), + vec![("person:alice".into(), 2)] + ); + assert_eq!(count_edges(&config).unwrap(), 1); +} diff --git a/src/memory/people/address_book.rs b/src/memory/people/address_book.rs index b02b2867..b211c25c 100644 --- a/src/memory/people/address_book.rs +++ b/src/memory/people/address_book.rs @@ -291,92 +291,5 @@ mod imp { // ── tests ───────────────────────────────────────────────────────────────────── #[cfg(test)] -pub mod tests { - use super::*; - - /// Test double that returns a canned list without any FFI calls. - pub struct MockContactsSource { - pub result: Result, AddressBookError>, - } - - impl MockContactsSource { - pub fn ok(contacts: Vec) -> Self { - Self { - result: Ok(contacts), - } - } - - pub fn permission_denied() -> Self { - Self { - result: Err(AddressBookError::PermissionDenied), - } - } - } - - impl ContactsSource for MockContactsSource { - fn fetch_contacts(&self) -> Result, AddressBookError> { - match &self.result { - Ok(v) => Ok(v.clone()), - Err(AddressBookError::PermissionDenied) => Err(AddressBookError::PermissionDenied), - Err(AddressBookError::Other(s)) => Err(AddressBookError::Other(s.clone())), - } - } - } - - fn mk_contact(name: &str, email: &str) -> AddressBookContact { - AddressBookContact { - display_name: Some(name.into()), - emails: vec![email.into()], - phones: vec![], - } - } - - #[test] - fn mock_source_returns_canned_contacts() { - let source = MockContactsSource::ok(vec![ - mk_contact("Alice", "alice@example.com"), - mk_contact("Bob", "bob@example.com"), - ]); - let result = read_with(&source).unwrap(); - assert_eq!(result.len(), 2); - assert_eq!(result[0].display_name.as_deref(), Some("Alice")); - assert_eq!(result[1].emails[0], "bob@example.com"); - } - - #[test] - fn mock_source_permission_denied_is_distinguished() { - let source = MockContactsSource::permission_denied(); - let err = read_with(&source).unwrap_err(); - assert_eq!(err, AddressBookError::PermissionDenied); - } - - #[test] - fn system_source_non_mac_returns_empty() { - // Mirrors the `imp` cfgs above: the stub is what compiles whenever the - // real CNContactStore path is absent, whether by target or by gate. - #[cfg(not(all(target_os = "macos", feature = "contacts")))] - { - let source = SystemContactsSource; - let result = read_with(&source).unwrap(); - assert!(result.is_empty()); - } - #[cfg(all(target_os = "macos", feature = "contacts"))] - { - // TCC state is environment-dependent; just verify no panic. - let source = SystemContactsSource; - let _ = read_with(&source); - } - } - - #[test] - fn contact_with_no_fields_is_excluded_by_mock() { - let source = MockContactsSource::ok(vec![AddressBookContact { - display_name: Some("Sarah Lee".into()), - emails: vec![], - phones: vec!["+1 555 000 0001".into()], - }]); - let result = read_with(&source).unwrap(); - assert_eq!(result.len(), 1); - assert_eq!(result[0].phones[0], "+1 555 000 0001"); - } -} +#[path = "address_book_tests.rs"] +pub mod tests; diff --git a/src/memory/people/address_book_tests.rs b/src/memory/people/address_book_tests.rs new file mode 100644 index 00000000..b94be973 --- /dev/null +++ b/src/memory/people/address_book_tests.rs @@ -0,0 +1,87 @@ +use super::*; + +/// Test double that returns a canned list without any FFI calls. +pub struct MockContactsSource { + pub result: Result, AddressBookError>, +} + +impl MockContactsSource { + pub fn ok(contacts: Vec) -> Self { + Self { + result: Ok(contacts), + } + } + + pub fn permission_denied() -> Self { + Self { + result: Err(AddressBookError::PermissionDenied), + } + } +} + +impl ContactsSource for MockContactsSource { + fn fetch_contacts(&self) -> Result, AddressBookError> { + match &self.result { + Ok(v) => Ok(v.clone()), + Err(AddressBookError::PermissionDenied) => Err(AddressBookError::PermissionDenied), + Err(AddressBookError::Other(s)) => Err(AddressBookError::Other(s.clone())), + } + } +} + +fn mk_contact(name: &str, email: &str) -> AddressBookContact { + AddressBookContact { + display_name: Some(name.into()), + emails: vec![email.into()], + phones: vec![], + } +} + +#[test] +fn mock_source_returns_canned_contacts() { + let source = MockContactsSource::ok(vec![ + mk_contact("Alice", "alice@example.com"), + mk_contact("Bob", "bob@example.com"), + ]); + let result = read_with(&source).unwrap(); + assert_eq!(result.len(), 2); + assert_eq!(result[0].display_name.as_deref(), Some("Alice")); + assert_eq!(result[1].emails[0], "bob@example.com"); +} + +#[test] +fn mock_source_permission_denied_is_distinguished() { + let source = MockContactsSource::permission_denied(); + let err = read_with(&source).unwrap_err(); + assert_eq!(err, AddressBookError::PermissionDenied); +} + +#[test] +fn system_source_non_mac_returns_empty() { + // Mirrors the `imp` cfgs above: the stub is what compiles whenever the + // real CNContactStore path is absent, whether by target or by gate. + #[cfg(not(all(target_os = "macos", feature = "contacts")))] + { + let source = SystemContactsSource; + let result = read_with(&source).unwrap(); + assert!(result.is_empty()); + } + #[cfg(all(target_os = "macos", feature = "contacts"))] + { + // TCC state is environment-dependent; just verify no panic. + let source = SystemContactsSource; + let _ = read_with(&source); + } +} + +#[test] +fn contact_with_no_fields_is_excluded_by_mock() { + let source = MockContactsSource::ok(vec![AddressBookContact { + display_name: Some("Sarah Lee".into()), + emails: vec![], + phones: vec!["+1 555 000 0001".into()], + }]); + let result = read_with(&source).unwrap(); + assert_eq!(result.len(), 1); + assert_eq!(result[0].phones[0], "+1 555 000 0001"); +} diff --git a/src/memory/people/migrations.rs b/src/memory/people/migrations.rs index 57d153ec..75e406a3 100644 --- a/src/memory/people/migrations.rs +++ b/src/memory/people/migrations.rs @@ -46,48 +46,5 @@ pub fn run(conn: &Connection) -> Result<()> { } #[cfg(test)] -mod tests { - use super::*; - - fn fresh() -> Connection { - Connection::open_in_memory().unwrap() - } - - #[test] - fn migrations_create_expected_tables() { - let conn = fresh(); - run(&conn).unwrap(); - let mut stmt = conn - .prepare("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name") - .unwrap(); - let names: Vec = stmt - .query_map([], |row| row.get(0)) - .unwrap() - .map(|r| r.unwrap()) - .collect(); - for expected in [ - "people", - "handle_aliases", - "interactions", - "_people_migrations", - ] { - assert!( - names.iter().any(|n| n == expected), - "missing {expected}: {names:?}" - ); - } - } - - #[test] - fn migrations_are_idempotent() { - let conn = fresh(); - run(&conn).unwrap(); - run(&conn).unwrap(); - let count: i64 = conn - .query_row("SELECT count(*) FROM _people_migrations", [], |row| { - row.get(0) - }) - .unwrap(); - assert_eq!(count, MIGRATIONS.len() as i64); - } -} +#[path = "migrations_tests.rs"] +mod tests; diff --git a/src/memory/people/migrations_tests.rs b/src/memory/people/migrations_tests.rs new file mode 100644 index 00000000..bc14e367 --- /dev/null +++ b/src/memory/people/migrations_tests.rs @@ -0,0 +1,43 @@ +use super::*; + +fn fresh() -> Connection { + Connection::open_in_memory().unwrap() +} + +#[test] +fn migrations_create_expected_tables() { + let conn = fresh(); + run(&conn).unwrap(); + let mut stmt = conn + .prepare("SELECT name FROM sqlite_master WHERE type='table' ORDER BY name") + .unwrap(); + let names: Vec = stmt + .query_map([], |row| row.get(0)) + .unwrap() + .map(|r| r.unwrap()) + .collect(); + for expected in [ + "people", + "handle_aliases", + "interactions", + "_people_migrations", + ] { + assert!( + names.iter().any(|n| n == expected), + "missing {expected}: {names:?}" + ); + } +} + +#[test] +fn migrations_are_idempotent() { + let conn = fresh(); + run(&conn).unwrap(); + run(&conn).unwrap(); + let count: i64 = conn + .query_row("SELECT count(*) FROM _people_migrations", [], |row| { + row.get(0) + }) + .unwrap(); + assert_eq!(count, MIGRATIONS.len() as i64); +} diff --git a/src/memory/people/resolver.rs b/src/memory/people/resolver.rs index 44283459..8e38cef9 100644 --- a/src/memory/people/resolver.rs +++ b/src/memory/people/resolver.rs @@ -174,354 +174,5 @@ impl<'a> HandleResolver<'a> { } #[cfg(test)] -mod tests { - use super::*; - use crate::memory::people::address_book::tests::MockContactsSource; - use crate::memory::people::types::AddressBookContact; - - #[tokio::test] - async fn resolve_returns_none_for_unknown_handle() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - let got = r.resolve(&Handle::Email("x@y.z".into())).await.unwrap(); - assert!(got.is_none()); - } - - #[tokio::test] - async fn resolve_or_create_is_deterministic_across_case_and_whitespace() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - let a = r - .resolve_or_create(&Handle::Email("Sarah@Example.COM".into())) - .await - .unwrap(); - let b = r - .resolve_or_create(&Handle::Email(" sarah@example.com ".into())) - .await - .unwrap(); - assert_eq!(a, b, "canonicalization must collapse case+whitespace"); - } - - #[tokio::test] - async fn concurrent_resolve_or_create_returns_one_database_id() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - let handles: Vec<_> = (0..16) - .map(|_| Handle::Email("Race@Example.COM".into())) - .collect(); - - let ids = futures::future::join_all(handles.iter().map(|h| r.resolve_or_create(h))).await; - let first = ids[0].as_ref().unwrap(); - for id in &ids { - assert_eq!(id.as_ref().unwrap(), first); - } - - let people = s.list().await.unwrap(); - assert_eq!(people.len(), 1); - assert_eq!(people[0].id, *first); - } - - #[tokio::test] - async fn same_email_different_display_name_resolve_same_id() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - let via_email = r - .resolve_or_create(&Handle::Email("a@b.c".into())) - .await - .unwrap(); - // Linking a display name to the same email must not mint a second id. - let via_linked = r - .link( - &Handle::Email("a@b.c".into()), - Handle::DisplayName("Alice".into()), - ) - .await - .unwrap(); - assert_eq!(via_email, via_linked); - // And now resolving the display name returns the same id. - let via_name = r - .resolve(&Handle::DisplayName("Alice".into())) - .await - .unwrap(); - assert_eq!(via_name, Some(via_email)); - } - - #[tokio::test] - async fn distinct_handles_without_linking_produce_distinct_ids() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - let a = r - .resolve_or_create(&Handle::Email("a@b.c".into())) - .await - .unwrap(); - let b = r - .resolve_or_create(&Handle::Email("x@y.z".into())) - .await - .unwrap(); - assert_ne!(a, b); - } - - #[tokio::test] - async fn seed_from_address_book_populates_store() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - let source = MockContactsSource::ok(vec![ - AddressBookContact { - display_name: Some("Alice Smith".into()), - emails: vec!["alice@example.com".into()], - phones: vec!["+1 555 000 0001".into()], - }, - AddressBookContact { - display_name: Some("Bob Jones".into()), - emails: vec!["bob@example.com".into()], - phones: vec![], - }, - ]); - - let (seeded, skipped) = r.seed_from_address_book(&source).await.unwrap(); - assert_eq!(seeded, 2, "both contacts should be seeded"); - assert_eq!(skipped, 0); - - // Alice is resolvable by email - let alice_id = r - .resolve(&Handle::Email("alice@example.com".into())) - .await - .unwrap(); - assert!(alice_id.is_some(), "alice must be resolvable after seed"); - - // Alice is also resolvable by phone (linked as alias) - let alice_via_phone = r - .resolve(&Handle::IMessage("+1 555 000 0001".into())) - .await - .unwrap(); - assert_eq!( - alice_id, alice_via_phone, - "email and phone must resolve to same person" - ); - - // Bob is resolvable - let bob_id = r - .resolve(&Handle::Email("bob@example.com".into())) - .await - .unwrap(); - assert!(bob_id.is_some()); - assert_ne!(alice_id, bob_id, "distinct contacts must have distinct ids"); - } - - #[tokio::test] - async fn seed_from_address_book_permission_denied_is_propagated() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - let source = MockContactsSource::permission_denied(); - let err = r.seed_from_address_book(&source).await.unwrap_err(); - assert_eq!(err, AddressBookError::PermissionDenied); - - // Store must still be empty — no partial writes. - let people = s.list().await.unwrap(); - assert!( - people.is_empty(), - "no people should be inserted on permission denied" - ); - } - - #[tokio::test] - async fn seed_is_idempotent() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - let source = MockContactsSource::ok(vec![AddressBookContact { - display_name: Some("Carol".into()), - emails: vec!["carol@example.com".into()], - phones: vec![], - }]); - - let (s1, _) = r.seed_from_address_book(&source).await.unwrap(); - let (s2, _) = r.seed_from_address_book(&source).await.unwrap(); - assert_eq!(s1, 1); - assert_eq!(s2, 1, "second seed call should still report 1 (upsert)"); - - // Only one person in store. - let people = s.list().await.unwrap(); - assert_eq!(people.len(), 1, "idempotent — must not duplicate"); - } - - #[tokio::test] - async fn contact_with_only_display_name_is_seeded() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - let source = MockContactsSource::ok(vec![AddressBookContact { - display_name: Some("No Email Person".into()), - emails: vec![], - phones: vec![], - }]); - let (seeded, skipped) = r.seed_from_address_book(&source).await.unwrap(); - assert_eq!(seeded, 1); - assert_eq!(skipped, 0); - } - - #[tokio::test] - async fn contact_with_no_fields_is_skipped() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - let source = MockContactsSource::ok(vec![AddressBookContact { - display_name: None, - emails: vec![], - phones: vec![], - }]); - let (seeded, skipped) = r.seed_from_address_book(&source).await.unwrap(); - assert_eq!(seeded, 0); - assert_eq!(skipped, 1); - } - - // ── Cross-source merge safety tests (issue#1538) ────────────────────────── - // - // The people resolver must NOT silently merge two distinct identities that - // happen to share only a display name or only an unverified handle from - // different sources. These tests lock in the "ambiguous cross-source" - // contract: two handles from unrelated sources remain distinct unless - // explicitly linked via `link()`. - - /// Two contacts that share only a display name (no email or phone overlap) - /// must NOT be merged — they may be homonymous individuals. - #[tokio::test] - async fn same_display_name_from_different_sources_does_not_merge() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - // Source A — email-backed identity - let id_a = r - .resolve_or_create(&Handle::Email("alice@company-a.com".into())) - .await - .unwrap(); - r.link( - &Handle::Email("alice@company-a.com".into()), - Handle::DisplayName("Alice Smith".into()), - ) - .await - .unwrap(); - - // Source B — different email; the same display name surfaces again, - // but as a *separate* DisplayName-backed mint (NOT linked to either - // email). This is the actual collision scenario: two ingestion paths - // both encounter "Alice Smith" without any cross-source identifier. - let id_b = r - .resolve_or_create(&Handle::Email("alice@company-b.com".into())) - .await - .unwrap(); - // The display-name resolver must already pin to id_a (linked above), - // so a second mint of the same DisplayName does NOT spawn a third - // identity — but crucially it also does NOT silently merge id_b into id_a. - let id_name_again = r - .resolve_or_create(&Handle::DisplayName("Alice Smith".into())) - .await - .unwrap(); - - // The two email-backed identities must be distinct. - assert_ne!( - id_a, id_b, - "two email handles with identical display names must not be merged without explicit link" - ); - - // The repeated DisplayName mint resolves to the linked identity (id_a), - // NOT to id_b. If display names auto-merged, id_b would have collapsed - // into id_a; if they minted fresh on every call, this would be a third id. - assert_eq!( - id_name_again, id_a, - "repeated DisplayName mint should resolve to the existing linked identity" - ); - assert_ne!( - id_name_again, id_b, - "DisplayName collision must not silently merge id_b into id_a" - ); - - // Resolving the display name returns the ONE identity that was explicitly linked. - let via_name = r - .resolve(&Handle::DisplayName("Alice Smith".into())) - .await - .unwrap(); - assert_eq!( - via_name, - Some(id_a), - "display name resolves to the explicitly linked identity" - ); - - // company-b Alice is still addressable by email only. - let via_b_email = r - .resolve(&Handle::Email("alice@company-b.com".into())) - .await - .unwrap(); - assert_eq!(via_b_email, Some(id_b)); - } - - /// Minting the same email handle from two logically distinct call sites - /// must always collapse to one `PersonId` (idempotent mint). This is the - /// safe side of cross-source: we never mint duplicates for an identical - /// canonical handle. - #[tokio::test] - async fn same_email_from_two_sources_collapses_to_one_person() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - // Simulate two different ingestion paths (gmail vs slack) that both - // surface the same email address. - let from_gmail = r - .resolve_or_create(&Handle::Email("shared@example.com".into())) - .await - .unwrap(); - let from_slack = r - .resolve_or_create(&Handle::Email("shared@example.com".into())) - .await - .unwrap(); - - assert_eq!( - from_gmail, from_slack, - "identical canonical email from two ingestion paths must resolve to one PersonId" - ); - - // Exactly one person in the store. - let people = s.list().await.unwrap(); - assert_eq!( - people.len(), - 1, - "no duplicate person rows must exist for the same canonical email" - ); - } - - /// An iMessage phone handle from one source and an email from a different - /// source for the SAME real person must stay distinct until explicitly linked. - /// Memory must not unsafely merge the same person's identities across sources - /// (issue#1538). - #[tokio::test] - async fn phone_and_email_from_different_sources_are_not_merged_without_link() { - let s = PeopleStore::open_in_memory().unwrap(); - let r = HandleResolver::new(&s); - - // iMessage source sees only a phone. - let id_phone = r - .resolve_or_create(&Handle::IMessage("+15550001234".into())) - .await - .unwrap(); - - // Gmail source sees only an email. - let id_email = r - .resolve_or_create(&Handle::Email("sam@example.com".into())) - .await - .unwrap(); - - // Without an explicit link these are separate identities. This is the - // contract under test — cross-source handles for the same real person - // must NOT auto-merge. Asserting post-link merge semantics is out of - // scope: link()'s exact propagation rule (does the email handle - // afterwards canonically resolve to the phone PersonId, or remain - // independent with only the link table updated?) is a separate - // behavior tested in store_tests.rs. - assert_ne!( - id_phone, id_email, - "phone and email from unrelated sources must not be auto-merged" - ); - } -} +#[path = "resolver_tests.rs"] +mod tests; diff --git a/src/memory/people/resolver_tests.rs b/src/memory/people/resolver_tests.rs new file mode 100644 index 00000000..cf5a6842 --- /dev/null +++ b/src/memory/people/resolver_tests.rs @@ -0,0 +1,349 @@ +use super::*; +use crate::memory::people::address_book::tests::MockContactsSource; +use crate::memory::people::types::AddressBookContact; + +#[tokio::test] +async fn resolve_returns_none_for_unknown_handle() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + let got = r.resolve(&Handle::Email("x@y.z".into())).await.unwrap(); + assert!(got.is_none()); +} + +#[tokio::test] +async fn resolve_or_create_is_deterministic_across_case_and_whitespace() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + let a = r + .resolve_or_create(&Handle::Email("Sarah@Example.COM".into())) + .await + .unwrap(); + let b = r + .resolve_or_create(&Handle::Email(" sarah@example.com ".into())) + .await + .unwrap(); + assert_eq!(a, b, "canonicalization must collapse case+whitespace"); +} + +#[tokio::test] +async fn concurrent_resolve_or_create_returns_one_database_id() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + let handles: Vec<_> = (0..16) + .map(|_| Handle::Email("Race@Example.COM".into())) + .collect(); + + let ids = futures::future::join_all(handles.iter().map(|h| r.resolve_or_create(h))).await; + let first = ids[0].as_ref().unwrap(); + for id in &ids { + assert_eq!(id.as_ref().unwrap(), first); + } + + let people = s.list().await.unwrap(); + assert_eq!(people.len(), 1); + assert_eq!(people[0].id, *first); +} + +#[tokio::test] +async fn same_email_different_display_name_resolve_same_id() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + let via_email = r + .resolve_or_create(&Handle::Email("a@b.c".into())) + .await + .unwrap(); + // Linking a display name to the same email must not mint a second id. + let via_linked = r + .link( + &Handle::Email("a@b.c".into()), + Handle::DisplayName("Alice".into()), + ) + .await + .unwrap(); + assert_eq!(via_email, via_linked); + // And now resolving the display name returns the same id. + let via_name = r + .resolve(&Handle::DisplayName("Alice".into())) + .await + .unwrap(); + assert_eq!(via_name, Some(via_email)); +} + +#[tokio::test] +async fn distinct_handles_without_linking_produce_distinct_ids() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + let a = r + .resolve_or_create(&Handle::Email("a@b.c".into())) + .await + .unwrap(); + let b = r + .resolve_or_create(&Handle::Email("x@y.z".into())) + .await + .unwrap(); + assert_ne!(a, b); +} + +#[tokio::test] +async fn seed_from_address_book_populates_store() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + let source = MockContactsSource::ok(vec![ + AddressBookContact { + display_name: Some("Alice Smith".into()), + emails: vec!["alice@example.com".into()], + phones: vec!["+1 555 000 0001".into()], + }, + AddressBookContact { + display_name: Some("Bob Jones".into()), + emails: vec!["bob@example.com".into()], + phones: vec![], + }, + ]); + + let (seeded, skipped) = r.seed_from_address_book(&source).await.unwrap(); + assert_eq!(seeded, 2, "both contacts should be seeded"); + assert_eq!(skipped, 0); + + // Alice is resolvable by email + let alice_id = r + .resolve(&Handle::Email("alice@example.com".into())) + .await + .unwrap(); + assert!(alice_id.is_some(), "alice must be resolvable after seed"); + + // Alice is also resolvable by phone (linked as alias) + let alice_via_phone = r + .resolve(&Handle::IMessage("+1 555 000 0001".into())) + .await + .unwrap(); + assert_eq!( + alice_id, alice_via_phone, + "email and phone must resolve to same person" + ); + + // Bob is resolvable + let bob_id = r + .resolve(&Handle::Email("bob@example.com".into())) + .await + .unwrap(); + assert!(bob_id.is_some()); + assert_ne!(alice_id, bob_id, "distinct contacts must have distinct ids"); +} + +#[tokio::test] +async fn seed_from_address_book_permission_denied_is_propagated() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + let source = MockContactsSource::permission_denied(); + let err = r.seed_from_address_book(&source).await.unwrap_err(); + assert_eq!(err, AddressBookError::PermissionDenied); + + // Store must still be empty — no partial writes. + let people = s.list().await.unwrap(); + assert!( + people.is_empty(), + "no people should be inserted on permission denied" + ); +} + +#[tokio::test] +async fn seed_is_idempotent() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + let source = MockContactsSource::ok(vec![AddressBookContact { + display_name: Some("Carol".into()), + emails: vec!["carol@example.com".into()], + phones: vec![], + }]); + + let (s1, _) = r.seed_from_address_book(&source).await.unwrap(); + let (s2, _) = r.seed_from_address_book(&source).await.unwrap(); + assert_eq!(s1, 1); + assert_eq!(s2, 1, "second seed call should still report 1 (upsert)"); + + // Only one person in store. + let people = s.list().await.unwrap(); + assert_eq!(people.len(), 1, "idempotent — must not duplicate"); +} + +#[tokio::test] +async fn contact_with_only_display_name_is_seeded() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + let source = MockContactsSource::ok(vec![AddressBookContact { + display_name: Some("No Email Person".into()), + emails: vec![], + phones: vec![], + }]); + let (seeded, skipped) = r.seed_from_address_book(&source).await.unwrap(); + assert_eq!(seeded, 1); + assert_eq!(skipped, 0); +} + +#[tokio::test] +async fn contact_with_no_fields_is_skipped() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + let source = MockContactsSource::ok(vec![AddressBookContact { + display_name: None, + emails: vec![], + phones: vec![], + }]); + let (seeded, skipped) = r.seed_from_address_book(&source).await.unwrap(); + assert_eq!(seeded, 0); + assert_eq!(skipped, 1); +} + +// ── Cross-source merge safety tests (issue#1538) ────────────────────────── +// +// The people resolver must NOT silently merge two distinct identities that +// happen to share only a display name or only an unverified handle from +// different sources. These tests lock in the "ambiguous cross-source" +// contract: two handles from unrelated sources remain distinct unless +// explicitly linked via `link()`. + +/// Two contacts that share only a display name (no email or phone overlap) +/// must NOT be merged — they may be homonymous individuals. +#[tokio::test] +async fn same_display_name_from_different_sources_does_not_merge() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + // Source A — email-backed identity + let id_a = r + .resolve_or_create(&Handle::Email("alice@company-a.com".into())) + .await + .unwrap(); + r.link( + &Handle::Email("alice@company-a.com".into()), + Handle::DisplayName("Alice Smith".into()), + ) + .await + .unwrap(); + + // Source B — different email; the same display name surfaces again, + // but as a *separate* DisplayName-backed mint (NOT linked to either + // email). This is the actual collision scenario: two ingestion paths + // both encounter "Alice Smith" without any cross-source identifier. + let id_b = r + .resolve_or_create(&Handle::Email("alice@company-b.com".into())) + .await + .unwrap(); + // The display-name resolver must already pin to id_a (linked above), + // so a second mint of the same DisplayName does NOT spawn a third + // identity — but crucially it also does NOT silently merge id_b into id_a. + let id_name_again = r + .resolve_or_create(&Handle::DisplayName("Alice Smith".into())) + .await + .unwrap(); + + // The two email-backed identities must be distinct. + assert_ne!( + id_a, id_b, + "two email handles with identical display names must not be merged without explicit link" + ); + + // The repeated DisplayName mint resolves to the linked identity (id_a), + // NOT to id_b. If display names auto-merged, id_b would have collapsed + // into id_a; if they minted fresh on every call, this would be a third id. + assert_eq!( + id_name_again, id_a, + "repeated DisplayName mint should resolve to the existing linked identity" + ); + assert_ne!( + id_name_again, id_b, + "DisplayName collision must not silently merge id_b into id_a" + ); + + // Resolving the display name returns the ONE identity that was explicitly linked. + let via_name = r + .resolve(&Handle::DisplayName("Alice Smith".into())) + .await + .unwrap(); + assert_eq!( + via_name, + Some(id_a), + "display name resolves to the explicitly linked identity" + ); + + // company-b Alice is still addressable by email only. + let via_b_email = r + .resolve(&Handle::Email("alice@company-b.com".into())) + .await + .unwrap(); + assert_eq!(via_b_email, Some(id_b)); +} + +/// Minting the same email handle from two logically distinct call sites +/// must always collapse to one `PersonId` (idempotent mint). This is the +/// safe side of cross-source: we never mint duplicates for an identical +/// canonical handle. +#[tokio::test] +async fn same_email_from_two_sources_collapses_to_one_person() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + // Simulate two different ingestion paths (gmail vs slack) that both + // surface the same email address. + let from_gmail = r + .resolve_or_create(&Handle::Email("shared@example.com".into())) + .await + .unwrap(); + let from_slack = r + .resolve_or_create(&Handle::Email("shared@example.com".into())) + .await + .unwrap(); + + assert_eq!( + from_gmail, from_slack, + "identical canonical email from two ingestion paths must resolve to one PersonId" + ); + + // Exactly one person in the store. + let people = s.list().await.unwrap(); + assert_eq!( + people.len(), + 1, + "no duplicate person rows must exist for the same canonical email" + ); +} + +/// An iMessage phone handle from one source and an email from a different +/// source for the SAME real person must stay distinct until explicitly linked. +/// Memory must not unsafely merge the same person's identities across sources +/// (issue#1538). +#[tokio::test] +async fn phone_and_email_from_different_sources_are_not_merged_without_link() { + let s = PeopleStore::open_in_memory().unwrap(); + let r = HandleResolver::new(&s); + + // iMessage source sees only a phone. + let id_phone = r + .resolve_or_create(&Handle::IMessage("+15550001234".into())) + .await + .unwrap(); + + // Gmail source sees only an email. + let id_email = r + .resolve_or_create(&Handle::Email("sam@example.com".into())) + .await + .unwrap(); + + // Without an explicit link these are separate identities. This is the + // contract under test — cross-source handles for the same real person + // must NOT auto-merge. Asserting post-link merge semantics is out of + // scope: link()'s exact propagation rule (does the email handle + // afterwards canonically resolve to the phone PersonId, or remain + // independent with only the link table updated?) is a separate + // behavior tested in store_tests.rs. + assert_ne!( + id_phone, id_email, + "phone and email from unrelated sources must not be auto-merged" + ); +} diff --git a/src/memory/people/scorer.rs b/src/memory/people/scorer.rs index e9a455ef..ce175244 100644 --- a/src/memory/people/scorer.rs +++ b/src/memory/people/scorer.rs @@ -94,117 +94,5 @@ pub fn score(interactions: &[Interaction], now: DateTime) -> ScoreComponent } #[cfg(test)] -mod tests { - use super::*; - use crate::memory::people::types::PersonId; - use chrono::Duration; - - fn mk(ts: DateTime, outbound: bool, length: u32) -> Interaction { - Interaction { - person_id: PersonId::new(), - ts, - is_outbound: outbound, - length, - } - } - - #[test] - fn empty_interactions_score_zero() { - let s = score(&[], Utc::now()); - assert_eq!(s.score, 0.0); - assert_eq!(s.recency, 0.0); - assert_eq!(s.frequency, 0.0); - } - - #[test] - fn recency_half_life_matches_config() { - let now = Utc::now(); - let half_ago = now - Duration::days(RECENCY_HALF_LIFE_DAYS as i64); - let s = score(&[mk(half_ago, true, 100)], now); - // Half-life point → recency ≈ 0.5 (allow small float slack). - assert!((s.recency - 0.5).abs() < 0.05, "got {}", s.recency); - } - - #[test] - fn all_components_clamped_to_unit_interval() { - let now = Utc::now(); - let interactions: Vec = (0..200) - .map(|i| mk(now - Duration::hours(i), i % 2 == 0, 10_000)) - .collect(); - let s = score(&interactions, now); - for c in [s.recency, s.frequency, s.reciprocity, s.depth, s.score] { - assert!((0.0..=1.0).contains(&c), "component out of range: {c}"); - } - // 200 interactions all within a few days → window_count ≥ FREQUENCY_CAP - assert_eq!(s.frequency, 1.0); - assert_eq!(s.depth, 1.0); - } - - #[test] - fn one_sided_conversation_has_zero_reciprocity() { - let now = Utc::now(); - let v: Vec<_> = (0..5) - .map(|i| mk(now - Duration::hours(i), true, 100)) - .collect(); - let s = score(&v, now); - assert_eq!(s.reciprocity, 0.0); - assert_eq!( - s.score, 0.0, - "composite must be zero when any factor is zero" - ); - } - - #[test] - fn deterministic_given_same_inputs() { - let now = Utc::now(); - let v = vec![ - mk(now - Duration::days(1), true, 100), - mk(now - Duration::days(2), false, 150), - mk(now - Duration::days(3), true, 200), - ]; - let a = score(&v, now); - let b = score(&v, now); - assert_eq!(a.score, b.score); - assert_eq!(a.recency, b.recency); - } - - #[test] - fn old_burst_does_not_inflate_frequency_score() { - // 100 interactions from 90 days ago (outside FREQUENCY_WINDOW_DAYS=30) - // should contribute 0 to frequency; 1 interaction today should give - // 1/FREQUENCY_CAP. - let now = Utc::now(); - let mut v: Vec = (0..100) - .map(|i| mk(now - Duration::days(90 + i), true, 100)) - .collect(); - // Add one recent interaction to avoid zero reciprocity forcing score=0 - v.push(mk(now - Duration::hours(1), false, 100)); - let s = score(&v, now); - // Only 1 interaction falls within the 30-day window. - let expected_frequency = 1.0 / FREQUENCY_CAP; - assert!( - (s.frequency - expected_frequency).abs() < 0.001, - "frequency should be {expected_frequency}, got {}", - s.frequency - ); - } - - #[test] - fn interactions_exactly_at_window_boundary_are_included() { - let now = Utc::now(); - // Interaction exactly FREQUENCY_WINDOW_DAYS ago — should be included - // (boundary is inclusive via >=). - let boundary = now - Duration::days(FREQUENCY_WINDOW_DAYS as i64); - let v = vec![ - mk(boundary, true, 100), - mk(now - Duration::hours(1), false, 100), - ]; - let s = score(&v, now); - let expected = 2.0 / FREQUENCY_CAP; - assert!( - (s.frequency - expected).abs() < 0.001, - "expected {expected} got {}", - s.frequency - ); - } -} +#[path = "scorer_tests.rs"] +mod tests; diff --git a/src/memory/people/scorer_tests.rs b/src/memory/people/scorer_tests.rs new file mode 100644 index 00000000..2b3161be --- /dev/null +++ b/src/memory/people/scorer_tests.rs @@ -0,0 +1,112 @@ +use super::*; +use crate::memory::people::types::PersonId; +use chrono::Duration; + +fn mk(ts: DateTime, outbound: bool, length: u32) -> Interaction { + Interaction { + person_id: PersonId::new(), + ts, + is_outbound: outbound, + length, + } +} + +#[test] +fn empty_interactions_score_zero() { + let s = score(&[], Utc::now()); + assert_eq!(s.score, 0.0); + assert_eq!(s.recency, 0.0); + assert_eq!(s.frequency, 0.0); +} + +#[test] +fn recency_half_life_matches_config() { + let now = Utc::now(); + let half_ago = now - Duration::days(RECENCY_HALF_LIFE_DAYS as i64); + let s = score(&[mk(half_ago, true, 100)], now); + // Half-life point → recency ≈ 0.5 (allow small float slack). + assert!((s.recency - 0.5).abs() < 0.05, "got {}", s.recency); +} + +#[test] +fn all_components_clamped_to_unit_interval() { + let now = Utc::now(); + let interactions: Vec = (0..200) + .map(|i| mk(now - Duration::hours(i), i % 2 == 0, 10_000)) + .collect(); + let s = score(&interactions, now); + for c in [s.recency, s.frequency, s.reciprocity, s.depth, s.score] { + assert!((0.0..=1.0).contains(&c), "component out of range: {c}"); + } + // 200 interactions all within a few days → window_count ≥ FREQUENCY_CAP + assert_eq!(s.frequency, 1.0); + assert_eq!(s.depth, 1.0); +} + +#[test] +fn one_sided_conversation_has_zero_reciprocity() { + let now = Utc::now(); + let v: Vec<_> = (0..5) + .map(|i| mk(now - Duration::hours(i), true, 100)) + .collect(); + let s = score(&v, now); + assert_eq!(s.reciprocity, 0.0); + assert_eq!( + s.score, 0.0, + "composite must be zero when any factor is zero" + ); +} + +#[test] +fn deterministic_given_same_inputs() { + let now = Utc::now(); + let v = vec![ + mk(now - Duration::days(1), true, 100), + mk(now - Duration::days(2), false, 150), + mk(now - Duration::days(3), true, 200), + ]; + let a = score(&v, now); + let b = score(&v, now); + assert_eq!(a.score, b.score); + assert_eq!(a.recency, b.recency); +} + +#[test] +fn old_burst_does_not_inflate_frequency_score() { + // 100 interactions from 90 days ago (outside FREQUENCY_WINDOW_DAYS=30) + // should contribute 0 to frequency; 1 interaction today should give + // 1/FREQUENCY_CAP. + let now = Utc::now(); + let mut v: Vec = (0..100) + .map(|i| mk(now - Duration::days(90 + i), true, 100)) + .collect(); + // Add one recent interaction to avoid zero reciprocity forcing score=0 + v.push(mk(now - Duration::hours(1), false, 100)); + let s = score(&v, now); + // Only 1 interaction falls within the 30-day window. + let expected_frequency = 1.0 / FREQUENCY_CAP; + assert!( + (s.frequency - expected_frequency).abs() < 0.001, + "frequency should be {expected_frequency}, got {}", + s.frequency + ); +} + +#[test] +fn interactions_exactly_at_window_boundary_are_included() { + let now = Utc::now(); + // Interaction exactly FREQUENCY_WINDOW_DAYS ago — should be included + // (boundary is inclusive via >=). + let boundary = now - Duration::days(FREQUENCY_WINDOW_DAYS as i64); + let v = vec![ + mk(boundary, true, 100), + mk(now - Duration::hours(1), false, 100), + ]; + let s = score(&v, now); + let expected = 2.0 / FREQUENCY_CAP; + assert!( + (s.frequency - expected).abs() < 0.001, + "expected {expected} got {}", + s.frequency + ); +} diff --git a/src/memory/people/store.rs b/src/memory/people/store.rs index 3855d939..084e4d7f 100644 --- a/src/memory/people/store.rs +++ b/src/memory/people/store.rs @@ -579,75 +579,5 @@ fn ts_to_dt(ts: i64) -> DateTime { } #[cfg(test)] -mod tests { - use super::*; - - #[tokio::test] - async fn insert_list_and_lookup_round_trip() { - let s = PeopleStore::open_in_memory().unwrap(); - let now = Utc::now(); - let p = Person { - id: PersonId::new(), - display_name: Some("Sarah Lee".into()), - primary_email: Some("sarah@example.com".into()), - primary_phone: None, - handles: vec![], - created_at: now, - updated_at: now, - }; - s.insert_person( - &p, - &[ - Handle::Email("Sarah@Example.com".into()), - Handle::DisplayName("Sarah Lee".into()), - ], - ) - .await - .unwrap(); - - let got = s - .lookup(&Handle::Email("sarah@example.com".into())) - .await - .unwrap(); - assert_eq!(got, Some(p.id)); - - let list = s.list().await.unwrap(); - assert_eq!(list.len(), 1); - assert_eq!(list[0].handles.len(), 2); - } - - #[tokio::test] - async fn interactions_round_trip() { - let s = PeopleStore::open_in_memory().unwrap(); - let now = Utc::now(); - let pid = PersonId::new(); - let p = Person { - id: pid, - display_name: Some("X".into()), - primary_email: None, - primary_phone: None, - handles: vec![], - created_at: now, - updated_at: now, - }; - s.insert_person(&p, &[]).await.unwrap(); - s.record_interaction(Interaction { - person_id: pid, - ts: now, - is_outbound: true, - length: 100, - }) - .await - .unwrap(); - s.record_interaction(Interaction { - person_id: pid, - ts: now, - is_outbound: false, - length: 50, - }) - .await - .unwrap(); - let ints = s.interactions_for(pid).await.unwrap(); - assert_eq!(ints.len(), 2); - } -} +#[path = "store_tests.rs"] +mod tests; diff --git a/src/memory/people/store_tests.rs b/src/memory/people/store_tests.rs new file mode 100644 index 00000000..73876400 --- /dev/null +++ b/src/memory/people/store_tests.rs @@ -0,0 +1,70 @@ +use super::*; + +#[tokio::test] +async fn insert_list_and_lookup_round_trip() { + let s = PeopleStore::open_in_memory().unwrap(); + let now = Utc::now(); + let p = Person { + id: PersonId::new(), + display_name: Some("Sarah Lee".into()), + primary_email: Some("sarah@example.com".into()), + primary_phone: None, + handles: vec![], + created_at: now, + updated_at: now, + }; + s.insert_person( + &p, + &[ + Handle::Email("Sarah@Example.com".into()), + Handle::DisplayName("Sarah Lee".into()), + ], + ) + .await + .unwrap(); + + let got = s + .lookup(&Handle::Email("sarah@example.com".into())) + .await + .unwrap(); + assert_eq!(got, Some(p.id)); + + let list = s.list().await.unwrap(); + assert_eq!(list.len(), 1); + assert_eq!(list[0].handles.len(), 2); +} + +#[tokio::test] +async fn interactions_round_trip() { + let s = PeopleStore::open_in_memory().unwrap(); + let now = Utc::now(); + let pid = PersonId::new(); + let p = Person { + id: pid, + display_name: Some("X".into()), + primary_email: None, + primary_phone: None, + handles: vec![], + created_at: now, + updated_at: now, + }; + s.insert_person(&p, &[]).await.unwrap(); + s.record_interaction(Interaction { + person_id: pid, + ts: now, + is_outbound: true, + length: 100, + }) + .await + .unwrap(); + s.record_interaction(Interaction { + person_id: pid, + ts: now, + is_outbound: false, + length: 50, + }) + .await + .unwrap(); + let ints = s.interactions_for(pid).await.unwrap(); + assert_eq!(ints.len(), 2); +} diff --git a/src/memory/people/types.rs b/src/memory/people/types.rs index 34a0ec70..49504461 100644 --- a/src/memory/people/types.rs +++ b/src/memory/people/types.rs @@ -120,40 +120,5 @@ pub struct AddressBookContact { } #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn handle_canonicalize_lowercases_emails_and_imessage() { - assert_eq!( - Handle::Email(" Foo@Example.COM ".into()).canonicalize(), - Handle::Email("foo@example.com".into()) - ); - assert_eq!( - Handle::IMessage("+1 (555) 123".into()).canonicalize(), - Handle::IMessage("+1 (555) 123".into()) - ); - assert_eq!( - Handle::IMessage(" Foo@Bar.com ".into()).canonicalize(), - Handle::IMessage("foo@bar.com".into()) - ); - } - - #[test] - fn handle_canonicalize_collapses_display_name_whitespace() { - assert_eq!( - Handle::DisplayName(" Sarah Lee ".into()).canonicalize(), - Handle::DisplayName("Sarah Lee".into()) - ); - } - - #[test] - fn handle_as_key_returns_correct_kind() { - assert_eq!(Handle::Email("a@b.c".into()).as_key(), ("email", "a@b.c")); - assert_eq!(Handle::IMessage("+1".into()).as_key(), ("imessage", "+1")); - assert_eq!( - Handle::DisplayName("X".into()).as_key(), - ("display_name", "X") - ); - } -} +#[path = "types_tests.rs"] +mod tests; diff --git a/src/memory/people/types_tests.rs b/src/memory/people/types_tests.rs new file mode 100644 index 00000000..cd55284b --- /dev/null +++ b/src/memory/people/types_tests.rs @@ -0,0 +1,35 @@ +use super::*; + +#[test] +fn handle_canonicalize_lowercases_emails_and_imessage() { + assert_eq!( + Handle::Email(" Foo@Example.COM ".into()).canonicalize(), + Handle::Email("foo@example.com".into()) + ); + assert_eq!( + Handle::IMessage("+1 (555) 123".into()).canonicalize(), + Handle::IMessage("+1 (555) 123".into()) + ); + assert_eq!( + Handle::IMessage(" Foo@Bar.com ".into()).canonicalize(), + Handle::IMessage("foo@bar.com".into()) + ); +} + +#[test] +fn handle_canonicalize_collapses_display_name_whitespace() { + assert_eq!( + Handle::DisplayName(" Sarah Lee ".into()).canonicalize(), + Handle::DisplayName("Sarah Lee".into()) + ); +} + +#[test] +fn handle_as_key_returns_correct_kind() { + assert_eq!(Handle::Email("a@b.c".into()).as_key(), ("email", "a@b.c")); + assert_eq!(Handle::IMessage("+1".into()).as_key(), ("imessage", "+1")); + assert_eq!( + Handle::DisplayName("X".into()).as_key(), + ("display_name", "X") + ); +} diff --git a/src/memory/sync/audit.rs b/src/memory/sync/audit.rs index 3a4711bb..6bb4db15 100644 --- a/src/memory/sync/audit.rs +++ b/src/memory/sync/audit.rs @@ -167,63 +167,5 @@ impl RealCostAccumulator { } #[cfg(test)] -mod tests { - use super::*; - - fn entry(id: &str, timestamp: DateTime) -> SyncAuditEntry { - SyncAuditEntry { - timestamp, - source_id: id.into(), - source_kind: "test".into(), - scope: "all".into(), - items_fetched: 1, - batches: 1, - input_tokens: 10, - output_tokens: 2, - estimated_cost_usd: 0.1, - composio_actions_called: 1, - composio_cost_usd: 0.02, - actual_charged_usd: None, - duration_ms: 5, - success: true, - error: None, - } - } - - #[test] - fn audit_round_trip_is_newest_first_and_skips_malformed_lines() { - let temp = tempfile::tempdir().unwrap(); - let config = MemoryConfig::new(temp.path()); - let first = entry("first", Utc::now()); - let second = entry("second", Utc::now()); - append_audit_entry(&config, &first).unwrap(); - std::fs::OpenOptions::new() - .append(true) - .open(temp.path().join("memory_tree/sync_audit.jsonl")) - .unwrap() - .write_all(b"not-json\n") - .unwrap(); - append_audit_entry(&config, &second).unwrap(); - let entries = read_audit_log(&config).unwrap(); - assert_eq!(entries.len(), 2); - assert_eq!(entries[0].source_id, "second"); - assert_eq!(entries[1].source_id, "first"); - } - - #[test] - fn accumulator_uses_real_values_only_when_every_batch_reports_them() { - let mut complete = RealCostAccumulator::new(); - complete.add_batch(100, 10, 80, 8, Some(0.01)); - complete.add_batch(100, 10, 90, 9, Some(0.02)); - assert_eq!(complete.audit_input_tokens(), 170); - assert_eq!(complete.actual_charged_usd(), Some(0.03)); - assert!(complete.cost_is_actual()); - - let mut partial = RealCostAccumulator::new(); - partial.add_batch(100, 10, 80, 8, Some(0.01)); - partial.add_batch(200, 20, 0, 0, None); - assert_eq!(partial.audit_input_tokens(), 300); - assert_eq!(partial.actual_charged_usd(), None); - assert!(!partial.usage_is_real()); - } -} +#[path = "audit_tests.rs"] +mod tests; diff --git a/src/memory/sync/audit_tests.rs b/src/memory/sync/audit_tests.rs new file mode 100644 index 00000000..fbbea383 --- /dev/null +++ b/src/memory/sync/audit_tests.rs @@ -0,0 +1,58 @@ +use super::*; + +fn entry(id: &str, timestamp: DateTime) -> SyncAuditEntry { + SyncAuditEntry { + timestamp, + source_id: id.into(), + source_kind: "test".into(), + scope: "all".into(), + items_fetched: 1, + batches: 1, + input_tokens: 10, + output_tokens: 2, + estimated_cost_usd: 0.1, + composio_actions_called: 1, + composio_cost_usd: 0.02, + actual_charged_usd: None, + duration_ms: 5, + success: true, + error: None, + } +} + +#[test] +fn audit_round_trip_is_newest_first_and_skips_malformed_lines() { + let temp = tempfile::tempdir().unwrap(); + let config = MemoryConfig::new(temp.path()); + let first = entry("first", Utc::now()); + let second = entry("second", Utc::now()); + append_audit_entry(&config, &first).unwrap(); + std::fs::OpenOptions::new() + .append(true) + .open(temp.path().join("memory_tree/sync_audit.jsonl")) + .unwrap() + .write_all(b"not-json\n") + .unwrap(); + append_audit_entry(&config, &second).unwrap(); + let entries = read_audit_log(&config).unwrap(); + assert_eq!(entries.len(), 2); + assert_eq!(entries[0].source_id, "second"); + assert_eq!(entries[1].source_id, "first"); +} + +#[test] +fn accumulator_uses_real_values_only_when_every_batch_reports_them() { + let mut complete = RealCostAccumulator::new(); + complete.add_batch(100, 10, 80, 8, Some(0.01)); + complete.add_batch(100, 10, 90, 9, Some(0.02)); + assert_eq!(complete.audit_input_tokens(), 170); + assert_eq!(complete.actual_charged_usd(), Some(0.03)); + assert!(complete.cost_is_actual()); + + let mut partial = RealCostAccumulator::new(); + partial.add_batch(100, 10, 80, 8, Some(0.01)); + partial.add_batch(200, 20, 0, 0, None); + assert_eq!(partial.audit_input_tokens(), 300); + assert_eq!(partial.actual_charged_usd(), None); + assert!(!partial.usage_is_real()); +} diff --git a/src/memory/sync/composio/client.rs b/src/memory/sync/composio/client.rs index d2c4ce57..0d993188 100644 --- a/src/memory/sync/composio/client.rs +++ b/src/memory/sync/composio/client.rs @@ -271,34 +271,5 @@ async fn decode_response( } #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn proxied_backend_envelope_decodes_provider_response() { - let response = decode_proxy_response(serde_json::json!({ - "success": true, - "data": { - "successful": true, - "data": {"messages": [{"messageId": "message-1"}]}, - "error": null - } - })) - .unwrap(); - - assert!(response.successful); - assert_eq!(response.data["messages"][0]["messageId"], "message-1"); - } - - #[test] - fn flat_proxy_response_remains_supported() { - let response = decode_proxy_response(serde_json::json!({ - "successful": true, - "data": {"items": [1]} - })) - .unwrap(); - - assert!(response.successful); - assert_eq!(response.data["items"], serde_json::json!([1])); - } -} +#[path = "client_tests.rs"] +mod tests; diff --git a/src/memory/sync/composio/client_tests.rs b/src/memory/sync/composio/client_tests.rs new file mode 100644 index 00000000..df83f979 --- /dev/null +++ b/src/memory/sync/composio/client_tests.rs @@ -0,0 +1,29 @@ +use super::*; + +#[test] +fn proxied_backend_envelope_decodes_provider_response() { + let response = decode_proxy_response(serde_json::json!({ + "success": true, + "data": { + "successful": true, + "data": {"messages": [{"messageId": "message-1"}]}, + "error": null + } + })) + .unwrap(); + + assert!(response.successful); + assert_eq!(response.data["messages"][0]["messageId"], "message-1"); +} + +#[test] +fn flat_proxy_response_remains_supported() { + let response = decode_proxy_response(serde_json::json!({ + "successful": true, + "data": {"items": [1]} + })) + .unwrap(); + + assert!(response.successful); + assert_eq!(response.data["items"], serde_json::json!([1])); +} diff --git a/src/memory/sync/github.rs b/src/memory/sync/github.rs index 8b5c2d79..2ec52fb4 100644 --- a/src/memory/sync/github.rs +++ b/src/memory/sync/github.rs @@ -173,27 +173,5 @@ fn raw_coordinates(item_id: &str) -> Option<(RawKind, String)> { } #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn parses_supported_repo_urls() { - assert_eq!( - parse_repo("https://github.com/tinyhumansai/openhuman.git").unwrap(), - ("tinyhumansai".into(), "openhuman".into()) - ); - assert_eq!( - parse_repo("git@github.com:tinyhumansai/openhuman").unwrap(), - ("tinyhumansai".into(), "openhuman".into()) - ); - } - - #[test] - fn maps_raw_item_coordinates() { - assert_eq!( - raw_coordinates("issue:42"), - Some((RawKind::Issue, "42".into())) - ); - assert_eq!(raw_coordinates("unknown:42"), None); - } -} +#[path = "github_tests.rs"] +mod tests; diff --git a/src/memory/sync/github_tests.rs b/src/memory/sync/github_tests.rs new file mode 100644 index 00000000..0bb93680 --- /dev/null +++ b/src/memory/sync/github_tests.rs @@ -0,0 +1,22 @@ +use super::*; + +#[test] +fn parses_supported_repo_urls() { + assert_eq!( + parse_repo("https://github.com/tinyhumansai/openhuman.git").unwrap(), + ("tinyhumansai".into(), "openhuman".into()) + ); + assert_eq!( + parse_repo("git@github.com:tinyhumansai/openhuman").unwrap(), + ("tinyhumansai".into(), "openhuman".into()) + ); +} + +#[test] +fn maps_raw_item_coordinates() { + assert_eq!( + raw_coordinates("issue:42"), + Some((RawKind::Issue, "42".into())) + ); + assert_eq!(raw_coordinates("unknown:42"), None); +} diff --git a/src/memory/sync/periodic.rs b/src/memory/sync/periodic.rs index 7bd5210b..ed4ca957 100644 --- a/src/memory/sync/periodic.rs +++ b/src/memory/sync/periodic.rs @@ -69,79 +69,5 @@ fn elapsed_since(timestamp: DateTime, now: DateTime) -> Duration { } #[cfg(test)] -mod tests { - use super::*; - - fn source(id: &str, kind: SourceKind, enabled: bool) -> MemorySourceEntry { - MemorySourceEntry { - id: id.into(), - kind, - label: id.into(), - enabled, - toolkit: None, - connection_id: None, - path: Some("/tmp".into()), - glob: None, - url: Some("https://example.com".into()), - branch: None, - paths: Vec::new(), - max_commits: None, - max_issues: None, - max_prs: None, - query: None, - since_days: None, - max_items: None, - selector: None, - max_tokens_per_sync: None, - max_cost_per_sync_usd: None, - sync_depth_days: None, - } - } - - fn audit(id: &str, timestamp: DateTime, success: bool) -> SyncAuditEntry { - SyncAuditEntry { - timestamp, - source_id: id.into(), - source_kind: "folder".into(), - scope: id.into(), - items_fetched: 0, - batches: 0, - input_tokens: 0, - output_tokens: 0, - estimated_cost_usd: 0.0, - composio_actions_called: 0, - composio_cost_usd: 0.0, - actual_charged_usd: None, - duration_ms: 0, - success, - error: None, - } - } - - #[test] - fn cadence_handles_manual_minimum_and_persisted_success() { - assert_eq!(effective_interval_secs(Some(0)), None); - assert_eq!(effective_interval_secs(Some(60)), Some(60)); - let now = Utc::now(); - let sources = vec![ - source("new", SourceKind::Folder, true), - source("recent", SourceKind::Folder, true), - source("old", SourceKind::Folder, true), - source("disabled", SourceKind::Folder, false), - source("conversation", SourceKind::Conversation, true), - ]; - let history = vec![ - audit("recent", now - chrono::Duration::hours(1), true), - audit("old", now - chrono::Duration::hours(25), true), - audit("recent", now - chrono::Duration::hours(30), true), - ]; - let due = due_workspace_sources(&sources, &history, None, now); - assert_eq!( - due.iter() - .map(|source| source.id.as_str()) - .collect::>(), - vec!["new", "old"] - ); - assert!(due_workspace_sources(&sources, &history, Some(0), now).is_empty()); - } -} +#[path = "periodic_tests.rs"] +mod tests; diff --git a/src/memory/sync/periodic_tests.rs b/src/memory/sync/periodic_tests.rs new file mode 100644 index 00000000..b9818072 --- /dev/null +++ b/src/memory/sync/periodic_tests.rs @@ -0,0 +1,74 @@ +use super::*; + +fn source(id: &str, kind: SourceKind, enabled: bool) -> MemorySourceEntry { + MemorySourceEntry { + id: id.into(), + kind, + label: id.into(), + enabled, + toolkit: None, + connection_id: None, + path: Some("/tmp".into()), + glob: None, + url: Some("https://example.com".into()), + branch: None, + paths: Vec::new(), + max_commits: None, + max_issues: None, + max_prs: None, + query: None, + since_days: None, + max_items: None, + selector: None, + max_tokens_per_sync: None, + max_cost_per_sync_usd: None, + sync_depth_days: None, + } +} + +fn audit(id: &str, timestamp: DateTime, success: bool) -> SyncAuditEntry { + SyncAuditEntry { + timestamp, + source_id: id.into(), + source_kind: "folder".into(), + scope: id.into(), + items_fetched: 0, + batches: 0, + input_tokens: 0, + output_tokens: 0, + estimated_cost_usd: 0.0, + composio_actions_called: 0, + composio_cost_usd: 0.0, + actual_charged_usd: None, + duration_ms: 0, + success, + error: None, + } +} + +#[test] +fn cadence_handles_manual_minimum_and_persisted_success() { + assert_eq!(effective_interval_secs(Some(0)), None); + assert_eq!(effective_interval_secs(Some(60)), Some(60)); + let now = Utc::now(); + let sources = vec![ + source("new", SourceKind::Folder, true), + source("recent", SourceKind::Folder, true), + source("old", SourceKind::Folder, true), + source("disabled", SourceKind::Folder, false), + source("conversation", SourceKind::Conversation, true), + ]; + let history = vec![ + audit("recent", now - chrono::Duration::hours(1), true), + audit("old", now - chrono::Duration::hours(25), true), + audit("recent", now - chrono::Duration::hours(30), true), + ]; + let due = due_workspace_sources(&sources, &history, None, now); + assert_eq!( + due.iter() + .map(|source| source.id.as_str()) + .collect::>(), + vec!["new", "old"] + ); + assert!(due_workspace_sources(&sources, &history, Some(0), now).is_empty()); +} diff --git a/src/memory/sync/rebuild.rs b/src/memory/sync/rebuild.rs index 52245d3b..174f9e0c 100644 --- a/src/memory/sync/rebuild.rs +++ b/src/memory/sync/rebuild.rs @@ -409,83 +409,5 @@ fn backfill_coverage_from_summaries( } #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn coverage_is_incremental_and_ignores_source_metadata() { - let temp = tempfile::tempdir().unwrap(); - let mut config = MemoryConfig::new(temp.path()); - let custom_root = temp.path().join("custom-content"); - config.content_root = Some(custom_root.clone()); - let root = custom_root.join("raw/github-com-org-repo/issues"); - std::fs::create_dir_all(&root).unwrap(); - std::fs::write(root.join("100_one.md"), "one").unwrap(); - std::fs::write(root.join("200_two.md"), "two").unwrap(); - std::fs::write(root.parent().unwrap().join("_source.md"), "metadata").unwrap(); - - let first = raw_coverage(&config, "github:org/repo", "github.com/org/repo").unwrap(); - assert_eq!(first.total, 2); - assert_eq!(first.pending.len(), 2); - mark_raw_paths_ingested(&config, &[first.pending[0].rel.clone()]).unwrap(); - let second = raw_coverage(&config, "github:org/repo", "github.com/org/repo").unwrap(); - assert_eq!(second.covered, 1); - assert_eq!(second.pending.len(), 1); - assert!(needs_rebuild( - &config, - "github:org/repo", - "github.com/org/repo" - )); - } - - #[tokio::test] - async fn rebuild_ingests_l1_summaries_marks_coverage_and_is_idempotent() { - let temp = tempfile::tempdir().unwrap(); - let mut config = MemoryConfig::new(temp.path()); - config.tree.input_token_budget = 2; - let root = config - .workspace - .join("memory_tree/content/raw/github-com-org-repo/issues"); - std::fs::create_dir_all(&root).unwrap(); - std::fs::write(root.join("100_one.md"), "first body").unwrap(); - std::fs::write(root.join("200_two.md"), "second body").unwrap(); - - let first = rebuild_tree_from_raw( - &config, - "github:org/repo", - "github.com/org/repo", - &crate::memory::tree::ConcatSummariser, - ) - .await - .unwrap(); - assert_eq!(first.files_read, 2); - assert_eq!(first.batches, 2); - assert!(!needs_rebuild( - &config, - "github:org/repo", - "github.com/org/repo" - )); - assert_eq!(crate::memory::chunks::count_chunks(&config).unwrap(), 0); - let tree = TreeFactory::source("github:org/repo") - .get_or_create(&config) - .unwrap(); - assert_eq!( - list_summaries_at_level(&config, &tree.id, 1).unwrap().len(), - 2 - ); - - let second = rebuild_tree_from_raw( - &config, - "github:org/repo", - "github.com/org/repo", - &crate::memory::tree::ConcatSummariser, - ) - .await - .unwrap(); - assert_eq!(second, RebuildOutcome::default()); - assert_eq!( - list_summaries_at_level(&config, &tree.id, 1).unwrap().len(), - 2 - ); - } -} +#[path = "rebuild_tests.rs"] +mod tests; diff --git a/src/memory/sync/rebuild_tests.rs b/src/memory/sync/rebuild_tests.rs new file mode 100644 index 00000000..6bb66f9b --- /dev/null +++ b/src/memory/sync/rebuild_tests.rs @@ -0,0 +1,78 @@ +use super::*; + +#[test] +fn coverage_is_incremental_and_ignores_source_metadata() { + let temp = tempfile::tempdir().unwrap(); + let mut config = MemoryConfig::new(temp.path()); + let custom_root = temp.path().join("custom-content"); + config.content_root = Some(custom_root.clone()); + let root = custom_root.join("raw/github-com-org-repo/issues"); + std::fs::create_dir_all(&root).unwrap(); + std::fs::write(root.join("100_one.md"), "one").unwrap(); + std::fs::write(root.join("200_two.md"), "two").unwrap(); + std::fs::write(root.parent().unwrap().join("_source.md"), "metadata").unwrap(); + + let first = raw_coverage(&config, "github:org/repo", "github.com/org/repo").unwrap(); + assert_eq!(first.total, 2); + assert_eq!(first.pending.len(), 2); + mark_raw_paths_ingested(&config, &[first.pending[0].rel.clone()]).unwrap(); + let second = raw_coverage(&config, "github:org/repo", "github.com/org/repo").unwrap(); + assert_eq!(second.covered, 1); + assert_eq!(second.pending.len(), 1); + assert!(needs_rebuild( + &config, + "github:org/repo", + "github.com/org/repo" + )); +} + +#[tokio::test] +async fn rebuild_ingests_l1_summaries_marks_coverage_and_is_idempotent() { + let temp = tempfile::tempdir().unwrap(); + let mut config = MemoryConfig::new(temp.path()); + config.tree.input_token_budget = 2; + let root = config + .workspace + .join("memory_tree/content/raw/github-com-org-repo/issues"); + std::fs::create_dir_all(&root).unwrap(); + std::fs::write(root.join("100_one.md"), "first body").unwrap(); + std::fs::write(root.join("200_two.md"), "second body").unwrap(); + + let first = rebuild_tree_from_raw( + &config, + "github:org/repo", + "github.com/org/repo", + &crate::memory::tree::ConcatSummariser, + ) + .await + .unwrap(); + assert_eq!(first.files_read, 2); + assert_eq!(first.batches, 2); + assert!(!needs_rebuild( + &config, + "github:org/repo", + "github.com/org/repo" + )); + assert_eq!(crate::memory::chunks::count_chunks(&config).unwrap(), 0); + let tree = TreeFactory::source("github:org/repo") + .get_or_create(&config) + .unwrap(); + assert_eq!( + list_summaries_at_level(&config, &tree.id, 1).unwrap().len(), + 2 + ); + + let second = rebuild_tree_from_raw( + &config, + "github:org/repo", + "github.com/org/repo", + &crate::memory::tree::ConcatSummariser, + ) + .await + .unwrap(); + assert_eq!(second, RebuildOutcome::default()); + assert_eq!( + list_summaries_at_level(&config, &tree.id, 1).unwrap().len(), + 2 + ); +} diff --git a/src/memory/sync/state.rs b/src/memory/sync/state.rs index 9a11227f..aa708ce2 100644 --- a/src/memory/sync/state.rs +++ b/src/memory/sync/state.rs @@ -182,69 +182,5 @@ fn today() -> String { } #[cfg(test)] -mod tests { - use std::collections::HashMap; - use std::sync::Mutex; - - use super::*; - - #[derive(Default)] - struct MemoryStateStore(Mutex>); - - #[async_trait] - impl SyncStateStore for MemoryStateStore { - async fn get( - &self, - namespace: &str, - key: &str, - ) -> anyhow::Result> { - Ok(self - .0 - .lock() - .unwrap() - .get(&format!("{namespace}:{key}")) - .cloned()) - } - - async fn set( - &self, - namespace: &str, - key: &str, - value: &serde_json::Value, - ) -> anyhow::Result<()> { - self.0 - .lock() - .unwrap() - .insert(format!("{namespace}:{key}"), value.clone()); - Ok(()) - } - } - - #[tokio::test] - async fn state_round_trips_cursor_dedup_and_budget() { - let store = MemoryStateStore::default(); - let mut state = SyncState::new("gmail", "conn-1"); - state.advance_cursor("cursor-2"); - state.mark_synced("message-1"); - state.record_requests(3); - state.save(&store).await.unwrap(); - - let loaded = SyncState::load(&store, "gmail", "conn-1").await.unwrap(); - assert_eq!(loaded.cursor.as_deref(), Some("cursor-2")); - assert!(loaded.is_synced("message-1")); - assert_eq!(loaded.daily_budget.requests_used, 3); - } - - #[test] - fn stale_budget_reports_full_and_resets_on_record() { - let mut budget = DailyBudget { - date: "2000-01-01".into(), - requests_used: 499, - limit: 500, - }; - assert_eq!(budget.remaining(), 500); - budget.record_requests(1); - assert_eq!(budget.requests_used, 1); - assert_eq!(budget.remaining(), 499); - } -} +#[path = "state_tests.rs"] +mod tests; diff --git a/src/memory/sync/state_tests.rs b/src/memory/sync/state_tests.rs new file mode 100644 index 00000000..b3cabaf1 --- /dev/null +++ b/src/memory/sync/state_tests.rs @@ -0,0 +1,60 @@ +use std::collections::HashMap; +use std::sync::Mutex; + +use super::*; + +#[derive(Default)] +struct MemoryStateStore(Mutex>); + +#[async_trait] +impl SyncStateStore for MemoryStateStore { + async fn get(&self, namespace: &str, key: &str) -> anyhow::Result> { + Ok(self + .0 + .lock() + .unwrap() + .get(&format!("{namespace}:{key}")) + .cloned()) + } + + async fn set( + &self, + namespace: &str, + key: &str, + value: &serde_json::Value, + ) -> anyhow::Result<()> { + self.0 + .lock() + .unwrap() + .insert(format!("{namespace}:{key}"), value.clone()); + Ok(()) + } +} + +#[tokio::test] +async fn state_round_trips_cursor_dedup_and_budget() { + let store = MemoryStateStore::default(); + let mut state = SyncState::new("gmail", "conn-1"); + state.advance_cursor("cursor-2"); + state.mark_synced("message-1"); + state.record_requests(3); + state.save(&store).await.unwrap(); + + let loaded = SyncState::load(&store, "gmail", "conn-1").await.unwrap(); + assert_eq!(loaded.cursor.as_deref(), Some("cursor-2")); + assert!(loaded.is_synced("message-1")); + assert_eq!(loaded.daily_budget.requests_used, 3); +} + +#[test] +fn stale_budget_reports_full_and_resets_on_record() { + let mut budget = DailyBudget { + date: "2000-01-01".into(), + requests_used: 499, + limit: 500, + }; + assert_eq!(budget.remaining(), 500); + budget.record_requests(1); + assert_eq!(budget.requests_used, 1); + assert_eq!(budget.remaining(), 499); +} diff --git a/src/memory/sync/status.rs b/src/memory/sync/status.rs index 11a9e998..ba8d33c1 100644 --- a/src/memory/sync/status.rs +++ b/src/memory/sync/status.rs @@ -103,73 +103,5 @@ fn nonnegative(value: i64) -> u64 { } #[cfg(test)] -mod tests { - use chrono::{TimeZone, Utc}; - - use super::*; - use crate::memory::chunks::{upsert_chunks, Chunk, Metadata, SourceKind}; - - fn chunk(id: &str, source_id: &str, created_ms: i64) -> Chunk { - let timestamp = Utc.timestamp_millis_opt(created_ms).unwrap(); - Chunk { - id: id.into(), - content: "content".into(), - token_count: 1, - seq_in_source: 0, - created_at: timestamp, - partial_message: false, - metadata: Metadata { - source_kind: SourceKind::Document, - source_id: source_id.into(), - path_scope: None, - source_ref: None, - owner: "test".into(), - timestamp, - time_range: (timestamp, timestamp), - tags: Vec::new(), - }, - } - } - - #[test] - fn status_groups_provider_and_tracks_active_wave_resolution() { - let temp = tempfile::tempdir().unwrap(); - let config = MemoryConfig::new(temp.path()); - let now = 1_777_000_000_000i64; - upsert_chunks( - &config, - &[ - chunk("a", "gmail:conn", now - 2_000), - chunk("b", "gmail:conn", now - 1_000), - chunk("c", "slack:conn", now - 600_000), - ], - ) - .unwrap(); - with_connection(&config, |connection| { - connection.execute( - "INSERT INTO mem_tree_chunk_embeddings (chunk_id, model_signature, vector, dim, created_at) VALUES ('a', 'test', X'00', 1, 0)", - [], - )?; - connection.execute("UPDATE mem_tree_chunks SET lifecycle_status = 'dropped' WHERE id = 'c'", [])?; - Ok(()) - }).unwrap(); - - let statuses = list_sync_statuses_at(&config, now).unwrap(); - let gmail = statuses - .iter() - .find(|status| status.provider == "gmail") - .unwrap(); - assert_eq!(gmail.chunks_synced, 2); - assert_eq!(gmail.chunks_pending, 1); - assert_eq!(gmail.batch_total, 2); - assert_eq!(gmail.batch_processed, 1); - assert_eq!(gmail.freshness, FreshnessLabel::Active); - let slack = statuses - .iter() - .find(|status| status.provider == "slack") - .unwrap(); - assert_eq!(slack.chunks_pending, 0); - assert_eq!(slack.batch_total, 0); - assert_eq!(slack.freshness, FreshnessLabel::Idle); - } -} +#[path = "status_tests.rs"] +mod tests; diff --git a/src/memory/sync/status_tests.rs b/src/memory/sync/status_tests.rs new file mode 100644 index 00000000..42b2ed4a --- /dev/null +++ b/src/memory/sync/status_tests.rs @@ -0,0 +1,68 @@ +use chrono::{TimeZone, Utc}; + +use super::*; +use crate::memory::chunks::{upsert_chunks, Chunk, Metadata, SourceKind}; + +fn chunk(id: &str, source_id: &str, created_ms: i64) -> Chunk { + let timestamp = Utc.timestamp_millis_opt(created_ms).unwrap(); + Chunk { + id: id.into(), + content: "content".into(), + token_count: 1, + seq_in_source: 0, + created_at: timestamp, + partial_message: false, + metadata: Metadata { + source_kind: SourceKind::Document, + source_id: source_id.into(), + path_scope: None, + source_ref: None, + owner: "test".into(), + timestamp, + time_range: (timestamp, timestamp), + tags: Vec::new(), + }, + } +} + +#[test] +fn status_groups_provider_and_tracks_active_wave_resolution() { + let temp = tempfile::tempdir().unwrap(); + let config = MemoryConfig::new(temp.path()); + let now = 1_777_000_000_000i64; + upsert_chunks( + &config, + &[ + chunk("a", "gmail:conn", now - 2_000), + chunk("b", "gmail:conn", now - 1_000), + chunk("c", "slack:conn", now - 600_000), + ], + ) + .unwrap(); + with_connection(&config, |connection| { + connection.execute( + "INSERT INTO mem_tree_chunk_embeddings (chunk_id, model_signature, vector, dim, created_at) VALUES ('a', 'test', X'00', 1, 0)", + [], + )?; + connection.execute("UPDATE mem_tree_chunks SET lifecycle_status = 'dropped' WHERE id = 'c'", [])?; + Ok(()) + }).unwrap(); + + let statuses = list_sync_statuses_at(&config, now).unwrap(); + let gmail = statuses + .iter() + .find(|status| status.provider == "gmail") + .unwrap(); + assert_eq!(gmail.chunks_synced, 2); + assert_eq!(gmail.chunks_pending, 1); + assert_eq!(gmail.batch_total, 2); + assert_eq!(gmail.batch_processed, 1); + assert_eq!(gmail.freshness, FreshnessLabel::Active); + let slack = statuses + .iter() + .find(|status| status.provider == "slack") + .unwrap(); + assert_eq!(slack.chunks_pending, 0); + assert_eq!(slack.batch_total, 0); + assert_eq!(slack.freshness, FreshnessLabel::Idle); +} diff --git a/src/memory/sync/workspace.rs b/src/memory/sync/workspace.rs index 5010435c..b7d1580b 100644 --- a/src/memory/sync/workspace.rs +++ b/src/memory/sync/workspace.rs @@ -176,304 +176,5 @@ impl SyncPipeline for WorkspaceSourcePipeline { } #[cfg(test)] -mod tests { - use std::sync::{Arc, Mutex}; - - use super::*; - use crate::memory::sources::SourceKind; - use crate::memory::sync::state::SyncStateStore; - use crate::memory::sync::traits::{ - ExternalSourceReader, LocalDocumentSink, SkillDocSink, SkillDocument, SyncEventSink, - }; - - #[derive(Default)] - struct Host { - documents: Mutex>, - state: Mutex>, - events: Mutex>, - deletes: Mutex>, - external_items: Mutex>, - external_bodies: Mutex>, - } - - #[async_trait] - impl SkillDocSink for Host { - async fn store(&self, _: SkillDocument) -> anyhow::Result<()> { - Ok(()) - } - async fn delete(&self, _: &str, _: &str) -> anyhow::Result<()> { - Ok(()) - } - } - - #[async_trait] - impl LocalDocumentSink for Host { - async fn upsert(&self, document: LocalDocument) -> anyhow::Result<()> { - self.documents - .lock() - .unwrap() - .insert(document.source_id.clone(), document); - Ok(()) - } - - async fn delete(&self, source_id: &str) -> anyhow::Result<()> { - self.documents.lock().unwrap().remove(source_id); - self.deletes.lock().unwrap().push(source_id.into()); - Ok(()) - } - } - - #[async_trait] - impl SyncEventSink for Host { - async fn emit(&self, event: SyncEvent) -> anyhow::Result<()> { - self.events.lock().unwrap().push(event); - Ok(()) - } - } - - #[async_trait] - impl SyncStateStore for Host { - async fn get( - &self, - namespace: &str, - key: &str, - ) -> anyhow::Result> { - Ok(self - .state - .lock() - .unwrap() - .get(&format!("{namespace}:{key}")) - .cloned()) - } - async fn set( - &self, - namespace: &str, - key: &str, - value: &serde_json::Value, - ) -> anyhow::Result<()> { - self.state - .lock() - .unwrap() - .insert(format!("{namespace}:{key}"), value.clone()); - Ok(()) - } - } - - #[async_trait] - impl ExternalSourceReader for Host { - async fn list_items( - &self, - _: &MemorySourceEntry, - ) -> anyhow::Result> { - Ok(self.external_items.lock().unwrap().clone()) - } - - async fn read_item( - &self, - _: &MemorySourceEntry, - item_id: &str, - ) -> anyhow::Result { - self.external_bodies - .lock() - .unwrap() - .get(item_id) - .cloned() - .ok_or_else(|| anyhow::anyhow!("missing external item {item_id}")) - } - } - - fn folder_source(path: &std::path::Path) -> MemorySourceEntry { - MemorySourceEntry { - id: "folder-1".into(), - kind: SourceKind::Folder, - label: "Notes".into(), - enabled: true, - toolkit: None, - connection_id: None, - path: Some(path.to_string_lossy().into_owned()), - glob: Some("**/*.md".into()), - url: None, - branch: None, - paths: Vec::new(), - max_commits: None, - max_issues: None, - max_prs: None, - query: None, - since_days: None, - max_items: None, - selector: None, - max_tokens_per_sync: None, - max_cost_per_sync_usd: None, - sync_depth_days: None, - } - } - - #[tokio::test] - async fn folder_pipeline_tracks_create_update_noop_and_remove() { - let temp = tempfile::tempdir().unwrap(); - let notes = temp.path().join("notes"); - std::fs::create_dir_all(¬es).unwrap(); - let file = notes.join("daily.md"); - std::fs::write(&file, "first").unwrap(); - let config = MemoryConfig::new(temp.path().join("workspace")); - let pipeline = WorkspaceSourcePipeline::new(folder_source(¬es)).unwrap(); - let host = Arc::new(Host::default()); - let context = SyncContext { - events: host.clone(), - documents: host.clone(), - state: host.clone(), - local_documents: Some(host.clone()), - external_sources: None, - summariser: None, - }; - - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 1 - ); - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 0 - ); - std::thread::sleep(std::time::Duration::from_millis(5)); - std::fs::write(&file, "second").unwrap(); - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 1 - ); - assert_eq!( - host.documents.lock().unwrap()["mem_src:folder-1:daily.md"].body, - "second" - ); - std::fs::remove_file(&file).unwrap(); - let removed = pipeline.tick(&config, &context).await.unwrap(); - assert_eq!(removed.records_ingested, 0); - assert_eq!(removed.note.as_deref(), Some("1 removed")); - assert!(host.documents.lock().unwrap().is_empty()); - assert_eq!( - host.deletes.lock().unwrap().as_slice(), - ["mem_src:folder-1:daily.md"] - ); - } - - #[tokio::test] - async fn external_pipeline_tracks_versions_without_destructive_absence() { - use crate::memory::sources::{ContentType, SourceContent, SourceItem}; - - let config = MemoryConfig::new(tempfile::tempdir().unwrap().path().join("workspace")); - let mut source = folder_source(std::path::Path::new("unused")); - source.id = "rss-1".into(); - source.kind = SourceKind::RssFeed; - source.path = None; - source.glob = None; - source.url = Some("https://example.test/feed.xml".into()); - let pipeline = WorkspaceSourcePipeline::new(source).unwrap(); - let host = Arc::new(Host::default()); - host.external_items.lock().unwrap().push(SourceItem { - id: "post-1".into(), - title: "Post".into(), - updated_at_ms: Some(1), - }); - host.external_bodies.lock().unwrap().insert( - "post-1".into(), - SourceContent { - id: "post-1".into(), - title: "Post".into(), - body: "first".into(), - content_type: ContentType::Plaintext, - metadata: serde_json::Value::Null, - }, - ); - let context = SyncContext { - events: host.clone(), - documents: host.clone(), - state: host.clone(), - local_documents: Some(host.clone()), - external_sources: Some(host.clone()), - summariser: None, - }; - - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 1 - ); - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 0 - ); - host.external_items.lock().unwrap()[0].updated_at_ms = Some(2); - host.external_bodies.lock().unwrap().remove("post-1"); - assert!(pipeline.tick(&config, &context).await.is_err()); - assert_eq!( - host.documents.lock().unwrap()["mem_src:rss-1:post-1"].body, - "first" - ); - host.external_bodies.lock().unwrap().insert( - "post-1".into(), - SourceContent { - id: "post-1".into(), - title: "Post".into(), - body: "second".into(), - content_type: ContentType::Plaintext, - metadata: serde_json::Value::Null, - }, - ); - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 1 - ); - host.external_items.lock().unwrap()[0].updated_at_ms = None; - host.external_bodies - .lock() - .unwrap() - .get_mut("post-1") - .unwrap() - .body = "timestamp-less refresh".into(); - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 1 - ); - host.external_items.lock().unwrap().clear(); - assert_eq!( - pipeline - .tick(&config, &context) - .await - .unwrap() - .records_ingested, - 0 - ); - assert!(host - .documents - .lock() - .unwrap() - .contains_key("mem_src:rss-1:post-1")); - } -} +#[path = "workspace_tests.rs"] +mod tests; diff --git a/src/memory/sync/workspace_tests.rs b/src/memory/sync/workspace_tests.rs new file mode 100644 index 00000000..27a39ffc --- /dev/null +++ b/src/memory/sync/workspace_tests.rs @@ -0,0 +1,295 @@ +use std::sync::{Arc, Mutex}; + +use super::*; +use crate::memory::sources::SourceKind; +use crate::memory::sync::state::SyncStateStore; +use crate::memory::sync::traits::{ + ExternalSourceReader, LocalDocumentSink, SkillDocSink, SkillDocument, SyncEventSink, +}; + +#[derive(Default)] +struct Host { + documents: Mutex>, + state: Mutex>, + events: Mutex>, + deletes: Mutex>, + external_items: Mutex>, + external_bodies: Mutex>, +} + +#[async_trait] +impl SkillDocSink for Host { + async fn store(&self, _: SkillDocument) -> anyhow::Result<()> { + Ok(()) + } + async fn delete(&self, _: &str, _: &str) -> anyhow::Result<()> { + Ok(()) + } +} + +#[async_trait] +impl LocalDocumentSink for Host { + async fn upsert(&self, document: LocalDocument) -> anyhow::Result<()> { + self.documents + .lock() + .unwrap() + .insert(document.source_id.clone(), document); + Ok(()) + } + + async fn delete(&self, source_id: &str) -> anyhow::Result<()> { + self.documents.lock().unwrap().remove(source_id); + self.deletes.lock().unwrap().push(source_id.into()); + Ok(()) + } +} + +#[async_trait] +impl SyncEventSink for Host { + async fn emit(&self, event: SyncEvent) -> anyhow::Result<()> { + self.events.lock().unwrap().push(event); + Ok(()) + } +} + +#[async_trait] +impl SyncStateStore for Host { + async fn get(&self, namespace: &str, key: &str) -> anyhow::Result> { + Ok(self + .state + .lock() + .unwrap() + .get(&format!("{namespace}:{key}")) + .cloned()) + } + async fn set( + &self, + namespace: &str, + key: &str, + value: &serde_json::Value, + ) -> anyhow::Result<()> { + self.state + .lock() + .unwrap() + .insert(format!("{namespace}:{key}"), value.clone()); + Ok(()) + } +} + +#[async_trait] +impl ExternalSourceReader for Host { + async fn list_items( + &self, + _: &MemorySourceEntry, + ) -> anyhow::Result> { + Ok(self.external_items.lock().unwrap().clone()) + } + + async fn read_item( + &self, + _: &MemorySourceEntry, + item_id: &str, + ) -> anyhow::Result { + self.external_bodies + .lock() + .unwrap() + .get(item_id) + .cloned() + .ok_or_else(|| anyhow::anyhow!("missing external item {item_id}")) + } +} + +fn folder_source(path: &std::path::Path) -> MemorySourceEntry { + MemorySourceEntry { + id: "folder-1".into(), + kind: SourceKind::Folder, + label: "Notes".into(), + enabled: true, + toolkit: None, + connection_id: None, + path: Some(path.to_string_lossy().into_owned()), + glob: Some("**/*.md".into()), + url: None, + branch: None, + paths: Vec::new(), + max_commits: None, + max_issues: None, + max_prs: None, + query: None, + since_days: None, + max_items: None, + selector: None, + max_tokens_per_sync: None, + max_cost_per_sync_usd: None, + sync_depth_days: None, + } +} + +#[tokio::test] +async fn folder_pipeline_tracks_create_update_noop_and_remove() { + let temp = tempfile::tempdir().unwrap(); + let notes = temp.path().join("notes"); + std::fs::create_dir_all(¬es).unwrap(); + let file = notes.join("daily.md"); + std::fs::write(&file, "first").unwrap(); + let config = MemoryConfig::new(temp.path().join("workspace")); + let pipeline = WorkspaceSourcePipeline::new(folder_source(¬es)).unwrap(); + let host = Arc::new(Host::default()); + let context = SyncContext { + events: host.clone(), + documents: host.clone(), + state: host.clone(), + local_documents: Some(host.clone()), + external_sources: None, + summariser: None, + }; + + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 1 + ); + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 0 + ); + std::thread::sleep(std::time::Duration::from_millis(5)); + std::fs::write(&file, "second").unwrap(); + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 1 + ); + assert_eq!( + host.documents.lock().unwrap()["mem_src:folder-1:daily.md"].body, + "second" + ); + std::fs::remove_file(&file).unwrap(); + let removed = pipeline.tick(&config, &context).await.unwrap(); + assert_eq!(removed.records_ingested, 0); + assert_eq!(removed.note.as_deref(), Some("1 removed")); + assert!(host.documents.lock().unwrap().is_empty()); + assert_eq!( + host.deletes.lock().unwrap().as_slice(), + ["mem_src:folder-1:daily.md"] + ); +} + +#[tokio::test] +async fn external_pipeline_tracks_versions_without_destructive_absence() { + use crate::memory::sources::{ContentType, SourceContent, SourceItem}; + + let config = MemoryConfig::new(tempfile::tempdir().unwrap().path().join("workspace")); + let mut source = folder_source(std::path::Path::new("unused")); + source.id = "rss-1".into(); + source.kind = SourceKind::RssFeed; + source.path = None; + source.glob = None; + source.url = Some("https://example.test/feed.xml".into()); + let pipeline = WorkspaceSourcePipeline::new(source).unwrap(); + let host = Arc::new(Host::default()); + host.external_items.lock().unwrap().push(SourceItem { + id: "post-1".into(), + title: "Post".into(), + updated_at_ms: Some(1), + }); + host.external_bodies.lock().unwrap().insert( + "post-1".into(), + SourceContent { + id: "post-1".into(), + title: "Post".into(), + body: "first".into(), + content_type: ContentType::Plaintext, + metadata: serde_json::Value::Null, + }, + ); + let context = SyncContext { + events: host.clone(), + documents: host.clone(), + state: host.clone(), + local_documents: Some(host.clone()), + external_sources: Some(host.clone()), + summariser: None, + }; + + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 1 + ); + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 0 + ); + host.external_items.lock().unwrap()[0].updated_at_ms = Some(2); + host.external_bodies.lock().unwrap().remove("post-1"); + assert!(pipeline.tick(&config, &context).await.is_err()); + assert_eq!( + host.documents.lock().unwrap()["mem_src:rss-1:post-1"].body, + "first" + ); + host.external_bodies.lock().unwrap().insert( + "post-1".into(), + SourceContent { + id: "post-1".into(), + title: "Post".into(), + body: "second".into(), + content_type: ContentType::Plaintext, + metadata: serde_json::Value::Null, + }, + ); + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 1 + ); + host.external_items.lock().unwrap()[0].updated_at_ms = None; + host.external_bodies + .lock() + .unwrap() + .get_mut("post-1") + .unwrap() + .body = "timestamp-less refresh".into(); + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 1 + ); + host.external_items.lock().unwrap().clear(); + assert_eq!( + pipeline + .tick(&config, &context) + .await + .unwrap() + .records_ingested, + 0 + ); + assert!(host + .documents + .lock() + .unwrap() + .contains_key("mem_src:rss-1:post-1")); +} diff --git a/src/memory/tree/direct_ingest.rs b/src/memory/tree/direct_ingest.rs index e8bca97d..06321906 100644 --- a/src/memory/tree/direct_ingest.rs +++ b/src/memory/tree/direct_ingest.rs @@ -145,42 +145,5 @@ pub async fn ingest_summary( } #[cfg(test)] -mod tests { - use super::*; - use crate::memory::tree::{store::get_buffer, ConcatSummariser}; - - #[tokio::test] - async fn direct_summary_lands_at_l1_without_creating_chunks() { - let temp = tempfile::tempdir().unwrap(); - let config = MemoryConfig::new(temp.path()); - let tree = TreeFactory::source("github:org/repo") - .get_or_create(&config) - .unwrap(); - let now = Utc::now(); - let outcome = ingest_summary( - &config, - &tree, - SummaryIngestInput { - content: "summary".into(), - token_count: 2, - entities: Vec::new(), - topics: Vec::new(), - time_range_start: now, - time_range_end: now, - score: 0.5, - child_labels: vec!["100_issue-1".into()], - child_basenames: Vec::new(), - }, - &ConcatSummariser, - ) - .await - .unwrap(); - assert_eq!(crate::memory::chunks::count_chunks(&config).unwrap(), 0); - let summary = store::get_summary(&config, &outcome.summary_id) - .unwrap() - .unwrap(); - assert_eq!(summary.level, 1); - assert_eq!(summary.child_ids, vec!["100_issue-1"]); - assert_eq!(get_buffer(&config, &tree.id, 1).unwrap().item_ids.len(), 1); - } -} +#[path = "direct_ingest_tests.rs"] +mod tests; diff --git a/src/memory/tree/direct_ingest_tests.rs b/src/memory/tree/direct_ingest_tests.rs new file mode 100644 index 00000000..ef013ea8 --- /dev/null +++ b/src/memory/tree/direct_ingest_tests.rs @@ -0,0 +1,37 @@ +use super::*; +use crate::memory::tree::{store::get_buffer, ConcatSummariser}; + +#[tokio::test] +async fn direct_summary_lands_at_l1_without_creating_chunks() { + let temp = tempfile::tempdir().unwrap(); + let config = MemoryConfig::new(temp.path()); + let tree = TreeFactory::source("github:org/repo") + .get_or_create(&config) + .unwrap(); + let now = Utc::now(); + let outcome = ingest_summary( + &config, + &tree, + SummaryIngestInput { + content: "summary".into(), + token_count: 2, + entities: Vec::new(), + topics: Vec::new(), + time_range_start: now, + time_range_end: now, + score: 0.5, + child_labels: vec!["100_issue-1".into()], + child_basenames: Vec::new(), + }, + &ConcatSummariser, + ) + .await + .unwrap(); + assert_eq!(crate::memory::chunks::count_chunks(&config).unwrap(), 0); + let summary = store::get_summary(&config, &outcome.summary_id) + .unwrap() + .unwrap(); + assert_eq!(summary.level, 1); + assert_eq!(summary.child_ids, vec!["100_issue-1"]); + assert_eq!(get_buffer(&config, &tree.id, 1).unwrap().item_ids.len(), 1); +} From c2faf478831a2a58bf0b91596e8e921b63b25c9f Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Thu, 1 Oct 2026 18:34:48 +0300 Subject: [PATCH 2/3] fix(kv): correct safety module doc link Changed the intra-doc link for the safety guard from an absolute path to a relative one using `super::safety`, fixing a broken documentation reference that would have failed to resolve in the rendered docs. Auto-committed-on: dragonfly --- src/memory/store/kv.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/memory/store/kv.rs b/src/memory/store/kv.rs index 20323416..d45469b9 100644 --- a/src/memory/store/kv.rs +++ b/src/memory/store/kv.rs @@ -5,7 +5,7 @@ //! from OpenHuman's `memory_store::kv` (lifted off the `UnifiedMemory` //! connection into a standalone [`KvStore`] with its own connection). //! -//! Writes run through the [`safety`](crate::memory::store::safety) guard: +//! Writes run through the [`safety`](super::safety) guard: //! secret-like or PII-like keys/namespaces are rejected outright, and values //! are sanitized before they land in the store. From d27a8c14f768d4e29331991d8922948329bb0459 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Thu, 1 Oct 2026 18:34:53 +0300 Subject: [PATCH 3/3] fix(kv): remove unnecessary module path in doc link The doc comment for the KvStore module used a fully qualified path to reference the `safety` guard, but the link is within the same parent module so the shorter form works and avoids a stale reference if the module is ever reorganised. Auto-committed-on: dragonfly --- src/memory/store/kv.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/memory/store/kv.rs b/src/memory/store/kv.rs index d45469b9..d3aa2fb4 100644 --- a/src/memory/store/kv.rs +++ b/src/memory/store/kv.rs @@ -5,7 +5,7 @@ //! from OpenHuman's `memory_store::kv` (lifted off the `UnifiedMemory` //! connection into a standalone [`KvStore`] with its own connection). //! -//! Writes run through the [`safety`](super::safety) guard: +//! Writes run through the [`safety`] guard: //! secret-like or PII-like keys/namespaces are rejected outright, and values //! are sanitized before they land in the store.