use kcode_k1_chat_chatend::{
AGENT_MESSAGE_TYPE, AGENT_RESPONSE_TYPE, BoxId, ChatBox, Chatend, ProviderCall,
ProviderGenerated, RecoveryError, SYSTEM_MESSAGE_TYPE, TOOL_CALL_TYPE, TOOL_MESSAGE_TYPE,
TOOL_RESULT_TYPE, ToolCallId, TransitionError, USER_MESSAGE_TYPE,
};
fn tool_id(byte: u8, sequence: u64) -> ToolCallId {
ToolCallId::new([byte; 12], sequence)
}
fn provider_call(tool_call_id: ToolCallId) -> ProviderCall {
ProviderCall {
tool_call_id,
name: "WebSearch".into(),
arguments: "{}".into(),
}
}
fn assert_box(value: &ChatBox, id: u64, box_type: &str, contents: &str) {
assert_eq!(value.id(), BoxId::new(id));
assert_eq!(value.box_type(), box_type);
assert_eq!(value.contents(), contents);
}
fn box_types(values: &[ChatBox]) -> Vec<&str> {
values.iter().map(ChatBox::box_type).collect()
}
#[test]
fn empty_agent_response_is_retained_at_call_wave() {
let call_id = tool_id(1, 7);
let mut chat = Chatend::new();
assert_eq!(chat.start_round(), Ok(None));
let dispatched = chat
.append_stage(
String::new(),
vec![ProviderGenerated::ToolCall(provider_call(call_id))],
)
.unwrap();
assert_eq!(chat.boxes().len(), 2);
assert_box(&chat.boxes()[0], 1, AGENT_RESPONSE_TYPE, "");
assert_box(
&chat.boxes()[1],
2,
TOOL_CALL_TYPE,
"Call ID: 010101010101010101010101/7\nCall Name: WebSearch\nArguments:\n{}",
);
assert_eq!(
chat.boxes()[1]
.tool_call_metadata()
.unwrap()
.unwrap()
.tool_call_id,
call_id
);
assert_eq!(dispatched.len(), 1);
assert_eq!(dispatched[0].tool_call_id, call_id);
assert_eq!(dispatched[0].call_box_id, BoxId::new(2));
}
#[test]
fn empty_agent_response_is_retained_at_done() {
let mut chat = Chatend::new();
assert_eq!(chat.start_round(), Ok(None));
let appended = chat.done(String::new()).unwrap();
assert_eq!(appended.len(), 1);
assert_box(&appended[0], 1, AGENT_RESPONSE_TYPE, "");
assert_eq!(chat.boxes(), appended.as_slice());
assert!(!chat.round_active());
}
#[test]
fn arrivals_do_not_take_the_promised_response_id() {
let mut chat = Chatend::new();
assert_eq!(
chat.accept_system("context".into()),
Ok(Some(BoxId::new(1)))
);
assert_eq!(chat.start_round(), Ok(Some(BoxId::new(1))));
assert_eq!(chat.accept_user("queued".into()), Ok(None));
chat.append_stage("answer".into(), Vec::new()).unwrap();
assert_eq!(chat.boxes().len(), 2);
assert_box(&chat.boxes()[1], 2, AGENT_RESPONSE_TYPE, "answer");
let arrivals = chat.flush_active_arrivals().unwrap();
assert_eq!(arrivals.len(), 1);
assert_box(&arrivals[0], 3, USER_MESSAGE_TYPE, "queued");
assert!(chat.abort().unwrap().is_empty());
}
#[test]
fn provider_order_is_preserved_and_only_calls_are_dispatched() {
let first = tool_id(2, 1);
let second = tool_id(3, 9);
let mut chat = Chatend::new();
chat.start_round().unwrap();
let dispatched = chat
.append_stage(
"response".into(),
vec![
ProviderGenerated::AgentMessage {
contents: "visible one".into(),
},
ProviderGenerated::ToolCall(provider_call(first)),
ProviderGenerated::AgentMessage {
contents: "visible two".into(),
},
ProviderGenerated::ToolCall(provider_call(second)),
],
)
.unwrap();
assert_eq!(
box_types(chat.boxes()),
vec![
AGENT_RESPONSE_TYPE,
AGENT_MESSAGE_TYPE,
TOOL_CALL_TYPE,
AGENT_MESSAGE_TYPE,
TOOL_CALL_TYPE,
]
);
assert_eq!(chat.boxes()[1].contents(), "visible one");
assert_eq!(chat.boxes()[3].contents(), "visible two");
assert_eq!(dispatched.len(), 2);
assert_eq!(dispatched[0].tool_call_id, first);
assert_eq!(dispatched[0].call_box_id, BoxId::new(3));
assert_eq!(dispatched[1].tool_call_id, second);
assert_eq!(dispatched[1].call_box_id, BoxId::new(5));
}
#[test]
fn flush_before_stage_fails_transactionally() {
let mut chat = Chatend::new();
chat.accept_user("existing".into()).unwrap();
chat.start_round().unwrap();
chat.accept_system("queued".into()).unwrap();
let before = chat.boxes().to_vec();
assert_eq!(
chat.flush_active_arrivals(),
Err(TransitionError::InvalidPhase)
);
assert_eq!(chat.boxes(), before.as_slice());
chat.append_stage("stage".into(), Vec::new()).unwrap();
assert_box(&chat.boxes()[1], 2, AGENT_RESPONSE_TYPE, "stage");
let arrivals = chat.flush_active_arrivals().unwrap();
assert_eq!(arrivals.len(), 1);
assert_box(&arrivals[0], 3, SYSTEM_MESSAGE_TYPE, "queued");
}
#[test]
fn marker_only_boundary_flush_starts_the_next_generation() {
let mut chat = Chatend::new();
chat.start_round().unwrap();
chat.append_stage("first".into(), Vec::new()).unwrap();
assert!(chat.flush_active_arrivals().unwrap().is_empty());
let appended = chat.done("second".into()).unwrap();
assert_eq!(appended.len(), 1);
assert_box(&appended[0], 2, AGENT_RESPONSE_TYPE, "second");
assert!(!chat.round_active());
}
#[test]
fn done_orders_trailing_arrivals_after_its_response() {
let mut chat = Chatend::new();
chat.start_round().unwrap();
chat.append_stage("wave".into(), Vec::new()).unwrap();
assert!(chat.flush_active_arrivals().unwrap().is_empty());
chat.accept_user("late user".into()).unwrap();
chat.accept_system("late system".into()).unwrap();
let appended = chat.done("final".into()).unwrap();
assert_eq!(
box_types(&appended),
vec![AGENT_RESPONSE_TYPE, USER_MESSAGE_TYPE, SYSTEM_MESSAGE_TYPE]
);
assert_box(&appended[0], 2, AGENT_RESPONSE_TYPE, "final");
assert_box(&appended[1], 3, USER_MESSAGE_TYPE, "late user");
assert_box(&appended[2], 4, SYSTEM_MESSAGE_TYPE, "late system");
}
#[test]
fn recovery_is_idle_preserves_history_and_rejects_gaps() {
let mut original = Chatend::new();
original.accept_user("question".into()).unwrap();
original.start_round().unwrap();
original
.append_stage(
"first".into(),
vec![ProviderGenerated::AgentMessage {
contents: "notice".into(),
}],
)
.unwrap();
assert!(original.flush_active_arrivals().unwrap().is_empty());
original.done("last".into()).unwrap();
let history = original.boxes().to_vec();
let mut recovered = Chatend::recover(history.clone()).unwrap();
assert!(!recovered.round_active());
assert_eq!(recovered.boxes(), history.as_slice());
assert_eq!(
recovered.accept_system("next".into()),
Ok(Some(BoxId::new(5)))
);
let gap = ChatBox::new(
BoxId::new(2),
USER_MESSAGE_TYPE.into(),
"gap".into(),
String::new(),
String::new(),
);
assert!(matches!(
Chatend::recover(vec![gap]),
Err(RecoveryError::NonContiguousBoxId)
));
}
#[test]
fn tool_correlation_behavior_is_retained() {
let call_id = tool_id(4, 12);
let unknown = tool_id(9, 1);
let mut chat = Chatend::new();
chat.start_round().unwrap();
let dispatched = chat
.append_stage(
"calling".into(),
vec![ProviderGenerated::ToolCall(provider_call(call_id))],
)
.unwrap();
let call_box_id = dispatched[0].call_box_id;
assert_eq!(
chat.accept_tool_message(unknown, "orphan".into()),
Err(TransitionError::UnknownToolCall)
);
assert_eq!(chat.accept_tool_message(call_id, "one".into()), Ok(None));
assert_eq!(chat.accept_tool_message(call_id, "two".into()), Ok(None));
assert_eq!(
chat.accept_async_return_v2(
call_id,
Ok("result".into()),
"test/metadata".into(),
"opaque".into(),
),
Ok(None)
);
assert_eq!(
chat.accept_async_return(call_id, Ok("duplicate".into())),
Err(TransitionError::DuplicateReturn)
);
assert_eq!(
chat.accept_tool_message(call_id, "late".into()),
Err(TransitionError::ToolMessageAfterResult)
);
let arrivals = chat.flush_active_arrivals().unwrap();
assert_eq!(
box_types(&arrivals),
vec![TOOL_MESSAGE_TYPE, TOOL_MESSAGE_TYPE, TOOL_RESULT_TYPE]
);
for (index, value) in arrivals[..2].iter().enumerate() {
let metadata = value.tool_message_metadata().unwrap().unwrap();
assert_eq!(metadata.tool_call_id, call_id);
assert_eq!(metadata.originating_call, call_box_id);
assert_eq!(metadata.message_index, index as u64 + 1);
}
let result = arrivals[2].tool_result_v2_metadata().unwrap().unwrap();
assert_eq!(result.tool_call_id, call_id);
assert_eq!(result.originating_call, call_box_id);
assert_eq!(result.result, Ok("result".into()));
assert_eq!(result.metadata_type, "test/metadata");
assert_eq!(result.metadata_contents, "opaque");
Chatend::recover(chat.boxes().to_vec()).unwrap();
}