Skip to main content

rig_core/test_utils/
streaming_conformance.rs

1//! Wire-sequence conformance scenarios for provider streaming pipelines.
2//!
3//! The streaming sibling of `rig-agent`'s `model_conformance`: each scenario
4//! drives raw wire bytes (SSE or NDJSON) through a provider's *complete*
5//! streaming path — bytes → decode → the writer's parts →
6//! [`CompletionStream`](crate::streaming::CompletionStream) → its response
7//! — and asserts the streaming contract: an error is the stream's last item,
8//! a reply the provider did not end has no response, and every part the
9//! stream emits reaches the response exactly once. Scenarios state the contract; a per-provider
10//! [`ProviderWireFixture`] supplies the frames, since each wire format spells
11//! the same event differently.
12//!
13//! Every sequence family here pins a shipped bug from the #2257 review rounds
14//! (`rig-2257-code-review-findings-*.md`); the per-scenario comments cite the
15//! specific finding.
16//!
17//! Suites are expanded per wire family by
18//! [`streaming_conformance_suite!`](crate::streaming_conformance_suite).
19//! Scenarios a wire cannot spell return an explicit
20//! [`ScenarioOutcome::Skipped`] that the macro cross-checks against the
21//! suite's declared [`SuiteCapabilities`] — a skip is always visible and can
22//! never masquerade as a pass, so the executed count is exactly the declared
23//! grid minus the named skips (#2258 review, F8 corpus honesty).
24
25use bytes::Bytes;
26use futures::StreamExt;
27use futures::future::BoxFuture;
28
29use crate::completion::{CompletionRequest, CompletionResponse};
30use crate::error::ProviderError;
31use crate::{
32    completion::FinishReason,
33    error::ErrorReport,
34    http_client,
35    message::AssistantContent,
36    streaming::{Item, StreamEvent},
37};
38
39/// Typed failure from a wire-conformance scenario.
40#[derive(Debug, thiserror::Error)]
41pub enum ConformanceError {
42    /// Opening the stream failed before any wire frame was consumed.
43    #[error(transparent)]
44    Completion(#[from] ProviderError),
45    /// The pipeline violated the streaming contract table.
46    #[error("{scenario} conformance failed for {provider}: {details}")]
47    Contract {
48        /// Stable scenario name.
49        scenario: &'static str,
50        /// Provider driver under test.
51        provider: &'static str,
52        /// Actionable observation explaining the failure.
53        details: String,
54    },
55}
56
57impl ConformanceError {
58    fn contract(
59        scenario: &'static str,
60        provider: &'static str,
61        details: impl Into<String>,
62    ) -> Self {
63        Self::Contract {
64            scenario,
65            provider,
66            details: details.into(),
67        }
68    }
69}
70
71/// Outcome of a passing wire-conformance scenario.
72#[derive(Debug)]
73pub struct ScenarioReport {
74    /// Stable scenario name.
75    pub name: &'static str,
76    /// Provider driver the scenario ran against.
77    pub provider: &'static str,
78    /// Human-readable observations, one per verified sub-case.
79    pub observations: Vec<String>,
80}
81
82/// What a capability-gated scenario did: ran its assertions, or skipped
83/// because the wire family cannot spell the sequence shape.
84///
85/// A skip is an explicit, named outcome — never a silent pass. The
86/// [`streaming_conformance_suite!`](crate::streaming_conformance_suite)
87/// macro cross-checks it against the suite's declared capability flags via
88/// [`check_gated_outcome`], so a fixture cannot vacuously pass a scenario its
89/// capabilities claim to cover (#2258 review, F8 corpus-honesty batch).
90#[derive(Debug)]
91pub enum ScenarioOutcome {
92    /// The scenario ran and its assertions held.
93    Ran(ScenarioReport),
94    /// The fixture lacks the sequence shape; nothing was asserted.
95    Skipped {
96        /// Stable scenario name.
97        name: &'static str,
98        /// Provider driver under test.
99        provider: &'static str,
100        /// Why the wire family cannot spell the shape.
101        reason: &'static str,
102    },
103}
104
105/// Streaming-relevant capability flags for one wire family's conformance
106/// suite: which optional sequence shapes the wire can spell.
107///
108/// Each flag mirrors an `Option` field on [`ProviderWireFixture`], and the
109/// only constructor is [`ProviderWireFixture::capabilities`] — suites never
110/// hand-write flags, so a flag structurally cannot drift from the wire
111/// fixture that backs it (it *is* the fixture's populated-field set).
112#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
113pub struct SuiteCapabilities {
114    /// The wire streams tool-call arguments incrementally
115    /// (`partial_tool_call_frames`).
116    pub partial_tool_args: bool,
117    /// The wire has a genuine terminal that can omit usage metrics
118    /// (`zero_usage_terminal_frames`).
119    pub zero_usage_terminal: bool,
120    /// The wire has a data-less terminal signal (`bare_terminal_frames`).
121    pub bare_terminal: bool,
122    /// A frame-level decode failure can be spelled (`malformed_frame`).
123    pub malformed_frame: bool,
124    /// An unknown event type can be spelled (`unknown_event_frame`).
125    pub unknown_event_frame: bool,
126    /// A known event with a schema-defective payload can be spelled
127    /// (`defective_known_frame`).
128    pub defective_known_frame: bool,
129    /// The wire has a delta-less choice prelude shape
130    /// (`delta_less_prelude_frame`).
131    pub delta_less_prelude: bool,
132    /// The wire has a refusal channel (`refusal`).
133    pub refusal: bool,
134    /// The wire mints a constant per-stream reasoning identity, so
135    /// interleaving output is its only reasoning boundary
136    /// (`interleaved_reasoning`).
137    pub interleaved_reasoning: bool,
138}
139
140impl SuiteCapabilities {
141    /// Build a capability set from manifest names. An unknown name is an
142    /// error so a typo in a suite's `manifest:` list fails loudly rather
143    /// than silently asserting an empty flag; the macro-expanded test
144    /// asserts on it.
145    pub fn from_names(names: &[&str]) -> Result<Self, String> {
146        let mut caps = Self::default();
147        for name in names {
148            match *name {
149                "partial_tool_args" => caps.partial_tool_args = true,
150                "zero_usage_terminal" => caps.zero_usage_terminal = true,
151                "bare_terminal" => caps.bare_terminal = true,
152                "malformed_frame" => caps.malformed_frame = true,
153                "unknown_event_frame" => caps.unknown_event_frame = true,
154                "defective_known_frame" => caps.defective_known_frame = true,
155                "delta_less_prelude" => caps.delta_less_prelude = true,
156                "refusal" => caps.refusal = true,
157                "interleaved_reasoning" => caps.interleaved_reasoning = true,
158                other => {
159                    return Err(format!(
160                        "unknown capability name in suite manifest: {other}"
161                    ));
162                }
163            }
164        }
165        Ok(caps)
166    }
167}
168
169/// The canonical fixture-driven scenario set every wire-family suite must
170/// expand — one named test each, compared against the macro's emitted list by
171/// its `suite_is_complete` test (langchain's anti-tamper precedent).
172pub const CANONICAL_SCENARIOS: &[&str] = &[
173    "truncation_preserves_content_without_terminal",
174    "transport_error_after_tool_call_yields_err_then_end",
175    "malformed_frame_ends_the_reply",
176    "unknown_event_is_skipped",
177    "defective_known_event_ends_the_reply",
178    "delta_less_choice_prelude_is_a_noop",
179    "refusal_frames_deliver_text_without_error",
180    "bare_terminal_after_only_unparseable_frames_fabricates_nothing",
181    "usage_variants_are_reported_or_absent",
182    "interleaved_constant_id_reasoning_preserves_order",
183];
184
185/// Every streaming wire family in the workspace. The workspace registry test
186/// (`all_wire_families_have_conformance_suites`) fails CI when any family
187/// lacks a [`streaming_conformance_suite!`](crate::streaming_conformance_suite)
188/// invocation naming it.
189pub const WIRE_FAMILIES: &[&str] = &[
190    "openai_chat",
191    "openai_responses",
192    "openai_responses_websocket",
193    "chatgpt",
194    "anthropic",
195    "gemini_rest",
196    "gemini_interactions",
197    "gemini_grpc",
198    "cohere",
199    "ollama",
200    "xai",
201    "copilot",
202    "bedrock",
203    "candle",
204];
205
206/// The sanctioned reason for an expected-failure scenario, from `xfail`
207/// entries of the form `"scenario_name: reason (finding reference)"`.
208pub fn xfail_reason<'a>(xfail: &[&'a str], scenario: &str) -> Option<&'a str> {
209    xfail.iter().find_map(|entry| {
210        let (name, reason) = entry.split_once(':')?;
211        (name.trim() == scenario).then(|| reason.trim())
212    })
213}
214
215/// `xfail` entries that do not name a canonical scenario or carry no reason.
216pub fn invalid_xfail_entries(xfail: &[&str]) -> Vec<String> {
217    xfail
218        .iter()
219        .filter(|entry| match entry.split_once(':') {
220            Some((name, reason)) => {
221                !CANONICAL_SCENARIOS.contains(&name.trim()) || reason.trim().is_empty()
222            }
223            None => true,
224        })
225        .map(std::string::ToString::to_string)
226        .collect()
227}
228
229/// Enforce a capability-gated scenario's outcome against the suite's declared
230/// capability flag and its `xfail` list.
231///
232/// A `Skipped` outcome passes only when the capability is disclaimed; a `Ran`
233/// outcome passes only when it is declared — so a vacuous pass (fixture lacks
234/// the shape but the suite claims to cover it) is impossible, and the skip is
235/// visible in the test output.
236pub fn check_gated_outcome(
237    scenario: &'static str,
238    capability: bool,
239    xfail: &[&str],
240    outcome: Result<ScenarioOutcome, ConformanceError>,
241) -> Result<(), String> {
242    match (xfail_reason(xfail, scenario), outcome) {
243        (Some(reason), Err(error)) => {
244            eprintln!("xfail {scenario}: {reason} ({error})");
245            Ok(())
246        }
247        (Some(reason), Ok(_)) => Err(format!(
248            "{scenario} passed but is listed as xfail ({reason}); remove the xfail entry"
249        )),
250        (None, Err(error)) => Err(format!("{scenario} failed: {error}")),
251        (None, Ok(ScenarioOutcome::Ran(_))) => {
252            if capability {
253                Ok(())
254            } else {
255                Err(format!(
256                    "{scenario} ran but the suite disclaims the capability; set the flag to true"
257                ))
258            }
259        }
260        (None, Ok(ScenarioOutcome::Skipped { reason, .. })) => {
261            if capability {
262                Err(format!(
263                    "{scenario} skipped ({reason}) but the suite declares the capability; \
264                     a declared capability's scenario must run"
265                ))
266            } else {
267                eprintln!("skipped {scenario}: {reason}");
268                Ok(())
269            }
270        }
271    }
272}
273
274/// Enforce an always-runnable scenario's result against the `xfail` list.
275pub fn check_ungated_outcome(
276    scenario: &'static str,
277    xfail: &[&str],
278    result: Result<ScenarioReport, ConformanceError>,
279) -> Result<(), String> {
280    match (xfail_reason(xfail, scenario), result) {
281        (Some(reason), Err(error)) => {
282            eprintln!("xfail {scenario}: {reason} ({error})");
283            Ok(())
284        }
285        (Some(reason), Ok(_)) => Err(format!(
286            "{scenario} passed but is listed as xfail ({reason}); remove the xfail entry"
287        )),
288        (None, Err(error)) => Err(format!("{scenario} failed: {error}")),
289        (None, Ok(_)) => Ok(()),
290    }
291}
292
293/// One scripted wire input frame.
294///
295/// Byte-transport wires (SSE, NDJSON, websocket) script raw bytes fed through
296/// the provider's HTTP layer; typed-event wires (bedrock, candle,
297/// gemini-grpc) script already-typed SDK events fed to the adapter directly —
298/// events-first, no mock transport — which the typed driver downcasts back.
299#[derive(Clone)]
300pub enum WireInput {
301    /// A raw wire byte frame.
302    Bytes(Bytes),
303    /// An already-typed SDK event for a typed-event wire.
304    Event(std::sync::Arc<dyn std::any::Any + Send + Sync>),
305}
306
307impl WireInput {
308    /// The frame's raw bytes, when it is a byte frame.
309    pub fn as_bytes(&self) -> Option<&Bytes> {
310        match self {
311            Self::Bytes(bytes) => Some(bytes),
312            Self::Event(_) => None,
313        }
314    }
315
316    /// The frame's typed event, when it is an event frame of type `T`.
317    pub fn downcast_event<T: 'static>(&self) -> Option<&T> {
318        match self {
319            Self::Bytes(_) => None,
320            Self::Event(event) => event.downcast_ref(),
321        }
322    }
323}
324
325impl From<Bytes> for WireInput {
326    fn from(bytes: Bytes) -> Self {
327        Self::Bytes(bytes)
328    }
329}
330
331impl std::fmt::Debug for WireInput {
332    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
333        match self {
334            Self::Bytes(bytes) => formatter.debug_tuple("Bytes").field(bytes).finish(),
335            Self::Event(_) => formatter.write_str("Event(..)"),
336        }
337    }
338}
339
340/// Build a typed-event fixture frame.
341pub fn event_frame<T: Send + Sync + 'static>(event: T) -> WireInput {
342    WireInput::Event(std::sync::Arc::new(event))
343}
344
345/// The wire frames a driver feeds into the provider's pipeline. An `Err`
346/// chunk models a mid-stream transport failure.
347pub type WireChunks = Vec<http_client::Result<WireInput>>;
348
349/// Build the chunk list for an all-delivered frame sequence.
350pub fn ok_chunks(frames: impl IntoIterator<Item = impl Into<WireInput>>) -> WireChunks {
351    frames.into_iter().map(|frame| Ok(frame.into())).collect()
352}
353
354/// A scripted mid-stream transport failure chunk: the connection dropped,
355/// so there is no reply and no status.
356pub fn transport_error_chunk() -> http_client::Result<WireInput> {
357    Err(http_client::Error::instance(std::io::Error::new(
358        std::io::ErrorKind::ConnectionReset,
359        "connection reset",
360    )))
361}
362
363/// Executable stream-lifecycle validator (#2258 C1).
364///
365/// The invariants every normalized stream must satisfy, stated once and run
366/// over every recorded cassette and corpus fixture that drains through
367/// [`fixtures::drain`] — the langchain `assert_valid_event_stream` move:
368/// prose contracts scattered across N adapters become one executable
369/// artifact. Panics with the violated law.
370///
371/// Laws (universal — they hold for truncated and errored streams too):
372///
373/// 1. **Terminal error.** An error is the stream's last item: nothing
374///    follows it.
375/// 2. **Text conservation.** The aggregated text is exactly the
376///    concatenation of the yielded text deltas: accumulated delta content
377///    equals the payload the aggregate delivers.
378/// 3. **Completed-call conservation.** Every completed tool call yielded on
379///    the stream appears in the aggregated choice exactly once, and vice
380///    versa (counts match; aggregation neither drops nor duplicates).
381/// 4. **Sequence.** The events read back through
382///    [`Transcript::parse_prefix`](crate::streaming::Transcript::parse_prefix):
383///    each part starts once, in position, grows only while open, and ends
384///    once.
385/// 5. **Reasoning provenance.** The response contains a reasoning part only
386///    if the stream ended one.
387pub fn assert_valid_event_stream(
388    items: &[Result<Item<StreamEvent>, ErrorReport>],
389    choice: &[AssistantContent],
390) {
391    use crate::message::AssistantContent;
392
393    // Law 1: an error is the last item.
394    if let Some(error_index) = items.iter().position(Result::is_err) {
395        assert_eq!(
396            error_index + 1,
397            items.len(),
398            "law 1 (terminal error): an item followed the stream's error"
399        );
400    }
401    let events: Vec<&StreamEvent> = items
402        .iter()
403        .filter_map(|item| match item {
404            Ok(Item::Event(event)) => Some(event),
405            _ => None,
406        })
407        .collect();
408
409    // Law 2: text conservation.
410    let streamed_text: String = events
411        .iter()
412        .filter_map(|event| match event {
413            StreamEvent::Text { text, .. } => Some(text.as_str()),
414            _ => None,
415        })
416        .collect();
417    let aggregated_text: String = choice
418        .iter()
419        .filter_map(|content| match content {
420            AssistantContent::Text(text) => Some(text.text.as_str()),
421            _ => None,
422        })
423        .collect();
424    assert_eq!(
425        aggregated_text, streamed_text,
426        "law 2 (text conservation): aggregated text differs from the streamed fragments"
427    );
428
429    // Law 3: completed-call conservation.
430    let yielded_calls = events
431        .iter()
432        .filter(|event| {
433            matches!(
434                event,
435                StreamEvent::End {
436                    content: AssistantContent::ToolCall(_),
437                    ..
438                }
439            )
440        })
441        .count();
442    let aggregated_calls = choice
443        .iter()
444        .filter(|content| matches!(content, AssistantContent::ToolCall(_)))
445        .count();
446    assert_eq!(
447        aggregated_calls, yielded_calls,
448        "law 3 (completed-call conservation): {yielded_calls} calls yielded, \
449         {aggregated_calls} aggregated"
450    );
451
452    // Law 4: the sequence reads back.
453    let serialized = serde_json::to_value(
454        items
455            .iter()
456            .filter_map(|item| item.as_ref().ok())
457            .collect::<Vec<_>>(),
458    )
459    .unwrap_or_default();
460    let read_back = crate::streaming::Transcript::parse_prefix(serialized);
461    assert!(
462        read_back.is_ok(),
463        "law 4 (sequence): the stream's events do not read back: {read_back:?}"
464    );
465
466    // Law 5: reasoning provenance.
467    let yielded_reasoning = events.iter().any(|event| {
468        matches!(
469            event,
470            StreamEvent::End {
471                content: AssistantContent::Reasoning(_),
472                ..
473            }
474        )
475    });
476    let aggregated_reasoning = choice
477        .iter()
478        .any(|content| matches!(content, AssistantContent::Reasoning(_)));
479    assert!(
480        yielded_reasoning || !aggregated_reasoning,
481        "law 5 (reasoning provenance): aggregated reasoning with no reasoning yielded"
482    );
483}
484
485/// Everything the consumer observed from one full pipeline run: the yielded
486/// items in order, plus the response they folded into.
487#[derive(Debug)]
488pub struct DrainedStream {
489    /// Every item the stream yielded, in order.
490    pub items: Vec<Result<Item<StreamEvent>, ErrorReport>>,
491    /// The parts that ended, in position: the response's choice, or what
492    /// arrived before the stream stopped.
493    pub choice: Vec<AssistantContent>,
494    /// The response, absent on truncation or a terminal error.
495    pub response: Option<CompletionResponse>,
496}
497
498impl DrainedStream {
499    fn events(&self) -> impl Iterator<Item = &StreamEvent> {
500        self.items.iter().filter_map(|item| match item {
501            Ok(Item::Event(event)) => Some(event),
502            _ => None,
503        })
504    }
505
506    /// Text fragments yielded to the consumer, in order.
507    pub fn texts(&self) -> Vec<&str> {
508        self.events()
509            .filter_map(|event| match event {
510                StreamEvent::Text { text, .. } => Some(text.as_str()),
511                _ => None,
512            })
513            .collect()
514    }
515
516    /// Names of the tool calls that ended on the stream, in order.
517    pub fn tool_call_names(&self) -> Vec<&str> {
518        self.events()
519            .filter_map(|event| match event {
520                StreamEvent::End {
521                    content: AssistantContent::ToolCall(tool_call),
522                    ..
523                } => Some(tool_call.function.name.as_str()),
524                _ => None,
525            })
526            .collect()
527    }
528
529    /// Raw payloads of the `Unknown` items the stream yielded, in order.
530    pub fn unknown_values(&self) -> Vec<&serde_json::Value> {
531        self.items
532            .iter()
533            .filter_map(|item| match item {
534                Ok(Item::Unknown(value)) => Some(value.value()),
535                _ => None,
536            })
537            .collect()
538    }
539
540    /// Number of `Err` items the stream yielded.
541    pub fn error_count(&self) -> usize {
542        self.items.iter().filter(|item| item.is_err()).count()
543    }
544
545    /// Whether the provider ended the reply: it has a response.
546    fn has_terminal(&self) -> bool {
547        self.response.is_some()
548    }
549
550    /// Whether the stream reached the provider's end with no `Err` item —
551    /// the clean-completion precondition the content-shape scenarios share.
552    fn completed_cleanly(&self) -> bool {
553        self.error_count() == 0 && self.response.is_some()
554    }
555
556    /// Index of the first `Err` item, if any.
557    fn first_error_index(&self) -> Option<usize> {
558        self.items.iter().position(std::result::Result::is_err)
559    }
560
561    /// Text parts in the aggregated choice, in order.
562    pub fn choice_texts(&self) -> Vec<&str> {
563        self.choice
564            .iter()
565            .filter_map(|content| match content {
566                AssistantContent::Text(text) => Some(text.text.as_str()),
567                _ => None,
568            })
569            .collect()
570    }
571
572    /// Reasoning parts in the aggregated choice, in order.
573    pub fn choice_reasoning(&self) -> Vec<&crate::message::Reasoning> {
574        self.choice
575            .iter()
576            .filter_map(|content| match content {
577                AssistantContent::Reasoning(reasoning) => reasoning.open(reasoning.issuer()),
578                _ => None,
579            })
580            .collect()
581    }
582}
583
584type DriveFn = Box<
585    dyn Fn(WireChunks) -> BoxFuture<'static, Result<DrainedStream, ProviderError>> + Send + Sync,
586>;
587
588/// One provider's full streaming pipeline over scripted wire chunks.
589///
590/// The closure builds a fresh provider client over a scripted HTTP double
591/// (`SequencedStreamingHttpClient`), opens `Model::stream`, drains
592/// it, and returns everything the consumer observed.
593pub struct WireDriver {
594    /// Stable descriptor name of the provider under test.
595    pub provider: &'static str,
596    drive: DriveFn,
597}
598
599impl WireDriver {
600    /// Wrap a provider pipeline closure.
601    pub fn new(
602        provider: &'static str,
603        drive: impl Fn(WireChunks) -> BoxFuture<'static, Result<DrainedStream, ProviderError>>
604        + Send
605        + Sync
606        + 'static,
607    ) -> Self {
608        Self {
609            provider,
610            drive: Box::new(drive),
611        }
612    }
613
614    /// Run the provider's full pipeline over `chunks` and drain it.
615    pub async fn drive(&self, chunks: WireChunks) -> Result<DrainedStream, ProviderError> {
616        (self.drive)(chunks).await
617    }
618}
619
620/// Refusal frames and the text the pipeline must deliver for them.
621pub struct RefusalFixture {
622    /// Frames carrying the refusal content.
623    pub frames: Vec<WireInput>,
624    /// Text the consumer must observe.
625    pub expected_text: &'static str,
626}
627
628/// The interleaving-boundary shape for a wire whose reasoning identity is
629/// a constant per-stream minted key: reasoning, an interleaved tool call,
630/// then more reasoning, which must aggregate as three ordered parts —
631/// never one merged item that misorders history on replay.
632pub struct InterleavedReasoningFixture {
633    /// Reasoning → tool call → reasoning frames, terminal included.
634    pub frames: Vec<WireInput>,
635    /// The reasoning content streamed before the boundary.
636    pub first_reasoning: &'static str,
637    /// The interleaved call's tool name.
638    pub tool_name: &'static str,
639    /// The reasoning content streamed after the boundary.
640    pub second_reasoning: &'static str,
641}
642
643type BufferedDriveFn = Box<
644    dyn Fn(String) -> BoxFuture<'static, Result<Vec<AssistantContent>, ProviderError>>
645        + Send
646        + Sync,
647>;
648
649/// A buffered-body pipeline (the ChatGPT backend shape): the full SSE body is
650/// re-parsed after the fact and merged with the terminal response body.
651pub struct BufferedBodyDriver {
652    /// Stable descriptor name of the provider under test.
653    pub provider: &'static str,
654    drive: BufferedDriveFn,
655}
656
657impl BufferedBodyDriver {
658    /// Wrap a buffered pipeline closure.
659    pub fn new(
660        provider: &'static str,
661        drive: impl Fn(String) -> BoxFuture<'static, Result<Vec<AssistantContent>, ProviderError>>
662        + Send
663        + Sync
664        + 'static,
665    ) -> Self {
666        Self {
667            provider,
668            drive: Box::new(drive),
669        }
670    }
671
672    /// Run the buffered pipeline over a complete SSE body.
673    pub async fn drive(&self, body: String) -> Result<Vec<AssistantContent>, ProviderError> {
674        (self.drive)(body).await
675    }
676}
677
678/// Per-provider wire frames for the shared scenario set.
679///
680/// `Option` fields cover sequence shapes a wire family cannot spell (e.g.
681/// ollama's NDJSON has no event types, so no "unknown event type" frame).
682pub struct ProviderWireFixture {
683    /// The provider's full pipeline.
684    pub driver: WireDriver,
685    /// Frames that deliver exactly the text deltas in `expected_texts`.
686    pub text_frames: Vec<WireInput>,
687    /// The text deltas `text_frames` delivers, in order.
688    pub expected_texts: Vec<&'static str>,
689    /// Frames that fully deliver one tool call (including any completion
690    /// signal the wire needs, but no stream terminal).
691    pub tool_call_frames: Vec<WireInput>,
692    /// Name of the tool call `tool_call_frames` delivers.
693    pub expected_tool_name: &'static str,
694    /// Frames that leave a tool call mid-arguments, where the wire streams
695    /// arguments incrementally.
696    pub partial_tool_call_frames: Option<Vec<WireInput>>,
697    /// The provider's genuine stream terminal, carrying usage.
698    pub terminal_frames: Vec<WireInput>,
699    /// Total tokens `terminal_frames` reports.
700    pub expected_usage_total: u64,
701    /// Finish reason `terminal_frames` reports.
702    pub expected_finish_reason: Option<FinishReason>,
703    /// A genuine terminal that reports no usage metrics at all.
704    pub zero_usage_terminal_frames: Option<Vec<WireInput>>,
705    /// A terminal signal that carries no data of its own (e.g. a bare
706    /// `[DONE]`), for wires that have one.
707    pub bare_terminal_frames: Option<Vec<WireInput>>,
708    /// A frame that fails the wire decode entirely. `None` only for
709    /// typed-event wires, whose SDK surfaces decode failures as transport
710    /// errors — a frame-level corrupt input cannot be spelled there.
711    pub malformed_frame: Option<WireInput>,
712    /// An event type this client does not know, for typed-event wires.
713    pub unknown_event_frame: Option<WireInput>,
714    /// A known event whose payload is schema-defective.
715    pub defective_known_frame: Option<WireInput>,
716    /// A delta-less choice prelude (the Azure `prompt_filter_results` shape).
717    pub delta_less_prelude_frame: Option<WireInput>,
718    /// Refusal content frames, where the wire has a refusal channel.
719    pub refusal: Option<RefusalFixture>,
720    /// The interleaving-boundary shape, where the wire's reasoning identity
721    /// is a constant per-stream minted key (its adapter synthesizes the
722    /// reasoning ends other output implies).
723    pub interleaved_reasoning: Option<InterleavedReasoningFixture>,
724}
725
726impl ProviderWireFixture {
727    /// The capability set this fixture's populated optional fields spell —
728    /// the descriptor the suite macro gates scenarios on.
729    ///
730    /// Deriving flags here (instead of hand-writing them per suite
731    /// invocation) makes flag/fixture drift structurally impossible: a shape
732    /// the fixture supplies is a declared capability, a shape it lacks is a
733    /// visible named skip, and there is nothing else to keep in sync.
734    pub fn capabilities(&self) -> SuiteCapabilities {
735        SuiteCapabilities {
736            partial_tool_args: self.partial_tool_call_frames.is_some(),
737            zero_usage_terminal: self.zero_usage_terminal_frames.is_some(),
738            bare_terminal: self.bare_terminal_frames.is_some(),
739            malformed_frame: self.malformed_frame.is_some(),
740            unknown_event_frame: self.unknown_event_frame.is_some(),
741            defective_known_frame: self.defective_known_frame.is_some(),
742            delta_less_prelude: self.delta_less_prelude_frame.is_some(),
743            refusal: self.refusal.is_some(),
744            interleaved_reasoning: self.interleaved_reasoning.is_some(),
745        }
746    }
747}
748
749fn concat_frames(parts: &[&[WireInput]]) -> Vec<WireInput> {
750    parts
751        .iter()
752        .flat_map(|frames| frames.iter().cloned())
753        .collect()
754}
755
756/// A scenario's identity and running notes, bound once.
757///
758/// The scenario name and the provider under test are facts of the *scenario*,
759/// not of each check inside it: every contract violation in a body reports
760/// the same pair. `Checks` holds them, so a check states only what must hold
761/// and a report states only what was observed.
762struct Checks {
763    name: &'static str,
764    provider: &'static str,
765    observations: Vec<String>,
766}
767
768impl Checks {
769    fn new(name: &'static str, provider: &'static str) -> Self {
770        Self {
771            name,
772            provider,
773            observations: Vec::new(),
774        }
775    }
776
777    /// This scenario's contract violation, for the `ok_or_else` shape.
778    fn fail(&self, details: impl Into<String>) -> ConformanceError {
779        ConformanceError::contract(self.name, self.provider, details)
780    }
781
782    /// Require a contract to hold. `details` names the violation and is
783    /// built only when the check fails.
784    fn require<D: Into<String>>(
785        &self,
786        held: bool,
787        details: impl FnOnce() -> D,
788    ) -> Result<(), ConformanceError> {
789        if held {
790            return Ok(());
791        }
792        Err(self.fail(details()))
793    }
794
795    /// Record a verified sub-case.
796    fn note(&mut self, observation: impl Into<String>) {
797        self.observations.push(observation.into());
798    }
799
800    /// The wire family cannot spell this sequence shape: an explicit, named
801    /// skip that [`check_gated_outcome`] cross-checks against the suite's
802    /// declared capabilities.
803    fn skip(&self, reason: &'static str) -> ScenarioOutcome {
804        ScenarioOutcome::Skipped {
805            name: self.name,
806            provider: self.provider,
807            reason,
808        }
809    }
810
811    fn report(self) -> ScenarioReport {
812        ScenarioReport {
813            name: self.name,
814            provider: self.provider,
815            observations: self.observations,
816        }
817    }
818
819    fn ran(self) -> ScenarioOutcome {
820        ScenarioOutcome::Ran(self.report())
821    }
822}
823
824/// Truncation at every position — EOF before content, mid-text, mid-tool-args,
825/// after a tool call's frames — must preserve the content that ended and
826/// never produce a terminal record.
827///
828/// Pins the truncation family from round one (`rig-2257-code-review-findings-ec9f2625.md`):
829/// EOF without the provider's end event must not synthesize a successful
830/// usage-less terminal.
831pub async fn truncation_preserves_content_without_terminal(
832    fixture: &ProviderWireFixture,
833) -> Result<ScenarioReport, ConformanceError> {
834    let mut checks = Checks::new(
835        "truncation_preserves_content_without_terminal",
836        fixture.driver.provider,
837    );
838
839    // EOF before any content.
840    let drained = fixture.driver.drive(Vec::new()).await?;
841    checks.require(
842        !drained.has_terminal(),
843        || "an empty stream must not synthesize a terminal record",
844    )?;
845    checks.note("EOF before content: no terminal");
846
847    // EOF after text deltas.
848    let drained = fixture
849        .driver
850        .drive(ok_chunks(fixture.text_frames.clone()))
851        .await?;
852    checks.require(drained.texts() == fixture.expected_texts, || {
853        format!(
854            "text delivered before truncation must be preserved: expected {:?}, observed {:?}",
855            fixture.expected_texts,
856            drained.texts()
857        )
858    })?;
859    checks.require(
860        !drained.has_terminal(),
861        || "EOF after text deltas must not synthesize a terminal record",
862    )?;
863    checks.note("EOF mid-text: content preserved, no terminal");
864
865    // EOF mid-tool-arguments, where the wire streams arguments.
866    if let Some(partial) = &fixture.partial_tool_call_frames {
867        let drained = fixture.driver.drive(ok_chunks(partial.clone())).await?;
868        checks.require(
869            !drained.has_terminal(),
870            || "EOF mid-tool-arguments must not synthesize a terminal record",
871        )?;
872        checks.note("EOF mid-tool-args: no terminal");
873    }
874
875    // EOF after a tool call's frames, before the provider's end. A call is
876    // visible once it closes: on a wire that closes calls at the end, the
877    // cut call never surfaces.
878    let drained = fixture
879        .driver
880        .drive(ok_chunks(fixture.tool_call_frames.clone()))
881        .await?;
882    checks.require(
883        drained
884            .tool_call_names()
885            .iter()
886            .all(|name| *name == fixture.expected_tool_name),
887        || {
888            format!(
889                "only the delivered call may surface: observed {:?}",
890                drained.tool_call_names()
891            )
892        },
893    )?;
894    checks.require(
895        !drained.has_terminal(),
896        || "EOF after a tool call must not synthesize a terminal record",
897    )?;
898    checks.note("EOF after a tool call: no terminal");
899
900    Ok(checks.report())
901}
902
903/// A transport failure after a tool call's frames ends the reply: the `Err`
904/// is the last item, and there is no terminal record. A call the wire closed
905/// before the failure precedes it; one still open never surfaces.
906pub async fn transport_error_after_tool_call_yields_err_then_end(
907    fixture: &ProviderWireFixture,
908) -> Result<ScenarioReport, ConformanceError> {
909    let mut checks = Checks::new(
910        "transport_error_after_tool_call_yields_err_then_end",
911        fixture.driver.provider,
912    );
913
914    let mut chunks = ok_chunks(fixture.tool_call_frames.clone());
915    chunks.push(transport_error_chunk());
916    let drained = fixture.driver.drive(chunks).await?;
917
918    checks.require(
919        drained
920            .tool_call_names()
921            .iter()
922            .all(|name| *name == fixture.expected_tool_name),
923        || {
924            format!(
925                "only the delivered call may precede the transport error: observed {:?}",
926                drained.tool_call_names()
927            )
928        },
929    )?;
930    let error_index = drained
931        .first_error_index()
932        .ok_or_else(|| checks.fail("the transport failure must reach the consumer"))?;
933    checks.require(
934        error_index + 1 == drained.items.len(),
935        || "nothing may follow the terminal transport error",
936    )?;
937    checks.require(
938        !drained.has_terminal(),
939        || "a transport failure must not be papered over with a terminal record",
940    )?;
941
942    checks.note("Err, then end; no terminal");
943    Ok(checks.report())
944}
945
946/// A malformed frame between valid content and the provider's end ends the
947/// reply: it surfaces as the stream's last item, content before it is kept,
948/// and there is no response.
949pub async fn malformed_frame_ends_the_reply(
950    fixture: &ProviderWireFixture,
951) -> Result<ScenarioOutcome, ConformanceError> {
952    let mut checks = Checks::new("malformed_frame_ends_the_reply", fixture.driver.provider);
953    let Some(malformed) = &fixture.malformed_frame else {
954        return Ok(checks.skip("wire family cannot spell a frame-level decode failure"));
955    };
956
957    let frames = concat_frames(&[
958        &fixture.text_frames,
959        std::slice::from_ref(malformed),
960        &fixture.terminal_frames,
961    ]);
962    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
963
964    checks.require(drained.error_count() == 1, || {
965        format!(
966            "the malformed frame must surface as exactly one Err item, observed {}",
967            drained.error_count()
968        )
969    })?;
970    checks.require(
971        drained.first_error_index() == Some(drained.items.len() - 1),
972        || "the malformed frame's error must be the stream's last item",
973    )?;
974    checks.require(
975        drained.texts() == fixture.expected_texts,
976        || "content before the malformed frame must be preserved",
977    )?;
978    checks.require(
979        !drained.has_terminal(),
980        || "a reply cut by a corrupt frame has no response",
981    )?;
982
983    checks.note("Err surfaced last; content kept; no response");
984    Ok(checks.ran())
985}
986
987/// An event type the client does not know must be skipped without an error,
988/// and the stream must still complete.
989///
990/// Pins the unknown-event forward-compatibility policy (round three,
991/// `rig-2257-code-review-findings-8a2f41c7.md`).
992pub async fn unknown_event_is_skipped(
993    fixture: &ProviderWireFixture,
994) -> Result<ScenarioOutcome, ConformanceError> {
995    let mut checks = Checks::new("unknown_event_is_skipped", fixture.driver.provider);
996    let Some(unknown) = &fixture.unknown_event_frame else {
997        return Ok(checks.skip("wire family cannot spell an unknown event type"));
998    };
999
1000    let frames = concat_frames(&[
1001        &fixture.text_frames,
1002        std::slice::from_ref(unknown),
1003        &fixture.terminal_frames,
1004    ]);
1005    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1006
1007    checks.require(
1008        drained.error_count() == 0,
1009        || "an unknown event type must be skipped, not surfaced as an error",
1010    )?;
1011    checks.require(
1012        drained.texts() == fixture.expected_texts && drained.response.is_some(),
1013        || "the stream must deliver its content and complete around the skipped event",
1014    )?;
1015    // The frame is skipped semantically but observable verbatim on the raw
1016    // passthrough channel (openai-agents' raw-event precedent, #2258 item 5).
1017    checks.require(drained.unknown_values().len() == 1, || {
1018        format!(
1019            "exactly one Unknown passthrough item must surface for the unknown frame, \
1020             observed {}",
1021            drained.unknown_values().len()
1022        )
1023    })?;
1024
1025    // Control run without the unknown frame: the aggregated assistant choice
1026    // must be byte-identical — the passthrough item is never folded in.
1027    let control_frames = concat_frames(&[&fixture.text_frames, &fixture.terminal_frames]);
1028    let control = fixture.driver.drive(ok_chunks(control_frames)).await?;
1029    checks.require(
1030        drained.choice == control.choice,
1031        || "the unknown frame must not perturb the aggregated assistant choice",
1032    )?;
1033
1034    checks.note(
1035        "unknown event skipped semantically, surfaced on the raw channel, \
1036         choice unchanged, stream completed",
1037    );
1038    Ok(checks.ran())
1039}
1040
1041/// A *known* event whose payload is schema-defective ends the reply: it
1042/// surfaces as the stream's last item, and there is no response.
1043///
1044/// Pins the round-5 known-type strictness policy (a defective known event is
1045/// never demoted to an unknown one).
1046pub async fn defective_known_event_ends_the_reply(
1047    fixture: &ProviderWireFixture,
1048) -> Result<ScenarioOutcome, ConformanceError> {
1049    let mut checks = Checks::new(
1050        "defective_known_event_ends_the_reply",
1051        fixture.driver.provider,
1052    );
1053    let Some(defective) = &fixture.defective_known_frame else {
1054        return Ok(
1055            checks.skip("wire family cannot spell a known event with a schema-defective payload")
1056        );
1057    };
1058
1059    let frames = concat_frames(&[
1060        &fixture.text_frames,
1061        std::slice::from_ref(defective),
1062        &fixture.terminal_frames,
1063    ]);
1064    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1065
1066    checks.require(drained.error_count() == 1, || {
1067        format!(
1068            "a known event with a schema defect must surface exactly one Err item, observed {}",
1069            drained.error_count()
1070        )
1071    })?;
1072    checks.require(
1073        !drained.has_terminal(),
1074        || "a reply cut by a defective event has no response",
1075    )?;
1076
1077    checks.note("defective known event ended the reply with its Err");
1078    Ok(checks.ran())
1079}
1080
1081/// A delta-less choice (the Azure `prompt_filter_results` prelude) must be a
1082/// no-op — no error, no content, and the rest of the stream unaffected.
1083///
1084/// Pins the Azure prelude no-op from round two
1085/// (`rig-2257-code-review-findings-b91d03aa.md`).
1086pub async fn delta_less_choice_prelude_is_a_noop(
1087    fixture: &ProviderWireFixture,
1088) -> Result<ScenarioOutcome, ConformanceError> {
1089    let mut checks = Checks::new(
1090        "delta_less_choice_prelude_is_a_noop",
1091        fixture.driver.provider,
1092    );
1093    let Some(prelude) = &fixture.delta_less_prelude_frame else {
1094        return Ok(checks.skip("wire family has no delta-less prelude shape"));
1095    };
1096
1097    let frames = concat_frames(&[
1098        std::slice::from_ref(prelude),
1099        &fixture.text_frames,
1100        &fixture.terminal_frames,
1101    ]);
1102    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1103
1104    checks.require(
1105        drained.error_count() == 0,
1106        || "the delta-less prelude must not surface an error",
1107    )?;
1108    checks.require(
1109        drained.texts() == fixture.expected_texts && drained.response.is_some(),
1110        || "the prelude must not perturb content delivery or the terminal",
1111    )?;
1112
1113    checks.note("delta-less prelude ignored; stream unaffected");
1114    Ok(checks.ran())
1115}
1116
1117/// Refusal frames must deliver their text to the consumer without an error.
1118///
1119/// Pins the refusal-delta handling from round three
1120/// (`rig-2257-code-review-findings-8a2f41c7.md`).
1121pub async fn refusal_frames_deliver_text_without_error(
1122    fixture: &ProviderWireFixture,
1123) -> Result<ScenarioOutcome, ConformanceError> {
1124    let mut checks = Checks::new(
1125        "refusal_frames_deliver_text_without_error",
1126        fixture.driver.provider,
1127    );
1128    let Some(refusal) = &fixture.refusal else {
1129        return Ok(checks.skip("wire family has no refusal channel"));
1130    };
1131
1132    let frames = concat_frames(&[&refusal.frames, &fixture.terminal_frames]);
1133    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1134
1135    checks.require(
1136        drained.error_count() == 0,
1137        || "refusal content must not surface as an error",
1138    )?;
1139    let delivered = drained.texts().concat();
1140    checks.require(delivered == refusal.expected_text, || {
1141        format!(
1142            "refusal text must be delivered: expected {:?}, observed {delivered:?}",
1143            refusal.expected_text
1144        )
1145    })?;
1146    checks.require(
1147        drained.response.is_some(),
1148        || "a refused turn still ends with the provider's genuine terminal",
1149    )?;
1150
1151    checks.note("refusal text delivered without error");
1152    Ok(checks.ran())
1153}
1154
1155/// On the buffered-body pipeline (the ChatGPT backend), a terminal whose body
1156/// carries text never seen as a delta must merge that text into the choice
1157/// exactly once, and a body restating streamed deltas must not duplicate them.
1158///
1159/// Pins the terminal-body/delta per-kind merge from round five
1160/// (`rig-2257-code-review-findings-5c73639c.md`) and the empty-delta merge
1161/// direction verified in round six (`rig-2257-code-review-findings-34ee8ba5.md`
1162/// P3-2).
1163pub async fn terminal_body_content_merges_per_kind(
1164    driver: &BufferedBodyDriver,
1165    cases: Vec<(&'static str, String)>,
1166    expected_text: &str,
1167) -> Result<ScenarioReport, ConformanceError> {
1168    let mut checks = Checks::new("terminal_body_content_merges_per_kind", driver.provider);
1169
1170    for (label, body) in cases {
1171        let choice = driver.drive(body).await?;
1172        let choice_text: String = choice
1173            .iter()
1174            .filter_map(|content| match content {
1175                AssistantContent::Text(text) => Some(text.text.as_str()),
1176                _ => None,
1177            })
1178            .collect();
1179        let occurrences = choice_text.matches(expected_text).count();
1180        checks.require(occurrences == 1, || {
1181            format!(
1182                "{label}: terminal-body text must appear exactly once in the choice, observed {occurrences} in {choice_text:?}"
1183            )
1184        })?;
1185        checks.note(format!("{label}: text merged exactly once"));
1186    }
1187
1188    Ok(checks.report())
1189}
1190
1191/// A bare terminal signal after only-unparseable frames must not fabricate a
1192/// successful terminal record: the parse errors were already surfaced, and a
1193/// default-usage terminal would dress the failure up as success.
1194///
1195/// Pins the bare-`[DONE]` guard from round six
1196/// (`rig-2257-code-review-findings-5c73639c.md`, carried into `34ee8ba5`).
1197pub async fn bare_terminal_after_only_unparseable_frames_fabricates_nothing(
1198    fixture: &ProviderWireFixture,
1199) -> Result<ScenarioOutcome, ConformanceError> {
1200    let mut checks = Checks::new(
1201        "bare_terminal_after_only_unparseable_frames_fabricates_nothing",
1202        fixture.driver.provider,
1203    );
1204    let Some(bare_terminal) = &fixture.bare_terminal_frames else {
1205        return Ok(checks.skip("wire family has no data-less terminal signal"));
1206    };
1207    let Some(malformed) = &fixture.malformed_frame else {
1208        return Ok(checks.skip("wire family cannot spell a frame-level decode failure"));
1209    };
1210
1211    let frames = concat_frames(&[std::slice::from_ref(malformed), bare_terminal]);
1212    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1213
1214    checks.require(
1215        drained.error_count() != 0,
1216        || "the unparseable frame must surface as an Err item",
1217    )?;
1218    checks.require(
1219        !drained.has_terminal(),
1220        || "a bare terminal with no decoded frame must not fabricate a terminal record",
1221    )?;
1222
1223    checks.note("no fabricated terminal after only-unparseable frames");
1224    Ok(checks.ran())
1225}
1226
1227/// The genuine terminal must report the provider's usage; a terminal without
1228/// usage metrics must complete with no counter reported rather than being
1229/// suppressed or invented.
1230///
1231/// Pins the absent-usage contract on [`CompletionResponse::usage`](crate::completion::CompletionResponse::usage)
1232/// (round one, `rig-2257-code-review-findings-ec9f2625.md`).
1233pub async fn usage_variants_are_reported_or_absent(
1234    fixture: &ProviderWireFixture,
1235) -> Result<ScenarioReport, ConformanceError> {
1236    let mut checks = Checks::new(
1237        "usage_variants_are_reported_or_absent",
1238        fixture.driver.provider,
1239    );
1240
1241    let frames = concat_frames(&[&fixture.text_frames, &fixture.terminal_frames]);
1242    let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1243    let response = drained
1244        .response
1245        .as_ref()
1246        .ok_or_else(|| checks.fail("the genuine terminal must produce a record"))?;
1247    checks.require(
1248        response.usage.total_tokens == Some(fixture.expected_usage_total),
1249        || {
1250            format!(
1251                "terminal usage must be preserved: expected total {}, observed {:?}",
1252                fixture.expected_usage_total, response.usage.total_tokens
1253            )
1254        },
1255    )?;
1256    checks.require(
1257        response.finish_reason() == fixture.expected_finish_reason,
1258        || {
1259            format!(
1260                "terminal finish reason must be normalized: expected {:?}, observed {:?}",
1261                fixture.expected_finish_reason,
1262                response.finish_reason()
1263            )
1264        },
1265    )?;
1266    checks.note(format!(
1267        "usage total {} and finish reason {:?} preserved",
1268        fixture.expected_usage_total, fixture.expected_finish_reason
1269    ));
1270
1271    if let Some(zero_usage) = &fixture.zero_usage_terminal_frames {
1272        let frames = concat_frames(&[&fixture.text_frames, zero_usage]);
1273        let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1274        let response = drained.response.as_ref().ok_or_else(|| {
1275            checks.fail("a usage-less genuine terminal must still complete the stream")
1276        })?;
1277        checks.require(!response.usage.is_reported(), || {
1278            format!(
1279                "missing usage metrics must leave every counter unreported, not invented: {:?}",
1280                response.usage
1281            )
1282        })?;
1283        checks.note("usage-less terminal completed with no counter reported");
1284    }
1285
1286    Ok(checks.report())
1287}
1288
1289/// Reasoning-summary deltas followed by the item's full `output_item.done`
1290/// block must aggregate to the summary exactly once — the full block
1291/// supersedes its own deltas, never duplicates them.
1292///
1293/// Pins the open P1 in `rig-2257-code-review-findings-34ee8ba5.md` ("OpenAI
1294/// Responses reasoning-summary streams duplicate reasoning content"):
1295/// `reasoning_summary_text.delta` drops `item_id`, so the strict same-item
1296/// table appends the full block beside the delta-built item.
1297pub async fn reasoning_summary_deltas_are_superseded_without_duplication(
1298    driver: &WireDriver,
1299    frames: Vec<WireInput>,
1300    summary_text: &str,
1301) -> Result<ScenarioReport, ConformanceError> {
1302    let mut checks = Checks::new(
1303        "reasoning_summary_deltas_are_superseded_without_duplication",
1304        driver.provider,
1305    );
1306
1307    let drained = driver.drive(ok_chunks(frames)).await?;
1308    checks.require(
1309        drained.completed_cleanly(),
1310        || "the reasoning stream must complete without errors",
1311    )?;
1312    let reasoning = drained.choice_reasoning();
1313    let occurrences: usize = reasoning
1314        .iter()
1315        .flat_map(|item| item.content.iter())
1316        .filter(|content| match content {
1317            crate::message::ReasoningContent::Summary(text)
1318            | crate::message::ReasoningContent::Text { text, .. } => text.contains(summary_text),
1319            _ => false,
1320        })
1321        .count();
1322    checks.require(occurrences == 1, || {
1323        format!(
1324            "the summary must appear exactly once in the aggregated choice, observed {occurrences} across {reasoning:?}"
1325        )
1326    })?;
1327    checks.require(reasoning.len() == 1, || {
1328        format!(
1329            "deltas and their full block must collapse to one reasoning item, observed {}",
1330            reasoning.len()
1331        )
1332    })?;
1333
1334    checks.note("summary aggregated exactly once");
1335    Ok(checks.report())
1336}
1337
1338/// A reasoning item whose `output_item.done` carries several parts under one
1339/// item id (summary parts, text, encrypted) must keep every part, in order —
1340/// same-id sibling blocks append, they never replace each other.
1341///
1342/// Pins the open P1 in `rig-2257-code-review-findings-34ee8ba5.md` ("The by-id
1343/// fallback collapses multi-part same-id reasoning items"): the `rposition`
1344/// fallback replaces the just-appended same-id sibling.
1345pub async fn multi_part_same_id_reasoning_keeps_every_part(
1346    driver: &WireDriver,
1347    frames: Vec<WireInput>,
1348    expected_parts: &[&str],
1349) -> Result<ScenarioReport, ConformanceError> {
1350    let mut checks = Checks::new(
1351        "multi_part_same_id_reasoning_keeps_every_part",
1352        driver.provider,
1353    );
1354
1355    let drained = driver.drive(ok_chunks(frames)).await?;
1356    checks.require(
1357        drained.completed_cleanly(),
1358        || "the reasoning stream must complete without errors",
1359    )?;
1360    let observed: Vec<String> = drained
1361        .choice_reasoning()
1362        .iter()
1363        .flat_map(|item| item.content.iter())
1364        .map(|content| match content {
1365            crate::message::ReasoningContent::Summary(text) => text.clone(),
1366            crate::message::ReasoningContent::Text { text, .. } => text.clone(),
1367            crate::message::ReasoningContent::Encrypted(data) => data.clone(),
1368            crate::message::ReasoningContent::Redacted { data } => data.clone(),
1369        })
1370        .collect();
1371    checks.require(observed == expected_parts, || {
1372        format!(
1373            "every same-id reasoning part must survive in order: expected {expected_parts:?}, observed {observed:?}"
1374        )
1375    })?;
1376
1377    checks.note(format!(
1378        "all {} reasoning parts survived",
1379        expected_parts.len()
1380    ));
1381    Ok(checks.report())
1382}
1383
1384/// Reasoning deltas interleaved with a tool call, then the item's completed
1385/// block, must aggregate to exactly one reasoning item carrying the block's
1386/// content.
1387///
1388/// Pins the interleaved-reasoning replacement contract on
1389/// completed reasoning block (round six,
1390/// `rig-2257-code-review-findings-34ee8ba5.md`, "Verified sound" section).
1391pub async fn interleaved_reasoning_aggregates_to_one_item(
1392    driver: &WireDriver,
1393    frames: Vec<WireInput>,
1394    expected_text: &str,
1395) -> Result<ScenarioReport, ConformanceError> {
1396    let mut checks = Checks::new(
1397        "interleaved_reasoning_aggregates_to_one_item",
1398        driver.provider,
1399    );
1400
1401    let drained = driver.drive(ok_chunks(frames)).await?;
1402    checks.require(
1403        drained.completed_cleanly(),
1404        || "the interleaved stream must complete without errors",
1405    )?;
1406    let reasoning = drained.choice_reasoning();
1407    checks.require(reasoning.len() == 1, || {
1408        format!(
1409            "interleaved deltas and their completed block must collapse to one reasoning item, observed {}",
1410            reasoning.len()
1411        )
1412    })?;
1413    let carries_text = reasoning
1414        .iter()
1415        .flat_map(|item| item.content.iter())
1416        .any(|content| match content {
1417            crate::message::ReasoningContent::Summary(text)
1418            | crate::message::ReasoningContent::Text { text, .. } => text == expected_text,
1419            _ => false,
1420        });
1421    checks.require(carries_text, || {
1422        format!("the reasoning item must carry the completed block's text {expected_text:?}")
1423    })?;
1424
1425    checks.note("exactly one reasoning item with the completed content");
1426    Ok(checks.report())
1427}
1428
1429/// On a constant-id wire (a boundary-minted per-stream reasoning id), other
1430/// output closes the open reasoning item: thought → tool call → thought must
1431/// aggregate as `[Reasoning(first), ToolCall, Reasoning(second)]` — two items
1432/// in arrival order, never one merged item that misorders history on replay.
1433///
1434/// Pins the F1b ordering dimension of the #2258 review (main's
1435/// "other output closes the reasoning item" boundary, lost when identity
1436/// became the per-stream constant).
1437pub async fn interleaved_constant_id_reasoning_preserves_order(
1438    fixture: &ProviderWireFixture,
1439) -> Result<ScenarioOutcome, ConformanceError> {
1440    let mut checks = Checks::new(
1441        "interleaved_constant_id_reasoning_preserves_order",
1442        fixture.driver.provider,
1443    );
1444    let Some(interleaved) = &fixture.interleaved_reasoning else {
1445        return Ok(checks.skip("wire fixture supplies no interleaved reasoning frames"));
1446    };
1447
1448    let drained = fixture
1449        .driver
1450        .drive(ok_chunks(interleaved.frames.clone()))
1451        .await?;
1452    checks.require(
1453        drained.completed_cleanly(),
1454        || "the interleaved stream must complete without errors",
1455    )?;
1456    assert_reasoning_tool_reasoning(
1457        &checks,
1458        &drained,
1459        interleaved.first_reasoning,
1460        interleaved.tool_name,
1461        interleaved.second_reasoning,
1462    )?;
1463
1464    checks.note("boundary kept: reasoning, tool call, reasoning in order");
1465    Ok(checks.ran())
1466}
1467
1468/// On a constant-id wire whose completed reasoning block arrives as a signed
1469/// full restatement (gemini `thoughtSignature`), a full block *after*
1470/// interleaved output must not replace-and-discard the thought accumulated
1471/// before the boundary: the choice keeps `[Reasoning(first), ToolCall,
1472/// Reasoning(second, signed)]`.
1473///
1474/// Pins the F1b erasure dimension of the #2258 review, on top of the F1
1475/// adapter fix (the signed chunk restates only post-boundary fragments).
1476pub async fn interleaved_signed_full_reasoning_does_not_erase_prior_thought(
1477    driver: &WireDriver,
1478    frames: Vec<WireInput>,
1479    first: &str,
1480    tool_name: &str,
1481    second: &str,
1482) -> Result<ScenarioReport, ConformanceError> {
1483    let mut checks = Checks::new(
1484        "interleaved_signed_full_reasoning_does_not_erase_prior_thought",
1485        driver.provider,
1486    );
1487
1488    let drained = driver.drive(ok_chunks(frames)).await?;
1489    checks.require(
1490        drained.completed_cleanly(),
1491        || "the interleaved stream must complete without errors",
1492    )?;
1493    assert_reasoning_tool_reasoning(&checks, &drained, first, tool_name, second)?;
1494    let signed = drained.choice_reasoning().last().is_some_and(|reasoning| {
1495        reasoning.content.iter().any(|content| {
1496            matches!(
1497                content,
1498                crate::message::ReasoningContent::Text {
1499                    signature: Some(_),
1500                    ..
1501                }
1502            )
1503        })
1504    });
1505    checks.require(signed, || "the post-boundary block must keep its signature")?;
1506
1507    checks.note("pre-boundary thought survived; signed block completed the post-boundary part");
1508    Ok(checks.report())
1509}
1510
1511/// Shared assertion: the aggregated choice is exactly
1512/// `[Reasoning(first), ToolCall(tool_name), Reasoning(second…)]`.
1513fn assert_reasoning_tool_reasoning(
1514    checks: &Checks,
1515    drained: &DrainedStream,
1516    first: &str,
1517    tool_name: &str,
1518    second: &str,
1519) -> Result<(), ConformanceError> {
1520    let shape: Vec<String> = drained
1521        .choice
1522        .iter()
1523        .map(|content| match content {
1524            AssistantContent::Reasoning(reasoning) => {
1525                let text: String = reasoning
1526                    .value()
1527                    .content
1528                    .iter()
1529                    .filter_map(|content| match content {
1530                        crate::message::ReasoningContent::Summary(text)
1531                        | crate::message::ReasoningContent::Text { text, .. } => {
1532                            Some(text.as_str())
1533                        }
1534                        _ => None,
1535                    })
1536                    .collect();
1537                format!("reasoning:{text}")
1538            }
1539            AssistantContent::ToolCall(tool_call) => {
1540                format!("tool:{}", tool_call.function.name)
1541            }
1542            AssistantContent::Text(text) => format!("text:{}", text.text),
1543            AssistantContent::Image(_) => "image".to_string(),
1544        })
1545        .collect();
1546    let expected = vec![
1547        format!("reasoning:{first}"),
1548        format!("tool:{tool_name}"),
1549        format!("reasoning:{second}"),
1550    ];
1551    checks.require(shape == expected, || {
1552        format!("the boundary must survive aggregation: expected {expected:?}, observed {shape:?}")
1553    })
1554}
1555
1556/// Per-provider wire fixtures for the shared scenario set.
1557pub mod fixtures {
1558    use super::*;
1559    use crate::driver::{Model, Transport};
1560    use crate::operation::Completion;
1561    use crate::test_utils::SequencedStreamingHttpClient;
1562    use crate::wire::Wire;
1563    use serde_json::json;
1564
1565    /// Drain a full normalized stream into everything the consumer observed.
1566    /// Public so provider-crate conformance suites (the typed-event wires)
1567    /// can reuse it in their drivers.
1568    pub async fn drain(mut stream: crate::streaming::CompletionStream) -> DrainedStream {
1569        let mut items = Vec::new();
1570        while let Some(item) = stream.next().await {
1571            items.push(item.map_err(|error| ErrorReport::from(&error)));
1572        }
1573        let partial = stream.partial().choice;
1574        let response = stream.finish().await.ok();
1575        let drained = DrainedStream {
1576            items,
1577            choice: response
1578                .as_ref()
1579                .map_or(partial, |response| response.choice.clone()),
1580            response,
1581        };
1582        // Every fixture and cassette that drains through this helper runs
1583        // the lifecycle validator — the prose invariants as one executable
1584        // artifact (#2258 C1).
1585        super::assert_valid_event_stream(&drained.items, &drained.choice);
1586        drained
1587    }
1588
1589    /// Drive `model` through the observed stream entry and drain it. Every
1590    /// wire threads the context it is handed: the trace must show the
1591    /// request it sent and how the attempt closed, or the wire has silently
1592    /// taken the context-discarding default.
1593    pub async fn drain_observed<W, T>(
1594        model: &Model<W, T>,
1595        request: crate::completion::CompletionRequest,
1596    ) -> Result<DrainedStream, ProviderError>
1597    where
1598        W: Wire<Op = Completion>,
1599        T: Transport<W>,
1600    {
1601        let log = std::sync::Arc::new(crate::observe::ObservationLog::default());
1602        let context = crate::observe::AdapterContext::new(
1603            log.clone(),
1604            crate::observe::Subject::default(),
1605            "conformance",
1606        );
1607        // A stream that fails to open still sent (or failed to send) a
1608        // request: the facts are asserted before the error propagates.
1609        let drained = match model.stream_observed(request, context) {
1610            Ok(stream) => Ok(drain(stream).await),
1611            Err(error) => Err(error),
1612        };
1613        let events: Vec<_> = log
1614            .trace()
1615            .observations
1616            .iter()
1617            .filter_map(|o| match &o.action {
1618                crate::observe::Action::Adapter { observation } => Some(observation.event.clone()),
1619                _ => None,
1620            })
1621            .collect();
1622        assert!(
1623            matches!(
1624                events.first(),
1625                Some(crate::observe::AdapterEvent::Started { .. })
1626            ),
1627            "the wire must attach the observation context it was handed: {events:?}"
1628        );
1629        assert!(
1630            matches!(
1631                events.last(),
1632                Some(crate::observe::AdapterEvent::Finished { .. })
1633            ),
1634            "the attempt must close: {events:?}"
1635        );
1636        drained
1637    }
1638
1639    /// Lower fixture frames onto the byte transport a `SequencedStreamingHttpClient`
1640    /// replays. Only byte frames are valid here — an event frame in a
1641    /// byte-driver fixture is a fixture authoring error.
1642    fn byte_chunks(chunks: WireChunks) -> Result<Vec<http_client::Result<Bytes>>, ProviderError> {
1643        chunks
1644            .into_iter()
1645            .map(|chunk| match chunk {
1646                Ok(WireInput::Bytes(bytes)) => Ok(Ok(bytes)),
1647                Ok(WireInput::Event(_)) => Err(ProviderError::Provider(
1648                    "typed-event frame fed to a byte-transport driver".to_string(),
1649                )),
1650                Err(error) => Ok(Err(error)),
1651            })
1652            .collect()
1653    }
1654
1655    /// Drive a byte-transport wire over the scripted chunks: bind the
1656    /// provider's wire to the replaying transport, open the model `bind`
1657    /// names, and drain the observed stream.
1658    ///
1659    /// Every byte-wire conformance driver is this walk; since the wire
1660    /// unification they differ only in which wire they name, so the walk is
1661    /// written once and each family supplies its own `bind`.
1662    fn byte_driver<W>(
1663        provider: &'static str,
1664        bind: fn(SequencedStreamingHttpClient) -> Model<W, SequencedStreamingHttpClient>,
1665    ) -> WireDriver
1666    where
1667        W: Wire<Op = Completion, Payload = crate::wire::Encoded, Frame = crate::wire::WireFrame>,
1668    {
1669        WireDriver::new(provider, move |chunks| {
1670            Box::pin(async move {
1671                let model = bind(SequencedStreamingHttpClient::new(byte_chunks(chunks)?));
1672                let request = CompletionRequest::new("hello");
1673                drain_observed(&model, request).await
1674            })
1675        })
1676    }
1677
1678    fn sse(frame: &serde_json::Value) -> WireInput {
1679        WireInput::Bytes(Bytes::from(format!("data: {frame}\n\n")))
1680    }
1681
1682    fn sse_raw(data: &str) -> WireInput {
1683        WireInput::Bytes(Bytes::from(format!("data: {data}\n\n")))
1684    }
1685
1686    fn ndjson(frame: &serde_json::Value) -> WireInput {
1687        WireInput::Bytes(Bytes::from(format!("{frame}\n")))
1688    }
1689
1690    /// The frame's SSE text, for buffered-body pipelines that re-parse a
1691    /// whole body string.
1692    fn frame_text(frame: &WireInput) -> String {
1693        frame
1694            .as_bytes()
1695            .map(|bytes| String::from_utf8_lossy(bytes).into_owned())
1696            .unwrap_or_default()
1697    }
1698
1699    /// OpenAI chat-completions wire (the shared OpenAI-compatible SSE path).
1700    pub mod openai_chat {
1701        use super::*;
1702
1703        fn driver() -> WireDriver {
1704            byte_driver("openai", |transport| {
1705                crate::driver::Model::new(
1706                    crate::providers::openai::wire::OpenAIConfig::with_key(
1707                        &crate::providers::openai::wire::OPENAI,
1708                        "test-key",
1709                    )
1710                    .chat("gpt-4o"),
1711                    transport,
1712                )
1713            })
1714        }
1715
1716        /// The chat-completions fixture.
1717        pub fn fixture() -> ProviderWireFixture {
1718            ProviderWireFixture {
1719                driver: driver(),
1720                text_frames: vec![sse(&json!({
1721                    "id": "chatcmpl-1",
1722                    "model": "gpt-4o-2024-08-06",
1723                    "choices": [{"index": 0, "delta": {"content": "hi"}, "finish_reason": null}],
1724                    "usage": null,
1725                }))],
1726                expected_texts: vec!["hi"],
1727                tool_call_frames: vec![
1728                    sse(&json!({
1729                        "choices": [{"index": 0, "delta": {"tool_calls": [{
1730                            "index": 0,
1731                            "id": "call_1",
1732                            "type": "function",
1733                            "function": {"name": "get_weather", "arguments": ""},
1734                        }]}, "finish_reason": null}],
1735                    })),
1736                    sse(&json!({
1737                        "choices": [{"index": 0, "delta": {"tool_calls": [{
1738                            "index": 0,
1739                            "function": {"arguments": "{\"city\":\"Tokyo\"}"},
1740                        }]}, "finish_reason": null}],
1741                    })),
1742                    // No `finish_reason` chunk: on the chat wire that IS the
1743                    // terminal signal, and these frames must stop short of it.
1744                    // EOF/error cleanup still flushes the completed call.
1745                ],
1746                expected_tool_name: "get_weather",
1747                partial_tool_call_frames: Some(vec![sse(&json!({
1748                    "choices": [{"index": 0, "delta": {"tool_calls": [{
1749                        "index": 0,
1750                        "id": "call_1",
1751                        "type": "function",
1752                        "function": {"name": "get_weather", "arguments": "{\"cit"},
1753                    }]}, "finish_reason": null}],
1754                }))]),
1755                terminal_frames: vec![
1756                    sse(&json!({
1757                        "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}],
1758                        "usage": null,
1759                    })),
1760                    sse(&json!({
1761                        "choices": [],
1762                        "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15},
1763                    })),
1764                    sse_raw("[DONE]"),
1765                ],
1766                expected_usage_total: 15,
1767                expected_finish_reason: Some(FinishReason::Stop),
1768                zero_usage_terminal_frames: Some(vec![
1769                    sse(&json!({
1770                        "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}],
1771                        "usage": null,
1772                    })),
1773                    sse_raw("[DONE]"),
1774                ]),
1775                bare_terminal_frames: Some(vec![sse_raw("[DONE]")]),
1776                malformed_frame: Some(sse_raw("{not json")),
1777                unknown_event_frame: None,
1778                // A wrongly-typed `content` is tolerated by the lenient delta
1779                // decode; a wrongly-typed `choices` is a genuine schema defect
1780                // of the known chunk shape.
1781                defective_known_frame: Some(sse_raw(r#"{"choices": 42}"#)),
1782                // The Azure `prompt_filter_results` prelude: a choice with no
1783                // `delta` at all.
1784                delta_less_prelude_frame: Some(sse_raw(
1785                    r#"{"id":"","object":"","choices":[{"prompt_index":0,"content_filter_results":{"hate":{"filtered":false,"severity":"safe"}}}]}"#,
1786                )),
1787                refusal: None,
1788                // Deliberately absent — a documented named skip, not a gap.
1789                // The chat wire streams tool calls as fragments that only
1790                // finalize at a boundary the wire itself signals (next slot,
1791                // finish_reason, terminal), so the AGGREGATED part order
1792                // cannot pin reasoning→tool→reasoning without risky early
1793                // finalization; and the chat request format erases part
1794                // order on replay regardless (`tool_calls` is a flat array
1795                // beside `content`). The boundary the adapter does own —
1796                // closing the open reasoning block before emitting tool
1797                // content — is pinned at emission level by the adapter's
1798                // unit tests and by the driver's debug-mode sequence laws.
1799                interleaved_reasoning: None,
1800            }
1801        }
1802    }
1803
1804    /// OpenAI Responses API wire.
1805    pub mod openai_responses {
1806        use super::*;
1807
1808        /// The driver alone, for the reasoning-specific scenarios.
1809        pub fn driver() -> WireDriver {
1810            byte_driver("openai", |transport| {
1811                crate::driver::Model::new(
1812                    crate::providers::openai::OpenAIConfig::new("test-key").responses("gpt-5.4"),
1813                    transport,
1814                )
1815            })
1816        }
1817
1818        fn completed_response(
1819            usage: Option<&serde_json::Value>,
1820            output: &serde_json::Value,
1821        ) -> serde_json::Value {
1822            json!({
1823                "id": "resp_1",
1824                "object": "response",
1825                "created_at": 0,
1826                "status": "completed",
1827                "model": "gpt-5.4",
1828                "output": output,
1829                "tools": [],
1830                "usage": usage,
1831            })
1832        }
1833
1834        fn terminal(usage: Option<&serde_json::Value>, output: &serde_json::Value) -> WireInput {
1835            sse(&json!({
1836                "type": "response.completed",
1837                "sequence_number": 99,
1838                "response": completed_response(usage, output),
1839            }))
1840        }
1841
1842        fn usage_json() -> serde_json::Value {
1843            json!({
1844                "input_tokens": 10,
1845                "output_tokens": 5,
1846                "output_tokens_details": {"reasoning_tokens": 0},
1847                "total_tokens": 15,
1848            })
1849        }
1850
1851        fn text_delta(text: &str) -> WireInput {
1852            sse(&json!({
1853                "type": "response.output_text.delta",
1854                "content_index": 0,
1855                "delta": text,
1856                "item_id": "msg_1",
1857                "output_index": 0,
1858                "sequence_number": 1,
1859            }))
1860        }
1861
1862        fn tool_call_done() -> WireInput {
1863            sse(&json!({
1864                "type": "response.output_item.done",
1865                "output_index": 0,
1866                "sequence_number": 2,
1867                "item": {
1868                    "type": "function_call",
1869                    "id": "fc_1",
1870                    "arguments": "{\"city\":\"Tokyo\"}",
1871                    "call_id": "call_1",
1872                    "name": "get_weather",
1873                    "status": "completed",
1874                },
1875            }))
1876        }
1877
1878        /// Synthetic twin of the recorded
1879        /// `openai/streaming_grammar/incomplete_mid_tool_call` cassette: a
1880        /// forced tool call cut by `max_output_tokens` mid-arguments. The wire
1881        /// restates the call on `response.output_item.done` with the
1882        /// arguments truncated mid-JSON and item status `incomplete`, then
1883        /// ends with a genuine `response.incomplete` terminal.
1884        pub fn incomplete_mid_tool_call_frames() -> Vec<WireInput> {
1885            vec![
1886                sse(&json!({
1887                    "type": "response.output_item.added",
1888                    "output_index": 0,
1889                    "sequence_number": 1,
1890                    "item": {
1891                        "type": "function_call",
1892                        "id": "fc_1",
1893                        "arguments": "",
1894                        "call_id": "call_1",
1895                        "name": "add",
1896                        "status": "in_progress",
1897                    },
1898                })),
1899                sse(&json!({
1900                    "type": "response.function_call_arguments.delta",
1901                    "item_id": "fc_1",
1902                    "output_index": 0,
1903                    "sequence_number": 2,
1904                    "delta": "{\"x",
1905                })),
1906                sse(&json!({
1907                    "type": "response.function_call_arguments.delta",
1908                    "item_id": "fc_1",
1909                    "output_index": 0,
1910                    "sequence_number": 3,
1911                    "delta": "\":48151",
1912                })),
1913                sse(&json!({
1914                    "type": "response.function_call_arguments.done",
1915                    "item_id": "fc_1",
1916                    "output_index": 0,
1917                    "sequence_number": 4,
1918                    "arguments": "{\"x\":48151",
1919                })),
1920                sse(&json!({
1921                    "type": "response.output_item.done",
1922                    "output_index": 0,
1923                    "sequence_number": 5,
1924                    "item": {
1925                        "type": "function_call",
1926                        "id": "fc_1",
1927                        "arguments": "{\"x\":48151",
1928                        "call_id": "call_1",
1929                        "name": "add",
1930                        "status": "incomplete",
1931                    },
1932                })),
1933                sse(&json!({
1934                    "type": "response.incomplete",
1935                    "sequence_number": 6,
1936                    "response": {
1937                        "id": "resp_1",
1938                        "object": "response",
1939                        "created_at": 0,
1940                        "status": "incomplete",
1941                        "incomplete_details": {"reason": "max_output_tokens"},
1942                        "model": "gpt-5.4",
1943                        "output": [{
1944                            "type": "function_call",
1945                            "id": "fc_1",
1946                            "arguments": "{\"x\":48151",
1947                            "call_id": "call_1",
1948                            "name": "add",
1949                            "status": "incomplete",
1950                        }],
1951                        "tools": [],
1952                        "usage": usage_json(),
1953                    },
1954                })),
1955            ]
1956        }
1957
1958        fn reasoning_done_item(
1959            id: &str,
1960            summary: &serde_json::Value,
1961            content: &serde_json::Value,
1962            encrypted: Option<&str>,
1963        ) -> WireInput {
1964            let mut item = json!({
1965                "type": "reasoning",
1966                "id": id,
1967                "summary": summary,
1968                "content": content,
1969                "status": "completed",
1970            });
1971            if let (Some(encrypted), Some(object)) = (encrypted, item.as_object_mut()) {
1972                object.insert("encrypted_content".to_string(), json!(encrypted));
1973            }
1974            sse(&json!({
1975                "type": "response.output_item.done",
1976                "output_index": 0,
1977                "sequence_number": 3,
1978                "item": item,
1979            }))
1980        }
1981
1982        /// The Responses-API fixture.
1983        pub fn fixture() -> ProviderWireFixture {
1984            ProviderWireFixture {
1985                driver: driver(),
1986                text_frames: vec![text_delta("hi")],
1987                expected_texts: vec!["hi"],
1988                tool_call_frames: vec![tool_call_done()],
1989                expected_tool_name: "get_weather",
1990                partial_tool_call_frames: Some(vec![
1991                    sse(&json!({
1992                        "type": "response.output_item.added",
1993                        "output_index": 0,
1994                        "sequence_number": 1,
1995                        "item": {
1996                            "type": "function_call",
1997                            "id": "fc_1",
1998                            "arguments": "",
1999                            "call_id": "call_1",
2000                            "name": "get_weather",
2001                            "status": "in_progress",
2002                        },
2003                    })),
2004                    sse(&json!({
2005                        "type": "response.function_call_arguments.delta",
2006                        "item_id": "fc_1",
2007                        "output_index": 0,
2008                        "sequence_number": 2,
2009                        "delta": "{\"cit",
2010                    })),
2011                ]),
2012                terminal_frames: vec![terminal(Some(&usage_json()), &json!([]))],
2013                expected_usage_total: 15,
2014                expected_finish_reason: Some(FinishReason::Stop),
2015                zero_usage_terminal_frames: Some(vec![terminal(None, &json!([]))]),
2016                bare_terminal_frames: None,
2017                malformed_frame: Some(sse_raw("{not json")),
2018                unknown_event_frame: Some(sse(&json!({
2019                    "type": "response.web_search_call.searching",
2020                    "output_index": 0,
2021                    "sequence_number": 4,
2022                    "item_id": "ws_1",
2023                }))),
2024                // The P2 probe shape from `rig-2257-code-review-findings-34ee8ba5.md`:
2025                // a known part tag (`output_text`) with a schema-defective payload.
2026                defective_known_frame: Some(sse(&json!({
2027                    "type": "response.content_part.added",
2028                    "item_id": "msg_1",
2029                    "output_index": 0,
2030                    "content_index": 0,
2031                    "sequence_number": 5,
2032                    "part": {"type": "output_text", "text": 42},
2033                }))),
2034                delta_less_prelude_frame: None,
2035                refusal: Some(RefusalFixture {
2036                    frames: vec![sse(&json!({
2037                        "type": "response.refusal.delta",
2038                        "content_index": 0,
2039                        "delta": "I cannot help with that.",
2040                        "item_id": "msg_1",
2041                        "output_index": 0,
2042                        "sequence_number": 1,
2043                    }))],
2044                    expected_text: "I cannot help with that.",
2045                }),
2046                interleaved_reasoning: None,
2047            }
2048        }
2049
2050        /// The buffered-body pipeline the ChatGPT backend uses: the SSE body
2051        /// is re-parsed after the fact and merged with the terminal response
2052        /// body, per content kind.
2053        ///
2054        /// Drives the *real* entry — `Model::call` on the
2055        /// ChatGPT wire whose HTTP double answers the `/responses` POST with
2056        /// the scripted SSE body — so the scenario exercises the buffered
2057        /// fold itself rather than a mirrored copy of it (#2258 review, F8
2058        /// drift risk).
2059        pub fn buffered_driver() -> BufferedBodyDriver {
2060            BufferedBodyDriver::new("chatgpt", |body| {
2061                Box::pin(async move {
2062                    let model = crate::driver::Model::new(
2063                        crate::providers::openai::OpenAIConfig::with_key(
2064                            &crate::providers::chatgpt::DIALECT,
2065                            "test-token",
2066                        )
2067                        .with_account_id("account-id")
2068                        .responses("gpt-5.4"),
2069                        crate::test_utils::RecordingHttpClient::new(body),
2070                    );
2071                    let request = CompletionRequest::new("hello");
2072                    let response = model.call(request).await?;
2073                    Ok(response.choice)
2074                })
2075            })
2076        }
2077
2078        fn message_output(text: &str) -> serde_json::Value {
2079            json!([{
2080                "type": "message",
2081                "id": "msg_1",
2082                "role": "assistant",
2083                "status": "completed",
2084                "content": [{"type": "output_text", "text": text, "annotations": []}],
2085            }])
2086        }
2087
2088        /// A terminal whose body carries text never seen as a delta.
2089        pub fn terminal_body_only_sse_body(text: &str) -> String {
2090            frame_text(&terminal(Some(&usage_json()), &message_output(text)))
2091        }
2092
2093        /// A streamed delta plus a terminal body restating the same text.
2094        pub fn terminal_body_and_delta_sse_body(text: &str) -> String {
2095            let frames = [
2096                text_delta(text),
2097                terminal(Some(&usage_json()), &message_output(text)),
2098            ];
2099            frames.iter().map(frame_text).collect()
2100        }
2101
2102        /// A streamed delta whose terminal body carries no output items — the
2103        /// gpt-5.x shape the buffered fallback exists for.
2104        pub fn delta_only_sse_body(text: &str) -> String {
2105            let frames = [text_delta(text), terminal(Some(&usage_json()), &json!([]))];
2106            frames.iter().map(frame_text).collect()
2107        }
2108
2109        /// The ChatGPT envelope-less replay shape (#2258 F3): a summary delta
2110        /// with NO envelope bookkeeping at all (repair injects
2111        /// `output_index: 0`, minting the `output-0` identity), then the
2112        /// item's envelope-full `output_item.done` restating the summary
2113        /// under its real `rs_*` id, then the terminal. The done item must
2114        /// adopt the minted per-slot identity and supersede the delta build.
2115        pub fn envelope_less_reasoning_supersede_sse_body() -> (String, &'static str) {
2116            let delta = json!({
2117                "type": "response.reasoning_summary_text.delta",
2118                "delta": "step 1",
2119            });
2120            let frames = [
2121                sse(&delta),
2122                reasoning_done_item(
2123                    "rs_1",
2124                    &json!([{"type": "summary_text", "text": "step 1"}]),
2125                    &json!([]),
2126                    None,
2127                ),
2128                terminal(Some(&usage_json()), &json!([])),
2129            ];
2130            (frames.iter().map(frame_text).collect(), "step 1")
2131        }
2132
2133        /// Summary deltas followed by their item's full `output_item.done`
2134        /// block, then the terminal. The deltas carry `item_id` on the wire;
2135        /// the full block restates the summary.
2136        pub fn reasoning_summary_supersede_frames() -> (Vec<WireInput>, &'static str) {
2137            let frames = vec![
2138                sse(&json!({
2139                    "type": "response.reasoning_summary_text.delta",
2140                    "item_id": "rs_1",
2141                    "output_index": 0,
2142                    "summary_index": 0,
2143                    "sequence_number": 1,
2144                    "delta": "step 1",
2145                })),
2146                reasoning_done_item(
2147                    "rs_1",
2148                    &json!([{"type": "summary_text", "text": "step 1"}]),
2149                    &json!([]),
2150                    None,
2151                ),
2152                terminal(Some(&usage_json()), &json!([])),
2153            ];
2154            (frames, "step 1")
2155        }
2156
2157        /// One reasoning item done-block carrying two summary parts, visible
2158        /// text, and encrypted content under a single item id.
2159        pub fn multi_part_reasoning_frames() -> (Vec<WireInput>, Vec<&'static str>) {
2160            let frames = vec![
2161                reasoning_done_item(
2162                    "rs_1",
2163                    &json!([
2164                        {"type": "summary_text", "text": "s1"},
2165                        {"type": "summary_text", "text": "s2"},
2166                    ]),
2167                    &json!([{"type": "reasoning_text", "text": "visible"}]),
2168                    Some("enc_blob"),
2169                ),
2170                terminal(Some(&usage_json()), &json!([])),
2171            ];
2172            (frames, vec!["s1", "s2", "visible", "enc_blob"])
2173        }
2174
2175        /// A reasoning delta, an interleaved tool call, then the reasoning
2176        /// item's completed block and the terminal.
2177        pub fn interleaved_reasoning_frames() -> (Vec<WireInput>, &'static str) {
2178            let frames = vec![
2179                sse(&json!({
2180                    "type": "response.reasoning_text.delta",
2181                    "item_id": "rs_2",
2182                    "output_index": 0,
2183                    "content_index": 0,
2184                    "sequence_number": 1,
2185                    "delta": "thinking",
2186                })),
2187                tool_call_done(),
2188                reasoning_done_item(
2189                    "rs_2",
2190                    &json!([]),
2191                    &json!([{"type": "reasoning_text", "text": "full reasoning"}]),
2192                    None,
2193                ),
2194                terminal(Some(&usage_json()), &json!([])),
2195            ];
2196            (frames, "full reasoning")
2197        }
2198    }
2199
2200    /// Gemini REST (`streamGenerateContent`) SSE wire.
2201    pub mod gemini_rest {
2202        use super::*;
2203
2204        fn driver() -> WireDriver {
2205            byte_driver("gemini", |transport| {
2206                crate::driver::Model::new(
2207                    crate::providers::gemini::GeminiConfig::new("test-key").completion(
2208                        crate::providers::gemini::completion::GEMINI_2_5_PRO_PREVIEW_06_05,
2209                    ),
2210                    transport,
2211                )
2212            })
2213        }
2214
2215        /// The Gemini REST fixture.
2216        pub fn fixture() -> ProviderWireFixture {
2217            ProviderWireFixture {
2218                driver: driver(),
2219                text_frames: vec![sse(&json!({
2220                    "candidates": [{"content": {"parts": [{"text": "hi"}], "role": "model"}}],
2221                    "responseId": "resp-1",
2222                    "modelVersion": "gemini-2.5-pro",
2223                }))],
2224                expected_texts: vec!["hi"],
2225                tool_call_frames: vec![sse(&json!({
2226                    "candidates": [{"content": {"parts": [{
2227                        "functionCall": {"name": "get_weather", "args": {"city": "Tokyo"}},
2228                    }], "role": "model"}}],
2229                    "responseId": "resp-1",
2230                    "modelVersion": "gemini-2.5-pro",
2231                }))],
2232                expected_tool_name: "get_weather",
2233                // Gemini delivers tool calls whole; arguments never stream.
2234                partial_tool_call_frames: None,
2235                terminal_frames: vec![sse(&json!({
2236                    "candidates": [{
2237                        "content": {"parts": [], "role": "model"},
2238                        "finishReason": "STOP",
2239                    }],
2240                    "usageMetadata": {
2241                        "promptTokenCount": 5,
2242                        "candidatesTokenCount": 2,
2243                        "totalTokenCount": 7,
2244                    },
2245                    "responseId": "resp-1",
2246                    "modelVersion": "gemini-2.5-pro",
2247                }))],
2248                expected_usage_total: 7,
2249                expected_finish_reason: Some(FinishReason::Stop),
2250                zero_usage_terminal_frames: Some(vec![sse(&json!({
2251                    "candidates": [{
2252                        "content": {"parts": [], "role": "model"},
2253                        "finishReason": "STOP",
2254                    }],
2255                    "responseId": "resp-1",
2256                    "modelVersion": "gemini-2.5-pro",
2257                }))]),
2258                bare_terminal_frames: None,
2259                malformed_frame: Some(sse_raw("{not json")),
2260                // The wire has no event tag; valid JSON carrying neither
2261                // `candidates` nor `usageMetadata` is unrecognizable and must
2262                // be warn-skipped, not silently decoded as an empty chunk.
2263                unknown_event_frame: Some(sse_raw(r#"{"noise":true}"#)),
2264                defective_known_frame: Some(sse_raw(r#"{"candidates": 42}"#)),
2265                delta_less_prelude_frame: None,
2266                refusal: None,
2267                interleaved_reasoning: Some(interleaved_thought_fixture()),
2268            }
2269        }
2270
2271        fn chunk(parts: &serde_json::Value) -> WireInput {
2272            sse(&json!({
2273                "candidates": [{"content": {"parts": parts, "role": "model"}}],
2274                "responseId": "resp-1",
2275                "modelVersion": "gemini-2.5-pro",
2276            }))
2277        }
2278
2279        fn terminal_frame() -> WireInput {
2280            sse(&json!({
2281                "candidates": [{
2282                    "content": {"parts": [], "role": "model"},
2283                    "finishReason": "STOP",
2284                }],
2285                "usageMetadata": {
2286                    "promptTokenCount": 5,
2287                    "candidatesTokenCount": 2,
2288                    "totalTokenCount": 7,
2289                },
2290                "responseId": "resp-1",
2291                "modelVersion": "gemini-2.5-pro",
2292            }))
2293        }
2294
2295        /// Thought delta, interleaved tool call, thought delta, terminal —
2296        /// the constant-id (`reasoning-0`) interleaving shape.
2297        fn interleaved_thought_fixture() -> InterleavedReasoningFixture {
2298            InterleavedReasoningFixture {
2299                frames: vec![
2300                    chunk(&json!([{"text": "before tool", "thought": true}])),
2301                    chunk(&json!([{
2302                        "functionCall": {"name": "get_weather", "args": {"city": "Tokyo"}},
2303                    }])),
2304                    chunk(&json!([{"text": "after tool", "thought": true}])),
2305                    terminal_frame(),
2306                ],
2307                first_reasoning: "before tool",
2308                tool_name: "get_weather",
2309                second_reasoning: "after tool",
2310            }
2311        }
2312
2313        /// Thought delta, interleaved tool call, then a signed full thought
2314        /// chunk carrying non-empty text — the F1 erasure shape.
2315        pub fn interleaved_signed_thought_frames()
2316        -> (Vec<WireInput>, &'static str, &'static str, &'static str) {
2317            let frames = vec![
2318                chunk(&json!([{"text": "before tool", "thought": true}])),
2319                chunk(&json!([{
2320                    "functionCall": {"name": "get_weather", "args": {"city": "Tokyo"}},
2321                }])),
2322                chunk(&json!([{
2323                    "text": "signed conclusion",
2324                    "thought": true,
2325                    "thoughtSignature": "sig-1",
2326                }])),
2327                terminal_frame(),
2328            ];
2329            (frames, "before tool", "get_weather", "signed conclusion")
2330        }
2331    }
2332
2333    /// Gemini Interactions SSE wire (`event_type`-tagged events).
2334    pub mod interactions {
2335        use super::*;
2336
2337        fn driver() -> WireDriver {
2338            byte_driver("gemini", |transport| {
2339                crate::driver::Model::new(
2340                    crate::providers::gemini::GeminiConfig::new("test-key")
2341                        .interactions("gemini-2.5-pro"),
2342                    transport,
2343                )
2344            })
2345        }
2346
2347        fn completed(usage: Option<serde_json::Value>) -> WireInput {
2348            let mut interaction = json!({
2349                "id": "int-1",
2350                "model": "gemini-2.5-pro",
2351                "status": "completed",
2352            });
2353            if let (Some(usage), Some(object)) = (usage, interaction.as_object_mut()) {
2354                object.insert("usage".to_string(), usage);
2355            }
2356            sse(&json!({
2357                "event_type": "interaction.completed",
2358                "interaction": interaction,
2359            }))
2360        }
2361
2362        /// The Interactions fixture.
2363        pub fn fixture() -> ProviderWireFixture {
2364            ProviderWireFixture {
2365                driver: driver(),
2366                text_frames: vec![sse(&json!({
2367                    "event_type": "step.delta",
2368                    "index": 0,
2369                    "delta": {"type": "text", "text": "hi"},
2370                }))],
2371                expected_texts: vec!["hi"],
2372                tool_call_frames: vec![sse(&json!({
2373                    "event_type": "step.delta",
2374                    "index": 0,
2375                    "delta": {
2376                        "type": "function_call",
2377                        "name": "get_weather",
2378                        "arguments": {"city": "Tokyo"},
2379                        "id": "call-1",
2380                    },
2381                }))],
2382                expected_tool_name: "get_weather",
2383                // The Interactions wire delivers function calls whole;
2384                // arguments never stream.
2385                partial_tool_call_frames: None,
2386                terminal_frames: vec![completed(Some(json!({
2387                    "total_input_tokens": 5,
2388                    "total_output_tokens": 2,
2389                    "total_tokens": 7,
2390                })))],
2391                expected_usage_total: 7,
2392                expected_finish_reason: Some(FinishReason::Stop),
2393                zero_usage_terminal_frames: Some(vec![completed(None)]),
2394                bare_terminal_frames: None,
2395                malformed_frame: Some(sse_raw("{not json")),
2396                unknown_event_frame: Some(sse(&json!({
2397                    "event_type": "future.event",
2398                    "index": 0,
2399                }))),
2400                // A known tag (`step.delta`) with a schema-defective payload
2401                // must classify `Corrupt`, never `Unknown`.
2402                defective_known_frame: Some(sse_raw(
2403                    r#"{"event_type":"step.delta","index":0,"delta":42}"#,
2404                )),
2405                delta_less_prelude_frame: None,
2406                refusal: None,
2407                interleaved_reasoning: Some(interleaved_thought_fixture()),
2408            }
2409        }
2410
2411        /// Thought-summary delta, interleaved function call, thought-summary
2412        /// delta, terminal — the constant-id (`reasoning-0`) interleaving
2413        /// shape on the Interactions wire.
2414        fn interleaved_thought_fixture() -> InterleavedReasoningFixture {
2415            let frames = vec![
2416                sse(&json!({
2417                    "event_type": "step.delta",
2418                    "index": 0,
2419                    "delta": {
2420                        "type": "thought_summary",
2421                        "content": {"type": "text", "text": "before tool"},
2422                    },
2423                })),
2424                sse(&json!({
2425                    "event_type": "step.delta",
2426                    "index": 0,
2427                    "delta": {
2428                        "type": "function_call",
2429                        "name": "get_weather",
2430                        "arguments": {"city": "Tokyo"},
2431                        "id": "call-1",
2432                    },
2433                })),
2434                sse(&json!({
2435                    "event_type": "step.delta",
2436                    "index": 0,
2437                    "delta": {
2438                        "type": "thought_summary",
2439                        "content": {"type": "text", "text": "after tool"},
2440                    },
2441                })),
2442                completed(Some(json!({
2443                    "total_input_tokens": 5,
2444                    "total_output_tokens": 2,
2445                    "total_tokens": 7,
2446                }))),
2447            ];
2448            InterleavedReasoningFixture {
2449                frames,
2450                first_reasoning: "before tool",
2451                tool_name: "get_weather",
2452                second_reasoning: "after tool",
2453            }
2454        }
2455    }
2456
2457    /// Anthropic Messages SSE wire (`type`-tagged events, index-as-id blocks).
2458    pub mod anthropic {
2459        use super::*;
2460
2461        fn driver() -> WireDriver {
2462            byte_driver("anthropic", |transport| {
2463                crate::driver::Model::new(
2464                    crate::providers::anthropic::wire::AnthropicConfig::new("test-key")
2465                        .completion(crate::providers::anthropic::completion::CLAUDE_SONNET_4_6),
2466                    transport,
2467                )
2468            })
2469        }
2470
2471        fn message_start() -> WireInput {
2472            sse(&json!({
2473                "type": "message_start",
2474                "message": {
2475                    "id": "msg_1",
2476                    "role": "assistant",
2477                    "content": [],
2478                    "model": "claude-sonnet-4-6",
2479                    "stop_reason": null,
2480                    "stop_sequence": null,
2481                    "usage": {"input_tokens": 5, "output_tokens": 0},
2482                },
2483            }))
2484        }
2485
2486        /// The Anthropic fixture.
2487        pub fn fixture() -> ProviderWireFixture {
2488            ProviderWireFixture {
2489                driver: driver(),
2490                text_frames: vec![
2491                    message_start(),
2492                    sse(&json!({
2493                        "type": "content_block_start",
2494                        "index": 0,
2495                        "content_block": {"type": "text", "text": ""},
2496                    })),
2497                    sse(&json!({
2498                        "type": "content_block_delta",
2499                        "index": 0,
2500                        "delta": {"type": "text_delta", "text": "hi"},
2501                    })),
2502                ],
2503                expected_texts: vec!["hi"],
2504                tool_call_frames: vec![
2505                    sse(&json!({
2506                        "type": "content_block_start",
2507                        "index": 0,
2508                        "content_block": {
2509                            "type": "tool_use",
2510                            "id": "toolu_1",
2511                            "name": "get_weather",
2512                            "input": {},
2513                        },
2514                    })),
2515                    sse(&json!({
2516                        "type": "content_block_delta",
2517                        "index": 0,
2518                        "delta": {"type": "input_json_delta", "partial_json": "{\"city\":\"Tokyo\"}"},
2519                    })),
2520                    // `content_block_stop` completes the call; the stream
2521                    // terminal (`message_delta`) is deliberately absent.
2522                    sse(&json!({"type": "content_block_stop", "index": 0})),
2523                ],
2524                expected_tool_name: "get_weather",
2525                partial_tool_call_frames: Some(vec![
2526                    sse(&json!({
2527                        "type": "content_block_start",
2528                        "index": 0,
2529                        "content_block": {
2530                            "type": "tool_use",
2531                            "id": "toolu_1",
2532                            "name": "get_weather",
2533                            "input": {},
2534                        },
2535                    })),
2536                    sse(&json!({
2537                        "type": "content_block_delta",
2538                        "index": 0,
2539                        "delta": {"type": "input_json_delta", "partial_json": "{\"cit"},
2540                    })),
2541                ]),
2542                terminal_frames: vec![sse(&json!({
2543                    "type": "message_delta",
2544                    "delta": {"stop_reason": "end_turn", "stop_sequence": null},
2545                    "usage": {"output_tokens": 4},
2546                }))],
2547                // input 5 (from message_start) + output 4.
2548                expected_usage_total: 9,
2549                expected_finish_reason: Some(FinishReason::Stop),
2550                // The Anthropic terminal always carries `usage`; there is no
2551                // usage-less genuine terminal to spell on this wire.
2552                zero_usage_terminal_frames: None,
2553                // `message_stop` carries no data of its own and must not
2554                // fabricate a terminal record.
2555                bare_terminal_frames: Some(vec![sse(&json!({"type": "message_stop"}))]),
2556                malformed_frame: Some(sse_raw("{not json")),
2557                unknown_event_frame: Some(sse(&json!({
2558                    "type": "content_block_heartbeat",
2559                    "index": 0,
2560                }))),
2561                // A known tag (`content_block_delta`) with a schema-defective
2562                // payload must classify `Corrupt`, never `Unknown`.
2563                defective_known_frame: Some(sse_raw(
2564                    r#"{"type":"content_block_delta","index":0,"delta":42}"#,
2565                )),
2566                delta_less_prelude_frame: None,
2567                refusal: None,
2568                interleaved_reasoning: None,
2569            }
2570        }
2571    }
2572
2573    /// Cohere v2 chat SSE wire.
2574    pub mod cohere {
2575        use super::*;
2576
2577        fn driver() -> WireDriver {
2578            byte_driver("cohere", |transport| {
2579                crate::driver::Model::new(
2580                    crate::providers::cohere::wire::CohereConfig::new("test-key")
2581                        .completion(crate::providers::cohere::COMMAND_R_08_2024),
2582                    transport,
2583                )
2584            })
2585        }
2586
2587        /// The Cohere fixture.
2588        pub fn fixture() -> ProviderWireFixture {
2589            ProviderWireFixture {
2590                driver: driver(),
2591                text_frames: vec![
2592                    sse(&json!({"type": "message-start", "id": "msg_1"})),
2593                    sse(&json!({
2594                        "type": "content-delta",
2595                        "delta": {"message": {"content": {"text": "hi"}}},
2596                    })),
2597                ],
2598                expected_texts: vec!["hi"],
2599                tool_call_frames: vec![
2600                    sse(&json!({
2601                        "type": "tool-call-start",
2602                        "delta": {"message": {"tool_calls": {
2603                            "id": "call_1",
2604                            "function": {"name": "get_weather", "arguments": ""},
2605                        }}},
2606                    })),
2607                    sse(&json!({
2608                        "type": "tool-call-delta",
2609                        "delta": {"message": {"tool_calls": {
2610                            "function": {"arguments": "{\"city\":\"Tokyo\"}"},
2611                        }}},
2612                    })),
2613                    sse(&json!({"type": "tool-call-end"})),
2614                ],
2615                expected_tool_name: "get_weather",
2616                partial_tool_call_frames: Some(vec![sse(&json!({
2617                    "type": "tool-call-start",
2618                    "delta": {"message": {"tool_calls": {
2619                        "id": "call_1",
2620                        "function": {"name": "get_weather", "arguments": "{\"cit"},
2621                    }}},
2622                }))]),
2623                terminal_frames: vec![sse(&json!({
2624                    "type": "message-end",
2625                    "delta": {
2626                        "finish_reason": "COMPLETE",
2627                        "usage": {"tokens": {"input_tokens": 10, "output_tokens": 4}},
2628                    },
2629                }))],
2630                expected_usage_total: 14,
2631                expected_finish_reason: Some(FinishReason::Stop),
2632                zero_usage_terminal_frames: Some(vec![sse(&json!({"type": "message-end"}))]),
2633                bare_terminal_frames: None,
2634                malformed_frame: Some(sse_raw("{not json")),
2635                unknown_event_frame: Some(sse(&json!({
2636                    "type": "citation-start",
2637                    "delta": {"message": {"citations": {}}},
2638                }))),
2639                defective_known_frame: Some(sse_raw(r#"{"type":"content-delta","delta":42}"#)),
2640                delta_less_prelude_frame: None,
2641                refusal: None,
2642                interleaved_reasoning: Some(interleaved_thinking_fixture()),
2643            }
2644        }
2645
2646        /// Thinking delta, interleaved tool call, thinking delta, terminal —
2647        /// the constant-id (`reasoning-0`) interleaving shape on the Cohere
2648        /// v2 SSE wire.
2649        fn interleaved_thinking_fixture() -> InterleavedReasoningFixture {
2650            let frames = vec![
2651                sse(&json!({"type": "message-start", "id": "msg_1"})),
2652                sse(&json!({
2653                    "type": "content-delta",
2654                    "delta": {"message": {"content": {"thinking": "before tool"}}},
2655                })),
2656                sse(&json!({
2657                    "type": "tool-call-start",
2658                    "delta": {"message": {"tool_calls": {
2659                        "id": "call_1",
2660                        "function": {"name": "get_weather", "arguments": "{\"city\":\"Tokyo\"}"},
2661                    }}},
2662                })),
2663                sse(&json!({"type": "tool-call-end"})),
2664                sse(&json!({
2665                    "type": "content-delta",
2666                    "delta": {"message": {"content": {"thinking": "after tool"}}},
2667                })),
2668                sse(&json!({
2669                    "type": "message-end",
2670                    "delta": {
2671                        "finish_reason": "COMPLETE",
2672                        "usage": {"tokens": {"input_tokens": 10, "output_tokens": 4}},
2673                    },
2674                })),
2675            ];
2676            InterleavedReasoningFixture {
2677                frames,
2678                first_reasoning: "before tool",
2679                tool_name: "get_weather",
2680                second_reasoning: "after tool",
2681            }
2682        }
2683    }
2684
2685    /// Ollama `/api/chat` NDJSON wire.
2686    pub mod ollama {
2687        use super::*;
2688
2689        fn driver() -> WireDriver {
2690            byte_driver("ollama", |transport| {
2691                crate::driver::Model::new(
2692                    crate::providers::ollama::wire::OllamaConfig::new().completion("llama3.2"),
2693                    transport,
2694                )
2695            })
2696        }
2697
2698        /// The Ollama fixture.
2699        pub fn fixture() -> ProviderWireFixture {
2700            ProviderWireFixture {
2701                driver: driver(),
2702                text_frames: vec![ndjson(&json!({
2703                    "model": "llama3.2",
2704                    "created_at": "2023-08-04T19:22:45.499127Z",
2705                    "message": {"role": "assistant", "content": "hi"},
2706                    "done": false,
2707                }))],
2708                expected_texts: vec!["hi"],
2709                tool_call_frames: vec![ndjson(&json!({
2710                    "model": "llama3.2",
2711                    "created_at": "2023-08-04T19:22:45.499127Z",
2712                    "message": {"role": "assistant", "content": "", "tool_calls": [{
2713                        "function": {"name": "get_weather", "arguments": {"city": "Tokyo"}},
2714                    }]},
2715                    "done": false,
2716                }))],
2717                expected_tool_name: "get_weather",
2718                // NDJSON delivers tool calls whole; arguments never stream.
2719                partial_tool_call_frames: None,
2720                terminal_frames: vec![ndjson(&json!({
2721                    "model": "llama3.2",
2722                    "created_at": "2023-08-04T19:22:47.499127Z",
2723                    "message": {"role": "assistant", "content": ""},
2724                    "done": true,
2725                    "done_reason": "stop",
2726                    "prompt_eval_count": 10,
2727                    "eval_count": 4,
2728                }))],
2729                expected_usage_total: 14,
2730                expected_finish_reason: Some(FinishReason::Stop),
2731                zero_usage_terminal_frames: Some(vec![ndjson(&json!({
2732                    "model": "llama3.2",
2733                    "created_at": "2023-08-04T19:22:47.499127Z",
2734                    "message": {"role": "assistant", "content": ""},
2735                    "done": true,
2736                    "done_reason": "stop",
2737                }))]),
2738                bare_terminal_frames: None,
2739                malformed_frame: Some(WireInput::Bytes(Bytes::from_static(b"{not json\n"))),
2740                unknown_event_frame: None,
2741                defective_known_frame: Some(ndjson(&json!({
2742                    "model": "llama3.2",
2743                    "created_at": "2023-08-04T19:22:46.499127Z",
2744                    "message": {"role": "assistant", "content": 42},
2745                    "done": false,
2746                }))),
2747                delta_less_prelude_frame: None,
2748                refusal: None,
2749                interleaved_reasoning: Some(interleaved_thinking_fixture()),
2750            }
2751        }
2752
2753        /// Thinking delta, interleaved tool call, thinking delta, terminal —
2754        /// the constant-id (`reasoning-0`) interleaving shape on NDJSON.
2755        fn interleaved_thinking_fixture() -> InterleavedReasoningFixture {
2756            let frames = vec![
2757                ndjson(&json!({
2758                    "model": "llama3.2",
2759                    "created_at": "2023-08-04T19:22:45.499127Z",
2760                    "message": {"role": "assistant", "content": "", "thinking": "before tool"},
2761                    "done": false,
2762                })),
2763                ndjson(&json!({
2764                    "model": "llama3.2",
2765                    "created_at": "2023-08-04T19:22:45.599127Z",
2766                    "message": {"role": "assistant", "content": "", "tool_calls": [{
2767                        "function": {"name": "get_weather", "arguments": {"city": "Tokyo"}},
2768                    }]},
2769                    "done": false,
2770                })),
2771                ndjson(&json!({
2772                    "model": "llama3.2",
2773                    "created_at": "2023-08-04T19:22:45.699127Z",
2774                    "message": {"role": "assistant", "content": "", "thinking": "after tool"},
2775                    "done": false,
2776                })),
2777                ndjson(&json!({
2778                    "model": "llama3.2",
2779                    "created_at": "2023-08-04T19:22:47.499127Z",
2780                    "message": {"role": "assistant", "content": ""},
2781                    "done": true,
2782                    "done_reason": "stop",
2783                    "prompt_eval_count": 10,
2784                    "eval_count": 4,
2785                })),
2786            ];
2787            InterleavedReasoningFixture {
2788                frames,
2789                first_reasoning: "before tool",
2790                tool_name: "get_weather",
2791                second_reasoning: "after tool",
2792            }
2793        }
2794    }
2795}