use crate::fold::{FoldedRecord, SyncState};
use crate::oplog::Hlc;
use car_inference_types::{Message, ToolCall};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
pub const DEFAULT_CONVERSATION: &str = "";
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Role {
User,
Assistant,
Tool,
}
impl Role {
fn parse(s: &str) -> Role {
match s {
"assistant" => Role::Assistant,
"tool" | "tool_result" => Role::Tool,
_ => Role::User,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Turn {
pub conversation_id: String,
pub role: Role,
pub content: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub tool_calls: Vec<ToolCall>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tool_use_id: Option<String>,
pub timestamp: u64,
pub hlc: Hlc,
pub op_id: String,
}
impl Turn {
pub fn user_payload(conversation_id: &str, content: &str, timestamp: u64) -> Value {
json!({
"conversation_id": conversation_id,
"role": "user",
"content": content,
"timestamp": timestamp,
})
}
pub fn assistant_payload(
conversation_id: &str,
content: &str,
tool_calls: Vec<Value>,
timestamp: u64,
) -> Value {
json!({
"conversation_id": conversation_id,
"role": "assistant",
"content": content,
"tool_calls": tool_calls,
"timestamp": timestamp,
})
}
pub fn tool_payload(
conversation_id: &str,
tool_use_id: &str,
content: &str,
timestamp: u64,
) -> Value {
json!({
"conversation_id": conversation_id,
"role": "tool",
"tool_use_id": tool_use_id,
"content": content,
"timestamp": timestamp,
})
}
pub fn from_record(record: &FoldedRecord) -> Option<Turn> {
if crate::compact::is_tombstone(record) {
return None;
}
let p = &record.payload;
let role = p
.get("role")
.or_else(|| p.get("speaker"))
.and_then(Value::as_str)
.map(Role::parse)
.unwrap_or(Role::User);
let content = p
.get("content")
.or_else(|| p.get("text"))
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
let conversation_id = p
.get("conversation_id")
.and_then(Value::as_str)
.unwrap_or(DEFAULT_CONVERSATION)
.to_string();
let tool_calls = p
.get("tool_calls")
.and_then(Value::as_array)
.map(|arr| {
arr.iter()
.filter_map(|v| serde_json::from_value::<ToolCall>(v.clone()).ok())
.collect()
})
.unwrap_or_default();
let tool_use_id = p
.get("tool_use_id")
.and_then(Value::as_str)
.map(str::to_string);
let timestamp = p
.get("timestamp")
.and_then(|v| v.as_u64().or_else(|| v.as_f64().map(|f| f as u64)))
.unwrap_or(0);
Some(Turn {
conversation_id,
role,
content,
tool_calls,
tool_use_id,
timestamp,
hlc: record.hlc.clone(),
op_id: record.op_id.clone(),
})
}
pub fn to_message(&self) -> Message {
match self.role {
Role::User => Message::User { content: self.content.clone() },
Role::Assistant => Message::Assistant {
content: self.content.clone(),
tool_calls: self.tool_calls.clone(),
},
Role::Tool => Message::ToolResult {
tool_use_id: self.tool_use_id.clone().unwrap_or_default(),
content: self.content.clone(),
},
}
}
}
fn join_content(a: &str, b: &str) -> String {
match (a.is_empty(), b.is_empty()) {
(true, _) => b.to_string(),
(_, true) => a.to_string(),
_ => format!("{a}\n\n{b}"),
}
}
pub fn repair(turns: Vec<Turn>) -> Vec<Message> {
let mut out: Vec<Turn> = Vec::new();
let mut tool_open = false;
for t in turns {
match t.role {
Role::Tool => {
if tool_open {
out.push(t); }
}
Role::User => {
if let Some(last) = out.last_mut() {
if last.role == Role::User {
last.content = join_content(&last.content, &t.content);
continue;
}
}
tool_open = false;
out.push(t);
}
Role::Assistant => {
if let Some(last) = out.last_mut() {
if last.role == Role::Assistant {
last.content = join_content(&last.content, &t.content);
last.tool_calls.extend(t.tool_calls);
tool_open = !last.tool_calls.is_empty();
continue;
}
}
tool_open = !t.tool_calls.is_empty();
out.push(t);
}
}
}
for i in 0..out.len() {
if out[i].role == Role::Assistant && !out[i].tool_calls.is_empty() {
let answered = out.get(i + 1).is_some_and(|n| n.role == Role::Tool);
if !answered {
out[i].tool_calls.clear();
}
}
}
out.iter().map(Turn::to_message).collect()
}
impl SyncState {
pub fn transcript(&self, conversation_id: &str) -> Vec<Turn> {
self.log_entries(&crate::oplog::Surface::Conversation.tag())
.into_iter()
.filter_map(Turn::from_record)
.filter(|t| t.conversation_id == conversation_id)
.collect()
}
pub fn conversation_ids(&self) -> Vec<String> {
let mut ids: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
for rec in self.log_entries(&crate::oplog::Surface::Conversation.tag()) {
if let Some(turn) = Turn::from_record(rec) {
ids.insert(turn.conversation_id);
}
}
ids.into_iter().collect()
}
pub fn resume_messages(&self, conversation_id: &str) -> Vec<Message> {
repair(self.transcript(conversation_id))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::compact::{apply_retention, plan_compaction, AckTable, RetentionPolicy, RetentionRule};
use crate::fold::{fold, fold_onto, state_hash};
use crate::oplog::{DeviceLog, Scope, Surface};
fn assert_provider_valid(ms: &[Message]) {
for (i, w) in ms.windows(2).enumerate() {
let dup = matches!(
(&w[0], &w[1]),
(Message::Assistant { .. }, Message::Assistant { .. })
| (Message::User { .. }, Message::User { .. })
);
assert!(!dup, "invalid adjacency at {i}: {:?} then {:?}", w[0], w[1]);
}
for (i, m) in ms.iter().enumerate() {
if matches!(m, Message::ToolResult { .. }) {
let opens_here = i > 0
&& (matches!(&ms[i - 1], Message::Assistant { tool_calls, .. } if !tool_calls.is_empty())
|| matches!(&ms[i - 1], Message::ToolResult { .. }));
assert!(opens_here, "orphan tool_result at index {i}");
}
if let Message::Assistant { tool_calls, .. } = m {
if !tool_calls.is_empty() {
let answered = ms.get(i + 1).is_some_and(|n| matches!(n, Message::ToolResult { .. }));
assert!(answered, "dangling assistant tool_call at index {i}");
}
}
}
}
fn one_device_conversation(conv: &str) -> Vec<crate::oplog::OpRecord> {
let mut d = DeviceLog::new("dev-a");
vec![
d.append(Scope::Personal, Surface::Conversation, Turn::user_payload(conv, "what's the weather?", 10)),
d.append(
Scope::Personal,
Surface::Conversation,
Turn::assistant_payload(
conv,
"",
vec![json!({"id": "call_0", "name": "get_weather", "arguments": {"city": "SF"}})],
11,
),
),
d.append(Scope::Personal, Surface::Conversation, Turn::tool_payload(conv, "call_0", "sunny, 72F", 12)),
d.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload(conv, "It's sunny and 72F in SF.", vec![], 13)),
]
}
#[test]
fn transcript_folds_in_causal_order_regardless_of_delivery() {
let ops = one_device_conversation("c1");
let expected = ["what's the weather?", "", "sunny, 72F", "It's sunny and 72F in SF."];
for delivery in [ops.clone(), ops.iter().rev().cloned().collect::<Vec<_>>()] {
let turns = fold(&delivery).transcript("c1");
let texts: Vec<&str> = turns.iter().map(|t| t.content.as_str()).collect();
assert_eq!(texts, expected);
assert_eq!(
turns.iter().map(|t| t.role).collect::<Vec<_>>(),
vec![Role::User, Role::Assistant, Role::Tool, Role::Assistant]
);
}
}
#[test]
fn resume_produces_a_valid_message_sequence_of_the_real_type() {
let ops = one_device_conversation("c1");
let messages: Vec<Message> = fold(&ops).resume_messages("c1");
assert_eq!(messages.len(), 4);
assert_provider_valid(&messages);
match &messages[0] {
Message::User { content } => assert_eq!(content, "what's the weather?"),
other => panic!("turn 0 should be a user message, got {other:?}"),
}
match &messages[1] {
Message::Assistant { tool_calls, .. } => {
assert_eq!(tool_calls.len(), 1, "assistant tool_calls round-trip as the real ToolCall");
assert_eq!(tool_calls[0].name, "get_weather");
}
other => panic!("turn 1 should be an assistant tool call, got {other:?}"),
}
match &messages[2] {
Message::ToolResult { tool_use_id, content } => {
assert_eq!(tool_use_id, "call_0");
assert_eq!(content, "sunny, 72F");
}
other => panic!("turn 2 should be a tool_result, got {other:?}"),
}
match &messages[3] {
Message::Assistant { content, tool_calls } => {
assert_eq!(content, "It's sunny and 72F in SF.");
assert!(tool_calls.is_empty());
}
other => panic!("turn 3 should be a plain assistant reply, got {other:?}"),
}
}
#[test]
fn crit2_two_genuine_same_payload_turns_do_not_collapse() {
let mut d = DeviceLog::new("dev-a");
let yes = || Turn::user_payload("c1", "yes", 5); let o1 = d.append(Scope::Personal, Surface::Conversation, yes());
let o2 = d.append(Scope::Personal, Surface::Conversation, yes());
assert_ne!(o1.op_id, o2.op_id, "distinct ops (different seq/hlc) → distinct op_id");
assert_eq!(fold(&[o1.clone(), o2]).transcript("c1").len(), 2, "both genuine turns survive (CRIT-2 fixed)");
assert_eq!(fold(&[o1.clone(), o1]).transcript("c1").len(), 1, "resent op dedups on op_id");
}
#[test]
fn same_payload_turns_on_two_devices_are_distinct_events() {
let mut a = DeviceLog::new("dev-a");
let mut b = DeviceLog::new("dev-b");
let payload = Turn::user_payload("c1", "hello", 5);
let oa = a.append(Scope::Personal, Surface::Conversation, payload.clone());
let ob = b.append(Scope::Personal, Surface::Conversation, payload);
assert_eq!(fold(&[oa, ob]).transcript("c1").len(), 2);
}
#[test]
fn crit1_concurrent_assistant_replies_repair_to_a_valid_sequence() {
let mut a = DeviceLog::new("dev-a");
let mut b = DeviceLog::new("dev-b");
let u = a.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "hi", 1));
b.observe(&u.hlc);
let ra = a.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "reply A", vec![], 2));
let rb = b.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "reply B", vec![], 2));
let ops = vec![u, ra, rb];
let raw = fold(&ops).transcript("c1");
assert_eq!(raw.len(), 3);
assert_eq!((raw[1].role, raw[2].role), (Role::Assistant, Role::Assistant));
let baseline = fold(&ops).resume_messages("c1");
assert_eq!(baseline.len(), 2, "the two concurrent replies coalesce");
assert_provider_valid(&baseline);
assert!(matches!(&baseline[0], Message::User { .. }));
assert!(matches!(&baseline[1], Message::Assistant { content, .. } if content.contains("reply A") && content.contains("reply B")));
for perm in permutations(&ops) {
assert_eq!(fold(&perm).resume_messages("c1"), baseline, "repair is delivery-order-independent");
}
}
#[test]
fn adjacent_user_turns_coalesce_before_a_reply() {
let mut a = DeviceLog::new("dev-a");
let mut b = DeviceLog::new("dev-b");
let ua = a.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "userA", 1));
b.observe(&ua.hlc);
let ub = b.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "userB", 2));
a.observe(&ub.hlc);
let ra = a.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "reply", vec![], 3));
let ops = vec![ua, ub, ra];
assert_eq!(fold(&ops).transcript("c1").iter().map(|t| t.role).collect::<Vec<_>>(),
vec![Role::User, Role::User, Role::Assistant], "raw has the invalid user/user adjacency");
let ms = fold(&ops).resume_messages("c1");
assert_eq!(ms.len(), 2, "the two users coalesce");
assert_provider_valid(&ms);
assert!(matches!(&ms[0], Message::User { content } if content.contains("userA") && content.contains("userB")));
assert!(matches!(&ms[1], Message::Assistant { .. }));
}
#[test]
fn crit3_lastn_orphan_tool_result_is_dropped_on_resume() {
let ops = one_device_conversation("c1"); let mut policy = RetentionPolicy::keep_all();
policy.rules.insert("conversation".to_string(), RetentionRule::LastN { n: 2 });
let (retained, _) = apply_retention(&fold(&ops), &policy, 1_000).unwrap();
let raw = retained.transcript("c1");
assert_eq!(raw.iter().map(|t| t.role).collect::<Vec<_>>(), vec![Role::Tool, Role::Assistant],
"the retained window is the orphan [tool_result, assistant]");
let ms = retained.resume_messages("c1");
assert!(!matches!(ms.first(), Some(Message::ToolResult { .. })), "leading orphan tool_result dropped");
assert_provider_valid(&ms);
assert_eq!(ms.len(), 1);
assert!(matches!(&ms[0], Message::Assistant { .. }));
}
#[test]
fn repair_drops_a_leading_orphan_tool_result_directly() {
let hlc = Hlc { wall_ms: 0, counter: 0, device_id: "d".into() };
let orphan = Turn { conversation_id: "c".into(), role: Role::Tool, content: "res".into(), tool_calls: vec![], tool_use_id: Some("call_0".into()), timestamp: 1, hlc: hlc.clone(), op_id: "op-x".into() };
let asst = Turn { conversation_id: "c".into(), role: Role::Assistant, content: "done".into(), tool_calls: vec![], tool_use_id: None, timestamp: 2, hlc, op_id: "op-y".into() };
let ms = repair(vec![orphan, asst]);
assert!(!matches!(ms.first(), Some(Message::ToolResult { .. })));
assert_provider_valid(&ms);
}
#[test]
fn concurrent_device_turns_interleave_deterministically_and_order_independently() {
let mut a = DeviceLog::new("dev-a");
let mut b = DeviceLog::new("dev-b");
let a1 = a.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "from A", 1));
b.observe(&a1.hlc);
let b1 = b.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "B replies to A", vec![], 2));
let a2 = a.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "A concurrent", 3));
let b2 = b.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "B concurrent", 3));
let ops = vec![a1, b1, a2, b2];
let baseline = fold(&ops).transcript("c1");
assert_eq!(baseline.len(), 4);
let texts: Vec<&str> = baseline.iter().map(|t| t.content.as_str()).collect();
let pos = |s: &str| texts.iter().position(|t| *t == s).unwrap();
assert_eq!(pos("from A"), 0);
assert!(pos("from A") < pos("B replies to A"), "causality survives");
let baseline_hash = state_hash(&fold(&ops));
for perm in permutations(&ops) {
let folded = fold(&perm);
assert_eq!(folded.transcript("c1"), baseline, "transcript is delivery-order-independent");
assert_eq!(state_hash(&folded), baseline_hash);
}
}
#[test]
fn transcripts_are_partitioned_by_conversation_id() {
let ops = {
let mut v = one_device_conversation("work");
let mut d = DeviceLog::new("dev-b");
v.push(d.append(Scope::Personal, Surface::Conversation, Turn::user_payload("home", "dinner?", 20)));
v
};
let state = fold(&ops);
assert_eq!(state.conversation_ids(), vec!["home".to_string(), "work".to_string()]);
assert_eq!(state.transcript("work").len(), 4);
assert_eq!(state.transcript("home").len(), 1);
assert!(state.transcript("nonexistent").is_empty());
}
#[test]
fn legacy_speaker_text_turns_project_and_resume() {
let mut d = DeviceLog::new("dev-a");
let ops = vec![
d.append(Scope::Personal, Surface::Conversation, json!({"speaker": "user", "text": "hi", "timestamp": 1})),
d.append(Scope::Personal, Surface::Conversation, json!({"speaker": "assistant", "text": "hello", "timestamp": 2})),
];
let state = fold(&ops);
let turns = state.transcript(DEFAULT_CONVERSATION);
assert_eq!(turns.len(), 2);
assert_eq!((turns[0].role, turns[1].role), (Role::User, Role::Assistant));
let ms = state.resume_messages(DEFAULT_CONVERSATION);
assert_provider_valid(&ms);
assert!(matches!(&ms[0], Message::User { content } if content == "hi"));
assert!(matches!(&ms[1], Message::Assistant { content, .. } if content == "hello"));
}
#[test]
fn lastn_compaction_keeps_the_last_n_in_order_and_round_trips() {
let mut a = DeviceLog::new("dev-a");
let mut b = DeviceLog::new("dev-b");
let mut ops = vec![
a.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "t0", 100)),
a.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "t1", vec![], 101)),
a.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "t2", 102)),
a.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "t3", vec![], 103)),
];
let split = ops.len();
for op in &ops {
b.observe(&op.hlc);
}
ops.push(b.append(Scope::Personal, Surface::Conversation, Turn::user_payload("c1", "t4", 104)));
ops.push(b.append(Scope::Personal, Surface::Conversation, Turn::assistant_payload("c1", "t5", vec![], 105)));
let frontier = ops[..split].iter().map(|o| o.hlc.clone()).max().unwrap();
let mut acks = AckTable::new();
for op in &ops {
acks.ack(op.device_id.clone(), frontier.clone());
}
let mut policy = RetentionPolicy::keep_all();
policy.rules.insert("conversation".to_string(), RetentionRule::LastN { n: 2 });
let plan = plan_compaction(&ops, &acks, &policy, Some(1_000)).unwrap();
let ckpt = plan.checkpoint.state.transcript("c1");
assert_eq!(ckpt.iter().map(|t| t.content.as_str()).collect::<Vec<_>>(), vec!["t2", "t3"]);
let ckpt_json = serde_json::to_string(&plan.checkpoint.state).unwrap();
let ckpt_back: SyncState = serde_json::from_str(&ckpt_json).unwrap();
assert_eq!(ckpt_back.transcript("c1"), plan.checkpoint.state.transcript("c1"));
let reconstructed = fold_onto(&plan.checkpoint.state, &plan.retained_ops);
let resumed: Vec<String> = reconstructed.transcript("c1").iter().map(|t| t.content.clone()).collect();
assert_eq!(resumed, vec!["t2", "t3", "t4", "t5"], "resume = retained window + live tail, in order");
assert_provider_valid(&reconstructed.resume_messages("c1"));
let (global, _) = apply_retention(&fold(&ops), &policy, 1_000).unwrap();
let (local, _) = apply_retention(&reconstructed, &policy, 1_000).unwrap();
assert_eq!(local.transcript("c1"), global.transcript("c1"));
assert_eq!(
global.transcript("c1").iter().map(|t| t.content.clone()).collect::<Vec<_>>(),
vec!["t4", "t5"],
"the last-N display window is the same on every device"
);
assert_eq!(state_hash(&local), state_hash(&global));
}
fn permutations<T: Clone>(items: &[T]) -> Vec<Vec<T>> {
fn heap<T: Clone>(k: usize, arr: &mut Vec<T>, out: &mut Vec<Vec<T>>) {
if k == 1 {
out.push(arr.clone());
return;
}
for i in 0..k {
heap(k - 1, arr, out);
if k.is_multiple_of(2) {
arr.swap(i, k - 1);
} else {
arr.swap(0, k - 1);
}
}
}
let mut arr = items.to_vec();
let mut out = Vec::new();
heap(arr.len(), &mut arr, &mut out);
out
}
}