use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use supercode_harness::configfile::{resolve, ResolveOptions};
use supercode_harness::session_journal::{self, JournalOp, PlanEntry, QueueKind};
use supercode_harness::store::SessionStore;
use supercode_harness::{
Agent, ChatMessage, ChatRequest, Config, FunctionCall, Provider, Role, ToolCall, Usage,
};
fn temp_root(tag: &str) -> std::path::PathBuf {
let p = std::env::temp_dir().join(format!(
"supercode-bp8-{tag}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&p).unwrap();
p
}
fn parity_config(preset: &str, cwd: &std::path::Path) -> Config {
let mut resolved = resolve(
&format!("extends = \"{preset}\"\n"),
None,
&ResolveOptions { strict: true },
)
.unwrap_or_else(|e| panic!("preset `{preset}` failed to resolve: {e}"));
resolved.config.cwd = cwd.to_path_buf();
resolved.config
}
struct Answering(&'static str);
#[async_trait]
impl Provider for Answering {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
Ok((ChatMessage::assistant(self.0), Usage::default()))
}
}
struct Planning {
plan: serde_json::Value,
calls: AtomicUsize,
}
#[async_trait]
impl Provider for Planning {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
let call = ChatMessage {
role: Role::Assistant,
content: None,
content_parts: None,
tool_calls: Some(vec![ToolCall {
id: "call_plan".into(),
kind: "function".into(),
function: FunctionCall {
name: "update_plan".into(),
arguments: self.plan.to_string(),
},
}]),
tool_call_id: None,
name: None,
metadata: Default::default(),
};
return Ok((call, Usage::default()));
}
Ok((ChatMessage::assistant("planned"), Usage::default()))
}
}
struct Hanging {
seen_first_turn: Arc<Mutex<bool>>,
}
#[async_trait]
impl Provider for Hanging {
async fn complete(
&self,
_req: &ChatRequest,
_on_delta: &(dyn for<'a> Fn(&'a str) + Send + Sync),
) -> supercode_harness::Result<(ChatMessage, Usage)> {
*self.seen_first_turn.lock().unwrap() = true;
Ok((ChatMessage::assistant("mid-turn reply"), Usage::default()))
}
}
fn persist(agent: &Agent, store: &SessionStore, name: &str) {
let jsonl: String = agent
.history()
.iter()
.filter_map(|m| serde_json::to_string(m).ok())
.collect::<Vec<_>>()
.join("\n");
store.save(name, "t", &jsonl).unwrap();
session_journal::checkpoint(agent, store, name, agent.history().len());
}
fn reopen(
preset: &str,
cwd: &std::path::Path,
store: &SessionStore,
name: &str,
provider: Box<dyn Provider>,
) -> (Agent, session_journal::RestoreReport) {
let mut agent = Agent::with_provider(parity_config(preset, cwd), provider);
if let Ok(Some(jsonl)) = store.load_if_present(name) {
let tmp = cwd.join(format!("reopen-{name}.jsonl"));
std::fs::write(&tmp, jsonl).unwrap();
agent.load_transcript(&tmp).unwrap();
let _ = std::fs::remove_file(&tmp);
}
let report = session_journal::arm(&mut agent, store, name);
(agent, report)
}
#[tokio::test]
async fn a_turn_is_on_disk_before_the_transcript_is_rewritten_under_both_presets() {
for preset in ["cc-parity", "cx-parity"] {
let root = temp_root("append-only");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-append";
let mut agent = Agent::with_provider(
parity_config(preset, &root),
Box::new(Answering("first answer")),
);
assert!(
agent.config().session_append_only,
"{preset} must set core.session.append_only"
);
let report = session_journal::arm(&mut agent, &store, name);
assert!(report.journal_armed, "{preset}: the journal must be armed");
agent.send("hello").await.unwrap();
assert!(
store.load_if_present(name).unwrap().is_none(),
"{preset}: the transcript must not exist yet"
);
let state = store
.load_journal(name)
.unwrap()
.expect("the journal must exist mid-turn");
let bodies: Vec<String> = state
.messages
.iter()
.map(|m| m.content.clone().unwrap_or_default())
.collect();
assert_eq!(
bodies,
["hello", "first answer"],
"{preset}: every message must already be flushed"
);
}
}
#[tokio::test]
async fn a_turn_lost_to_a_crash_comes_back_on_the_next_open() {
for preset in ["cc-parity", "cx-parity"] {
let root = temp_root("recover");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-recover";
let seen = Arc::new(Mutex::new(false));
let mut agent = Agent::with_provider(
parity_config(preset, &root),
Box::new(Answering("answer one")),
);
session_journal::arm(&mut agent, &store, name);
agent.send("turn one").await.unwrap();
persist(&agent, &store, name);
drop(agent);
let (mut agent, _) = reopen(
preset,
&root,
&store,
name,
Box::new(Hanging {
seen_first_turn: seen.clone(),
}),
);
agent.send("turn two").await.unwrap();
assert!(*seen.lock().unwrap());
let before_crash = agent.history().len();
drop(agent);
let persisted = store.load(name).unwrap();
assert!(!persisted.contains("turn two"));
let (agent, report) = reopen(preset, &root, &store, name, Box::new(Answering("later")));
assert_eq!(
report.recovered_messages, 2,
"{preset}: the user turn and the reply must both come back"
);
assert_eq!(agent.history().len(), before_crash, "{preset}");
assert!(
agent
.history()
.iter()
.any(|m| m.content.as_deref() == Some("turn two")),
"{preset}: the recovered user turn must be in the live conversation"
);
}
}
#[tokio::test]
async fn even_a_brand_new_sessions_first_turn_survives_a_crash() {
let root = temp_root("first-turn");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-first";
let mut agent = Agent::with_provider(
parity_config("cc-parity", &root),
Box::new(Answering("the only answer")),
);
let report = session_journal::arm(&mut agent, &store, name);
assert!(report.journal_armed);
assert!(
!store.journal_path(name).unwrap().exists(),
"the journal must not exist before the first message"
);
agent.send("the only question").await.unwrap();
drop(agent);
assert!(store.load_if_present(name).unwrap().is_none());
let (agent, report) = reopen(
"cc-parity",
&root,
&store,
name,
Box::new(Answering("later")),
);
assert_eq!(report.recovered_messages, 2);
assert!(agent
.history()
.iter()
.any(|m| m.content.as_deref() == Some("the only question")));
}
#[tokio::test]
async fn cc_parity_grows_a_conversation_tree_and_cx_parity_does_not() {
let root = temp_root("tree");
let store = SessionStore::open(root.join("sessions")).unwrap();
let mut agent =
Agent::with_provider(parity_config("cc-parity", &root), Box::new(Answering("a1")));
session_journal::arm(&mut agent, &store, "sess-tree");
agent.send("u1").await.unwrap();
agent.send("u2").await.unwrap();
let tree = agent
.session_tree()
.expect("cc-parity arms the tree module");
assert_eq!(tree.nodes.len(), 4, "one node per recorded message");
let projected: Vec<String> = tree
.linear_projection()
.unwrap()
.iter()
.map(|m| m.content.clone().unwrap_or_default())
.collect();
assert_eq!(projected, ["u1", "a1", "u2", "a1"]);
persist(&agent, &store, "sess-tree");
drop(agent);
let (agent, report) = reopen(
"cc-parity",
&root,
&store,
"sess-tree",
Box::new(Answering("a2")),
);
assert!(report.tree_loaded, "the tree must come back off disk");
assert_eq!(agent.session_tree().unwrap().nodes.len(), 4);
let mut cx = Agent::with_provider(parity_config("cx-parity", &root), Box::new(Answering("a1")));
session_journal::arm(&mut cx, &store, "sess-tree-cx");
cx.send("u1").await.unwrap();
assert!(
cx.session_tree().is_none(),
"cx-parity leaves capabilities.session_tree off"
);
}
#[tokio::test]
async fn rewinding_to_an_arbitrary_point_is_recorded_and_invertible() {
for preset in ["cc-parity", "cx-parity"] {
let root = temp_root("rewind");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-rewind";
let mut agent =
Agent::with_provider(parity_config(preset, &root), Box::new(Answering("a")));
session_journal::arm(&mut agent, &store, name);
for turn in ["u1", "u2", "u3"] {
agent.send(turn).await.unwrap();
}
assert_eq!(agent.history().len(), 7, "{preset}: system + 3 exchanges");
let outcome = agent.rewind_conversation(3);
assert_eq!(outcome.removed, 4, "{preset}");
assert_eq!(
agent
.history()
.iter()
.filter_map(|m| m.content.clone())
.filter(|c| c.starts_with('u'))
.collect::<Vec<_>>(),
["u1"],
"{preset}"
);
if preset == "cc-parity" {
let branch = outcome
.preserved_branch
.as_deref()
.expect("cc-parity preserves the old leaf under a sibling branch");
let kept: Vec<String> = agent
.session_tree()
.unwrap()
.linear_projection_of(branch)
.unwrap()
.iter()
.map(|m| m.content.clone().unwrap_or_default())
.collect();
assert_eq!(kept, ["u1", "a", "u2", "a", "u3", "a"]);
} else {
assert!(outcome.preserved_branch.is_none());
}
let raw = std::fs::read_to_string(store.journal_path(name).unwrap()).unwrap();
assert!(raw.contains("\"op\":\"rewind\""), "{preset}");
assert!(raw.contains("\"content\":\"u3\""), "{preset}");
assert!(agent.undo_rewind(), "{preset}");
assert_eq!(agent.history().len(), 7, "{preset}");
assert!(std::fs::read_to_string(store.journal_path(name).unwrap())
.unwrap()
.contains("\"op\":\"unrewind\""));
agent.rewind_conversation(3);
persist(&agent, &store, name);
drop(agent);
let (mut agent, report) = reopen(preset, &root, &store, name, Box::new(Answering("a")));
assert_eq!(report.restored_rewinds, 1, "{preset}");
assert!(agent.undo_rewind(), "{preset}: the undo crossed a restart");
assert_eq!(agent.history().len(), 7, "{preset}");
}
}
#[tokio::test]
async fn rename_moves_the_resume_handle_and_the_whole_family() {
let root = temp_root("rename");
let store = SessionStore::open(root.join("sessions")).unwrap();
let mut agent =
Agent::with_provider(parity_config("cc-parity", &root), Box::new(Answering("a1")));
session_journal::arm(&mut agent, &store, "old-handle");
agent.send("hello").await.unwrap();
persist(&agent, &store, "old-handle");
drop(agent);
let before: Vec<String> = std::fs::read_dir(store.root())
.unwrap()
.flatten()
.map(|e| e.file_name().to_string_lossy().into_owned())
.filter(|n| n.starts_with("old-handle."))
.collect();
assert!(
before.len() >= 3,
"the family should have transcript + meta + journal at least: {before:?}"
);
store.rename("old-handle", "new-handle").unwrap();
assert!(store.load_if_present("old-handle").unwrap().is_none());
assert!(store.load("new-handle").unwrap().contains("hello"));
for member in &before {
let moved = member.replacen("old-handle.", "new-handle.", 1);
assert!(
store.root().join(&moved).exists(),
"family member `{member}` did not follow the rename"
);
assert!(!store.root().join(member).exists());
}
assert!(store.list().iter().any(|s| s.name == "new-handle"));
assert!(!store.list().iter().any(|s| s.name == "old-handle"));
let (agent, _) = reopen(
"cc-parity",
&root,
&store,
"new-handle",
Box::new(Answering("a2")),
);
assert!(agent
.history()
.iter()
.any(|m| m.content.as_deref() == Some("hello")));
store.save("taken", "t", "").unwrap();
assert!(store.rename("new-handle", "taken").is_err());
assert!(store.load("new-handle").unwrap().contains("hello"));
}
#[test]
fn the_derived_index_survives_the_process_and_rebuilds_identically() {
let root = temp_root("index");
let store = SessionStore::open(root.join("sessions")).unwrap();
for (name, body) in [
(
"sess-a",
"{\"role\":\"user\",\"content\":\"first question\"}",
),
(
"sess-b",
"{\"role\":\"user\",\"content\":\"second question\"}",
),
] {
store.save(name, name, body).unwrap();
}
let cold = store.index();
assert_eq!(cold.rederived, 2, "a cold index derives every row");
assert_eq!(cold.reused, 0);
assert!(store.index_path().exists(), "the cache must be written");
assert!(cold
.entries
.iter()
.any(|e| e.preview == "first question" && e.messages == 1));
let other = SessionStore::at(store.root());
let warm = other.index();
assert_eq!(warm.reused, 2, "every unchanged row comes from the cache");
assert_eq!(warm.rederived, 0);
assert_eq!(warm.entries, cold.entries);
std::thread::sleep(std::time::Duration::from_millis(5));
store
.save(
"sess-a",
"sess-a",
"{\"role\":\"user\",\"content\":\"first question\"}\n{\"role\":\"assistant\",\"content\":\"ok\"}",
)
.unwrap();
let repaired = store.index();
assert_eq!(repaired.rederived, 1, "only the moved row is re-read");
assert_eq!(repaired.reused, 1);
assert_eq!(
repaired
.entries
.iter()
.find(|e| e.name == "sess-a")
.unwrap()
.messages,
2
);
let with_cache = store.index().entries;
store.invalidate_index();
assert!(!store.index_path().exists());
let rebuilt = store.index();
assert_eq!(rebuilt.rederived, 2);
assert_eq!(rebuilt.entries, with_cache);
}
#[test]
fn an_old_generation_is_upgraded_in_place_once_with_the_original_recoverable() {
let root = temp_root("upgrade");
let store = SessionStore::open(root.join("sessions")).unwrap();
let legacy = concat!(
"{\"role\":\"user\",\"content\":\"hi\",\"tool_calls\":null,\"timestamp\":\"2026-01-01\"}\n",
"{\"role\":\"assistant\",\"content\":\"hello\",\"legacy_id\":7}\n"
);
store.save("sess-old", "t", legacy).unwrap();
assert_eq!(store.format_version("sess-old"), 0);
let before: Vec<String> = store
.load("sess-old")
.unwrap()
.lines()
.map(|l| serde_json::from_str::<ChatMessage>(l).unwrap())
.map(|m| m.content.unwrap_or_default())
.collect();
let up = store
.upgrade_in_place("sess-old")
.unwrap()
.expect("an unmarked session must be upgraded");
assert_eq!(up.from_version, 0);
assert_eq!(up.to_version, SessionStore::FORMAT_VERSION);
assert!(up.rewritten);
assert_eq!(up.messages, 2);
assert_eq!(
store.format_version("sess-old"),
SessionStore::FORMAT_VERSION
);
assert!(store.upgrade_in_place("sess-old").unwrap().is_none());
let after = store.load("sess-old").unwrap();
let after_msgs: Vec<String> = after
.lines()
.map(|l| serde_json::from_str::<ChatMessage>(l).unwrap())
.map(|m| m.content.unwrap_or_default())
.collect();
assert_eq!(before, after_msgs);
assert!(!after.contains("legacy_id"));
assert!(!after.contains("\"tool_calls\":null"));
let original = std::fs::read_to_string(store.root().join(&up.original)).unwrap();
assert_eq!(original, legacy);
}
#[tokio::test]
async fn a_pending_input_survives_a_restart_under_cc_parity_only() {
let root = temp_root("queue");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-queue";
let mut agent =
Agent::with_provider(parity_config("cc-parity", &root), Box::new(Answering("a1")));
assert!(agent.config().session_queue_persist);
session_journal::arm(&mut agent, &store, name);
agent.send("first").await.unwrap();
agent.queue_follow_up("the queued follow-up");
persist(&agent, &store, name);
drop(agent);
let (mut agent, report) = reopen("cc-parity", &root, &store, name, Box::new(Answering("a2")));
assert_eq!(report.restored_queue, 1, "the pending input must come back");
assert_eq!(agent.queued_steer_count(), 1);
agent.send("second").await.unwrap();
assert!(agent
.history()
.iter()
.any(|m| m.content.as_deref() == Some("the queued follow-up")));
persist(&agent, &store, name);
drop(agent);
let (_, report) = reopen("cc-parity", &root, &store, name, Box::new(Answering("a3")));
assert_eq!(
report.restored_queue, 0,
"a consumed input must not be re-delivered"
);
let cx_root = temp_root("queue-cx");
let cx_store = SessionStore::open(cx_root.join("sessions")).unwrap();
let mut cx = Agent::with_provider(
parity_config("cx-parity", &cx_root),
Box::new(Answering("a1")),
);
assert!(!cx.config().session_queue_persist);
session_journal::arm(&mut cx, &cx_store, "sess-cx");
cx.queue_follow_up("not recorded");
let raw =
std::fs::read_to_string(cx_store.journal_path("sess-cx").unwrap()).unwrap_or_default();
assert!(!raw.contains("enqueue"), "cx-parity records no queue ops");
}
#[tokio::test]
async fn a_model_written_plan_survives_the_process_under_both_presets() {
for preset in ["cc-parity", "cx-parity"] {
let root = temp_root("plan");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-plan";
let mut agent = Agent::with_provider(
parity_config(preset, &root),
Box::new(Planning {
plan: serde_json::json!({
"plan": [
{"step": "read the ledger", "status": "completed"},
{"step": "write the test", "status": "in_progress"},
]
}),
calls: AtomicUsize::new(0),
}),
);
assert!(
agent.config().todos_persist,
"{preset} must set capabilities.todos.persist"
);
session_journal::arm(&mut agent, &store, name);
let reply = agent.send("plan the work").await.unwrap();
assert_eq!(reply, "planned", "{preset}: the tool round-trip must run");
assert_eq!(
agent.plan(),
vec![
PlanEntry {
step: "read the ledger".into(),
status: "completed".into()
},
PlanEntry {
step: "write the test".into(),
status: "in_progress".into()
},
],
"{preset}: the loop must see the plan the tool wrote"
);
persist(&agent, &store, name);
drop(agent);
assert_eq!(store.load_plan(name).unwrap().unwrap().len(), 2, "{preset}");
let (agent, report) = reopen(preset, &root, &store, name, Box::new(Answering("later")));
assert_eq!(report.restored_plan, 2, "{preset}");
assert_eq!(agent.plan()[1].status, "in_progress", "{preset}");
}
}
#[tokio::test]
async fn the_plan_is_journaled_before_the_turn_ends() {
let root = temp_root("plan-journal");
let store = SessionStore::open(root.join("sessions")).unwrap();
let mut agent = Agent::with_provider(
parity_config("cc-parity", &root),
Box::new(Planning {
plan: serde_json::json!({"plan": [{"step": "only step", "status": "pending"}]}),
calls: AtomicUsize::new(0),
}),
);
session_journal::arm(&mut agent, &store, "sess-pj");
agent.send("go").await.unwrap();
let state = store.load_journal("sess-pj").unwrap().unwrap();
assert_eq!(
state.plan_pairs(),
vec![("only step".to_string(), "pending".to_string())]
);
}
#[tokio::test]
async fn no_operation_ever_shortens_the_journal() {
let root = temp_root("append-invariant");
let store = SessionStore::open(root.join("sessions")).unwrap();
let name = "sess-inv";
let path = store.journal_path(name).unwrap();
let mut agent =
Agent::with_provider(parity_config("cc-parity", &root), Box::new(Answering("a")));
session_journal::arm(&mut agent, &store, name);
let mut sizes = vec![std::fs::metadata(&path).map(|m| m.len()).unwrap_or(0)];
for turn in ["u1", "u2"] {
agent.send(turn).await.unwrap();
sizes.push(std::fs::metadata(&path).unwrap().len());
}
let after_turns = std::fs::read_to_string(&path).unwrap();
agent.rewind_conversation(1);
sizes.push(std::fs::metadata(&path).unwrap().len());
agent.undo_rewind();
sizes.push(std::fs::metadata(&path).unwrap().len());
persist(&agent, &store, name);
sizes.push(std::fs::metadata(&path).unwrap().len());
for pair in sizes.windows(2) {
assert!(pair[1] >= pair[0], "the journal shrank: {sizes:?}");
}
let final_bytes = std::fs::read_to_string(&path).unwrap();
assert!(
final_bytes.starts_with(&after_turns),
"earlier bytes must be untouched by later operations"
);
}
#[test]
fn every_record_carries_the_discriminant() {
let dir = temp_root("discriminant");
let path = dir.join("j.journal.jsonl");
let mut journal = session_journal::SessionJournal::open_append(&path).unwrap();
journal.append_message(&ChatMessage::user("hi")).unwrap();
journal
.append(JournalOp::Enqueue {
queue: QueueKind::Steer,
text: "q".into(),
})
.unwrap();
journal.append(JournalOp::Rewind { to: 0 }).unwrap();
journal
.append(JournalOp::Checkpoint { messages: 1 })
.unwrap();
for line in std::fs::read_to_string(&path).unwrap().lines() {
let value: serde_json::Value = serde_json::from_str(line).unwrap();
assert_eq!(value["supercode_journal"], 1, "{line}");
}
}