1use 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#[derive(Debug, thiserror::Error)]
41pub enum ConformanceError {
42 #[error(transparent)]
44 Completion(#[from] ProviderError),
45 #[error("{scenario} conformance failed for {provider}: {details}")]
47 Contract {
48 scenario: &'static str,
50 provider: &'static str,
52 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#[derive(Debug)]
73pub struct ScenarioReport {
74 pub name: &'static str,
76 pub provider: &'static str,
78 pub observations: Vec<String>,
80}
81
82#[derive(Debug)]
91pub enum ScenarioOutcome {
92 Ran(ScenarioReport),
94 Skipped {
96 name: &'static str,
98 provider: &'static str,
100 reason: &'static str,
102 },
103}
104
105#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
113pub struct SuiteCapabilities {
114 pub partial_tool_args: bool,
117 pub zero_usage_terminal: bool,
120 pub bare_terminal: bool,
122 pub malformed_frame: bool,
124 pub unknown_event_frame: bool,
126 pub defective_known_frame: bool,
129 pub delta_less_prelude: bool,
132 pub refusal: bool,
134 pub interleaved_reasoning: bool,
138}
139
140impl SuiteCapabilities {
141 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
169pub 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
185pub 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
204pub 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
213pub 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
227pub 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
272pub 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#[derive(Clone)]
298pub enum WireInput {
299 Bytes(Bytes),
301 Event(std::sync::Arc<dyn std::any::Any + Send + Sync>),
303}
304
305impl WireInput {
306 pub fn as_bytes(&self) -> Option<&Bytes> {
308 match self {
309 Self::Bytes(bytes) => Some(bytes),
310 Self::Event(_) => None,
311 }
312 }
313
314 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
338pub fn event_frame<T: Send + Sync + 'static>(event: T) -> WireInput {
340 WireInput::Event(std::sync::Arc::new(event))
341}
342
343pub type WireChunks = Vec<http_client::Result<WireInput>>;
346
347pub fn ok_chunks(frames: impl IntoIterator<Item = impl Into<WireInput>>) -> WireChunks {
349 frames.into_iter().map(|frame| Ok(frame.into())).collect()
350}
351
352pub 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
361pub fn assert_valid_event_stream(
386 items: &[Result<Item<StreamEvent>, ErrorReport>],
387 choice: &[AssistantContent],
388) {
389 use crate::message::AssistantContent;
390
391 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 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 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 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 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#[derive(Debug)]
486pub struct DrainedStream {
487 pub items: Vec<Result<Item<StreamEvent>, ErrorReport>>,
489 pub choice: Vec<AssistantContent>,
492 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 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 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 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 pub fn error_count(&self) -> usize {
540 self.items.iter().filter(|item| item.is_err()).count()
541 }
542
543 fn has_terminal(&self) -> bool {
545 self.response.is_some()
546 }
547
548 fn completed_cleanly(&self) -> bool {
551 self.error_count() == 0 && self.response.is_some()
552 }
553
554 fn first_error_index(&self) -> Option<usize> {
556 self.items.iter().position(std::result::Result::is_err)
557 }
558
559 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 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
586pub struct WireDriver {
592 pub provider: &'static str,
594 drive: DriveFn,
595}
596
597impl WireDriver {
598 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 pub async fn drive(&self, chunks: WireChunks) -> Result<DrainedStream, ProviderError> {
614 (self.drive)(chunks).await
615 }
616}
617
618pub struct RefusalFixture {
620 pub frames: Vec<WireInput>,
622 pub expected_text: &'static str,
624}
625
626pub struct InterleavedReasoningFixture {
631 pub frames: Vec<WireInput>,
633 pub first_reasoning: &'static str,
635 pub tool_name: &'static str,
637 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
647pub struct BufferedBodyDriver {
650 pub provider: &'static str,
652 drive: BufferedDriveFn,
653}
654
655impl BufferedBodyDriver {
656 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 pub async fn drive(&self, body: String) -> Result<Vec<AssistantContent>, ProviderError> {
672 (self.drive)(body).await
673 }
674}
675
676pub struct ProviderWireFixture {
681 pub driver: WireDriver,
683 pub text_frames: Vec<WireInput>,
685 pub expected_texts: Vec<&'static str>,
687 pub tool_call_frames: Vec<WireInput>,
690 pub expected_tool_name: &'static str,
692 pub partial_tool_call_frames: Option<Vec<WireInput>>,
695 pub terminal_frames: Vec<WireInput>,
697 pub expected_usage_total: u64,
699 pub expected_finish_reason: Option<FinishReason>,
701 pub zero_usage_terminal_frames: Option<Vec<WireInput>>,
703 pub bare_terminal_frames: Option<Vec<WireInput>>,
706 pub malformed_frame: Option<WireInput>,
710 pub unknown_event_frame: Option<WireInput>,
712 pub defective_known_frame: Option<WireInput>,
714 pub delta_less_prelude_frame: Option<WireInput>,
716 pub refusal: Option<RefusalFixture>,
718 pub interleaved_reasoning: Option<InterleavedReasoningFixture>,
722}
723
724impl ProviderWireFixture {
725 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
754struct 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 fn fail(&self, details: impl Into<String>) -> ConformanceError {
777 ConformanceError::contract(self.name, self.provider, details)
778 }
779
780 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 fn note(&mut self, observation: impl Into<String>) {
795 self.observations.push(observation.into());
796 }
797
798 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
822pub 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 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 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 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 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
901pub 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
944pub 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
985pub 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 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 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
1039pub 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
1079pub 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
1115pub 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
1153pub 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
1189pub 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
1225pub 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
1287pub 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
1331pub 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
1379pub 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
1417pub 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
1456pub 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
1492fn 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
1524pub 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 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 super::assert_valid_event_stream(&drained.items, &drained.choice);
1554 drained
1555 }
1556
1557 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 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 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 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 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 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 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 ],
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 defective_known_frame: None,
1748 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 interleaved_reasoning: None,
1766 }
1767 }
1768 }
1769
1770 pub mod openai_responses {
1772 use super::*;
1773
1774 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 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 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 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 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 pub fn terminal_body_only_sse_body(text: &str) -> String {
2049 frame_text(&terminal(Some(&usage_json()), &message_output(text)))
2050 }
2051
2052 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 expected_usage_total: 9,
2520 expected_finish_reason: Some(FinishReason::Stop),
2521 zero_usage_terminal_frames: None,
2524 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 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}