use super::*;
use crate::embedder::HashEmbedder;
use crate::limits::MAX_FACT_BYTES;
use crate::mcp::dto::{
ExtractionJobStatusParams, ExtractionJobStatusResult, ListMemoriesParams, RecallFusedParams,
RememberExtractedParams, RememberExtractedResult, WhyParams,
};
use crate::model::{ColumnFilter, ColumnOp, Link};
use crate::service::Metadata;
use tempfile::TempDir;
const DECISION: &str = "we chose parking_lot to avoid lock poisoning";
fn server() -> (TempDir, McpServer) {
let dir = TempDir::new().expect("create tempdir");
let embedder: DynEmbedder = Box::new(HashEmbedder::new(crate::DEFAULT_DIMENSION));
let service = MemoryService::open(dir.path(), embedder).expect("open memory store");
let server = McpServer::new(service)
.with_extraction_jobs(dir.path())
.expect("start durable extraction jobs");
(dir, server)
}
async fn committed_extraction(
server: &McpServer,
receipt: RememberExtractedResult,
) -> ExtractionJobStatusResult {
tokio::time::timeout(Duration::from_secs(5), async {
loop {
let Json(status) = server
.extraction_status(Parameters(ExtractionJobStatusParams {
request_id: receipt.request_id.clone(),
}))
.await
.expect("query extraction status");
match status.state {
extraction_jobs::ExtractionJobState::Committed => return status,
extraction_jobs::ExtractionJobState::Failed => {
panic!("extraction failed: {:?}", status.error)
}
extraction_jobs::ExtractionJobState::Accepted
| extraction_jobs::ExtractionJobState::Running => tokio::task::yield_now().await,
}
}
})
.await
.expect("extraction job reaches a terminal state")
}
async fn why_one_hop(srv: &McpServer) -> (Vec<u64>, usize) {
let Json(why) = srv
.why(Parameters(WhyParams {
decision: DECISION.to_owned(),
max_hops: Some(1),
filter: None,
}))
.await
.expect("why");
let ids: Vec<u64> = why.nodes.iter().map(|n| n.id).collect();
(ids, why.edges.len())
}
#[tokio::test]
async fn remember_then_recall_roundtrips_through_the_server() {
let (_dir, srv) = server();
let Json(stored) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let Json(recalled) = srv
.recall(Parameters(RecallParams {
query: "parking_lot poisoning".to_owned(),
limit: None,
filter: None,
}))
.await
.expect("recall");
assert!(recalled.memories.iter().any(|m| m.id == stored.id));
}
#[tokio::test]
async fn feedback_tool_reinforces_and_returns_confidence() {
let (_dir, srv) = server();
let Json(stored) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let Json(up) = srv
.feedback(Parameters(FeedbackParams {
id: stored.id,
success: true,
}))
.await
.expect("feedback success");
assert_eq!(up.id, stored.id);
assert!(up.confidence > 0.5, "success raises confidence");
let Json(down) = srv
.feedback(Parameters(FeedbackParams {
id: stored.id,
success: false,
}))
.await
.expect("feedback failure");
assert!(down.confidence < up.confidence, "failure lowers confidence");
}
#[tokio::test]
async fn feedback_tool_errors_on_unknown_id() {
let (_dir, srv) = server();
let result = srv
.feedback(Parameters(FeedbackParams {
id: 4242,
success: true,
}))
.await;
assert!(result.is_err(), "feedback on an unknown id must error");
}
#[tokio::test]
async fn why_returns_the_connected_subgraph() {
let (_dir, srv) = server();
let Json(decision) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember decision");
let Json(pr) = srv
.remember(Parameters(RememberParams {
fact: "PR #42 swaps the mutex".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember pr");
srv.relate(Parameters(RelateParams {
from: decision.id,
to: pr.id,
relation: "decided_in".to_owned(),
}))
.await
.expect("relate");
let (ids, edges) = why_one_hop(&srv).await;
assert!(ids.contains(&decision.id) && ids.contains(&pr.id));
assert_eq!(edges, 1);
}
#[tokio::test]
async fn forget_removes_a_memory_through_the_server() {
let (_dir, srv) = server();
let Json(stored) = srv
.remember(Parameters(RememberParams {
fact: "ephemeral note about France".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let Json(forgotten) = srv
.forget(Parameters(ForgetParams { id: stored.id }))
.await
.expect("forget");
assert!(forgotten.found, "an existing memory must report found=true");
assert_eq!(forgotten.id, stored.id);
let Json(recalled) = srv
.recall(Parameters(RecallParams {
query: "France".to_owned(),
limit: None,
filter: None,
}))
.await
.expect("recall");
assert!(recalled.memories.iter().all(|m| m.id != stored.id));
}
#[tokio::test]
async fn forget_unknown_id_through_the_server_reports_not_found() {
let (_dir, srv) = server();
let Json(forgotten) = srv
.forget(Parameters(ForgetParams { id: 999_999 }))
.await
.expect("forget on an unknown id must not error");
assert!(
!forgotten.found,
"forget of an id that was never stored must report found=false"
);
}
#[tokio::test]
async fn remember_links_are_traversable_by_why() {
let (_dir, srv) = server();
let Json(pr) = srv
.remember(Parameters(RememberParams {
fact: "PR #99 refactors locks".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember pr");
let Json(decision) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: vec![Link {
target: pr.id,
relation: "decided_in".to_owned(),
}],
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember decision with link");
let (ids, _) = why_one_hop(&srv).await;
assert!(ids.contains(&decision.id) && ids.contains(&pr.id));
}
#[tokio::test]
async fn metadata_and_filter_flow_through_the_server() {
let (_dir, srv) = server();
let mut veles_meta = Metadata::new();
veles_meta.insert("project".to_owned(), serde_json::json!("veles"));
let mut acme_meta = Metadata::new();
acme_meta.insert("project".to_owned(), serde_json::json!("acme"));
let Json(kept) = srv
.remember(Parameters(RememberParams {
fact: "auth bug in login".to_owned(),
links: Vec::new(),
metadata: Some(veles_meta.clone()),
ttl_seconds: None,
}))
.await
.expect("remember veles");
let Json(dropped) = srv
.remember(Parameters(RememberParams {
fact: "auth bug in login too".to_owned(),
links: Vec::new(),
metadata: Some(acme_meta),
ttl_seconds: None,
}))
.await
.expect("remember acme");
let Json(recalled) = srv
.recall(Parameters(RecallParams {
query: "auth bug".to_owned(),
limit: None,
filter: Some(veles_meta),
}))
.await
.expect("recall filtered");
assert!(recalled.memories.iter().any(|m| m.id == kept.id));
assert!(recalled.memories.iter().all(|m| m.id != dropped.id));
}
#[tokio::test]
async fn remember_accepts_explicit_and_default_ttl() {
let (_dir, srv) = server();
let srv = srv.with_default_ttl(3_600);
let Json(explicit) = srv
.remember(Parameters(RememberParams {
fact: "explicit ttl fact".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: Some(3_600),
}))
.await
.expect("remember with explicit ttl");
let Json(defaulted) = srv
.remember(Parameters(RememberParams {
fact: "default ttl fact".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember with default ttl");
let Json(recalled) = srv
.recall(Parameters(RecallParams {
query: "ttl fact".to_owned(),
limit: None,
filter: None,
}))
.await
.expect("recall");
assert!(recalled.memories.iter().any(|m| m.id == explicit.id));
assert!(recalled.memories.iter().any(|m| m.id == defaulted.id));
}
fn ts_meta(ts: i64) -> Metadata {
let mut meta = Metadata::new();
meta.insert("ts".to_owned(), serde_json::json!(ts));
meta
}
#[tokio::test]
async fn recall_where_filters_by_range_through_the_server() {
let (_dir, srv) = server();
for (fact, ts) in [
("kickoff in january", 20_230_115),
("kickoff in june", 20_230_615),
] {
srv.remember(Parameters(RememberParams {
fact: fact.to_owned(),
links: Vec::new(),
metadata: Some(ts_meta(ts)),
ttl_seconds: None,
}))
.await
.expect("remember");
}
let Json(res) = srv
.recall_where(Parameters(RecallWhereParams {
query: "kickoff".to_owned(),
limit: None,
filters: vec![ColumnFilter {
field: "ts".to_owned(),
op: ColumnOp::Ge,
value: serde_json::json!(20_230_601),
}],
}))
.await
.expect("recall_where");
assert!(
res.memories.iter().any(|m| m.content.contains("june")),
"the june fact is within the ts range"
);
assert!(
res.memories.iter().all(|m| !m.content.contains("january")),
"the january fact is below the ts range and excluded"
);
}
#[tokio::test]
async fn recall_where_invalid_field_returns_invalid_params() {
let (_dir, srv) = server();
let err = srv
.recall_where(Parameters(RecallWhereParams {
query: "anything".to_owned(),
limit: None,
filters: vec![ColumnFilter {
field: "ts; DROP TABLE".to_owned(),
op: ColumnOp::Ge,
value: serde_json::json!(1),
}],
}))
.await
.map(|_| ())
.expect_err("a non-identifier filter field must be rejected");
assert_eq!(err.code, ErrorCode::INVALID_PARAMS);
}
#[tokio::test]
async fn empty_fact_returns_invalid_params_not_internal_error() {
let (_dir, srv) = server();
let err = srv
.remember(Parameters(RememberParams {
fact: " ".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.map(|_| ())
.expect_err("whitespace fact must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"EmptyFact must map to invalid_params so clients distinguish bad input from server faults"
);
}
#[tokio::test]
async fn unknown_link_target_returns_invalid_params_not_internal_error() {
let (_dir, srv) = server();
let err = srv
.remember(Parameters(RememberParams {
fact: "a decision".to_owned(),
links: vec![Link {
target: 9_999_999,
relation: "x".to_owned(),
}],
metadata: None,
ttl_seconds: None,
}))
.await
.map(|_| ())
.expect_err("unknown link target must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"UnknownMemory must map to invalid_params"
);
}
#[tokio::test]
async fn relate_to_unknown_endpoint_returns_invalid_params_not_internal_error() {
let (_dir, srv) = server();
let Json(stored) = srv
.remember(Parameters(RememberParams {
fact: "an existing memory".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let err = srv
.relate(Parameters(RelateParams {
from: stored.id,
to: 9_999_999,
relation: "references".to_owned(),
}))
.await
.map(|_| ())
.expect_err("relate to a missing endpoint must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"relate to an unknown endpoint must map to invalid_params"
);
}
#[tokio::test]
async fn oversized_fact_returns_invalid_params() {
let (_dir, srv) = server();
let huge = "x".repeat(MAX_FACT_BYTES + 1);
let err = srv
.remember(Parameters(RememberParams {
fact: huge,
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.map(|_| ())
.expect_err("oversized fact must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"oversized fact must map to invalid_params"
);
}
#[tokio::test]
async fn entity_miss_echoes_the_canonical_queried_name() {
let (_dir, srv) = server();
let Json(profile) = srv
.entity(Parameters(EntityParams {
name: " Zzz Personne Inexistante ".to_owned(),
}))
.await
.expect("entity lookup");
assert!(!profile.found, "nothing was ever stored under that name");
assert_eq!(
profile.name, "zzz personne inexistante",
"a miss must echo the queried name, canonicalized like a hit's"
);
assert_eq!(profile.id, 0, "a miss still reports no id");
}
#[tokio::test]
async fn recall_fused_folds_in_a_graph_reached_fact() {
let (_dir, srv) = server();
let Json(anchor) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember anchor");
let Json(linked) = srv
.remember(Parameters(RememberParams {
fact: "the on-call rotation moved to Tuesdays".to_owned(),
links: vec![Link {
target: anchor.id,
relation: "context".to_owned(),
}],
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember linked");
let Json(plain) = srv
.recall(Parameters(RecallParams {
query: DECISION.to_owned(),
limit: Some(1),
filter: None,
}))
.await
.expect("recall");
assert!(!plain.memories.iter().any(|m| m.id == linked.id));
let Json(fused) = srv
.recall_fused(Parameters(RecallFusedParams {
query: DECISION.to_owned(),
limit: Some(10),
filter: None,
hops: None,
graph_boost: None,
pool: None,
date_field: None,
}))
.await
.expect("recall_fused");
assert!(
fused.memories.iter().any(|m| m.id == anchor.id),
"anchor present in fused recall"
);
assert!(
fused.memories.iter().any(|m| m.id == linked.id),
"graph-reached fact folded into fused recall"
);
}
#[tokio::test]
async fn recall_fused_with_date_field_returns_a_dated_timeline() {
let (_dir, srv) = server();
for (fact, ts) in [
("the release shipped", 20_260_701_i64),
("the project kicked off", 20_260_103),
] {
srv.remember(Parameters(RememberParams {
fact: fact.to_owned(),
links: Vec::new(),
metadata: Some(ts_meta(ts)),
ttl_seconds: None,
}))
.await
.expect("remember dated fact");
}
let Json(res) = srv
.recall_fused(Parameters(RecallFusedParams {
query: "project release timeline".to_owned(),
limit: Some(10),
filter: None,
hops: None,
graph_boost: None,
pool: None,
date_field: Some("ts".to_owned()),
}))
.await
.expect("recall_fused dated");
let timeline = res
.dated_context
.expect("dated_context present when date_field set");
assert!(timeline.contains("- [2026-01-03] the project kicked off"));
assert!(timeline.contains("- [2026-07-01] the release shipped"));
assert!(
timeline.find("2026-01-03").unwrap() < timeline.find("2026-07-01").unwrap(),
"facts must be ordered oldest-first"
);
assert_eq!(res.now.as_deref(), Some("2026-07-01"));
}
#[tokio::test]
async fn recall_fused_without_date_field_omits_the_timeline() {
let (_dir, srv) = server();
srv.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let Json(res) = srv
.recall_fused(Parameters(RecallFusedParams {
query: DECISION.to_owned(),
limit: Some(5),
filter: None,
hops: None,
graph_boost: None,
pool: None,
date_field: None,
}))
.await
.expect("recall_fused");
assert!(res.dated_context.is_none());
assert!(res.now.is_none());
}
#[tokio::test]
async fn recall_fused_survives_a_non_finite_graph_boost() {
let (_dir, srv) = server();
let Json(anchor) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember anchor");
let Json(linked) = srv
.remember(Parameters(RememberParams {
fact: "the on-call rotation moved to Tuesdays".to_owned(),
links: vec![Link {
target: anchor.id,
relation: "context".to_owned(),
}],
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember linked");
let Json(fused) = srv
.recall_fused(Parameters(RecallFusedParams {
query: DECISION.to_owned(),
limit: Some(10),
filter: None,
hops: None,
graph_boost: Some(f64::NAN),
pool: None,
date_field: None,
}))
.await
.expect("recall_fused");
assert!(
fused.memories.iter().any(|m| m.id == linked.id),
"graph-reached fact must still surface despite a non-finite graph_boost"
);
}
#[tokio::test]
async fn recall_fused_limit_is_capped_at_max() {
let (_dir, srv) = server();
let Json(result) = srv
.recall_fused(Parameters(RecallFusedParams {
query: "anything".to_owned(),
limit: Some(usize::MAX),
filter: None,
hops: Some(usize::MAX),
graph_boost: None,
pool: Some(usize::MAX),
date_field: None,
}))
.await
.expect("recall_fused with huge limit/hops/pool must succeed (silently capped)");
let _ = result;
}
#[tokio::test]
async fn recall_limit_is_capped_at_max() {
let (_dir, srv) = server();
let Json(result) = srv
.recall(Parameters(RecallParams {
query: "anything".to_owned(),
limit: Some(usize::MAX),
filter: None,
}))
.await
.expect("recall with huge limit must succeed (silently capped)");
let _ = result;
}
#[tokio::test]
async fn why_hop_depth_is_capped_at_max() {
let (_dir, srv) = server();
srv.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let Json(_) = srv
.why(Parameters(WhyParams {
decision: DECISION.to_owned(),
max_hops: Some(usize::MAX),
filter: None,
}))
.await
.expect("why with huge max_hops must succeed (silently capped)");
}
struct GraphStub;
impl Extractor for GraphStub {
fn extract(&self, _text: &str) -> Result<Vec<ExtractedFact>, ExtractError> {
Ok(vec![
ExtractedFact {
text: "Alice ships the parser in Rust.".to_owned(),
entities: vec!["rust".to_owned()],
},
ExtractedFact {
text: "Bob maintains the Rust toolchain.".to_owned(),
entities: vec!["rust".to_owned()],
},
])
}
}
struct GenerativeStub;
impl Extractor for GenerativeStub {
fn extract(&self, _text: &str) -> Result<Vec<ExtractedFact>, ExtractError> {
Ok(vec![ExtractedFact {
text: "Generative interpretation".to_owned(),
entities: vec!["generative".to_owned()],
}])
}
}
struct CountingExtractionStub {
calls: AtomicUsize,
}
struct FailingExtractionStub {
calls: AtomicUsize,
}
impl Extractor for FailingExtractionStub {
fn extract(&self, _text: &str) -> Result<Vec<ExtractedFact>, ExtractError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Err(ExtractError::Backend("model unavailable".to_owned()))
}
}
impl Extractor for CountingExtractionStub {
fn extract(&self, _text: &str) -> Result<Vec<ExtractedFact>, ExtractError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(vec![ExtractedFact {
text: "One idempotent extraction result.".to_owned(),
entities: Vec::new(),
}])
}
}
#[tokio::test]
async fn remember_extracted_builds_a_graph_through_the_server() {
let (_dir, srv) = server();
let srv = srv
.with_named_extractor("ollama", Arc::new(GraphStub) as DynExtractor)
.expect("ollama is a supported extractor name");
let Json(receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "Alice and Bob work in Rust.".to_owned(),
metadata: None,
extractor: None,
idempotency_key: None,
}))
.await
.expect("remember_extracted");
let res = committed_extraction(&srv, receipt).await;
assert_eq!(res.ids.len(), 2, "both facts stored");
let Json(why) = srv
.why(Parameters(WhyParams {
decision: "parser in rust".to_owned(),
max_hops: Some(2),
filter: None,
}))
.await
.expect("why");
assert!(why.nodes.len() > 1, "graph is alive through the server");
assert!(
!why.nodes[0].content.starts_with("Entity:"),
"seed is a fact, not a hub"
);
}
#[tokio::test]
async fn two_successive_calls_can_select_different_extractors() {
let (_dir, srv) = server();
let srv = srv
.with_named_extractor("ollama", Arc::new(GenerativeStub) as DynExtractor)
.expect("ollama is a supported extractor name");
let source = "fact: Deterministic directive | outline";
let Json(outlined_receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: source.to_owned(),
metadata: None,
extractor: Some("outline".to_owned()),
idempotency_key: None,
}))
.await
.expect("the per-call outline backend must work");
let Json(generative_receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: source.to_owned(),
metadata: None,
extractor: Some("ollama".to_owned()),
idempotency_key: None,
}))
.await
.expect("the configured remote backend must remain selectable");
let outlined = committed_extraction(&srv, outlined_receipt).await;
let generative = committed_extraction(&srv, generative_receipt).await;
assert_ne!(outlined.ids, generative.ids, "the extractions must differ");
let Json(listed) = srv
.list_memories(Parameters(ListMemoriesParams {
cursor: None,
limit: None,
filter: None,
include_internal: false,
}))
.await
.expect("list the two extracted facts");
let contents: Vec<&str> = listed
.memories
.iter()
.map(|memory| memory.content.as_str())
.collect();
assert!(contents.contains(&"Deterministic directive"));
assert!(contents.contains(&"Generative interpretation"));
}
#[tokio::test]
async fn retry_key_runs_one_extraction_and_rejects_a_changed_payload() {
let (_dir, srv) = server();
let extractor = Arc::new(CountingExtractionStub {
calls: AtomicUsize::new(0),
});
let srv = srv.with_extractor(extractor.clone());
let submit = |text: &str| RememberExtractedParams {
text: text.to_owned(),
metadata: None,
extractor: None,
idempotency_key: Some("client-operation-42".to_owned()),
};
let Json(first) = srv
.remember_extracted(Parameters(submit("same passage")))
.await
.expect("first durable receipt");
let Json(retry) = srv
.remember_extracted(Parameters(submit("same passage")))
.await
.expect("retry must reuse the receipt");
assert_eq!(retry.request_id, first.request_id);
assert!(retry.reused);
let result = committed_extraction(&srv, first).await;
assert_eq!(result.ids.len(), 1);
assert_eq!(extractor.calls.load(Ordering::SeqCst), 1);
let error = srv
.remember_extracted(Parameters(submit("changed passage")))
.await
.map(|_| ())
.expect_err("one retry key cannot alias a changed payload");
assert_eq!(error.code, ErrorCode::INVALID_PARAMS);
}
#[tokio::test]
async fn failed_extraction_is_terminal_queryable_and_not_retried() {
let (_dir, srv) = server();
let extractor = Arc::new(FailingExtractionStub {
calls: AtomicUsize::new(0),
});
let srv = srv.with_extractor(extractor.clone());
let params = || RememberExtractedParams {
text: "passage whose model call fails".to_owned(),
metadata: None,
extractor: None,
idempotency_key: Some("terminal-failure-proof".to_owned()),
};
let Json(receipt) = srv
.remember_extracted(Parameters(params()))
.await
.expect("the durable request is accepted before generation fails");
let status = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let Json(status) = srv
.extraction_status(Parameters(ExtractionJobStatusParams {
request_id: receipt.request_id.clone(),
}))
.await
.expect("query failed extraction");
if status.state.is_terminal() {
return status;
}
tokio::task::yield_now().await;
}
})
.await
.expect("failed extraction reaches its terminal record");
assert_eq!(status.state, extraction_jobs::ExtractionJobState::Failed);
assert!(status.ids.is_empty());
assert!(status.skipped_over_cap.is_none());
assert!(status
.error
.as_deref()
.is_some_and(|error| error.contains("model unavailable")));
let Json(retry) = srv
.remember_extracted(Parameters(params()))
.await
.expect("same key and payload reuses terminal failure");
assert!(retry.reused);
assert_eq!(retry.state, extraction_jobs::ExtractionJobState::Failed);
assert_eq!(extractor.calls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn reserved_metadata_key_returns_invalid_params() {
let (_dir, srv) = server();
let mut bad_meta = Metadata::new();
bad_meta.insert("_veles_hub".to_owned(), serde_json::json!(true));
let err = srv
.remember(Parameters(RememberParams {
fact: "a fact".to_owned(),
links: Vec::new(),
metadata: Some(bad_meta),
ttl_seconds: None,
}))
.await
.map(|_| ())
.expect_err("reserved metadata key must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"ReservedKey must map to invalid_params, not internal_error"
);
}
#[tokio::test]
async fn recall_where_non_scalar_filter_value_returns_invalid_params() {
let (_dir, srv) = server();
let err = srv
.recall_where(Parameters(RecallWhereParams {
query: "query".to_owned(),
limit: None,
filters: vec![ColumnFilter {
field: "ts".to_owned(),
op: ColumnOp::Eq,
value: serde_json::json!([1, 2, 3]),
}],
}))
.await
.map(|_| ())
.expect_err("array filter value must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"non-scalar filter value must map to invalid_params"
);
}
#[tokio::test]
async fn relate_with_empty_relation_returns_invalid_params() {
let (_dir, srv) = server();
let Json(a) = srv
.remember(Parameters(RememberParams {
fact: "fact A".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember A");
let Json(b) = srv
.remember(Parameters(RememberParams {
fact: "fact B".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember B");
let err = srv
.relate(Parameters(RelateParams {
from: a.id,
to: b.id,
relation: String::new(),
}))
.await
.map(|_| ())
.expect_err("empty relation must be rejected");
assert_eq!(
err.code,
ErrorCode::INVALID_PARAMS,
"InvalidRelation must map to invalid_params"
);
}
#[tokio::test]
async fn remember_extracted_without_backend_returns_internal_error() {
let (_dir, srv) = server(); let err = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "anything".to_owned(),
metadata: None,
extractor: None,
idempotency_key: None,
}))
.await
.map(|_| ())
.expect_err("extraction with no backend must error");
assert_eq!(err.code, ErrorCode::INTERNAL_ERROR);
}
#[tokio::test]
async fn explicit_outline_works_without_a_daemon_default() {
let (_dir, srv) = server();
let Json(receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "fact: A deterministic fact | outline".to_owned(),
metadata: None,
extractor: Some("outline".to_owned()),
idempotency_key: None,
}))
.await
.expect("explicit outline must not need a daemon default");
let stored = committed_extraction(&srv, receipt).await;
assert_eq!(stored.ids.len(), 1);
}
#[tokio::test]
async fn outline_extraction_lands_a_typed_edge_in_the_graph() {
let (_dir, srv) = server();
let Json(receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "fact: Alice ships the parser.\nedge: Alice | works_with | Bob".to_owned(),
metadata: None,
extractor: Some("outline".to_owned()),
idempotency_key: None,
}))
.await
.expect("the offline reader must not need a daemon default");
let stored = committed_extraction(&srv, receipt).await;
assert_eq!(stored.ids.len(), 1, "the `fact:` line is the only fact");
let deadline = Instant::now() + Duration::from_secs(5);
loop {
let Json(profile) = srv
.entity(Parameters(EntityParams {
name: "alice".to_owned(),
}))
.await
.expect("entity lookup");
let edge = profile.relations.iter().find(|relation| {
relation.predicate == "works_with" && relation.target.to_lowercase().ends_with("bob")
});
if edge.is_some() {
assert!(profile.found, "an entity with an edge is a found entity");
return;
}
assert!(
Instant::now() < deadline,
"no `works_with` edge reached the graph from the offline reader; \
found={} relations=[{}]",
profile.found,
profile
.relations
.iter()
.map(|relation| format!("{} -> {}", relation.predicate, relation.target))
.collect::<Vec<_>>()
.join(", ")
);
tokio::task::yield_now().await;
}
}
#[tokio::test]
async fn unknown_per_call_extractor_returns_invalid_params() {
let (_dir, srv) = server();
let err = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "anything".to_owned(),
metadata: None,
extractor: Some("lmstudio".to_owned()),
idempotency_key: None,
}))
.await
.map(|_| ())
.expect_err("unknown extractor must be rejected");
assert_eq!(err.code, ErrorCode::INVALID_PARAMS);
assert!(err.message.contains("unknown extraction backend"));
}
#[test]
fn relate_params_accept_string_or_number_ids_on_the_wire() {
let numeric: RelateParams =
serde_json::from_value(serde_json::json!({"from": 1, "to": 2, "relation": "r"}))
.expect("numeric ids must still deserialize (0.9.x compat)");
assert_eq!((numeric.from, numeric.to), (1, 2));
let stringy: RelateParams =
serde_json::from_value(serde_json::json!({"from": "1", "to": "2", "relation": "r"}))
.expect("decimal-string ids must deserialize");
assert_eq!((stringy.from, stringy.to), (1, 2));
}
#[test]
fn forget_and_feedback_params_accept_string_ids_on_the_wire() {
let forget: ForgetParams = serde_json::from_value(serde_json::json!({"id": "42"}))
.expect("forget id must accept a decimal string");
assert_eq!(forget.id, 42);
let feedback: FeedbackParams =
serde_json::from_value(serde_json::json!({"id": "42", "success": true}))
.expect("feedback id must accept a decimal string");
assert_eq!(feedback.id, 42);
}
#[test]
fn remember_link_target_accepts_a_string_id_on_the_wire() {
let link: Link = serde_json::from_value(serde_json::json!({"target": "7", "relation": "r"}))
.expect("Link::target must accept a decimal string");
assert_eq!(link.target, 7);
}
#[tokio::test]
async fn remember_recall_relate_forget_feedback_responses_echo_an_id_str_twin() {
let (_dir, srv) = server();
let Json(remembered) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
assert_eq!(remembered.id_str, remembered.id.to_string());
let Json(recalled) = srv
.recall(Parameters(RecallParams {
query: "parking_lot poisoning".to_owned(),
limit: None,
filter: None,
}))
.await
.expect("recall");
let hit = recalled
.memories
.iter()
.find(|m| m.id == remembered.id)
.expect("recalled memory present");
assert_eq!(hit.id_str, hit.id.to_string());
let Json(pr) = srv
.remember(Parameters(RememberParams {
fact: "PR #42 swaps the mutex".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember pr");
let Json(relate_res) = srv
.relate(Parameters(RelateParams {
from: remembered.id,
to: pr.id,
relation: "decided_in".to_owned(),
}))
.await
.expect("relate");
assert_eq!(relate_res.edge_id_str, relate_res.edge_id.to_string());
let Json(feedback_res) = srv
.feedback(Parameters(FeedbackParams {
id: remembered.id,
success: true,
}))
.await
.expect("feedback");
assert_eq!(feedback_res.id_str, feedback_res.id.to_string());
let Json(forget_res) = srv
.forget(Parameters(ForgetParams { id: pr.id }))
.await
.expect("forget");
assert_eq!(forget_res.id_str, forget_res.id.to_string());
}
#[tokio::test]
async fn why_response_echoes_id_str_and_from_to_str_on_nodes_and_edges() {
let (_dir, srv) = server();
let Json(decision) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember decision");
let Json(pr) = srv
.remember(Parameters(RememberParams {
fact: "PR #42 swaps the mutex".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember pr");
srv.relate(Parameters(RelateParams {
from: decision.id,
to: pr.id,
relation: "decided_in".to_owned(),
}))
.await
.expect("relate");
let Json(why) = srv
.why(Parameters(WhyParams {
decision: DECISION.to_owned(),
max_hops: Some(1),
filter: None,
}))
.await
.expect("why");
assert!(!why.nodes.is_empty() && !why.edges.is_empty());
for node in &why.nodes {
assert_eq!(node.id_str, node.id.to_string());
}
for edge in &why.edges {
assert_eq!(edge.from_str, edge.from.to_string());
assert_eq!(edge.to_str, edge.to.to_string());
}
}
#[tokio::test]
async fn remember_extracted_response_echoes_ids_str() {
use crate::extract::{ExtractError, ExtractedFact, Extractor};
struct Stub;
impl Extractor for Stub {
fn extract(&self, _text: &str) -> Result<Vec<ExtractedFact>, ExtractError> {
Ok(vec![ExtractedFact {
text: "Alice ships the parser in Rust.".to_owned(),
entities: vec!["rust".to_owned()],
}])
}
}
let (_dir, srv) = server();
let srv = srv.with_extractor(Arc::new(Stub) as DynExtractor);
let Json(receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "Alice works in Rust.".to_owned(),
metadata: None,
extractor: None,
idempotency_key: None,
}))
.await
.expect("remember_extracted");
let res = committed_extraction(&srv, receipt).await;
assert_eq!(res.ids_str.len(), res.ids.len());
for (id, id_str) in res.ids.iter().zip(res.ids_str.iter()) {
assert_eq!(*id_str, id.to_string());
}
}
#[tokio::test]
async fn relate_by_wrong_numeric_id_fails_but_id_str_round_trip_succeeds() {
let (_dir, srv) = server();
let Json(decision) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember decision");
let Json(pr) = srv
.remember(Parameters(RememberParams {
fact: "PR #42 swaps the mutex".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember pr");
let wrong_from = decision.id + 1_000_003;
let err = srv
.relate(Parameters(RelateParams {
from: wrong_from,
to: pr.id,
relation: "decided_in".to_owned(),
}))
.await
.map(|_| ())
.expect_err("a rounded/wrong id must not silently resolve to the real memory");
assert_eq!(err.code, ErrorCode::INVALID_PARAMS);
let params: RelateParams = serde_json::from_value(serde_json::json!({
"from": decision.id_str,
"to": pr.id_str,
"relation": "decided_in",
}))
.expect("id_str values must deserialize as RelateParams");
srv.relate(Parameters(params))
.await
.expect("relate via id_str must succeed");
let (ids, edges) = why_one_hop(&srv).await;
assert!(ids.contains(&decision.id) && ids.contains(&pr.id));
assert_eq!(edges, 1);
}
#[test]
fn test_get_info_instructions_cover_memory_compiler_and_working_context() {
let (_dir, srv) = server();
let info = srv.get_info();
let instructions = info.instructions.expect("instructions must be set");
assert!(
instructions.contains("recall") && instructions.contains("relate"),
"must mention the memory family: {instructions}"
);
#[cfg(feature = "context")]
{
assert!(
instructions.contains("compile_context"),
"must mention the context compiler family: {instructions}"
);
assert!(
instructions.contains("working"),
"must mention working-context resumption: {instructions}"
);
assert!(
instructions.contains("compile_transcript"),
"must mention the compile_transcript shortcut (V2b-2/V2b-3): {instructions}"
);
}
}
#[test]
fn test_recall_where_description_documents_type_strict_comparisons() {
let tool = McpServer::recall_where_tool_attr();
let description = tool
.description
.expect("recall_where must declare a description");
let lower = description.to_lowercase();
assert!(
lower.contains("type-strict") || lower.contains("type strict"),
"recall_where's description must document type-strict comparisons: {description}"
);
assert!(
description.to_lowercase().contains("numeric"),
"recall_where's description must advise storing comparable values \
(e.g. dates) numerically: {description}"
);
}
fn description_names_field(haystack: &str, field: &str) -> bool {
let continues_word = |byte: u8| byte.is_ascii_alphanumeric() || byte == b'_';
let bytes = haystack.as_bytes();
haystack.match_indices(field).any(|(start, _)| {
let end = start + field.len();
let opens = start == 0 || !continues_word(bytes[start - 1]);
let closes = end == bytes.len() || !continues_word(bytes[end]);
opens && closes
})
}
#[test]
fn test_remember_extracted_description_names_every_published_output_field() {
let tool = McpServer::remember_extracted_tool_attr();
let description = tool
.description
.clone()
.expect("remember_extracted must declare a description");
let schema = tool
.output_schema
.as_ref()
.expect("remember_extracted declares an explicit output_schema");
let properties = schema
.get("properties")
.and_then(serde_json::Value::as_object)
.expect("its output schema is an object schema carrying `properties`");
assert!(
!properties.is_empty(),
"the published output schema must expose root properties, else this \
test would pass vacuously"
);
let missing: Vec<&str> = properties
.keys()
.filter(|field| !description_names_field(&description, field))
.map(String::as_str)
.collect();
assert!(
missing.is_empty(),
"remember_extracted's description must name every root field of its \
published outputSchema. Missing: {missing:?}. A field the description \
omits is a field the model never learns to read, and `skipped_over_cap` \
omitted is a silent data loss — the caller believes everything it sent \
was stored. Description was: {description}"
);
}
#[test]
fn test_recall_fused_input_schema_types_every_parameter_directly() {
let tool = McpServer::recall_fused_tool_attr();
let schema = serde_json::to_value(&tool.input_schema).expect("schema serializes");
let properties = schema["properties"]
.as_object()
.expect("recall_fused input schema must have properties");
for (name, subschema) in properties {
assert!(
subschema.get("type").is_some(),
"recall_fused parameter `{name}` must advertise a direct `type` \
keyword (anyOf/$ref-only schemas get stringified by real MCP \
harnesses); got: {subschema}"
);
}
}
#[test]
fn recall_family_advertises_k_as_the_canonical_count_parameter() {
let tools = [
McpServer::recall_tool_attr(),
McpServer::recall_where_tool_attr(),
McpServer::recall_fused_tool_attr(),
];
for tool in tools {
let schema = serde_json::to_value(&tool.input_schema).expect("schema serializes");
let properties = schema["properties"]
.as_object()
.expect("recall input schema must have properties");
assert!(
properties.contains_key("k"),
"{} must advertise `k`",
tool.name
);
assert!(
!properties.contains_key("limit"),
"{} must keep `limit` as a wire alias, not the canonical schema field",
tool.name
);
}
}
#[test]
fn test_recall_fused_params_accept_stringified_scalars_and_objects() {
let params: RecallFusedParams = serde_json::from_value(serde_json::json!({
"query": "q",
"k": "6",
"hops": "2",
"graph_boost": "0.15",
"pool": "128",
"filter": "{\"project\": \"velesdb\"}"
}))
.expect("stringified scalar/object arguments must deserialize");
assert_eq!(params.limit, Some(6));
assert_eq!(params.hops, Some(2));
assert_eq!(params.pool, Some(128));
assert!((params.graph_boost.unwrap() - 0.15).abs() < f64::EPSILON);
let filter = params.filter.expect("filter must parse from a JSON string");
assert_eq!(
filter.get("project").and_then(|v| v.as_str()),
Some("velesdb")
);
}
#[test]
fn recall_fused_accepts_k_and_the_deprecated_limit_alias() {
for (input, value) in [
(serde_json::json!({"query": "q", "k": 7}), 7),
(serde_json::json!({"query": "q", "limit": 8}), 8),
] {
let params: RecallFusedParams = serde_json::from_value(input)
.expect("canonical and deprecated spellings must deserialize");
assert_eq!(params.limit, Some(value));
}
}
#[test]
fn recall_and_recall_where_accept_the_deprecated_limit_alias() {
let recall: RecallParams = serde_json::from_value(serde_json::json!({
"query": "q",
"limit": 7
}))
.expect("recall must accept deprecated `limit`");
assert_eq!(recall.limit, Some(7));
let recall_where: RecallWhereParams = serde_json::from_value(serde_json::json!({
"query": "q",
"limit": 8,
"filters": []
}))
.expect("recall_where must accept deprecated `limit`");
assert_eq!(recall_where.limit, Some(8));
}
#[test]
fn recall_fused_rejects_k_and_limit_together() {
let result = serde_json::from_value::<RecallFusedParams>(serde_json::json!({
"query": "q",
"k": 7,
"limit": 8
}));
let Err(error) = result else {
panic!("two spellings for one parameter must be ambiguous");
};
assert!(error.to_string().contains("duplicate field"), "{error}");
}
#[tokio::test]
async fn unrelate_tool_removes_the_edge_and_is_idempotent() {
let (_dir, srv) = server();
let Json(a) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember a");
let Json(b) = srv
.remember(Parameters(RememberParams {
fact: "the cause behind the decision".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember b");
srv.relate(Parameters(RelateParams {
from: a.id,
to: b.id,
relation: "caused_by".to_owned(),
}))
.await
.expect("relate");
let Json(res) = srv
.unrelate(Parameters(UnrelateParams {
from: a.id,
to: b.id,
relation: "caused_by".to_owned(),
}))
.await
.expect("unrelate");
assert!(res.found, "the edge existed");
assert_eq!(res.removed, 1);
let Json(res) = srv
.unrelate(Parameters(UnrelateParams {
from: a.id,
to: b.id,
relation: "caused_by".to_owned(),
}))
.await
.expect("a second unrelate must not error — cleanups are replayable");
assert!(!res.found, "already removed");
assert_eq!(res.removed, 0);
let (_ids, edge_count) = why_one_hop(&srv).await;
assert_eq!(edge_count, 0, "the edge must be gone from traversal");
}
#[tokio::test]
async fn unrelate_refuses_a_self_loop_as_invalid_params() {
let (_dir, srv) = server();
let Json(a) = srv
.remember(Parameters(RememberParams {
fact: DECISION.to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let err = srv
.unrelate(Parameters(UnrelateParams {
from: a.id,
to: a.id,
relation: "supports".to_owned(),
}))
.await
.map(|_| ())
.expect_err("a self-loop is refused exactly like relate's");
assert_eq!(err.code, ErrorCode::INVALID_PARAMS);
}
#[test]
fn unrelate_params_accept_string_or_number_ids_on_the_wire() {
let params: UnrelateParams = serde_json::from_value(serde_json::json!({
"from": "18446744073709551615",
"to": 42,
"relation": "caused_by"
}))
.expect("string and number ids must both deserialize");
assert_eq!(params.from, u64::MAX);
assert_eq!(params.to, 42);
}
#[tokio::test]
async fn an_outline_configured_server_builds_a_hub_that_entity_finds() {
let dir = TempDir::new().expect("create tempdir");
let embedder: DynEmbedder = Box::new(HashEmbedder::new(crate::DEFAULT_DIMENSION));
let service = MemoryService::open(dir.path(), embedder).expect("open memory store");
let crate::ExtractorSelection::Ready(extractor) =
crate::select_extractor("outline").expect("`outline` must be an accepted backend name")
else {
panic!("`outline` must be usable as-is, with no remote configuration");
};
let srv = McpServer::new(service)
.with_extractor(extractor)
.with_extraction_jobs(dir.path())
.expect("start durable extraction jobs");
let Json(receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "fact: Theo has a sister called Camille | Theo, Camille\n\
edge: Camille | soeur de | Theo\n\
attr: Theo | age | 15"
.to_owned(),
metadata: None,
extractor: None,
idempotency_key: None,
}))
.await
.expect("remember_extracted must work on an outline-configured server");
let stored = committed_extraction(&srv, receipt).await;
assert_eq!(
stored.ids.len(),
1,
"the passage carries exactly one `fact:` directive, so exactly one \
fact must be stored (`edge:` and `attr:` build the graph around it, \
they are not facts of their own)"
);
let Json(theo) = srv
.entity(Parameters(EntityParams {
name: "Theo".to_owned(),
}))
.await
.expect("entity");
assert!(
theo.found,
"entity must find the hub that remember_extracted just created — \
this is the second of the two behaviours #1734 reported dead"
);
let outgoing: Vec<&str> = theo
.relations
.iter()
.map(|r| r.predicate.as_str())
.collect();
let incoming: Vec<&str> = theo
.relations_in
.iter()
.map(|r| r.predicate.as_str())
.collect();
assert!(
incoming.contains(&"soeur de") || outgoing.contains(&"soeur de"),
"the `edge:` directive must reach the graph — outgoing {outgoing:?}, \
incoming {incoming:?}"
);
}
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::mpsc::{sync_channel, Receiver, SyncSender};
use std::sync::Mutex;
use std::time::{Duration, Instant};
use crate::extract::{ExtractError, ExtractedFact, ExtractedRelation, Extraction, Extractor};
struct GatedServerExtractor {
gate: Mutex<Receiver<()>>,
entered: AtomicUsize,
completed: AtomicUsize,
}
impl GatedServerExtractor {
fn new() -> (Arc<Self>, SyncSender<()>) {
let (tx, rx) = sync_channel::<()>(64);
(
Arc::new(Self {
gate: Mutex::new(rx),
entered: AtomicUsize::new(0),
completed: AtomicUsize::new(0),
}),
tx,
)
}
}
impl Extractor for GatedServerExtractor {
fn extract(&self, _text: &str) -> Result<Vec<ExtractedFact>, ExtractError> {
Ok(Vec::new())
}
fn extract_graph(&self, _text: &str) -> Result<Extraction, ExtractError> {
self.entered.fetch_add(1, Ordering::SeqCst);
self.gate
.lock()
.expect("gate lock")
.recv()
.map_err(|_| ExtractError::Backend("gate closed".to_owned()))?;
self.completed.fetch_add(1, Ordering::SeqCst);
Ok(Extraction {
facts: Vec::new(),
relations: vec![ExtractedRelation {
subject: "alice martin".to_owned(),
predicate: "travaille chez".to_owned(),
object: "wiscale".to_owned(),
}],
attributes: Vec::new(),
})
}
}
fn gated_server() -> (
TempDir,
McpServer,
Arc<GatedServerExtractor>,
SyncSender<()>,
) {
let dir = TempDir::new().expect("create tempdir");
let embedder: DynEmbedder = Box::new(HashEmbedder::new(crate::DEFAULT_DIMENSION));
let (extractor, release) = GatedServerExtractor::new();
let service = MemoryService::open(dir.path(), embedder)
.expect("open memory store")
.with_autograph(extractor.clone());
(dir, McpServer::new(service), extractor, release)
}
#[tokio::test]
async fn remember_extracted_receipt_does_not_wait_for_generation() {
let (_dir, srv) = server();
let (extractor, release) = GatedServerExtractor::new();
let srv = srv.with_extractor(extractor.clone());
let Json(receipt) = srv
.remember_extracted(Parameters(RememberExtractedParams {
text: "Alice works at Wiscale.".to_owned(),
metadata: None,
extractor: None,
idempotency_key: Some("immediate-receipt-proof".to_owned()),
}))
.await
.expect("durable acceptance must not wait for the extractor");
assert_eq!(receipt.state, extraction_jobs::ExtractionJobState::Accepted);
assert_eq!(
extractor.completed.load(Ordering::SeqCst),
0,
"the receipt must arrive while generation is still blocked"
);
assert!(
wait_for(Duration::from_secs(5), || extractor
.entered
.load(Ordering::SeqCst)
== 1),
"the background worker must start the accepted job"
);
let Json(running) = srv
.extraction_status(Parameters(ExtractionJobStatusParams {
request_id: receipt.request_id.clone(),
}))
.await
.expect("read the in-flight job");
assert_eq!(running.state, extraction_jobs::ExtractionJobState::Running);
release.send(()).expect("release extraction");
let status = committed_extraction(&srv, receipt).await;
assert_eq!(status.state, extraction_jobs::ExtractionJobState::Committed);
assert_eq!(extractor.completed.load(Ordering::SeqCst), 1);
}
fn wait_for(deadline: Duration, mut probe: impl FnMut() -> bool) -> bool {
let end = Instant::now() + deadline;
loop {
if probe() {
return true;
}
if Instant::now() >= end {
return false;
}
std::thread::sleep(Duration::from_millis(20));
}
}
#[tokio::test]
async fn remember_tool_answers_while_the_autograph_gate_is_still_shut() {
let (_dir, srv, _extractor, release) = gated_server();
#[allow(clippy::used_underscore_binding)]
let handle_stored = srv._autograph_worker.is_some();
assert!(
handle_stored,
"constructing the server over an autograph-carrying service must \
spawn the background worker and STORE its handle — the handle is \
what makes the server's drop bound shutdown"
);
assert!(
srv.service
.inspect(MemoryService::autograph_queue_open)
.expect("active generation"),
"the spawned worker's queue must be open, so remember enqueues \
instead of running the enrichment on the response path"
);
let started = Instant::now();
let Json(_stored) = srv
.remember(Parameters(RememberParams {
fact: "Alice Martin travaille chez Wiscale.".to_owned(),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_secs(2),
"the remember TOOL must answer with the extractor gate still SHUT — \
it took {elapsed:?}, so the enrichment sat on the server's response \
path, the exact pre-#1851 behaviour"
);
release.send(()).expect("release the gated extraction");
let deadline = Instant::now() + Duration::from_secs(10);
loop {
let Json(profile) = srv
.entity(Parameters(EntityParams {
name: "alice martin".to_owned(),
}))
.await
.expect("entity");
if profile.found && !profile.relations.is_empty() {
break;
}
assert!(
Instant::now() < deadline,
"the deferred enrichment must eventually wire alice martin's \
edge — deferred is not dropped"
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
#[tokio::test]
async fn dropping_the_server_bounds_shutdown_and_skips_queued_jobs() {
let (_dir, srv, extractor, release) = gated_server();
let service = Arc::clone(&srv.service);
for i in 0..3 {
srv.remember(Parameters(RememberParams {
fact: format!("fait en rafale numero {i}"),
links: Vec::new(),
metadata: None,
ttl_seconds: None,
}))
.await
.expect("remember");
}
assert!(
wait_for(Duration::from_secs(5), || {
extractor.entered.load(Ordering::SeqCst) == 1
}),
"the worker must have dequeued job 0 and be blocked on the gate"
);
let started = Instant::now();
let joiner = std::thread::spawn(move || drop(srv));
assert!(
wait_for(Duration::from_secs(5), || {
service
.inspect(MemoryService::autograph_queue_open)
.is_ok_and(|open| !open)
}),
"dropping the server must close the autograph queue via the stored \
worker handle"
);
release.send(()).expect("release the in-flight job");
joiner.join().expect("join the dropping thread");
assert!(
started.elapsed() < Duration::from_secs(10),
"shutdown is BOUNDED: it waits for at most the ONE in-flight \
generation, never the queue behind it"
);
assert_eq!(
extractor.completed.load(Ordering::SeqCst),
1,
"only the in-flight job is wired on shutdown"
);
assert_eq!(
service
.inspect(MemoryService::autograph_dropped)
.expect("active generation"),
2,
"the two still-queued jobs are SKIPPED and counted, not waited out"
);
}