Skip to main content

kcode_k1_chat_codex_state/
lib.rs

1#![forbid(unsafe_code)]
2
3mod recovery;
4
5pub use kcode_k1_chat_codex_codec::{BoxValue, Call};
6pub use kcode_k1_chat_state::{
7    AGENT_ATTACHMENT_TYPE, AGENT_MESSAGE_TYPE, AGENT_RESPONSE_TYPE, ActorState, BoxId, ChatBox,
8    ProviderGenerated, SYSTEM_MESSAGE_TYPE, TOOL_ATTACHMENT_TYPE, TOOL_CALL_TYPE,
9    TOOL_MESSAGE_TYPE, TOOL_RESULT_TYPE, ToolCallId, USER_ATTACHMENT_TYPE, USER_MESSAGE_TYPE,
10};
11pub use kcode_k1_codex_adapter::{ShimItem, ShimOutput};
12
13use std::sync::Arc;
14use std::sync::atomic::AtomicU8;
15
16use kcode_k1_chat_codex_codec::{open_agent_response, project};
17use kcode_k1_chat_state::{PreflightGenerated, ProviderCall, StateError};
18use recovery::recovered_sequence;
19
20pub const MALFORMED_NATIVE_CALL_METADATA_TYPE: &str = "k1.malformed-native-call.v1";
21
22#[derive(Clone, Copy, Debug, Eq, PartialEq)]
23pub enum PreflightMode {
24    Blocking,
25    NonBlocking,
26}
27
28#[derive(Clone, Debug, Eq, PartialEq)]
29pub enum PreflightItem {
30    SystemMessage {
31        contents: String,
32    },
33    KtoolCall {
34        name: String,
35        arguments: String,
36        mode: PreflightMode,
37    },
38}
39
40#[derive(Clone, Debug, Eq, PartialEq)]
41pub struct PreparedPreflightCall {
42    pub tool_call_id: ToolCallId,
43    pub call_box_id: BoxId,
44    pub name: String,
45    pub arguments: String,
46    pub mode: PreflightMode,
47}
48
49#[derive(Clone, Debug)]
50pub struct Start {
51    pub job: u64,
52    pub values: Vec<BoxValue>,
53    pub attempt: Arc<AtomicU8>,
54}
55
56#[derive(Clone, Debug, Eq, PartialEq)]
57pub struct PreparedCall {
58    pub tool_call_id: ToolCallId,
59    pub name: String,
60    pub arguments: String,
61    disposition: PreparedCallDisposition,
62}
63
64impl PreparedCall {
65    pub fn disposition(&self) -> &PreparedCallDisposition {
66        &self.disposition
67    }
68}
69
70#[derive(Clone, Debug, Eq, PartialEq)]
71pub enum PreparedCallDisposition {
72    External,
73    ImmediateError(Box<ImmediateToolError>),
74}
75
76#[derive(Clone, Debug, Eq, PartialEq)]
77pub struct ImmediateToolError {
78    pub message: String,
79    pub native_tool: String,
80    pub attempted_ktool: Option<String>,
81    pub validation_code: String,
82    pub path: String,
83    pub expected: String,
84    pub received: String,
85    pub native_arguments: String,
86    pub metadata_type: String,
87    pub metadata_contents: String,
88}
89
90#[derive(Clone, Debug)]
91pub struct PreparedMailboxFlush(Arc<Prepared>);
92
93#[derive(Debug)]
94struct Prepared {
95    token: u64,
96    values: Vec<BoxValue>,
97    job: u64,
98    external_count: usize,
99}
100
101impl PreparedMailboxFlush {
102    pub fn values(&self) -> &[BoxValue] {
103        &self.0.values
104    }
105}
106
107#[derive(Clone, Debug, Eq, PartialEq)]
108pub enum Status {
109    Running,
110    Quiet,
111    Stalled { message: String, restartable: bool },
112}
113
114#[derive(Clone, Copy, Debug, Eq, PartialEq)]
115pub enum RestartError {
116    NotStalled,
117    ProviderActionAccepted,
118}
119
120#[derive(Clone, Copy, Debug, Eq, PartialEq)]
121enum Phase {
122    ProviderActive,
123    ChatendBoundary,
124    PendingGeneration,
125}
126
127struct Round {
128    job: u64,
129    accepted_provider_action: bool,
130    phase: Phase,
131    mailbox_flush_needed: bool,
132    restartable: bool,
133}
134
135enum Mode {
136    Idle,
137    Running(Round),
138    Stalled(Status, bool),
139}
140
141pub struct ConversationState {
142    state: ActorState,
143    session: [u8; 12],
144    sequence: u64,
145    unsubmitted: Vec<BoxValue>,
146    queued_trigger: bool,
147    token: u64,
148    prepared: Option<Arc<Prepared>>,
149    mode: Mode,
150}
151
152impl ConversationState {
153    pub fn new(session: [u8; 12]) -> Self {
154        Self {
155            state: ActorState::new(false),
156            session,
157            sequence: 0,
158            unsubmitted: Vec::new(),
159            queued_trigger: false,
160            token: 0,
161            prepared: None,
162            mode: Mode::Idle,
163        }
164    }
165
166    pub fn recover(session: [u8; 12], boxes: Vec<ChatBox>, force: bool) -> Result<Self, String> {
167        let sequence = recovered_sequence(session, &boxes)?;
168        let state = ActorState::recover(boxes, force).map_err(debug)?;
169        Ok(Self {
170            unsubmitted: state.boxes().iter().map(project).collect(),
171            state,
172            sequence,
173            ..Self::new(session)
174        })
175    }
176
177    pub fn boxes(&self) -> &[ChatBox] {
178        self.state.boxes()
179    }
180
181    pub fn status(&self) -> Status {
182        match &self.mode {
183            Mode::Running(_) => Status::Running,
184            Mode::Idle if self.state.quiet() => Status::Quiet,
185            Mode::Idle => Status::Running,
186            Mode::Stalled(status, _) => status.clone(),
187        }
188    }
189
190    pub fn prepare_preflight(
191        &mut self,
192        items: Vec<PreflightItem>,
193    ) -> Result<Vec<PreparedPreflightCall>, String> {
194        if items.is_empty()
195            || !matches!(self.mode, Mode::Idle)
196            || !self.state.boxes().is_empty()
197            || self.prepared.is_some()
198        {
199            return Err("preflight requires one fresh nonempty session".to_owned());
200        }
201        let mut sequence = self.sequence;
202        let mut generated = Vec::with_capacity(items.len());
203        let mut calls = Vec::new();
204        for item in items {
205            match item {
206                PreflightItem::SystemMessage { contents } => {
207                    generated.push(PreflightGenerated::SystemMessage { contents });
208                }
209                PreflightItem::KtoolCall {
210                    name,
211                    arguments,
212                    mode,
213                } => {
214                    sequence = next_sequence(sequence)?;
215                    let tool_call_id = ToolCallId::new(self.session, sequence);
216                    generated.push(PreflightGenerated::ToolCall(ProviderCall {
217                        tool_call_id,
218                        name: name.clone(),
219                        arguments: arguments.clone(),
220                    }));
221                    calls.push((tool_call_id, name, arguments, mode));
222                }
223            }
224        }
225        let before = self.state.boxes().len();
226        let dispatched = self.state.append_preflight(generated).map_err(debug)?;
227        if dispatched.len() != calls.len() {
228            return Err("preflight call correlation diverged".to_owned());
229        }
230        self.unsubmitted
231            .extend(self.state.boxes()[before..].iter().map(project));
232        self.sequence = sequence;
233        Ok(dispatched
234            .into_iter()
235            .zip(calls)
236            .map(|(dispatched, (tool_call_id, name, arguments, mode))| {
237                debug_assert_eq!(dispatched.tool_call_id, tool_call_id);
238                PreparedPreflightCall {
239                    tool_call_id,
240                    call_box_id: dispatched.call_box_id,
241                    name,
242                    arguments,
243                    mode,
244                }
245            })
246            .collect())
247    }
248
249    pub fn accept(
250        &mut self,
251        box_type: String,
252        contents: String,
253        hidden_type: String,
254        hidden_contents: String,
255    ) -> Result<(), String> {
256        self.accept_arrival(true, |state| {
257            state.accept_box(box_type, contents, hidden_type, hidden_contents)
258        })
259    }
260
261    pub fn accept_tool_message(
262        &mut self,
263        tool_call_id: ToolCallId,
264        message: String,
265    ) -> Result<(), String> {
266        self.accept_arrival(false, |state| {
267            state.accept_tool_message(tool_call_id, message)
268        })
269    }
270
271    pub fn accept_tool_return(
272        &mut self,
273        tool_call_id: ToolCallId,
274        result: Result<String, String>,
275    ) -> Result<(), String> {
276        self.accept_arrival(true, |state| {
277            state.accept_async_return(tool_call_id, result)
278        })
279    }
280
281    pub fn accept_tool_return_v2(
282        &mut self,
283        tool_call_id: ToolCallId,
284        result: Result<String, String>,
285        metadata_type: String,
286        metadata_contents: String,
287    ) -> Result<(), String> {
288        self.accept_arrival(true, |state| {
289            state.accept_async_return_v2(tool_call_id, result, metadata_type, metadata_contents)
290        })
291    }
292
293    pub fn begin(&mut self) -> Result<Option<Start>, String> {
294        if !matches!(self.mode, Mode::Idle) {
295            return Ok(None);
296        }
297        let Some(start) = self.state.begin_inference().map_err(debug)? else {
298            return Ok(None);
299        };
300        let promised_id = self.promised_id()?;
301        let mut values = std::mem::take(&mut self.unsubmitted);
302        values.push(open_agent_response(promised_id));
303        let mailbox_flush_needed = std::mem::take(&mut self.queued_trigger);
304        self.mode = Mode::Running(Round {
305            job: start.job,
306            accepted_provider_action: false,
307            phase: Phase::ProviderActive,
308            mailbox_flush_needed,
309            restartable: true,
310        });
311        Ok(Some(Start {
312            job: start.job,
313            values,
314            attempt: start.attempt,
315        }))
316    }
317
318    pub fn prepare_stage(
319        &mut self,
320        job: u64,
321        text: String,
322        values: Vec<BoxValue>,
323    ) -> Result<Vec<PreparedCall>, String> {
324        if !matches!(
325            &self.mode,
326            Mode::Running(round) if round.job == job && round.phase == Phase::ProviderActive
327        ) {
328            return Err("stale Codex inference stage".to_owned());
329        }
330        if self.prepared.is_some() {
331            return Err("a Codex mailbox flush remains uncommitted".to_owned());
332        }
333        let accepted_provider_action = !values.is_empty();
334        let mut sequence = self.sequence;
335        let mut generated = Vec::with_capacity(values.len());
336        let mut prepared = Vec::new();
337        for value in values {
338            match value {
339                BoxValue::AgentMessage(Ok(contents)) => {
340                    generated.push(ProviderGenerated::AgentMessage { contents });
341                }
342                BoxValue::Call(Ok(call)) => {
343                    sequence = next_sequence(sequence)?;
344                    let tool_call_id = ToolCallId::new(self.session, sequence);
345                    generated.push(ProviderGenerated::ToolCall(ProviderCall {
346                        tool_call_id,
347                        name: call.name.clone(),
348                        arguments: call.arguments.clone(),
349                    }));
350                    prepared.push(PreparedCall {
351                        tool_call_id,
352                        name: call.name,
353                        arguments: call.arguments,
354                        disposition: PreparedCallDisposition::External,
355                    });
356                }
357                BoxValue::MalformedNativeAction(action) => {
358                    sequence = next_sequence(sequence)?;
359                    let tool_call_id = ToolCallId::new(self.session, sequence);
360                    let name = action
361                        .attempted_ktool()
362                        .unwrap_or_else(|| action.native_tool())
363                        .to_owned();
364                    let arguments = action.native_arguments_json();
365                    let error = ImmediateToolError {
366                        message: action.diagnostic(),
367                        native_tool: action.native_tool().to_owned(),
368                        attempted_ktool: action.attempted_ktool().map(str::to_owned),
369                        validation_code: action.validation_code().to_owned(),
370                        path: action.path().to_owned(),
371                        expected: action.expected().to_owned(),
372                        received: action.received().to_owned(),
373                        native_arguments: arguments.clone(),
374                        metadata_type: MALFORMED_NATIVE_CALL_METADATA_TYPE.to_owned(),
375                        metadata_contents: action.diagnostic_json(),
376                    };
377                    generated.push(ProviderGenerated::ToolCall(ProviderCall {
378                        tool_call_id,
379                        name: name.clone(),
380                        arguments: arguments.clone(),
381                    }));
382                    prepared.push(PreparedCall {
383                        tool_call_id,
384                        name,
385                        arguments,
386                        disposition: PreparedCallDisposition::ImmediateError(Box::new(error)),
387                    });
388                }
389                _ => return Err("stage contains a malformed provider action".to_owned()),
390            }
391        }
392        let before = self.state.boxes().len();
393        self.state
394            .append_stage(job, text, generated)
395            .map_err(debug)?;
396        self.unsubmitted.extend(
397            self.state.boxes()[before..]
398                .iter()
399                .filter(|box_| box_.box_type() == TOOL_CALL_TYPE)
400                .map(project),
401        );
402        self.sequence = sequence;
403        if let Mode::Running(round) = &mut self.mode {
404            round.accepted_provider_action |= accepted_provider_action;
405            round.phase = Phase::ChatendBoundary;
406            round.mailbox_flush_needed = true;
407            round.restartable = false;
408        }
409        Ok(prepared)
410    }
411
412    pub fn mailbox_flush(&mut self, job: u64) -> Result<Vec<ChatBox>, String> {
413        if !matches!(
414            &self.mode,
415            Mode::Running(round) if round.job == job && round.phase == Phase::ChatendBoundary
416        ) {
417            return Err("stale Codex active-arrival mailbox flush".to_owned());
418        }
419        let boxes = self.state.flush_active_arrivals(job).map_err(debug)?;
420        self.unsubmitted.extend(boxes.iter().map(project));
421        if let Mode::Running(round) = &mut self.mode {
422            round.phase = Phase::PendingGeneration;
423        }
424        Ok(boxes)
425    }
426
427    pub fn prepare_mailbox_flush(
428        &mut self,
429        job: u64,
430    ) -> Result<Option<PreparedMailboxFlush>, String> {
431        let (mailbox_flush_needed, phase) = match &self.mode {
432            Mode::Running(round) if round.job == job => (round.mailbox_flush_needed, round.phase),
433            _ => return Err("stale Codex inference mailbox flush".to_owned()),
434        };
435        if let Some(prepared) = &self.prepared {
436            return Ok(Some(PreparedMailboxFlush(Arc::clone(prepared))));
437        }
438        if !mailbox_flush_needed {
439            return Ok(None);
440        }
441        if phase == Phase::ProviderActive {
442            return Ok(None);
443        }
444        if phase == Phase::ChatendBoundary {
445            self.mailbox_flush(job)?;
446        }
447        let token = self
448            .token
449            .checked_add(1)
450            .ok_or_else(|| "Codex mailbox-flush token space was exhausted".to_owned())?;
451        let external_count = self.unsubmitted.len();
452        let mut values = self.unsubmitted.clone();
453        values.push(open_agent_response(self.promised_id()?));
454        let prepared = Arc::new(Prepared {
455            token,
456            values,
457            job,
458            external_count,
459        });
460        self.token = token;
461        self.prepared = Some(Arc::clone(&prepared));
462        if let Mode::Running(round) = &mut self.mode {
463            round.mailbox_flush_needed = false;
464        }
465        Ok(Some(PreparedMailboxFlush(prepared)))
466    }
467
468    pub fn validate_mailbox_flush(&self, prepared: &PreparedMailboxFlush) -> Result<(), String> {
469        let prepared = &prepared.0;
470        let prefix_matches = self.unsubmitted.get(..prepared.external_count)
471            == Some(&prepared.values[..prepared.external_count]);
472        let valid = matches!(
473            &self.mode,
474            Mode::Running(round)
475                if round.job == prepared.job && round.phase == Phase::PendingGeneration
476        ) && self.token == prepared.token
477            && prefix_matches
478            && self
479                .prepared
480                .as_ref()
481                .is_some_and(|current| Arc::ptr_eq(current, prepared));
482        if valid {
483            Ok(())
484        } else {
485            Err("stale or invalid Codex mailbox flush".to_owned())
486        }
487    }
488
489    pub fn commit_mailbox_flush(&mut self, prepared: PreparedMailboxFlush) -> Result<(), String> {
490        self.validate_mailbox_flush(&prepared)?;
491        self.unsubmitted.drain(..prepared.0.external_count);
492        self.prepared = None;
493        if let Mode::Running(round) = &mut self.mode {
494            round.phase = Phase::ProviderActive;
495        }
496        Ok(())
497    }
498
499    pub fn complete(&mut self, job: u64, output: ShimOutput<BoxValue>) -> Result<(), String> {
500        let mut round = self.take_round(job)?;
501        if round.phase != Phase::ProviderActive {
502            let message = "Codex inference completed outside a provider generation".to_owned();
503            self.preserve(round, message.clone(), false);
504            return Err(message);
505        }
506        if self.prepared.is_some() {
507            let message = "Codex inference completed with an uncommitted mailbox flush".to_owned();
508            self.preserve(round, message.clone(), false);
509            return Err(message);
510        }
511        let mut text = String::new();
512        for item in output.items {
513            match item {
514                ShimItem::Text(value) => text.push_str(&value),
515                ShimItem::Box(_) => {
516                    round.accepted_provider_action = true;
517                    let message = "terminal Codex output contains a box".to_owned();
518                    self.preserve(round, message.clone(), false);
519                    return Err(message);
520                }
521            }
522        }
523        let before = self.state.boxes().len();
524        if let Err(error) = self.state.complete_inference(job, text).map_err(debug) {
525            self.preserve(round, error.clone(), false);
526            return Err(error);
527        }
528        self.unsubmitted
529            .extend(self.state.boxes()[before..].iter().skip(1).map(project));
530        self.queued_trigger = false;
531        self.mode = Mode::Idle;
532        Ok(())
533    }
534
535    pub fn fail(&mut self, job: u64, message: String, restartable_before_launch: bool) {
536        if let Ok(round) = self.take_round(job) {
537            self.preserve(round, message, restartable_before_launch);
538        }
539    }
540
541    pub fn restart(&mut self) -> Result<(), RestartError> {
542        match &self.mode {
543            Mode::Stalled(_, true) => return Err(RestartError::ProviderActionAccepted),
544            Mode::Stalled(
545                Status::Stalled {
546                    restartable: true, ..
547                },
548                false,
549            ) => {}
550            _ => return Err(RestartError::NotStalled),
551        }
552        self.state.restart().map_err(|_| RestartError::NotStalled)?;
553        self.unsubmitted = self.state.boxes().iter().map(project).collect();
554        self.queued_trigger = false;
555        self.prepared = None;
556        self.mode = Mode::Idle;
557        Ok(())
558    }
559
560    fn accept_arrival<F>(&mut self, triggering: bool, accept: F) -> Result<(), String>
561    where
562        F: FnOnce(&mut ActorState) -> Result<(), StateError>,
563    {
564        let before = self.state.boxes().len();
565        accept(&mut self.state).map_err(debug)?;
566        let appended = &self.state.boxes()[before..];
567        self.unsubmitted.extend(appended.iter().map(project));
568        match &mut self.mode {
569            Mode::Running(round) => round.mailbox_flush_needed |= triggering,
570            Mode::Idle if appended.is_empty() => self.queued_trigger |= triggering,
571            Mode::Idle => self.queued_trigger = false,
572            Mode::Stalled(_, _) => {}
573        }
574        Ok(())
575    }
576
577    fn promised_id(&self) -> Result<BoxId, String> {
578        let previous = self.state.boxes().last().map_or(0, |box_| box_.id().get());
579        let value = previous
580            .checked_add(1)
581            .ok_or_else(|| "BoxId space was exhausted".to_owned())?;
582        Ok(BoxId::new(value))
583    }
584
585    fn take_round(&mut self, job: u64) -> Result<Round, String> {
586        match std::mem::replace(&mut self.mode, Mode::Idle) {
587            Mode::Running(round) if round.job == job => Ok(round),
588            other => {
589                self.mode = other;
590                Err("stale Codex inference completion".to_owned())
591            }
592        }
593    }
594
595    fn preserve(&mut self, round: Round, message: String, restartable_before_launch: bool) {
596        self.prepared = None;
597        let restartable = restartable_before_launch
598            && round.restartable
599            && !round.accepted_provider_action
600            && self
601                .state
602                .stall_inference(round.job, message.clone())
603                .is_ok();
604        if !restartable {
605            let _ = self.state.halt(message.clone());
606        }
607        self.mode = Mode::Stalled(
608            Status::Stalled {
609                message,
610                restartable,
611            },
612            round.accepted_provider_action,
613        );
614    }
615}
616
617fn next_sequence(sequence: u64) -> Result<u64, String> {
618    sequence
619        .checked_add(1)
620        .ok_or_else(|| "ToolCallId space was exhausted".to_owned())
621}
622
623fn debug(error: impl std::fmt::Debug) -> String {
624    format!("{error:?}")
625}
626
627#[cfg(test)]
628mod tests {
629    use super::*;
630    use kcode_k1_chat_codex_codec::Codec;
631    use kcode_k1_codex_adapter::{BoxCodec, ToolCall};
632
633    fn active() -> (ConversationState, u64) {
634        let mut state = ConversationState::new([7; 12]);
635        state
636            .accept(
637                USER_MESSAGE_TYPE.into(),
638                "start".into(),
639                String::new(),
640                String::new(),
641            )
642            .unwrap();
643        let job = state.begin().unwrap().unwrap().job;
644        (state, job)
645    }
646
647    fn native(name: &str, arguments: &str) -> BoxValue {
648        let mut codec = Codec;
649        codec.tool_call_box(&ToolCall {
650            call_id: "native".into(),
651            name: name.into(),
652            arguments: arguments.parse().unwrap(),
653        })
654    }
655
656    fn active_after_mailbox_flush() -> (ConversationState, u64, ToolCallId) {
657        let (mut state, job) = active();
658        let calls = state
659            .prepare_stage(
660                job,
661                String::new(),
662                vec![BoxValue::Call(Ok(Call {
663                    name: "tool".into(),
664                    arguments: "{}".into(),
665                }))],
666            )
667            .unwrap();
668        let prepared = state.prepare_mailbox_flush(job).unwrap().unwrap();
669        state.commit_mailbox_flush(prepared).unwrap();
670        (state, job, calls[0].tool_call_id)
671    }
672
673    fn empty_output() -> ShimOutput<BoxValue> {
674        ShimOutput { items: Vec::new() }
675    }
676
677    #[test]
678    fn fresh_preflight_is_ordered_unscheduled_and_part_of_first_input() {
679        let mut state = ConversationState::new([4; 12]);
680        let calls = state
681            .prepare_preflight(vec![
682                PreflightItem::SystemMessage {
683                    contents: "context".into(),
684                },
685                PreflightItem::KtoolCall {
686                    name: "CurrentTime".into(),
687                    arguments: "{}".into(),
688                    mode: PreflightMode::Blocking,
689                },
690            ])
691            .unwrap();
692        assert_eq!(calls.len(), 1);
693        assert_eq!(calls[0].call_box_id.get(), 2);
694        assert_eq!(calls[0].tool_call_id.sequence(), 1);
695        assert!(state.begin().unwrap().is_none());
696        state
697            .accept(
698                USER_MESSAGE_TYPE.into(),
699                "kickoff".into(),
700                String::new(),
701                String::new(),
702            )
703            .unwrap();
704        let start = state.begin().unwrap().unwrap();
705        assert_eq!(start.values.len(), 4);
706        assert!(
707            state
708                .prepare_preflight(vec![PreflightItem::SystemMessage {
709                    contents: "again".into()
710                }])
711                .is_err()
712        );
713    }
714
715    #[test]
716    fn malformed_only_is_persisted_and_prepared_as_deterministic_immediate_error() {
717        let (mut state, job) = active();
718        let malformed = native(
719            "call_ktool",
720            r#"{"z":0,"name":"recoverable","arguments":{"b":2,"a":1}}"#,
721        );
722        let (diagnostic, diagnostic_json) = match &malformed {
723            BoxValue::MalformedNativeAction(action) => {
724                (action.diagnostic(), action.diagnostic_json())
725            }
726            _ => panic!("test value must be malformed"),
727        };
728        let calls = state
729            .prepare_stage(job, String::new(), vec![malformed])
730            .unwrap();
731        assert_eq!(state.status(), Status::Running);
732        assert_eq!(calls.len(), 1);
733        assert_eq!(calls[0].name, "recoverable");
734        assert_eq!(
735            calls[0].arguments,
736            r#"{"arguments":{"a":1,"b":2},"name":"recoverable","z":0}"#
737        );
738        let PreparedCallDisposition::ImmediateError(error) = calls[0].disposition() else {
739            panic!("malformed call must be local");
740        };
741        assert_eq!(error.message, diagnostic);
742        assert_eq!(error.native_tool, "call_ktool");
743        assert_eq!(error.attempted_ktool.as_deref(), Some("recoverable"));
744        assert_eq!(error.validation_code, "invalid_wrapper_fields");
745        assert_eq!(error.path, "$");
746        assert_eq!(error.expected, "object with exactly name and arguments");
747        assert_eq!(error.received, calls[0].arguments);
748        assert_eq!(error.native_arguments, calls[0].arguments);
749        assert_eq!(error.metadata_type, MALFORMED_NATIVE_CALL_METADATA_TYPE);
750        assert_eq!(error.metadata_contents, diagnostic_json);
751        let persisted = state
752            .boxes()
753            .last()
754            .unwrap()
755            .tool_call_metadata()
756            .unwrap()
757            .unwrap();
758        assert_eq!(persisted.name, calls[0].name);
759        assert_eq!(persisted.arguments, calls[0].arguments);
760    }
761
762    #[test]
763    fn mixed_stage_preserves_order_contiguous_ids_name_selection_and_external_calls() {
764        let (mut state, job) = active();
765        let before = state.boxes().len();
766        let values = vec![
767            BoxValue::AgentMessage(Ok("first".into())),
768            native(
769                "call_ktool",
770                r#"{"name":"attempted","arguments":{},"extra":1}"#,
771            ),
772            BoxValue::Call(Ok(Call {
773                name: "valid".into(),
774                arguments: "{\"ok\":true}".into(),
775            })),
776            native("future_native", r#"{"z":0,"a":true}"#),
777        ];
778        let calls = state.prepare_stage(job, String::new(), values).unwrap();
779        assert_eq!(
780            calls
781                .iter()
782                .map(|call| call.tool_call_id.sequence())
783                .collect::<Vec<_>>(),
784            vec![1, 2, 3]
785        );
786        assert_eq!(
787            calls
788                .iter()
789                .map(|call| call.name.as_str())
790                .collect::<Vec<_>>(),
791            vec!["attempted", "valid", "future_native"]
792        );
793        assert!(matches!(
794            calls[0].disposition(),
795            PreparedCallDisposition::ImmediateError(_)
796        ));
797        assert_eq!(calls[1].disposition(), &PreparedCallDisposition::External);
798        assert!(matches!(
799            calls[2].disposition(),
800            PreparedCallDisposition::ImmediateError(_)
801        ));
802        assert_eq!(calls[2].arguments, r#"{"a":true,"z":0}"#);
803        assert_eq!(
804            state.boxes()[before..]
805                .iter()
806                .map(ChatBox::box_type)
807                .collect::<Vec<_>>(),
808            vec![
809                AGENT_RESPONSE_TYPE,
810                AGENT_MESSAGE_TYPE,
811                TOOL_CALL_TYPE,
812                TOOL_CALL_TYPE,
813                TOOL_CALL_TYPE,
814            ]
815        );
816    }
817
818    #[test]
819    fn normal_call_precedes_correlated_result_and_open_response_in_stable_flush() {
820        let (mut state, job) = active();
821        let calls = state
822            .prepare_stage(
823                job,
824                String::new(),
825                vec![BoxValue::Call(Ok(Call {
826                    name: "tool".into(),
827                    arguments: "{\"input\":1}".into(),
828                }))],
829            )
830            .unwrap();
831        state
832            .accept_tool_return(calls[0].tool_call_id, Ok("done".into()))
833            .unwrap();
834
835        let prepared = state.prepare_mailbox_flush(job).unwrap().unwrap();
836        let expected = vec![
837            project(
838                state
839                    .boxes()
840                    .iter()
841                    .find(|box_| box_.box_type() == TOOL_CALL_TYPE)
842                    .unwrap(),
843            ),
844            project(
845                state
846                    .boxes()
847                    .iter()
848                    .find(|box_| box_.box_type() == TOOL_RESULT_TYPE)
849                    .unwrap(),
850            ),
851            open_agent_response(state.promised_id().unwrap()),
852        ];
853        assert_eq!(prepared.values(), expected);
854
855        let repeated = state.prepare_mailbox_flush(job).unwrap().unwrap();
856        assert_eq!(repeated.values(), prepared.values());
857        state.validate_mailbox_flush(&prepared).unwrap();
858        state.commit_mailbox_flush(repeated).unwrap();
859        assert!(state.unsubmitted.is_empty());
860    }
861
862    #[test]
863    fn malformed_native_call_is_queued_as_its_canonical_tool_call() {
864        let (mut state, job) = active();
865        let malformed = native(
866            "call_ktool",
867            r#"{"name":"recoverable","arguments":{},"extra":1}"#,
868        );
869        state
870            .prepare_stage(job, String::new(), vec![malformed])
871            .unwrap();
872
873        let expected_call = project(
874            state
875                .boxes()
876                .iter()
877                .find(|box_| box_.box_type() == TOOL_CALL_TYPE)
878                .unwrap(),
879        );
880        let expected_response = open_agent_response(state.promised_id().unwrap());
881        let prepared = state.prepare_mailbox_flush(job).unwrap().unwrap();
882        assert_eq!(prepared.values(), &[expected_call, expected_response]);
883    }
884
885    #[test]
886    fn provider_agent_message_and_stage_response_are_not_requeued() {
887        let (mut state, job) = active();
888        state
889            .prepare_stage(
890                job,
891                "stage response".into(),
892                vec![BoxValue::AgentMessage(Ok("provider message".into()))],
893            )
894            .unwrap();
895
896        let expected_response = open_agent_response(state.promised_id().unwrap());
897        let prepared = state.prepare_mailbox_flush(job).unwrap().unwrap();
898        assert_eq!(prepared.values(), &[expected_response]);
899    }
900
901    #[test]
902    fn committed_mailbox_flush_does_not_schedule_a_fresh_turn() {
903        let (mut state, job, _) = active_after_mailbox_flush();
904        state.complete(job, empty_output()).unwrap();
905        assert_eq!(state.status(), Status::Quiet);
906        assert!(state.begin().unwrap().is_none());
907    }
908
909    #[test]
910    fn tool_messages_after_a_mailbox_flush_remain_inert() {
911        let (mut state, job, tool_call_id) = active_after_mailbox_flush();
912        state
913            .accept_tool_message(tool_call_id, "still running".into())
914            .unwrap();
915        state.complete(job, empty_output()).unwrap();
916        assert_eq!(state.status(), Status::Quiet);
917        assert!(state.begin().unwrap().is_none());
918    }
919
920    #[test]
921    fn one_tool_result_schedules_exactly_one_fresh_turn() {
922        let (mut state, job, tool_call_id) = active_after_mailbox_flush();
923        state
924            .accept_tool_return(tool_call_id, Ok("done".into()))
925            .unwrap();
926        state.complete(job, empty_output()).unwrap();
927        let followup = state.begin().unwrap().unwrap();
928        assert!(state.begin().unwrap().is_none());
929        state.complete(followup.job, empty_output()).unwrap();
930        assert_eq!(state.status(), Status::Quiet);
931        assert!(state.begin().unwrap().is_none());
932    }
933}