use std::io::Write;
use std::process::{Child, Command, Stdio};
mod support;
struct Client {
child: Child,
responses: support::BoundedLineReader,
}
impl Client {
fn spawn(dir: &std::path::Path) -> Self {
let mut cmd = Command::new(env!("CARGO_BIN_EXE_exocortex-mcp-client"));
cmd.args([
"--org",
"readback",
"--user",
"tester",
"--data-dir",
dir.to_str().unwrap(),
]);
let mut child = cmd
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null())
.spawn()
.expect("spawn exocortex-mcp-client");
let responses = support::BoundedLineReader::new(child.stdout.take().expect("stdout"));
Self { child, responses }
}
fn send_all(&mut self, msgs: &[serde_json::Value]) {
let stdin = self.child.stdin.as_mut().expect("stdin");
for m in msgs {
writeln!(stdin, "{m}").unwrap();
}
stdin.flush().unwrap();
}
fn read_line(&mut self) -> serde_json::Value {
self.responses.read_json(&mut self.child)
}
fn init_msgs() -> Vec<serde_json::Value> {
vec![
serde_json::json!({
"jsonrpc": "2.0", "id": 1, "method": "initialize",
"params": {
"protocolVersion": "2024-11-05",
"capabilities": {},
"clientInfo": { "name": "readback-test", "version": "0" }
}
}),
serde_json::json!({
"jsonrpc": "2.0", "method": "notifications/initialized", "params": {}
}),
]
}
fn search(query: &str) -> serde_json::Value {
serde_json::json!({
"jsonrpc": "2.0", "id": 90, "method": "tools/call",
"params": {
"name": "exocortex.search_memories",
"arguments": { "query": query, "limit": 50 }
}
})
}
fn end_session(memories: serde_json::Value, edges: serde_json::Value) -> serde_json::Value {
serde_json::json!({
"jsonrpc": "2.0", "id": 80, "method": "tools/call",
"params": {
"name": "exocortex.end_session",
"arguments": {
"session_id": "s-readback",
"project_id": "proj",
"memories": memories,
"edges": edges
}
}
})
}
}
impl Drop for Client {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
fn tempdir() -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!(
"exocortex-readback-{}-{}",
std::process::id(),
uuid_v4_hex()
));
std::fs::create_dir_all(&dir).unwrap();
dir
}
fn uuid_v4_hex() -> String {
uuid::Uuid::new_v4().simple().to_string()
}
fn one_memory(draft_key: &str, type_label: &str, title: &str) -> serde_json::Value {
serde_json::json!({
"draft_key": draft_key,
"memory_type": type_label,
"title": title,
"content": "body",
"visibility": "org"
})
}
fn search_hits(payload: &serde_json::Value) -> Vec<serde_json::Value> {
let text = payload["result"]["content"][0]["text"]
.as_str()
.unwrap_or("");
let inner: serde_json::Value = serde_json::from_str(text).expect("tool payload JSON");
inner["memories"].as_array().cloned().unwrap_or_default()
}
#[test]
fn in_session_read_back_after_offline_write() {
let dir = tempdir();
let mut c = Client::spawn(&dir);
let mut msgs = Client::init_msgs();
msgs.push(Client::end_session(
serde_json::json!([
{ "draft_key": "p", "memory_type": "Problem", "title": "Zebra race condition in cache swap", "content": "swap races under load", "visibility": "org", "tags": ["concurrency"] }
]),
serde_json::json!([]),
));
msgs.push(Client::search("zebra"));
c.send_all(&msgs);
let _init = c.read_line();
let ack = c.read_line();
assert!(ack.get("result").is_some(), "end_session ok: {ack}");
let ack_text = ack["result"]["content"][0]["text"].as_str().unwrap();
let ack_json: serde_json::Value = serde_json::from_str(ack_text).unwrap();
let acked_lsn = ack_json["local_lsns"][0].as_u64().expect("acked local lsn");
let search = c.read_line();
let hits = search_hits(&search);
assert_eq!(hits.len(), 1, "exactly the written memory: {hits:?}");
assert!(hits[0]["title"]
.as_str()
.unwrap()
.contains("Zebra race condition"));
let search_text = search["result"]["content"][0]["text"].as_str().unwrap();
let search_inner: serde_json::Value = serde_json::from_str(search_text).unwrap();
let local_lsn = search_inner["snapshot_version"]["local_lsn"].as_u64();
assert!(
local_lsn.is_some_and(|n| n >= acked_lsn),
"R-M7 stamp reflects the offline write (got {local_lsn:?}, acked {acked_lsn})"
);
}
#[test]
fn exact_offline_retry_reuses_the_original_wal_result_and_memory_ids() {
let dir = tempdir();
{
let mut client = Client::spawn(&dir);
let request = Client::end_session(
serde_json::json!([
{ "draft_key": "retry", "memory_type": "Problem", "title": "Retry-safe offline row", "content": "same request after response loss", "visibility": "org", "tags": [] }
]),
serde_json::json!([]),
);
let mut messages = Client::init_msgs();
messages.push(request.clone());
messages.push(request);
messages.push(Client::search("retry-safe"));
client.send_all(&messages);
let _init = client.read_line();
let first = client.read_line();
let retry = client.read_line();
assert_eq!(
first["result"]["content"][0]["text"], retry["result"]["content"][0]["text"],
"response-loss retry returns the original local LSN"
);
let hits = search_hits(&client.read_line());
assert_eq!(hits.len(), 1, "retry publishes one semantic memory");
}
let wal = exocortex_client::wal::Wal::open(&dir.join("wal")).unwrap();
assert_eq!(wal.db_len(), 1, "retry appends no second durable row");
}
#[test]
fn cross_restart_read_back_with_stable_ids() {
let dir = tempdir();
let id = {
let mut first = Client::spawn(&dir);
let mut msgs = Client::init_msgs();
msgs.push(Client::end_session(
serde_json::json!([one_memory("d", "CodePattern", "Adopt blake3 for ids")]),
serde_json::json!([]),
));
msgs.push(Client::search("blake3"));
first.send_all(&msgs);
let _init = first.read_line();
let ack = first.read_line();
assert!(ack.get("result").is_some(), "write ok: {ack}");
let search = first.read_line();
let hits = search_hits(&search);
assert_eq!(hits.len(), 1, "in-session hit (F2): {hits:?}");
hits[0]["id"].as_str().unwrap().to_string()
};
let mut second = Client::spawn(&dir);
let mut msgs = Client::init_msgs();
msgs.push(Client::search("blake3"));
second.send_all(&msgs);
let _init = second.read_line();
let search = second.read_line();
let hits = search_hits(&search);
assert_eq!(hits.len(), 1, "write survives restart: {hits:?}");
assert_eq!(
hits[0]["id"].as_str().unwrap(),
id,
"WAL-stored id byte-stable across restart"
);
}
#[test]
fn boot_seeds_pending_synced_and_failed_entries() {
let dir = tempdir();
{
let mut first = Client::spawn(&dir);
let mut msgs = Client::init_msgs();
for title in [
"Palladium purge panics",
"Silver sync stalls",
"Flax failure flakes",
] {
msgs.push(Client::end_session(
serde_json::json!([one_memory("d", "Problem", title)]),
serde_json::json!([]),
));
}
first.send_all(&msgs);
let _init = first.read_line();
for i in 0..3 {
let ack = first.read_line();
assert!(ack.get("result").is_some(), "write {i} ok: {ack}");
}
}
{
let wal = exocortex_client::wal::Wal::open(&dir.join("wal")).unwrap();
let states = wal.states_for_test().unwrap();
assert_eq!(states.len(), 3);
wal.mark_synced(states[0].0, 500).unwrap();
wal.mark_failed(states[1].0).unwrap();
drop(wal);
}
let mut second = Client::spawn(&dir);
let mut msgs = Client::init_msgs();
for q in ["palladium", "silver", "flax"] {
msgs.push(Client::search(q));
}
second.send_all(&msgs);
let _init = second.read_line();
for (i, title_part) in ["palladium", "silver", "flax"].into_iter().enumerate() {
let search = second.read_line();
let hits = search_hits(&search);
assert_eq!(
hits.len(),
1,
"search {i} (`{title_part}`) must hit its write in ANY wal state: {hits:?}"
);
}
}
#[test]
fn edges_and_grouping_readable_dangling_edge_harmless() {
let dir = tempdir();
let mut c = Client::spawn(&dir);
let mut msgs = Client::init_msgs();
msgs.push(Client::end_session(
serde_json::json!([
one_memory("prob", "Problem", "Oak overflow in parser"),
one_memory("fix", "Fix", "Raise the oak ceiling"),
]),
serde_json::json!([
{ "from_draft_key": "fix", "to_draft_key": "prob", "kind": "Fixes", "strength": 0.9 },
{ "from_draft_key": "fix", "to_memory_id": "0123456789abcdef0123456789abcdef", "kind": "Fixes" }
]),
));
msgs.push(Client::search("ceiling"));
c.send_all(&msgs);
let _init = c.read_line();
let ack = c.read_line();
assert!(
ack.get("result").is_some(),
"dangling edge must not fail the write: {ack}"
);
let search = c.read_line();
let hits = search_hits(&search);
assert_eq!(hits.len(), 1, "the fix is searchable: {hits:?}");
let fix_id = hits[0]["id"].as_str().unwrap().to_string();
let mut q = Vec::new();
q.push(serde_json::json!({
"jsonrpc": "2.0", "id": 70, "method": "tools/call",
"params": {
"name": "exocortex.find_related",
"arguments": { "anchor": fix_id, "k": 1 }
}
}));
c.send_all(&q);
let rel = c.read_line();
assert!(
rel.get("result").is_some(),
"dangling edge must not fail reads: {rel}"
);
let text = rel["result"]["content"][0]["text"].as_str().unwrap();
let inner: serde_json::Value = serde_json::from_str(text).unwrap();
let titles: Vec<&str> = inner["memories"]
.as_array()
.unwrap()
.iter()
.map(|m| m["title"].as_str().unwrap())
.collect();
assert!(
titles.contains(&"Oak overflow in parser"),
"Fixes edge traversable: {titles:?}"
);
assert!(
titles.iter().any(|t| t.starts_with("Session ")),
"Conversation node reachable via InSession: {titles:?}"
);
}
#[test]
fn grouping_parity_with_backend_commit_path() {
let _ = std::hint::black_box(exocortex_pack_dev_v1::pack_def().name.clone());
let ontology = exocortex_kernel::pack::load_registered_packs().unwrap();
let now = chrono::Utc::now();
let org = "parity-org";
let key = "session-key-1";
let rule = &exocortex_ingest::grouping::grouping_rules()[0];
assert_eq!(rule.flavor, "session", "the rule under parity");
assert_eq!(rule.node_type, "Conversation");
assert_eq!(rule.edge_kind, "InSession");
let backend_node = exocortex_ingest::grouping::grouping_node(&ontology, org, rule, key, now)
.expect("backend node");
let local_node = exocortex_client::materialize::grouping_node_local(&ontology, org, key, now)
.expect("local node");
assert_eq!(
serde_json::to_value(&backend_node).unwrap(),
serde_json::to_value(&local_node).unwrap(),
"grouping node identical to the backend mint"
);
let member = exocortex_client::materialize::grouping_node_local(&ontology, org, "other", now)
.expect("member");
let backend_edge =
exocortex_ingest::grouping::grouping_edge(&ontology, rule, &member, &backend_node, now)
.expect("backend edge");
let local_edge =
exocortex_client::materialize::grouping_edge_local(&ontology, &member, &local_node, now)
.expect("local edge");
assert_eq!(
serde_json::to_value(&backend_edge).unwrap(),
serde_json::to_value(&local_edge).unwrap(),
"InSession edge identical to the backend mint (same derivation, same id)"
);
}