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