use crate::fold::{FoldedRecord, SyncState};
use crate::oplog::Hlc;
use car_inference_types::{Message, Provenance, 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>,
#[serde(default, skip_serializing_if = "Provenance::is_internal")]
pub provenance: Provenance,
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 {
Self::tool_payload_with_provenance(
conversation_id,
tool_use_id,
content,
timestamp,
Provenance::Internal,
)
}
pub fn tool_payload_with_provenance(
conversation_id: &str,
tool_use_id: &str,
content: &str,
timestamp: u64,
provenance: Provenance,
) -> Value {
let mut v = json!({
"conversation_id": conversation_id,
"role": "tool",
"tool_use_id": tool_use_id,
"content": content,
"timestamp": timestamp,
});
if provenance.is_external() {
v["provenance"] = json!("external");
}
v
}
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 provenance = match p.get("provenance").and_then(|v| v.as_str()) {
Some("external") => Provenance::External,
_ => Provenance::Internal,
};
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,
provenance,
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(),
thinking: Vec::new(),
},
Role::Tool => Message::ToolResult {
tool_use_id: self.tool_use_id.clone().unwrap_or_default(),
content: self.content.clone(),
provenance: self.provenance,
},
}
}
}
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};
#[test]
fn external_provenance_survives_the_oplog_round_trip() {
let mut d = DeviceLog::new("dev-a");
let ops = vec![
d.append(
Scope::Personal,
Surface::Conversation,
Turn::user_payload("c1", "what does that page say?", 10),
),
d.append(
Scope::Personal,
Surface::Conversation,
Turn::assistant_payload(
"c1",
"",
vec![serde_json::json!({"id": "call_0", "name": "web_search",
"arguments": {"q": "x"}})],
11,
),
),
d.append(
Scope::Personal,
Surface::Conversation,
Turn::tool_payload_with_provenance(
"c1",
"call_0",
"fetched page text",
12,
Provenance::External,
),
),
];
let messages: Vec<Message> = fold(&ops).resume_messages("c1");
let tool = messages
.iter()
.find(|m| matches!(m, Message::ToolResult { .. }))
.expect("the tool result must survive resume");
match tool {
Message::ToolResult {
content,
provenance,
..
} => {
assert_eq!(content, "fetched page text");
assert_eq!(
*provenance,
Provenance::External,
"resume downgraded external content to trusted"
);
}
other => panic!("expected ToolResult, got {other:?}"),
}
}
#[test]
fn plain_tool_payload_is_internal_and_omits_the_field() {
let payload = Turn::tool_payload("c1", "call_0", "exit 0", 10);
assert!(
payload.get("provenance").is_none(),
"internal must add no bytes: {payload}"
);
let external =
Turn::tool_payload_with_provenance("c1", "call_0", "x", 10, Provenance::External);
assert_eq!(external["provenance"], "external");
}
#[test]
fn unknown_provenance_value_folds_as_internal_rather_than_failing() {
let mut d = DeviceLog::new("dev-a");
let mut payload = Turn::tool_payload("c1", "call_0", "x", 10);
payload["provenance"] = serde_json::json!("from-the-future");
let ops = vec![d.append(Scope::Personal, Surface::Conversation, payload)];
let turns = fold(&ops).transcript("c1");
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].provenance, Provenance::Internal);
}
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,
provenance,
} => {
assert_eq!(tool_use_id, "call_0");
assert_eq!(content, "sunny, 72F");
assert_eq!(*provenance, Provenance::Internal);
}
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()),
provenance: Provenance::Internal,
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,
provenance: Provenance::Internal,
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
}
}