mirror of
https://github.com/akitaonrails/ai-memory.git
synced 2026-10-02 03:24:46 +08:00
test: five adversarial security-boundary regressions across handoff, page-sharing, message, and lifecycle-guard invariants
Each test actively attempts a boundary violation and asserts refusal, with a legitimate control case, so a future regression that removes the guard fails CI: - handoff_admission: a non-admin proxied caller's `any_owner: true` on memory_handoff_accept is refused by require_admin_capability and does not claim the baton; root's any_owner accept is the control. - multi_session: a page with a non-null author_id is still readable by a different operator via both search_pages_for_project and page_body_by_ids, pinning that author_id is attribution and never a read filter (invariant #16). - agent_messages: the recipient cannot cancel the sender's message by exact id (the specific-id sibling of the existing whole-outbox test), and two concurrent pop_message calls on one message deliver it exactly once, pinning the state='pending' CAS in pop_message_in_transaction. - removal: reset/restore/reindex/uninstall --purge-data each refuse and leave the data dir untouched when a sibling ai-memory process is detected. Adds a minimal, test-only injection seam to process_guard::sibling_processes() (AI_MEMORY_TEST_FORCE_SIBLING_PIDS) so the refusal path is deterministically testable without spawning a real sibling process; production behavior is unchanged since the variable is never set outside a test harness. Each test was verified to fail when its guard is removed and pass when restored. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01MDbhmszrjG9s5MrPrTuNtm
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
fbef4b6f9a
commit
e297f9d3b6
@@ -17,6 +17,24 @@ pub const BIN_NAME: &str = "ai-memory";
|
||||
/// process and any threads of it).
|
||||
#[must_use]
|
||||
pub fn sibling_processes() -> Vec<sysinfo::Pid> {
|
||||
// Test injection: a comma-separated list of fake PIDs to report as alive
|
||||
// siblings, bypassing both the real scan AND the `cfg!(test)` opt-out
|
||||
// below. This is what lets the guard's REFUSAL path be exercised by an
|
||||
// in-process test (reset / reindex / restore / uninstall --purge-data all
|
||||
// call `sibling_processes()` directly, and `cfg!(test)` alone would
|
||||
// otherwise force every in-process test onto the "no siblings" branch).
|
||||
// Checked first, and not itself gated by `cfg!(test)`, matching the
|
||||
// existing `AI_MEMORY_TEST_NO_PROCESS_GUARD` opt-out below: neither is
|
||||
// reachable in a normal shipped run because neither is ever set outside
|
||||
// a test harness's own env.
|
||||
if let Ok(raw) = std::env::var("AI_MEMORY_TEST_FORCE_SIBLING_PIDS") {
|
||||
return raw
|
||||
.split(',')
|
||||
.filter(|s| !s.trim().is_empty())
|
||||
.filter_map(|s| s.trim().parse::<u32>().ok())
|
||||
.map(sysinfo::Pid::from_u32)
|
||||
.collect();
|
||||
}
|
||||
// Test opt-out. The destructive-command tests would otherwise flake
|
||||
// non-deterministically: a dev box (and a parallel test run) almost always
|
||||
// has some *other* `ai-memory` process alive, which the real scan rightly
|
||||
|
||||
@@ -943,9 +943,12 @@ fn uninstall_dry_run_previews_purge() {
|
||||
}
|
||||
}
|
||||
|
||||
/// Best-effort, NOT in the default run (sysinfo reads the real process table;
|
||||
/// no injection seam). Spawns a real sibling `ai-memory` process and asserts
|
||||
/// `--purge-data` refuses up front, leaving the wiring intact. Run with:
|
||||
/// Best-effort, NOT in the default run (sysinfo reads the real process
|
||||
/// table). Spawns a REAL sibling `ai-memory` process — as opposed to the
|
||||
/// `AI_MEMORY_TEST_FORCE_SIBLING_PIDS` injection seam the tests below use —
|
||||
/// so the actual sysinfo scan itself stays covered, not just the refusal
|
||||
/// logic downstream of it. Asserts `--purge-data` refuses up front, leaving
|
||||
/// the wiring intact. Run with:
|
||||
/// `cargo test -p ai-memory-cli --test removal -- --ignored`.
|
||||
#[test]
|
||||
#[ignore]
|
||||
@@ -988,6 +991,273 @@ fn purge_data_refuses_when_sibling_alive() {
|
||||
);
|
||||
}
|
||||
|
||||
// ----------------------------------------------------------------
|
||||
// Process-guard injection seam (`AI_MEMORY_TEST_FORCE_SIBLING_PIDS`).
|
||||
//
|
||||
// `process_guard::sibling_processes()` normally either does a real sysinfo
|
||||
// scan, or (only under `cfg!(test)` / `AI_MEMORY_TEST_NO_PROCESS_GUARD`)
|
||||
// short-circuits to "no siblings" — which made the guard's REFUSAL branch
|
||||
// untestable from an in-process test and left `purge_data_refuses_when_
|
||||
// sibling_alive` above `#[ignore]`d, needing a real spawned sibling. These
|
||||
// tests exercise the refusal deterministically via the injection seam,
|
||||
// against every direct-disk lifecycle command the guard protects: `reset`,
|
||||
// `restore`, `reindex`, and `uninstall --purge-data`.
|
||||
// ----------------------------------------------------------------
|
||||
|
||||
/// `reset --confirm` refuses while a sibling is reported alive, and leaves
|
||||
/// the data dir completely untouched.
|
||||
#[test]
|
||||
fn reset_refuses_when_sibling_pid_is_injected() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
std::fs::create_dir_all(data.path().join(sub)).unwrap();
|
||||
std::fs::write(data.path().join(sub).join("f.txt"), b"x").unwrap();
|
||||
}
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["reset", "--confirm"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "424242")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
!out.status.success(),
|
||||
"should refuse while a sibling is alive"
|
||||
);
|
||||
let stderr = String::from_utf8_lossy(&out.stderr);
|
||||
assert!(
|
||||
stderr.contains("refusing to reset") && stderr.contains("424242"),
|
||||
"stderr was: {stderr}"
|
||||
);
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
assert!(
|
||||
data.path().join(sub).join("f.txt").exists(),
|
||||
"{sub}/f.txt must survive a refused reset"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// The control: an explicitly EMPTY injected sibling list must not block a
|
||||
/// legitimate `reset --confirm`.
|
||||
#[test]
|
||||
fn reset_proceeds_when_injected_sibling_list_is_empty() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
std::fs::create_dir_all(data.path().join(sub)).unwrap();
|
||||
std::fs::write(data.path().join(sub).join("f.txt"), b"x").unwrap();
|
||||
}
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["reset", "--confirm"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"stderr: {}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
assert!(
|
||||
!data.path().join(sub).join("f.txt").exists(),
|
||||
"{sub} must be wiped once the guard reports no siblings"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// `restore --from` refuses while a sibling is reported alive, before ever
|
||||
/// checking whether the source tarball exists.
|
||||
#[test]
|
||||
fn restore_refuses_when_sibling_pid_is_injected() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["restore", "--from", "/does/not/exist.tar.gz"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "424242")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
!out.status.success(),
|
||||
"should refuse while a sibling is alive"
|
||||
);
|
||||
let stderr = String::from_utf8_lossy(&out.stderr);
|
||||
assert!(
|
||||
stderr.contains("refusing to restore") && stderr.contains("424242"),
|
||||
"stderr was: {stderr}"
|
||||
);
|
||||
assert!(
|
||||
!stderr.contains("not found"),
|
||||
"the guard must refuse BEFORE the missing-tarball check runs: {stderr}"
|
||||
);
|
||||
}
|
||||
|
||||
/// The control: an explicitly EMPTY injected sibling list lets `restore` past
|
||||
/// the guard — it then fails on the next check (missing tarball) instead of
|
||||
/// the busy message, proving the guard itself is not what stopped it.
|
||||
#[test]
|
||||
fn restore_proceeds_when_injected_sibling_list_is_empty() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["restore", "--from", "/does/not/exist.tar.gz"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(!out.status.success(), "still fails, but past the guard");
|
||||
let stderr = String::from_utf8_lossy(&out.stderr);
|
||||
assert!(
|
||||
stderr.contains("not found"),
|
||||
"expected the missing-tarball error once no sibling is reported: {stderr}"
|
||||
);
|
||||
assert!(
|
||||
!stderr.contains("refusing to restore"),
|
||||
"stderr was: {stderr}"
|
||||
);
|
||||
}
|
||||
|
||||
/// `reindex` refuses while a sibling is reported alive, before ever opening
|
||||
/// the SQLite store.
|
||||
#[test]
|
||||
fn reindex_refuses_when_sibling_pid_is_injected() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
std::fs::create_dir_all(data.path().join("wiki")).unwrap();
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["reindex"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "424242")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
!out.status.success(),
|
||||
"should refuse while a sibling is alive"
|
||||
);
|
||||
let stderr = String::from_utf8_lossy(&out.stderr);
|
||||
assert!(
|
||||
stderr.contains("refusing to reindex") && stderr.contains("424242"),
|
||||
"stderr was: {stderr}"
|
||||
);
|
||||
// Nothing should have been created; the guard runs before `Store::open`.
|
||||
assert!(
|
||||
!data.path().join("db").join("memory.sqlite").exists(),
|
||||
"the guard must refuse before the store is opened/created"
|
||||
);
|
||||
}
|
||||
|
||||
/// The control: an explicitly EMPTY injected sibling list lets `reindex`
|
||||
/// proceed and actually rebuild the (empty) index.
|
||||
#[test]
|
||||
fn reindex_proceeds_when_injected_sibling_list_is_empty() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
std::fs::create_dir_all(data.path().join("wiki")).unwrap();
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["reindex"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"stderr: {}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
assert!(
|
||||
data.path().join("db").join("memory.sqlite").exists(),
|
||||
"reindex must have opened/created the store once no sibling is reported"
|
||||
);
|
||||
}
|
||||
|
||||
/// `uninstall --purge-data` refuses while a sibling is reported alive, and
|
||||
/// leaves the data dir untouched — the sibling of the ignored real-process
|
||||
/// test above, but deterministic and in the default run.
|
||||
#[test]
|
||||
fn uninstall_purge_data_refuses_when_sibling_pid_is_injected() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
std::fs::create_dir_all(data.path().join(sub)).unwrap();
|
||||
std::fs::write(data.path().join(sub).join("f.txt"), b"x").unwrap();
|
||||
}
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["uninstall", "--apply", "--yes", "--purge-data"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "424242")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
!out.status.success(),
|
||||
"should refuse while a sibling is alive"
|
||||
);
|
||||
let stderr = String::from_utf8_lossy(&out.stderr);
|
||||
assert!(
|
||||
stderr.contains("refusing to purge data") && stderr.contains("424242"),
|
||||
"stderr was: {stderr}"
|
||||
);
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
assert!(
|
||||
data.path().join(sub).join("f.txt").exists(),
|
||||
"{sub}/f.txt must survive a refused purge"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// The control: an explicitly EMPTY injected sibling list must not block a
|
||||
/// legitimate `uninstall --apply --yes --purge-data`.
|
||||
#[test]
|
||||
fn uninstall_purge_data_proceeds_when_injected_sibling_list_is_empty() {
|
||||
let _guard = cli_test_lock();
|
||||
let home = tempfile::tempdir().unwrap();
|
||||
let data = tempfile::tempdir().unwrap();
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
std::fs::create_dir_all(data.path().join(sub)).unwrap();
|
||||
std::fs::write(data.path().join(sub).join("f.txt"), b"x").unwrap();
|
||||
}
|
||||
|
||||
let out = command_with_home(home.path())
|
||||
.args(["uninstall", "--apply", "--yes", "--purge-data"])
|
||||
.env("AI_MEMORY_DATA_DIR", data.path())
|
||||
.env("AI_MEMORY_TEST_FORCE_SIBLING_PIDS", "")
|
||||
.output()
|
||||
.unwrap();
|
||||
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"stderr: {}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
for sub in ["wiki", "db", "raw"] {
|
||||
assert!(
|
||||
!data.path().join(sub).join("f.txt").exists(),
|
||||
"{sub} must be purged once the guard reports no siblings"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uninstall_devin_hooks_preserves_user_entries() {
|
||||
let _guard = cli_test_lock();
|
||||
|
||||
@@ -431,3 +431,126 @@ async fn unauthorised_any_owner_cancel_reaches_no_webhook() {
|
||||
"the refused cancel must not have discarded the baton",
|
||||
);
|
||||
}
|
||||
|
||||
/// The sibling of `unauthorised_any_owner_cancel_reaches_no_webhook` for the
|
||||
/// accept path: a non-admin proxied caller ("alice") asking to accept a
|
||||
/// baton owned by "bob" via `any_owner: true` must be refused by
|
||||
/// `require_admin_capability`, must not reach the admission chain, and must
|
||||
/// NOT have claimed the handoff — accept is a compare-and-set consume, so a
|
||||
/// silently-successful bypass here would hand a stranger's session state to
|
||||
/// anyone who could reach the tool. An admin (root) doing the exact same
|
||||
/// call is the control: `any_owner` recovery must still work for the
|
||||
/// operator it exists for.
|
||||
#[tokio::test]
|
||||
async fn unauthorised_any_owner_accept_reaches_no_webhook_and_does_not_claim() {
|
||||
let (addr, calls) = recording_webhook_host().await;
|
||||
let h = proxied_harness(guarded_chain(addr)).await;
|
||||
|
||||
// Inserted directly, as the cancel test does: going through the tool
|
||||
// would put a legitimate `handoff_begin` in the recording, and what is
|
||||
// under test is that the refused accept adds nothing to it.
|
||||
let id = h
|
||||
.store
|
||||
.writer
|
||||
.insert_handoff(ai_memory_core::NewHandoff {
|
||||
workspace_id: h.ws,
|
||||
project_id: h.proj,
|
||||
from_session_id: None,
|
||||
from_agent: ai_memory_core::AgentKind::ClaudeCode,
|
||||
to_agent: None,
|
||||
cwd: None,
|
||||
summary: "bobs-baton".to_string(),
|
||||
open_questions: Vec::new(),
|
||||
next_steps: Vec::new(),
|
||||
files_touched: Vec::new(),
|
||||
owner_user: ai_memory_core::owner_stamp(
|
||||
Some(&ai_memory_core::IdentityKey::User("bob".into())),
|
||||
true,
|
||||
),
|
||||
})
|
||||
.await
|
||||
.expect("insert handoff");
|
||||
|
||||
let response = call_tool_raw(
|
||||
&h.router,
|
||||
"memory_handoff_accept",
|
||||
json!({
|
||||
"workspace": "default",
|
||||
"project": "scratch",
|
||||
"handoff_id": id.to_string(),
|
||||
"any_owner": true,
|
||||
}),
|
||||
&proxied_as("alice"),
|
||||
)
|
||||
.await;
|
||||
// Asserted on the message, not merely the presence of an error: a
|
||||
// rejection from the auth layer itself would also produce one, and would
|
||||
// make this test pass without the capability gate ever running.
|
||||
let message = response
|
||||
.pointer("/error/message")
|
||||
.and_then(|m| m.as_str())
|
||||
.unwrap_or_else(|| panic!("a non-admin caller must be refused any_owner: {response}"));
|
||||
assert!(
|
||||
message.contains("admin"),
|
||||
"expected the admin capability gate to refuse this, got: {message}",
|
||||
);
|
||||
|
||||
settle().await;
|
||||
assert!(
|
||||
calls.lock().unwrap().is_empty(),
|
||||
"an operation the caller was never permitted must reach no webhook: {:?}",
|
||||
calls.lock().unwrap(),
|
||||
);
|
||||
assert_eq!(
|
||||
h.store
|
||||
.reader
|
||||
.handoff_by_id(id)
|
||||
.await
|
||||
.expect("read back")
|
||||
.expect("row still there")
|
||||
.lifecycle
|
||||
.state,
|
||||
ai_memory_core::HandoffState::Open,
|
||||
"the refused accept must not have claimed bob's baton",
|
||||
);
|
||||
|
||||
// Control: an admin (root) CAN accept via `any_owner`, because the
|
||||
// recovery path this gate exists for must still work for the operator it
|
||||
// is meant for.
|
||||
let accepted = call_tool_raw(
|
||||
&h.router,
|
||||
"memory_handoff_accept",
|
||||
json!({
|
||||
"workspace": "default",
|
||||
"project": "scratch",
|
||||
"handoff_id": id.to_string(),
|
||||
"any_owner": true,
|
||||
}),
|
||||
&[("authorization", "Bearer the-root-token")],
|
||||
)
|
||||
.await;
|
||||
assert!(
|
||||
accepted.get("error").is_none(),
|
||||
"root's any_owner accept must succeed: {accepted}",
|
||||
);
|
||||
let claimed_summary = accepted
|
||||
.pointer("/result/content/0/text")
|
||||
.and_then(|t| t.as_str())
|
||||
.unwrap_or_else(|| panic!("missing tool text: {accepted}"));
|
||||
assert!(
|
||||
claimed_summary.contains("bobs-baton"),
|
||||
"root's recovery accept must actually claim bob's baton: {claimed_summary}",
|
||||
);
|
||||
assert_eq!(
|
||||
h.store
|
||||
.reader
|
||||
.handoff_by_id(id)
|
||||
.await
|
||||
.expect("read back")
|
||||
.expect("row still there")
|
||||
.lifecycle
|
||||
.state,
|
||||
ai_memory_core::HandoffState::Accepted,
|
||||
"the admin's accept must have claimed the baton",
|
||||
);
|
||||
}
|
||||
|
||||
@@ -219,6 +219,77 @@ async fn recipient_cannot_cancel_the_senders_outbox() {
|
||||
);
|
||||
}
|
||||
|
||||
/// The sibling of `recipient_cannot_cancel_the_senders_outbox` for the
|
||||
/// specific-id path: the existing test only exercises whole-outbox cancel
|
||||
/// (`specific_id: None`), which is scoped by `from_workspace_id`/
|
||||
/// `from_project_id` alone and would still refuse B even if the id-scoped
|
||||
/// branch dropped that guard. This targets B's exact knowledge of A's
|
||||
/// message id directly, so it can only pass if the id-scoped `UPDATE` also
|
||||
/// carries the sender-coordinate `AND from_workspace_id = ... AND
|
||||
/// from_project_id = ...` guard.
|
||||
#[tokio::test]
|
||||
async fn recipient_cannot_cancel_a_specific_message_by_id() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let a = project(&store, "default", "project-a").await;
|
||||
let b = project(&store, "default", "project-b").await;
|
||||
|
||||
let a_msg_id = store
|
||||
.writer
|
||||
.insert_message(message(a, b, "A's own message"))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// B knows the exact id (e.g. from its own inbox listing) but must not be
|
||||
// able to cancel A's outbound message by targeting it directly.
|
||||
let cancelled = store
|
||||
.writer
|
||||
.cancel_messages(b.0, b.1, Some(a_msg_id))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
cancelled, 0,
|
||||
"B cancelling A's message by exact id must not touch it"
|
||||
);
|
||||
|
||||
let popped = store
|
||||
.writer
|
||||
.pop_message(claim(b), Some(a_msg_id))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
popped
|
||||
.expect("A's message must still be pending and poppable by B")
|
||||
.body,
|
||||
"A's own message"
|
||||
);
|
||||
|
||||
// Control: A cancelling its OWN message by the same exact id succeeds.
|
||||
let a_msg_id_2 = store
|
||||
.writer
|
||||
.insert_message(message(a, b, "A's second message"))
|
||||
.await
|
||||
.unwrap();
|
||||
let cancelled_by_owner = store
|
||||
.writer
|
||||
.cancel_messages(a.0, a.1, Some(a_msg_id_2))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
cancelled_by_owner, 1,
|
||||
"A cancelling its own message by exact id must succeed"
|
||||
);
|
||||
assert!(
|
||||
store
|
||||
.writer
|
||||
.pop_message(claim(b), Some(a_msg_id_2))
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none(),
|
||||
"a cancelled message must not be delivered"
|
||||
);
|
||||
}
|
||||
|
||||
/// Popping by an explicit id claims that specific message; FIFO otherwise.
|
||||
#[tokio::test]
|
||||
async fn pop_can_target_a_specific_message_id() {
|
||||
@@ -257,6 +328,45 @@ async fn pop_can_target_a_specific_message_id() {
|
||||
assert_eq!(remaining[0].body, "first");
|
||||
}
|
||||
|
||||
/// Two racing pops of the SAME pending message must not both claim it: the
|
||||
/// baton claim-once discipline this module's own doc comment promises. Both
|
||||
/// calls are dispatched together via `tokio::join!` so their `WriteCmd`s land
|
||||
/// on the single-writer actor back-to-back, exercising the real
|
||||
/// `state = 'pending'` guards in `pop_message_in_transaction` rather than
|
||||
/// relying on test-side sequencing to keep them apart.
|
||||
#[tokio::test]
|
||||
async fn concurrent_pops_of_one_message_deliver_it_exactly_once() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let a = project(&store, "default", "project-a").await;
|
||||
let b = project(&store, "default", "project-b").await;
|
||||
|
||||
let id = store
|
||||
.writer
|
||||
.insert_message(message(a, b, "only one winner"))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Targeted by exact id (rather than `None`/FIFO) so the race exercises
|
||||
// Guard 1 + Guard 2 in `pop_message_in_transaction` directly, instead of
|
||||
// being pre-filtered by the FIFO candidate selection's own `state =
|
||||
// 'pending'` clause before either guard runs.
|
||||
let writer1 = store.writer.clone();
|
||||
let writer2 = store.writer.clone();
|
||||
let (r1, r2) = tokio::join!(
|
||||
writer1.pop_message(claim(b), Some(id)),
|
||||
writer2.pop_message(claim(b), Some(id)),
|
||||
);
|
||||
let r1 = r1.unwrap();
|
||||
let r2 = r2.unwrap();
|
||||
|
||||
let winners = [&r1, &r2].into_iter().filter(|r| r.is_some()).count();
|
||||
assert_eq!(
|
||||
winners, 1,
|
||||
"exactly one of two concurrent pops must claim the message: {r1:?} / {r2:?}",
|
||||
);
|
||||
}
|
||||
|
||||
/// A full recipient inbox rejects new sends so a flood cannot exhaust the
|
||||
/// recipient's context or storage.
|
||||
#[tokio::test]
|
||||
|
||||
@@ -18,7 +18,7 @@
|
||||
|
||||
use ai_memory_core::{
|
||||
ActorContext, AgentKind, HandoffAcceptance, IdentityKey, NewHandoff, NewPage, NewSession,
|
||||
OwnerFilter, PagePath, ProjectId, SessionId, Tier, WorkspaceId, owner_stamp,
|
||||
NewUser, OwnerFilter, PagePath, ProjectId, SessionId, Tier, UserRole, WorkspaceId, owner_stamp,
|
||||
};
|
||||
use ai_memory_store::Store;
|
||||
|
||||
@@ -130,6 +130,74 @@ async fn one_operators_page_is_readable_by_another_in_the_same_project() {
|
||||
assert!(body.body.contains("We picked SQLite"));
|
||||
}
|
||||
|
||||
/// The stronger form of the collaboration guarantee: the existing sibling
|
||||
/// test writes with `author_id: None`, so it cannot tell an "authored pages
|
||||
/// are private" regression from a genuine bug — a filter keyed on the
|
||||
/// caller's identity would happily let a NULL-authored page through. This
|
||||
/// stamps a real, non-null `author_id` (operator A) and asserts operator B —
|
||||
/// a *different* identity, reading with no owner coordinate at all, exactly
|
||||
/// as `search_pages_for_project` and `page_body_by_ids` are shaped — still
|
||||
/// sees the page in full, through both the search path and the direct-body
|
||||
/// path. `pages.author_id` is attribution, never a read filter.
|
||||
#[tokio::test]
|
||||
async fn an_authored_page_is_readable_by_a_different_operator_via_search_and_body() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let store = Store::open(tmp.path()).unwrap();
|
||||
let (ws, proj) = scope(&store).await;
|
||||
|
||||
let operator_a = store
|
||||
.writer
|
||||
.create_human_user(
|
||||
NewUser {
|
||||
username: "operator-a".into(),
|
||||
name: Some("Operator A".into()),
|
||||
email: Some("operator-a@example.com".into()),
|
||||
},
|
||||
UserRole::User,
|
||||
None,
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
store
|
||||
.writer
|
||||
.upsert_page(NewPage {
|
||||
author_id: Some(operator_a),
|
||||
..page(
|
||||
ws,
|
||||
proj,
|
||||
"decisions/0002.md",
|
||||
"Chose SQLite Again",
|
||||
"We picked SQLite for the derived index, authored by operator A.",
|
||||
)
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Operator B's read: no owner coordinate passed anywhere, because none of
|
||||
// these signatures accept one — that absence IS the invariant.
|
||||
let hits = store
|
||||
.reader
|
||||
.search_pages_for_project(ws, proj, "SQLite Again".to_string(), 10, None)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
hits.iter().any(|h| h.path.as_str() == "decisions/0002.md"),
|
||||
"an authored page must be visible to a different operator's search; \
|
||||
got {:?}",
|
||||
hits.iter().map(|h| h.path.as_str()).collect::<Vec<_>>()
|
||||
);
|
||||
|
||||
let body = store
|
||||
.reader
|
||||
.page_body_by_ids(ws, proj, "decisions/0002.md")
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("a different operator can still resolve the page by path");
|
||||
assert!(body.body.contains("authored by operator A"));
|
||||
}
|
||||
|
||||
/// Two harnesses editing the same page keep both versions.
|
||||
///
|
||||
/// The latest write wins the `is_latest` flag — there is no merge, and none is
|
||||
|
||||
Reference in New Issue
Block a user