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 "cohere",
199 "ollama",
200 "xai",
201 "copilot",
202 "bedrock",
203 "candle",
204];
205
206pub fn xfail_reason<'a>(xfail: &[&'a str], scenario: &str) -> Option<&'a str> {
209 xfail.iter().find_map(|entry| {
210 let (name, reason) = entry.split_once(':')?;
211 (name.trim() == scenario).then(|| reason.trim())
212 })
213}
214
215pub fn invalid_xfail_entries(xfail: &[&str]) -> Vec<String> {
217 xfail
218 .iter()
219 .filter(|entry| match entry.split_once(':') {
220 Some((name, reason)) => {
221 !CANONICAL_SCENARIOS.contains(&name.trim()) || reason.trim().is_empty()
222 }
223 None => true,
224 })
225 .map(std::string::ToString::to_string)
226 .collect()
227}
228
229pub fn check_gated_outcome(
237 scenario: &'static str,
238 capability: bool,
239 xfail: &[&str],
240 outcome: Result<ScenarioOutcome, ConformanceError>,
241) -> Result<(), String> {
242 match (xfail_reason(xfail, scenario), outcome) {
243 (Some(reason), Err(error)) => {
244 eprintln!("xfail {scenario}: {reason} ({error})");
245 Ok(())
246 }
247 (Some(reason), Ok(_)) => Err(format!(
248 "{scenario} passed but is listed as xfail ({reason}); remove the xfail entry"
249 )),
250 (None, Err(error)) => Err(format!("{scenario} failed: {error}")),
251 (None, Ok(ScenarioOutcome::Ran(_))) => {
252 if capability {
253 Ok(())
254 } else {
255 Err(format!(
256 "{scenario} ran but the suite disclaims the capability; set the flag to true"
257 ))
258 }
259 }
260 (None, Ok(ScenarioOutcome::Skipped { reason, .. })) => {
261 if capability {
262 Err(format!(
263 "{scenario} skipped ({reason}) but the suite declares the capability; \
264 a declared capability's scenario must run"
265 ))
266 } else {
267 eprintln!("skipped {scenario}: {reason}");
268 Ok(())
269 }
270 }
271 }
272}
273
274pub fn check_ungated_outcome(
276 scenario: &'static str,
277 xfail: &[&str],
278 result: Result<ScenarioReport, ConformanceError>,
279) -> Result<(), String> {
280 match (xfail_reason(xfail, scenario), result) {
281 (Some(reason), Err(error)) => {
282 eprintln!("xfail {scenario}: {reason} ({error})");
283 Ok(())
284 }
285 (Some(reason), Ok(_)) => Err(format!(
286 "{scenario} passed but is listed as xfail ({reason}); remove the xfail entry"
287 )),
288 (None, Err(error)) => Err(format!("{scenario} failed: {error}")),
289 (None, Ok(_)) => Ok(()),
290 }
291}
292
293#[derive(Clone)]
300pub enum WireInput {
301 Bytes(Bytes),
303 Event(std::sync::Arc<dyn std::any::Any + Send + Sync>),
305}
306
307impl WireInput {
308 pub fn as_bytes(&self) -> Option<&Bytes> {
310 match self {
311 Self::Bytes(bytes) => Some(bytes),
312 Self::Event(_) => None,
313 }
314 }
315
316 pub fn downcast_event<T: 'static>(&self) -> Option<&T> {
318 match self {
319 Self::Bytes(_) => None,
320 Self::Event(event) => event.downcast_ref(),
321 }
322 }
323}
324
325impl From<Bytes> for WireInput {
326 fn from(bytes: Bytes) -> Self {
327 Self::Bytes(bytes)
328 }
329}
330
331impl std::fmt::Debug for WireInput {
332 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
333 match self {
334 Self::Bytes(bytes) => formatter.debug_tuple("Bytes").field(bytes).finish(),
335 Self::Event(_) => formatter.write_str("Event(..)"),
336 }
337 }
338}
339
340pub fn event_frame<T: Send + Sync + 'static>(event: T) -> WireInput {
342 WireInput::Event(std::sync::Arc::new(event))
343}
344
345pub type WireChunks = Vec<http_client::Result<WireInput>>;
348
349pub fn ok_chunks(frames: impl IntoIterator<Item = impl Into<WireInput>>) -> WireChunks {
351 frames.into_iter().map(|frame| Ok(frame.into())).collect()
352}
353
354pub fn transport_error_chunk() -> http_client::Result<WireInput> {
357 Err(http_client::Error::instance(std::io::Error::new(
358 std::io::ErrorKind::ConnectionReset,
359 "connection reset",
360 )))
361}
362
363pub fn assert_valid_event_stream(
388 items: &[Result<Item<StreamEvent>, ErrorReport>],
389 choice: &[AssistantContent],
390) {
391 use crate::message::AssistantContent;
392
393 if let Some(error_index) = items.iter().position(Result::is_err) {
395 assert_eq!(
396 error_index + 1,
397 items.len(),
398 "law 1 (terminal error): an item followed the stream's error"
399 );
400 }
401 let events: Vec<&StreamEvent> = items
402 .iter()
403 .filter_map(|item| match item {
404 Ok(Item::Event(event)) => Some(event),
405 _ => None,
406 })
407 .collect();
408
409 let streamed_text: String = events
411 .iter()
412 .filter_map(|event| match event {
413 StreamEvent::Text { text, .. } => Some(text.as_str()),
414 _ => None,
415 })
416 .collect();
417 let aggregated_text: String = choice
418 .iter()
419 .filter_map(|content| match content {
420 AssistantContent::Text(text) => Some(text.text.as_str()),
421 _ => None,
422 })
423 .collect();
424 assert_eq!(
425 aggregated_text, streamed_text,
426 "law 2 (text conservation): aggregated text differs from the streamed fragments"
427 );
428
429 let yielded_calls = events
431 .iter()
432 .filter(|event| {
433 matches!(
434 event,
435 StreamEvent::End {
436 content: AssistantContent::ToolCall(_),
437 ..
438 }
439 )
440 })
441 .count();
442 let aggregated_calls = choice
443 .iter()
444 .filter(|content| matches!(content, AssistantContent::ToolCall(_)))
445 .count();
446 assert_eq!(
447 aggregated_calls, yielded_calls,
448 "law 3 (completed-call conservation): {yielded_calls} calls yielded, \
449 {aggregated_calls} aggregated"
450 );
451
452 let serialized = serde_json::to_value(
454 items
455 .iter()
456 .filter_map(|item| item.as_ref().ok())
457 .collect::<Vec<_>>(),
458 )
459 .unwrap_or_default();
460 let read_back = crate::streaming::Transcript::parse_prefix(serialized);
461 assert!(
462 read_back.is_ok(),
463 "law 4 (sequence): the stream's events do not read back: {read_back:?}"
464 );
465
466 let yielded_reasoning = events.iter().any(|event| {
468 matches!(
469 event,
470 StreamEvent::End {
471 content: AssistantContent::Reasoning(_),
472 ..
473 }
474 )
475 });
476 let aggregated_reasoning = choice
477 .iter()
478 .any(|content| matches!(content, AssistantContent::Reasoning(_)));
479 assert!(
480 yielded_reasoning || !aggregated_reasoning,
481 "law 5 (reasoning provenance): aggregated reasoning with no reasoning yielded"
482 );
483}
484
485#[derive(Debug)]
488pub struct DrainedStream {
489 pub items: Vec<Result<Item<StreamEvent>, ErrorReport>>,
491 pub choice: Vec<AssistantContent>,
494 pub response: Option<CompletionResponse>,
496}
497
498impl DrainedStream {
499 fn events(&self) -> impl Iterator<Item = &StreamEvent> {
500 self.items.iter().filter_map(|item| match item {
501 Ok(Item::Event(event)) => Some(event),
502 _ => None,
503 })
504 }
505
506 pub fn texts(&self) -> Vec<&str> {
508 self.events()
509 .filter_map(|event| match event {
510 StreamEvent::Text { text, .. } => Some(text.as_str()),
511 _ => None,
512 })
513 .collect()
514 }
515
516 pub fn tool_call_names(&self) -> Vec<&str> {
518 self.events()
519 .filter_map(|event| match event {
520 StreamEvent::End {
521 content: AssistantContent::ToolCall(tool_call),
522 ..
523 } => Some(tool_call.function.name.as_str()),
524 _ => None,
525 })
526 .collect()
527 }
528
529 pub fn unknown_values(&self) -> Vec<&serde_json::Value> {
531 self.items
532 .iter()
533 .filter_map(|item| match item {
534 Ok(Item::Unknown(value)) => Some(value.value()),
535 _ => None,
536 })
537 .collect()
538 }
539
540 pub fn error_count(&self) -> usize {
542 self.items.iter().filter(|item| item.is_err()).count()
543 }
544
545 fn has_terminal(&self) -> bool {
547 self.response.is_some()
548 }
549
550 fn completed_cleanly(&self) -> bool {
553 self.error_count() == 0 && self.response.is_some()
554 }
555
556 fn first_error_index(&self) -> Option<usize> {
558 self.items.iter().position(std::result::Result::is_err)
559 }
560
561 pub fn choice_texts(&self) -> Vec<&str> {
563 self.choice
564 .iter()
565 .filter_map(|content| match content {
566 AssistantContent::Text(text) => Some(text.text.as_str()),
567 _ => None,
568 })
569 .collect()
570 }
571
572 pub fn choice_reasoning(&self) -> Vec<&crate::message::Reasoning> {
574 self.choice
575 .iter()
576 .filter_map(|content| match content {
577 AssistantContent::Reasoning(reasoning) => reasoning.open(reasoning.issuer()),
578 _ => None,
579 })
580 .collect()
581 }
582}
583
584type DriveFn = Box<
585 dyn Fn(WireChunks) -> BoxFuture<'static, Result<DrainedStream, ProviderError>> + Send + Sync,
586>;
587
588pub struct WireDriver {
594 pub provider: &'static str,
596 drive: DriveFn,
597}
598
599impl WireDriver {
600 pub fn new(
602 provider: &'static str,
603 drive: impl Fn(WireChunks) -> BoxFuture<'static, Result<DrainedStream, ProviderError>>
604 + Send
605 + Sync
606 + 'static,
607 ) -> Self {
608 Self {
609 provider,
610 drive: Box::new(drive),
611 }
612 }
613
614 pub async fn drive(&self, chunks: WireChunks) -> Result<DrainedStream, ProviderError> {
616 (self.drive)(chunks).await
617 }
618}
619
620pub struct RefusalFixture {
622 pub frames: Vec<WireInput>,
624 pub expected_text: &'static str,
626}
627
628pub struct InterleavedReasoningFixture {
633 pub frames: Vec<WireInput>,
635 pub first_reasoning: &'static str,
637 pub tool_name: &'static str,
639 pub second_reasoning: &'static str,
641}
642
643type BufferedDriveFn = Box<
644 dyn Fn(String) -> BoxFuture<'static, Result<Vec<AssistantContent>, ProviderError>>
645 + Send
646 + Sync,
647>;
648
649pub struct BufferedBodyDriver {
652 pub provider: &'static str,
654 drive: BufferedDriveFn,
655}
656
657impl BufferedBodyDriver {
658 pub fn new(
660 provider: &'static str,
661 drive: impl Fn(String) -> BoxFuture<'static, Result<Vec<AssistantContent>, ProviderError>>
662 + Send
663 + Sync
664 + 'static,
665 ) -> Self {
666 Self {
667 provider,
668 drive: Box::new(drive),
669 }
670 }
671
672 pub async fn drive(&self, body: String) -> Result<Vec<AssistantContent>, ProviderError> {
674 (self.drive)(body).await
675 }
676}
677
678pub struct ProviderWireFixture {
683 pub driver: WireDriver,
685 pub text_frames: Vec<WireInput>,
687 pub expected_texts: Vec<&'static str>,
689 pub tool_call_frames: Vec<WireInput>,
692 pub expected_tool_name: &'static str,
694 pub partial_tool_call_frames: Option<Vec<WireInput>>,
697 pub terminal_frames: Vec<WireInput>,
699 pub expected_usage_total: u64,
701 pub expected_finish_reason: Option<FinishReason>,
703 pub zero_usage_terminal_frames: Option<Vec<WireInput>>,
705 pub bare_terminal_frames: Option<Vec<WireInput>>,
708 pub malformed_frame: Option<WireInput>,
712 pub unknown_event_frame: Option<WireInput>,
714 pub defective_known_frame: Option<WireInput>,
716 pub delta_less_prelude_frame: Option<WireInput>,
718 pub refusal: Option<RefusalFixture>,
720 pub interleaved_reasoning: Option<InterleavedReasoningFixture>,
724}
725
726impl ProviderWireFixture {
727 pub fn capabilities(&self) -> SuiteCapabilities {
735 SuiteCapabilities {
736 partial_tool_args: self.partial_tool_call_frames.is_some(),
737 zero_usage_terminal: self.zero_usage_terminal_frames.is_some(),
738 bare_terminal: self.bare_terminal_frames.is_some(),
739 malformed_frame: self.malformed_frame.is_some(),
740 unknown_event_frame: self.unknown_event_frame.is_some(),
741 defective_known_frame: self.defective_known_frame.is_some(),
742 delta_less_prelude: self.delta_less_prelude_frame.is_some(),
743 refusal: self.refusal.is_some(),
744 interleaved_reasoning: self.interleaved_reasoning.is_some(),
745 }
746 }
747}
748
749fn concat_frames(parts: &[&[WireInput]]) -> Vec<WireInput> {
750 parts
751 .iter()
752 .flat_map(|frames| frames.iter().cloned())
753 .collect()
754}
755
756struct Checks {
763 name: &'static str,
764 provider: &'static str,
765 observations: Vec<String>,
766}
767
768impl Checks {
769 fn new(name: &'static str, provider: &'static str) -> Self {
770 Self {
771 name,
772 provider,
773 observations: Vec::new(),
774 }
775 }
776
777 fn fail(&self, details: impl Into<String>) -> ConformanceError {
779 ConformanceError::contract(self.name, self.provider, details)
780 }
781
782 fn require<D: Into<String>>(
785 &self,
786 held: bool,
787 details: impl FnOnce() -> D,
788 ) -> Result<(), ConformanceError> {
789 if held {
790 return Ok(());
791 }
792 Err(self.fail(details()))
793 }
794
795 fn note(&mut self, observation: impl Into<String>) {
797 self.observations.push(observation.into());
798 }
799
800 fn skip(&self, reason: &'static str) -> ScenarioOutcome {
804 ScenarioOutcome::Skipped {
805 name: self.name,
806 provider: self.provider,
807 reason,
808 }
809 }
810
811 fn report(self) -> ScenarioReport {
812 ScenarioReport {
813 name: self.name,
814 provider: self.provider,
815 observations: self.observations,
816 }
817 }
818
819 fn ran(self) -> ScenarioOutcome {
820 ScenarioOutcome::Ran(self.report())
821 }
822}
823
824pub async fn truncation_preserves_content_without_terminal(
832 fixture: &ProviderWireFixture,
833) -> Result<ScenarioReport, ConformanceError> {
834 let mut checks = Checks::new(
835 "truncation_preserves_content_without_terminal",
836 fixture.driver.provider,
837 );
838
839 let drained = fixture.driver.drive(Vec::new()).await?;
841 checks.require(
842 !drained.has_terminal(),
843 || "an empty stream must not synthesize a terminal record",
844 )?;
845 checks.note("EOF before content: no terminal");
846
847 let drained = fixture
849 .driver
850 .drive(ok_chunks(fixture.text_frames.clone()))
851 .await?;
852 checks.require(drained.texts() == fixture.expected_texts, || {
853 format!(
854 "text delivered before truncation must be preserved: expected {:?}, observed {:?}",
855 fixture.expected_texts,
856 drained.texts()
857 )
858 })?;
859 checks.require(
860 !drained.has_terminal(),
861 || "EOF after text deltas must not synthesize a terminal record",
862 )?;
863 checks.note("EOF mid-text: content preserved, no terminal");
864
865 if let Some(partial) = &fixture.partial_tool_call_frames {
867 let drained = fixture.driver.drive(ok_chunks(partial.clone())).await?;
868 checks.require(
869 !drained.has_terminal(),
870 || "EOF mid-tool-arguments must not synthesize a terminal record",
871 )?;
872 checks.note("EOF mid-tool-args: no terminal");
873 }
874
875 let drained = fixture
879 .driver
880 .drive(ok_chunks(fixture.tool_call_frames.clone()))
881 .await?;
882 checks.require(
883 drained
884 .tool_call_names()
885 .iter()
886 .all(|name| *name == fixture.expected_tool_name),
887 || {
888 format!(
889 "only the delivered call may surface: observed {:?}",
890 drained.tool_call_names()
891 )
892 },
893 )?;
894 checks.require(
895 !drained.has_terminal(),
896 || "EOF after a tool call must not synthesize a terminal record",
897 )?;
898 checks.note("EOF after a tool call: no terminal");
899
900 Ok(checks.report())
901}
902
903pub async fn transport_error_after_tool_call_yields_err_then_end(
907 fixture: &ProviderWireFixture,
908) -> Result<ScenarioReport, ConformanceError> {
909 let mut checks = Checks::new(
910 "transport_error_after_tool_call_yields_err_then_end",
911 fixture.driver.provider,
912 );
913
914 let mut chunks = ok_chunks(fixture.tool_call_frames.clone());
915 chunks.push(transport_error_chunk());
916 let drained = fixture.driver.drive(chunks).await?;
917
918 checks.require(
919 drained
920 .tool_call_names()
921 .iter()
922 .all(|name| *name == fixture.expected_tool_name),
923 || {
924 format!(
925 "only the delivered call may precede the transport error: observed {:?}",
926 drained.tool_call_names()
927 )
928 },
929 )?;
930 let error_index = drained
931 .first_error_index()
932 .ok_or_else(|| checks.fail("the transport failure must reach the consumer"))?;
933 checks.require(
934 error_index + 1 == drained.items.len(),
935 || "nothing may follow the terminal transport error",
936 )?;
937 checks.require(
938 !drained.has_terminal(),
939 || "a transport failure must not be papered over with a terminal record",
940 )?;
941
942 checks.note("Err, then end; no terminal");
943 Ok(checks.report())
944}
945
946pub async fn malformed_frame_ends_the_reply(
950 fixture: &ProviderWireFixture,
951) -> Result<ScenarioOutcome, ConformanceError> {
952 let mut checks = Checks::new("malformed_frame_ends_the_reply", fixture.driver.provider);
953 let Some(malformed) = &fixture.malformed_frame else {
954 return Ok(checks.skip("wire family cannot spell a frame-level decode failure"));
955 };
956
957 let frames = concat_frames(&[
958 &fixture.text_frames,
959 std::slice::from_ref(malformed),
960 &fixture.terminal_frames,
961 ]);
962 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
963
964 checks.require(drained.error_count() == 1, || {
965 format!(
966 "the malformed frame must surface as exactly one Err item, observed {}",
967 drained.error_count()
968 )
969 })?;
970 checks.require(
971 drained.first_error_index() == Some(drained.items.len() - 1),
972 || "the malformed frame's error must be the stream's last item",
973 )?;
974 checks.require(
975 drained.texts() == fixture.expected_texts,
976 || "content before the malformed frame must be preserved",
977 )?;
978 checks.require(
979 !drained.has_terminal(),
980 || "a reply cut by a corrupt frame has no response",
981 )?;
982
983 checks.note("Err surfaced last; content kept; no response");
984 Ok(checks.ran())
985}
986
987pub async fn unknown_event_is_skipped(
993 fixture: &ProviderWireFixture,
994) -> Result<ScenarioOutcome, ConformanceError> {
995 let mut checks = Checks::new("unknown_event_is_skipped", fixture.driver.provider);
996 let Some(unknown) = &fixture.unknown_event_frame else {
997 return Ok(checks.skip("wire family cannot spell an unknown event type"));
998 };
999
1000 let frames = concat_frames(&[
1001 &fixture.text_frames,
1002 std::slice::from_ref(unknown),
1003 &fixture.terminal_frames,
1004 ]);
1005 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1006
1007 checks.require(
1008 drained.error_count() == 0,
1009 || "an unknown event type must be skipped, not surfaced as an error",
1010 )?;
1011 checks.require(
1012 drained.texts() == fixture.expected_texts && drained.response.is_some(),
1013 || "the stream must deliver its content and complete around the skipped event",
1014 )?;
1015 checks.require(drained.unknown_values().len() == 1, || {
1018 format!(
1019 "exactly one Unknown passthrough item must surface for the unknown frame, \
1020 observed {}",
1021 drained.unknown_values().len()
1022 )
1023 })?;
1024
1025 let control_frames = concat_frames(&[&fixture.text_frames, &fixture.terminal_frames]);
1028 let control = fixture.driver.drive(ok_chunks(control_frames)).await?;
1029 checks.require(
1030 drained.choice == control.choice,
1031 || "the unknown frame must not perturb the aggregated assistant choice",
1032 )?;
1033
1034 checks.note(
1035 "unknown event skipped semantically, surfaced on the raw channel, \
1036 choice unchanged, stream completed",
1037 );
1038 Ok(checks.ran())
1039}
1040
1041pub async fn defective_known_event_ends_the_reply(
1047 fixture: &ProviderWireFixture,
1048) -> Result<ScenarioOutcome, ConformanceError> {
1049 let mut checks = Checks::new(
1050 "defective_known_event_ends_the_reply",
1051 fixture.driver.provider,
1052 );
1053 let Some(defective) = &fixture.defective_known_frame else {
1054 return Ok(
1055 checks.skip("wire family cannot spell a known event with a schema-defective payload")
1056 );
1057 };
1058
1059 let frames = concat_frames(&[
1060 &fixture.text_frames,
1061 std::slice::from_ref(defective),
1062 &fixture.terminal_frames,
1063 ]);
1064 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1065
1066 checks.require(drained.error_count() == 1, || {
1067 format!(
1068 "a known event with a schema defect must surface exactly one Err item, observed {}",
1069 drained.error_count()
1070 )
1071 })?;
1072 checks.require(
1073 !drained.has_terminal(),
1074 || "a reply cut by a defective event has no response",
1075 )?;
1076
1077 checks.note("defective known event ended the reply with its Err");
1078 Ok(checks.ran())
1079}
1080
1081pub async fn delta_less_choice_prelude_is_a_noop(
1087 fixture: &ProviderWireFixture,
1088) -> Result<ScenarioOutcome, ConformanceError> {
1089 let mut checks = Checks::new(
1090 "delta_less_choice_prelude_is_a_noop",
1091 fixture.driver.provider,
1092 );
1093 let Some(prelude) = &fixture.delta_less_prelude_frame else {
1094 return Ok(checks.skip("wire family has no delta-less prelude shape"));
1095 };
1096
1097 let frames = concat_frames(&[
1098 std::slice::from_ref(prelude),
1099 &fixture.text_frames,
1100 &fixture.terminal_frames,
1101 ]);
1102 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1103
1104 checks.require(
1105 drained.error_count() == 0,
1106 || "the delta-less prelude must not surface an error",
1107 )?;
1108 checks.require(
1109 drained.texts() == fixture.expected_texts && drained.response.is_some(),
1110 || "the prelude must not perturb content delivery or the terminal",
1111 )?;
1112
1113 checks.note("delta-less prelude ignored; stream unaffected");
1114 Ok(checks.ran())
1115}
1116
1117pub async fn refusal_frames_deliver_text_without_error(
1122 fixture: &ProviderWireFixture,
1123) -> Result<ScenarioOutcome, ConformanceError> {
1124 let mut checks = Checks::new(
1125 "refusal_frames_deliver_text_without_error",
1126 fixture.driver.provider,
1127 );
1128 let Some(refusal) = &fixture.refusal else {
1129 return Ok(checks.skip("wire family has no refusal channel"));
1130 };
1131
1132 let frames = concat_frames(&[&refusal.frames, &fixture.terminal_frames]);
1133 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1134
1135 checks.require(
1136 drained.error_count() == 0,
1137 || "refusal content must not surface as an error",
1138 )?;
1139 let delivered = drained.texts().concat();
1140 checks.require(delivered == refusal.expected_text, || {
1141 format!(
1142 "refusal text must be delivered: expected {:?}, observed {delivered:?}",
1143 refusal.expected_text
1144 )
1145 })?;
1146 checks.require(
1147 drained.response.is_some(),
1148 || "a refused turn still ends with the provider's genuine terminal",
1149 )?;
1150
1151 checks.note("refusal text delivered without error");
1152 Ok(checks.ran())
1153}
1154
1155pub async fn terminal_body_content_merges_per_kind(
1164 driver: &BufferedBodyDriver,
1165 cases: Vec<(&'static str, String)>,
1166 expected_text: &str,
1167) -> Result<ScenarioReport, ConformanceError> {
1168 let mut checks = Checks::new("terminal_body_content_merges_per_kind", driver.provider);
1169
1170 for (label, body) in cases {
1171 let choice = driver.drive(body).await?;
1172 let choice_text: String = choice
1173 .iter()
1174 .filter_map(|content| match content {
1175 AssistantContent::Text(text) => Some(text.text.as_str()),
1176 _ => None,
1177 })
1178 .collect();
1179 let occurrences = choice_text.matches(expected_text).count();
1180 checks.require(occurrences == 1, || {
1181 format!(
1182 "{label}: terminal-body text must appear exactly once in the choice, observed {occurrences} in {choice_text:?}"
1183 )
1184 })?;
1185 checks.note(format!("{label}: text merged exactly once"));
1186 }
1187
1188 Ok(checks.report())
1189}
1190
1191pub async fn bare_terminal_after_only_unparseable_frames_fabricates_nothing(
1198 fixture: &ProviderWireFixture,
1199) -> Result<ScenarioOutcome, ConformanceError> {
1200 let mut checks = Checks::new(
1201 "bare_terminal_after_only_unparseable_frames_fabricates_nothing",
1202 fixture.driver.provider,
1203 );
1204 let Some(bare_terminal) = &fixture.bare_terminal_frames else {
1205 return Ok(checks.skip("wire family has no data-less terminal signal"));
1206 };
1207 let Some(malformed) = &fixture.malformed_frame else {
1208 return Ok(checks.skip("wire family cannot spell a frame-level decode failure"));
1209 };
1210
1211 let frames = concat_frames(&[std::slice::from_ref(malformed), bare_terminal]);
1212 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1213
1214 checks.require(
1215 drained.error_count() != 0,
1216 || "the unparseable frame must surface as an Err item",
1217 )?;
1218 checks.require(
1219 !drained.has_terminal(),
1220 || "a bare terminal with no decoded frame must not fabricate a terminal record",
1221 )?;
1222
1223 checks.note("no fabricated terminal after only-unparseable frames");
1224 Ok(checks.ran())
1225}
1226
1227pub async fn usage_variants_are_reported_or_absent(
1234 fixture: &ProviderWireFixture,
1235) -> Result<ScenarioReport, ConformanceError> {
1236 let mut checks = Checks::new(
1237 "usage_variants_are_reported_or_absent",
1238 fixture.driver.provider,
1239 );
1240
1241 let frames = concat_frames(&[&fixture.text_frames, &fixture.terminal_frames]);
1242 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1243 let response = drained
1244 .response
1245 .as_ref()
1246 .ok_or_else(|| checks.fail("the genuine terminal must produce a record"))?;
1247 checks.require(
1248 response.usage.total_tokens == Some(fixture.expected_usage_total),
1249 || {
1250 format!(
1251 "terminal usage must be preserved: expected total {}, observed {:?}",
1252 fixture.expected_usage_total, response.usage.total_tokens
1253 )
1254 },
1255 )?;
1256 checks.require(
1257 response.finish_reason() == fixture.expected_finish_reason,
1258 || {
1259 format!(
1260 "terminal finish reason must be normalized: expected {:?}, observed {:?}",
1261 fixture.expected_finish_reason,
1262 response.finish_reason()
1263 )
1264 },
1265 )?;
1266 checks.note(format!(
1267 "usage total {} and finish reason {:?} preserved",
1268 fixture.expected_usage_total, fixture.expected_finish_reason
1269 ));
1270
1271 if let Some(zero_usage) = &fixture.zero_usage_terminal_frames {
1272 let frames = concat_frames(&[&fixture.text_frames, zero_usage]);
1273 let drained = fixture.driver.drive(ok_chunks(frames)).await?;
1274 let response = drained.response.as_ref().ok_or_else(|| {
1275 checks.fail("a usage-less genuine terminal must still complete the stream")
1276 })?;
1277 checks.require(!response.usage.is_reported(), || {
1278 format!(
1279 "missing usage metrics must leave every counter unreported, not invented: {:?}",
1280 response.usage
1281 )
1282 })?;
1283 checks.note("usage-less terminal completed with no counter reported");
1284 }
1285
1286 Ok(checks.report())
1287}
1288
1289pub async fn reasoning_summary_deltas_are_superseded_without_duplication(
1298 driver: &WireDriver,
1299 frames: Vec<WireInput>,
1300 summary_text: &str,
1301) -> Result<ScenarioReport, ConformanceError> {
1302 let mut checks = Checks::new(
1303 "reasoning_summary_deltas_are_superseded_without_duplication",
1304 driver.provider,
1305 );
1306
1307 let drained = driver.drive(ok_chunks(frames)).await?;
1308 checks.require(
1309 drained.completed_cleanly(),
1310 || "the reasoning stream must complete without errors",
1311 )?;
1312 let reasoning = drained.choice_reasoning();
1313 let occurrences: usize = reasoning
1314 .iter()
1315 .flat_map(|item| item.content.iter())
1316 .filter(|content| match content {
1317 crate::message::ReasoningContent::Summary(text)
1318 | crate::message::ReasoningContent::Text { text, .. } => text.contains(summary_text),
1319 _ => false,
1320 })
1321 .count();
1322 checks.require(occurrences == 1, || {
1323 format!(
1324 "the summary must appear exactly once in the aggregated choice, observed {occurrences} across {reasoning:?}"
1325 )
1326 })?;
1327 checks.require(reasoning.len() == 1, || {
1328 format!(
1329 "deltas and their full block must collapse to one reasoning item, observed {}",
1330 reasoning.len()
1331 )
1332 })?;
1333
1334 checks.note("summary aggregated exactly once");
1335 Ok(checks.report())
1336}
1337
1338pub async fn multi_part_same_id_reasoning_keeps_every_part(
1346 driver: &WireDriver,
1347 frames: Vec<WireInput>,
1348 expected_parts: &[&str],
1349) -> Result<ScenarioReport, ConformanceError> {
1350 let mut checks = Checks::new(
1351 "multi_part_same_id_reasoning_keeps_every_part",
1352 driver.provider,
1353 );
1354
1355 let drained = driver.drive(ok_chunks(frames)).await?;
1356 checks.require(
1357 drained.completed_cleanly(),
1358 || "the reasoning stream must complete without errors",
1359 )?;
1360 let observed: Vec<String> = drained
1361 .choice_reasoning()
1362 .iter()
1363 .flat_map(|item| item.content.iter())
1364 .map(|content| match content {
1365 crate::message::ReasoningContent::Summary(text) => text.clone(),
1366 crate::message::ReasoningContent::Text { text, .. } => text.clone(),
1367 crate::message::ReasoningContent::Encrypted(data) => data.clone(),
1368 crate::message::ReasoningContent::Redacted { data } => data.clone(),
1369 })
1370 .collect();
1371 checks.require(observed == expected_parts, || {
1372 format!(
1373 "every same-id reasoning part must survive in order: expected {expected_parts:?}, observed {observed:?}"
1374 )
1375 })?;
1376
1377 checks.note(format!(
1378 "all {} reasoning parts survived",
1379 expected_parts.len()
1380 ));
1381 Ok(checks.report())
1382}
1383
1384pub async fn interleaved_reasoning_aggregates_to_one_item(
1392 driver: &WireDriver,
1393 frames: Vec<WireInput>,
1394 expected_text: &str,
1395) -> Result<ScenarioReport, ConformanceError> {
1396 let mut checks = Checks::new(
1397 "interleaved_reasoning_aggregates_to_one_item",
1398 driver.provider,
1399 );
1400
1401 let drained = driver.drive(ok_chunks(frames)).await?;
1402 checks.require(
1403 drained.completed_cleanly(),
1404 || "the interleaved stream must complete without errors",
1405 )?;
1406 let reasoning = drained.choice_reasoning();
1407 checks.require(reasoning.len() == 1, || {
1408 format!(
1409 "interleaved deltas and their completed block must collapse to one reasoning item, observed {}",
1410 reasoning.len()
1411 )
1412 })?;
1413 let carries_text = reasoning
1414 .iter()
1415 .flat_map(|item| item.content.iter())
1416 .any(|content| match content {
1417 crate::message::ReasoningContent::Summary(text)
1418 | crate::message::ReasoningContent::Text { text, .. } => text == expected_text,
1419 _ => false,
1420 });
1421 checks.require(carries_text, || {
1422 format!("the reasoning item must carry the completed block's text {expected_text:?}")
1423 })?;
1424
1425 checks.note("exactly one reasoning item with the completed content");
1426 Ok(checks.report())
1427}
1428
1429pub async fn interleaved_constant_id_reasoning_preserves_order(
1438 fixture: &ProviderWireFixture,
1439) -> Result<ScenarioOutcome, ConformanceError> {
1440 let mut checks = Checks::new(
1441 "interleaved_constant_id_reasoning_preserves_order",
1442 fixture.driver.provider,
1443 );
1444 let Some(interleaved) = &fixture.interleaved_reasoning else {
1445 return Ok(checks.skip("wire fixture supplies no interleaved reasoning frames"));
1446 };
1447
1448 let drained = fixture
1449 .driver
1450 .drive(ok_chunks(interleaved.frames.clone()))
1451 .await?;
1452 checks.require(
1453 drained.completed_cleanly(),
1454 || "the interleaved stream must complete without errors",
1455 )?;
1456 assert_reasoning_tool_reasoning(
1457 &checks,
1458 &drained,
1459 interleaved.first_reasoning,
1460 interleaved.tool_name,
1461 interleaved.second_reasoning,
1462 )?;
1463
1464 checks.note("boundary kept: reasoning, tool call, reasoning in order");
1465 Ok(checks.ran())
1466}
1467
1468pub async fn interleaved_signed_full_reasoning_does_not_erase_prior_thought(
1477 driver: &WireDriver,
1478 frames: Vec<WireInput>,
1479 first: &str,
1480 tool_name: &str,
1481 second: &str,
1482) -> Result<ScenarioReport, ConformanceError> {
1483 let mut checks = Checks::new(
1484 "interleaved_signed_full_reasoning_does_not_erase_prior_thought",
1485 driver.provider,
1486 );
1487
1488 let drained = driver.drive(ok_chunks(frames)).await?;
1489 checks.require(
1490 drained.completed_cleanly(),
1491 || "the interleaved stream must complete without errors",
1492 )?;
1493 assert_reasoning_tool_reasoning(&checks, &drained, first, tool_name, second)?;
1494 let signed = drained.choice_reasoning().last().is_some_and(|reasoning| {
1495 reasoning.content.iter().any(|content| {
1496 matches!(
1497 content,
1498 crate::message::ReasoningContent::Text {
1499 signature: Some(_),
1500 ..
1501 }
1502 )
1503 })
1504 });
1505 checks.require(signed, || "the post-boundary block must keep its signature")?;
1506
1507 checks.note("pre-boundary thought survived; signed block completed the post-boundary part");
1508 Ok(checks.report())
1509}
1510
1511fn assert_reasoning_tool_reasoning(
1514 checks: &Checks,
1515 drained: &DrainedStream,
1516 first: &str,
1517 tool_name: &str,
1518 second: &str,
1519) -> Result<(), ConformanceError> {
1520 let shape: Vec<String> = drained
1521 .choice
1522 .iter()
1523 .map(|content| match content {
1524 AssistantContent::Reasoning(reasoning) => {
1525 let text: String = reasoning
1526 .value()
1527 .content
1528 .iter()
1529 .filter_map(|content| match content {
1530 crate::message::ReasoningContent::Summary(text)
1531 | crate::message::ReasoningContent::Text { text, .. } => {
1532 Some(text.as_str())
1533 }
1534 _ => None,
1535 })
1536 .collect();
1537 format!("reasoning:{text}")
1538 }
1539 AssistantContent::ToolCall(tool_call) => {
1540 format!("tool:{}", tool_call.function.name)
1541 }
1542 AssistantContent::Text(text) => format!("text:{}", text.text),
1543 AssistantContent::Image(_) => "image".to_string(),
1544 })
1545 .collect();
1546 let expected = vec![
1547 format!("reasoning:{first}"),
1548 format!("tool:{tool_name}"),
1549 format!("reasoning:{second}"),
1550 ];
1551 checks.require(shape == expected, || {
1552 format!("the boundary must survive aggregation: expected {expected:?}, observed {shape:?}")
1553 })
1554}
1555
1556pub mod fixtures {
1558 use super::*;
1559 use crate::driver::{Model, Transport};
1560 use crate::operation::Completion;
1561 use crate::test_utils::SequencedStreamingHttpClient;
1562 use crate::wire::Wire;
1563 use serde_json::json;
1564
1565 pub async fn drain(mut stream: crate::streaming::CompletionStream) -> DrainedStream {
1569 let mut items = Vec::new();
1570 while let Some(item) = stream.next().await {
1571 items.push(item.map_err(|error| ErrorReport::from(&error)));
1572 }
1573 let partial = stream.partial().choice;
1574 let response = stream.finish().await.ok();
1575 let drained = DrainedStream {
1576 items,
1577 choice: response
1578 .as_ref()
1579 .map_or(partial, |response| response.choice.clone()),
1580 response,
1581 };
1582 super::assert_valid_event_stream(&drained.items, &drained.choice);
1586 drained
1587 }
1588
1589 pub async fn drain_observed<W, T>(
1594 model: &Model<W, T>,
1595 request: crate::completion::CompletionRequest,
1596 ) -> Result<DrainedStream, ProviderError>
1597 where
1598 W: Wire<Op = Completion>,
1599 T: Transport<W>,
1600 {
1601 let log = std::sync::Arc::new(crate::observe::ObservationLog::default());
1602 let context = crate::observe::AdapterContext::new(
1603 log.clone(),
1604 crate::observe::Subject::default(),
1605 "conformance",
1606 );
1607 let drained = match model.stream_observed(request, context) {
1610 Ok(stream) => Ok(drain(stream).await),
1611 Err(error) => Err(error),
1612 };
1613 let events: Vec<_> = log
1614 .trace()
1615 .observations
1616 .iter()
1617 .filter_map(|o| match &o.action {
1618 crate::observe::Action::Adapter { observation } => Some(observation.event.clone()),
1619 _ => None,
1620 })
1621 .collect();
1622 assert!(
1623 matches!(
1624 events.first(),
1625 Some(crate::observe::AdapterEvent::Started { .. })
1626 ),
1627 "the wire must attach the observation context it was handed: {events:?}"
1628 );
1629 assert!(
1630 matches!(
1631 events.last(),
1632 Some(crate::observe::AdapterEvent::Finished { .. })
1633 ),
1634 "the attempt must close: {events:?}"
1635 );
1636 drained
1637 }
1638
1639 fn byte_chunks(chunks: WireChunks) -> Result<Vec<http_client::Result<Bytes>>, ProviderError> {
1643 chunks
1644 .into_iter()
1645 .map(|chunk| match chunk {
1646 Ok(WireInput::Bytes(bytes)) => Ok(Ok(bytes)),
1647 Ok(WireInput::Event(_)) => Err(ProviderError::Provider(
1648 "typed-event frame fed to a byte-transport driver".to_string(),
1649 )),
1650 Err(error) => Ok(Err(error)),
1651 })
1652 .collect()
1653 }
1654
1655 fn byte_driver<W>(
1663 provider: &'static str,
1664 bind: fn(SequencedStreamingHttpClient) -> Model<W, SequencedStreamingHttpClient>,
1665 ) -> WireDriver
1666 where
1667 W: Wire<Op = Completion, Payload = crate::wire::Encoded, Frame = crate::wire::WireFrame>,
1668 {
1669 WireDriver::new(provider, move |chunks| {
1670 Box::pin(async move {
1671 let model = bind(SequencedStreamingHttpClient::new(byte_chunks(chunks)?));
1672 let request = CompletionRequest::new("hello");
1673 drain_observed(&model, request).await
1674 })
1675 })
1676 }
1677
1678 fn sse(frame: &serde_json::Value) -> WireInput {
1679 WireInput::Bytes(Bytes::from(format!("data: {frame}\n\n")))
1680 }
1681
1682 fn sse_raw(data: &str) -> WireInput {
1683 WireInput::Bytes(Bytes::from(format!("data: {data}\n\n")))
1684 }
1685
1686 fn ndjson(frame: &serde_json::Value) -> WireInput {
1687 WireInput::Bytes(Bytes::from(format!("{frame}\n")))
1688 }
1689
1690 fn frame_text(frame: &WireInput) -> String {
1693 frame
1694 .as_bytes()
1695 .map(|bytes| String::from_utf8_lossy(bytes).into_owned())
1696 .unwrap_or_default()
1697 }
1698
1699 pub mod openai_chat {
1701 use super::*;
1702
1703 fn driver() -> WireDriver {
1704 byte_driver("openai", |transport| {
1705 crate::driver::Model::new(
1706 crate::providers::openai::wire::OpenAIConfig::with_key(
1707 &crate::providers::openai::wire::OPENAI,
1708 "test-key",
1709 )
1710 .chat("gpt-4o"),
1711 transport,
1712 )
1713 })
1714 }
1715
1716 pub fn fixture() -> ProviderWireFixture {
1718 ProviderWireFixture {
1719 driver: driver(),
1720 text_frames: vec![sse(&json!({
1721 "id": "chatcmpl-1",
1722 "model": "gpt-4o-2024-08-06",
1723 "choices": [{"index": 0, "delta": {"content": "hi"}, "finish_reason": null}],
1724 "usage": null,
1725 }))],
1726 expected_texts: vec!["hi"],
1727 tool_call_frames: vec![
1728 sse(&json!({
1729 "choices": [{"index": 0, "delta": {"tool_calls": [{
1730 "index": 0,
1731 "id": "call_1",
1732 "type": "function",
1733 "function": {"name": "get_weather", "arguments": ""},
1734 }]}, "finish_reason": null}],
1735 })),
1736 sse(&json!({
1737 "choices": [{"index": 0, "delta": {"tool_calls": [{
1738 "index": 0,
1739 "function": {"arguments": "{\"city\":\"Tokyo\"}"},
1740 }]}, "finish_reason": null}],
1741 })),
1742 ],
1746 expected_tool_name: "get_weather",
1747 partial_tool_call_frames: Some(vec![sse(&json!({
1748 "choices": [{"index": 0, "delta": {"tool_calls": [{
1749 "index": 0,
1750 "id": "call_1",
1751 "type": "function",
1752 "function": {"name": "get_weather", "arguments": "{\"cit"},
1753 }]}, "finish_reason": null}],
1754 }))]),
1755 terminal_frames: vec![
1756 sse(&json!({
1757 "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}],
1758 "usage": null,
1759 })),
1760 sse(&json!({
1761 "choices": [],
1762 "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15},
1763 })),
1764 sse_raw("[DONE]"),
1765 ],
1766 expected_usage_total: 15,
1767 expected_finish_reason: Some(FinishReason::Stop),
1768 zero_usage_terminal_frames: Some(vec![
1769 sse(&json!({
1770 "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}],
1771 "usage": null,
1772 })),
1773 sse_raw("[DONE]"),
1774 ]),
1775 bare_terminal_frames: Some(vec![sse_raw("[DONE]")]),
1776 malformed_frame: Some(sse_raw("{not json")),
1777 unknown_event_frame: None,
1778 defective_known_frame: Some(sse_raw(r#"{"choices": 42}"#)),
1782 delta_less_prelude_frame: Some(sse_raw(
1785 r#"{"id":"","object":"","choices":[{"prompt_index":0,"content_filter_results":{"hate":{"filtered":false,"severity":"safe"}}}]}"#,
1786 )),
1787 refusal: None,
1788 interleaved_reasoning: None,
1800 }
1801 }
1802 }
1803
1804 pub mod openai_responses {
1806 use super::*;
1807
1808 pub fn driver() -> WireDriver {
1810 byte_driver("openai", |transport| {
1811 crate::driver::Model::new(
1812 crate::providers::openai::OpenAIConfig::new("test-key").responses("gpt-5.4"),
1813 transport,
1814 )
1815 })
1816 }
1817
1818 fn completed_response(
1819 usage: Option<&serde_json::Value>,
1820 output: &serde_json::Value,
1821 ) -> serde_json::Value {
1822 json!({
1823 "id": "resp_1",
1824 "object": "response",
1825 "created_at": 0,
1826 "status": "completed",
1827 "model": "gpt-5.4",
1828 "output": output,
1829 "tools": [],
1830 "usage": usage,
1831 })
1832 }
1833
1834 fn terminal(usage: Option<&serde_json::Value>, output: &serde_json::Value) -> WireInput {
1835 sse(&json!({
1836 "type": "response.completed",
1837 "sequence_number": 99,
1838 "response": completed_response(usage, output),
1839 }))
1840 }
1841
1842 fn usage_json() -> serde_json::Value {
1843 json!({
1844 "input_tokens": 10,
1845 "output_tokens": 5,
1846 "output_tokens_details": {"reasoning_tokens": 0},
1847 "total_tokens": 15,
1848 })
1849 }
1850
1851 fn text_delta(text: &str) -> WireInput {
1852 sse(&json!({
1853 "type": "response.output_text.delta",
1854 "content_index": 0,
1855 "delta": text,
1856 "item_id": "msg_1",
1857 "output_index": 0,
1858 "sequence_number": 1,
1859 }))
1860 }
1861
1862 fn tool_call_done() -> WireInput {
1863 sse(&json!({
1864 "type": "response.output_item.done",
1865 "output_index": 0,
1866 "sequence_number": 2,
1867 "item": {
1868 "type": "function_call",
1869 "id": "fc_1",
1870 "arguments": "{\"city\":\"Tokyo\"}",
1871 "call_id": "call_1",
1872 "name": "get_weather",
1873 "status": "completed",
1874 },
1875 }))
1876 }
1877
1878 pub fn incomplete_mid_tool_call_frames() -> Vec<WireInput> {
1885 vec![
1886 sse(&json!({
1887 "type": "response.output_item.added",
1888 "output_index": 0,
1889 "sequence_number": 1,
1890 "item": {
1891 "type": "function_call",
1892 "id": "fc_1",
1893 "arguments": "",
1894 "call_id": "call_1",
1895 "name": "add",
1896 "status": "in_progress",
1897 },
1898 })),
1899 sse(&json!({
1900 "type": "response.function_call_arguments.delta",
1901 "item_id": "fc_1",
1902 "output_index": 0,
1903 "sequence_number": 2,
1904 "delta": "{\"x",
1905 })),
1906 sse(&json!({
1907 "type": "response.function_call_arguments.delta",
1908 "item_id": "fc_1",
1909 "output_index": 0,
1910 "sequence_number": 3,
1911 "delta": "\":48151",
1912 })),
1913 sse(&json!({
1914 "type": "response.function_call_arguments.done",
1915 "item_id": "fc_1",
1916 "output_index": 0,
1917 "sequence_number": 4,
1918 "arguments": "{\"x\":48151",
1919 })),
1920 sse(&json!({
1921 "type": "response.output_item.done",
1922 "output_index": 0,
1923 "sequence_number": 5,
1924 "item": {
1925 "type": "function_call",
1926 "id": "fc_1",
1927 "arguments": "{\"x\":48151",
1928 "call_id": "call_1",
1929 "name": "add",
1930 "status": "incomplete",
1931 },
1932 })),
1933 sse(&json!({
1934 "type": "response.incomplete",
1935 "sequence_number": 6,
1936 "response": {
1937 "id": "resp_1",
1938 "object": "response",
1939 "created_at": 0,
1940 "status": "incomplete",
1941 "incomplete_details": {"reason": "max_output_tokens"},
1942 "model": "gpt-5.4",
1943 "output": [{
1944 "type": "function_call",
1945 "id": "fc_1",
1946 "arguments": "{\"x\":48151",
1947 "call_id": "call_1",
1948 "name": "add",
1949 "status": "incomplete",
1950 }],
1951 "tools": [],
1952 "usage": usage_json(),
1953 },
1954 })),
1955 ]
1956 }
1957
1958 fn reasoning_done_item(
1959 id: &str,
1960 summary: &serde_json::Value,
1961 content: &serde_json::Value,
1962 encrypted: Option<&str>,
1963 ) -> WireInput {
1964 let mut item = json!({
1965 "type": "reasoning",
1966 "id": id,
1967 "summary": summary,
1968 "content": content,
1969 "status": "completed",
1970 });
1971 if let (Some(encrypted), Some(object)) = (encrypted, item.as_object_mut()) {
1972 object.insert("encrypted_content".to_string(), json!(encrypted));
1973 }
1974 sse(&json!({
1975 "type": "response.output_item.done",
1976 "output_index": 0,
1977 "sequence_number": 3,
1978 "item": item,
1979 }))
1980 }
1981
1982 pub fn fixture() -> ProviderWireFixture {
1984 ProviderWireFixture {
1985 driver: driver(),
1986 text_frames: vec![text_delta("hi")],
1987 expected_texts: vec!["hi"],
1988 tool_call_frames: vec![tool_call_done()],
1989 expected_tool_name: "get_weather",
1990 partial_tool_call_frames: Some(vec![
1991 sse(&json!({
1992 "type": "response.output_item.added",
1993 "output_index": 0,
1994 "sequence_number": 1,
1995 "item": {
1996 "type": "function_call",
1997 "id": "fc_1",
1998 "arguments": "",
1999 "call_id": "call_1",
2000 "name": "get_weather",
2001 "status": "in_progress",
2002 },
2003 })),
2004 sse(&json!({
2005 "type": "response.function_call_arguments.delta",
2006 "item_id": "fc_1",
2007 "output_index": 0,
2008 "sequence_number": 2,
2009 "delta": "{\"cit",
2010 })),
2011 ]),
2012 terminal_frames: vec![terminal(Some(&usage_json()), &json!([]))],
2013 expected_usage_total: 15,
2014 expected_finish_reason: Some(FinishReason::Stop),
2015 zero_usage_terminal_frames: Some(vec![terminal(None, &json!([]))]),
2016 bare_terminal_frames: None,
2017 malformed_frame: Some(sse_raw("{not json")),
2018 unknown_event_frame: Some(sse(&json!({
2019 "type": "response.web_search_call.searching",
2020 "output_index": 0,
2021 "sequence_number": 4,
2022 "item_id": "ws_1",
2023 }))),
2024 defective_known_frame: Some(sse(&json!({
2027 "type": "response.content_part.added",
2028 "item_id": "msg_1",
2029 "output_index": 0,
2030 "content_index": 0,
2031 "sequence_number": 5,
2032 "part": {"type": "output_text", "text": 42},
2033 }))),
2034 delta_less_prelude_frame: None,
2035 refusal: Some(RefusalFixture {
2036 frames: vec![sse(&json!({
2037 "type": "response.refusal.delta",
2038 "content_index": 0,
2039 "delta": "I cannot help with that.",
2040 "item_id": "msg_1",
2041 "output_index": 0,
2042 "sequence_number": 1,
2043 }))],
2044 expected_text: "I cannot help with that.",
2045 }),
2046 interleaved_reasoning: None,
2047 }
2048 }
2049
2050 pub fn buffered_driver() -> BufferedBodyDriver {
2060 BufferedBodyDriver::new("chatgpt", |body| {
2061 Box::pin(async move {
2062 let model = crate::driver::Model::new(
2063 crate::providers::openai::OpenAIConfig::with_key(
2064 &crate::providers::chatgpt::DIALECT,
2065 "test-token",
2066 )
2067 .with_account_id("account-id")
2068 .responses("gpt-5.4"),
2069 crate::test_utils::RecordingHttpClient::new(body),
2070 );
2071 let request = CompletionRequest::new("hello");
2072 let response = model.call(request).await?;
2073 Ok(response.choice)
2074 })
2075 })
2076 }
2077
2078 fn message_output(text: &str) -> serde_json::Value {
2079 json!([{
2080 "type": "message",
2081 "id": "msg_1",
2082 "role": "assistant",
2083 "status": "completed",
2084 "content": [{"type": "output_text", "text": text, "annotations": []}],
2085 }])
2086 }
2087
2088 pub fn terminal_body_only_sse_body(text: &str) -> String {
2090 frame_text(&terminal(Some(&usage_json()), &message_output(text)))
2091 }
2092
2093 pub fn terminal_body_and_delta_sse_body(text: &str) -> String {
2095 let frames = [
2096 text_delta(text),
2097 terminal(Some(&usage_json()), &message_output(text)),
2098 ];
2099 frames.iter().map(frame_text).collect()
2100 }
2101
2102 pub fn delta_only_sse_body(text: &str) -> String {
2105 let frames = [text_delta(text), terminal(Some(&usage_json()), &json!([]))];
2106 frames.iter().map(frame_text).collect()
2107 }
2108
2109 pub fn envelope_less_reasoning_supersede_sse_body() -> (String, &'static str) {
2116 let delta = json!({
2117 "type": "response.reasoning_summary_text.delta",
2118 "delta": "step 1",
2119 });
2120 let frames = [
2121 sse(&delta),
2122 reasoning_done_item(
2123 "rs_1",
2124 &json!([{"type": "summary_text", "text": "step 1"}]),
2125 &json!([]),
2126 None,
2127 ),
2128 terminal(Some(&usage_json()), &json!([])),
2129 ];
2130 (frames.iter().map(frame_text).collect(), "step 1")
2131 }
2132
2133 pub fn reasoning_summary_supersede_frames() -> (Vec<WireInput>, &'static str) {
2137 let frames = vec![
2138 sse(&json!({
2139 "type": "response.reasoning_summary_text.delta",
2140 "item_id": "rs_1",
2141 "output_index": 0,
2142 "summary_index": 0,
2143 "sequence_number": 1,
2144 "delta": "step 1",
2145 })),
2146 reasoning_done_item(
2147 "rs_1",
2148 &json!([{"type": "summary_text", "text": "step 1"}]),
2149 &json!([]),
2150 None,
2151 ),
2152 terminal(Some(&usage_json()), &json!([])),
2153 ];
2154 (frames, "step 1")
2155 }
2156
2157 pub fn multi_part_reasoning_frames() -> (Vec<WireInput>, Vec<&'static str>) {
2160 let frames = vec![
2161 reasoning_done_item(
2162 "rs_1",
2163 &json!([
2164 {"type": "summary_text", "text": "s1"},
2165 {"type": "summary_text", "text": "s2"},
2166 ]),
2167 &json!([{"type": "reasoning_text", "text": "visible"}]),
2168 Some("enc_blob"),
2169 ),
2170 terminal(Some(&usage_json()), &json!([])),
2171 ];
2172 (frames, vec!["s1", "s2", "visible", "enc_blob"])
2173 }
2174
2175 pub fn interleaved_reasoning_frames() -> (Vec<WireInput>, &'static str) {
2178 let frames = vec![
2179 sse(&json!({
2180 "type": "response.reasoning_text.delta",
2181 "item_id": "rs_2",
2182 "output_index": 0,
2183 "content_index": 0,
2184 "sequence_number": 1,
2185 "delta": "thinking",
2186 })),
2187 tool_call_done(),
2188 reasoning_done_item(
2189 "rs_2",
2190 &json!([]),
2191 &json!([{"type": "reasoning_text", "text": "full reasoning"}]),
2192 None,
2193 ),
2194 terminal(Some(&usage_json()), &json!([])),
2195 ];
2196 (frames, "full reasoning")
2197 }
2198 }
2199
2200 pub mod gemini_rest {
2202 use super::*;
2203
2204 fn driver() -> WireDriver {
2205 byte_driver("gemini", |transport| {
2206 crate::driver::Model::new(
2207 crate::providers::gemini::GeminiConfig::new("test-key").completion(
2208 crate::providers::gemini::completion::GEMINI_2_5_PRO_PREVIEW_06_05,
2209 ),
2210 transport,
2211 )
2212 })
2213 }
2214
2215 pub fn fixture() -> ProviderWireFixture {
2217 ProviderWireFixture {
2218 driver: driver(),
2219 text_frames: vec![sse(&json!({
2220 "candidates": [{"content": {"parts": [{"text": "hi"}], "role": "model"}}],
2221 "responseId": "resp-1",
2222 "modelVersion": "gemini-2.5-pro",
2223 }))],
2224 expected_texts: vec!["hi"],
2225 tool_call_frames: vec![sse(&json!({
2226 "candidates": [{"content": {"parts": [{
2227 "functionCall": {"name": "get_weather", "args": {"city": "Tokyo"}},
2228 }], "role": "model"}}],
2229 "responseId": "resp-1",
2230 "modelVersion": "gemini-2.5-pro",
2231 }))],
2232 expected_tool_name: "get_weather",
2233 partial_tool_call_frames: None,
2235 terminal_frames: vec![sse(&json!({
2236 "candidates": [{
2237 "content": {"parts": [], "role": "model"},
2238 "finishReason": "STOP",
2239 }],
2240 "usageMetadata": {
2241 "promptTokenCount": 5,
2242 "candidatesTokenCount": 2,
2243 "totalTokenCount": 7,
2244 },
2245 "responseId": "resp-1",
2246 "modelVersion": "gemini-2.5-pro",
2247 }))],
2248 expected_usage_total: 7,
2249 expected_finish_reason: Some(FinishReason::Stop),
2250 zero_usage_terminal_frames: Some(vec![sse(&json!({
2251 "candidates": [{
2252 "content": {"parts": [], "role": "model"},
2253 "finishReason": "STOP",
2254 }],
2255 "responseId": "resp-1",
2256 "modelVersion": "gemini-2.5-pro",
2257 }))]),
2258 bare_terminal_frames: None,
2259 malformed_frame: Some(sse_raw("{not json")),
2260 unknown_event_frame: Some(sse_raw(r#"{"noise":true}"#)),
2264 defective_known_frame: Some(sse_raw(r#"{"candidates": 42}"#)),
2265 delta_less_prelude_frame: None,
2266 refusal: None,
2267 interleaved_reasoning: Some(interleaved_thought_fixture()),
2268 }
2269 }
2270
2271 fn chunk(parts: &serde_json::Value) -> WireInput {
2272 sse(&json!({
2273 "candidates": [{"content": {"parts": parts, "role": "model"}}],
2274 "responseId": "resp-1",
2275 "modelVersion": "gemini-2.5-pro",
2276 }))
2277 }
2278
2279 fn terminal_frame() -> WireInput {
2280 sse(&json!({
2281 "candidates": [{
2282 "content": {"parts": [], "role": "model"},
2283 "finishReason": "STOP",
2284 }],
2285 "usageMetadata": {
2286 "promptTokenCount": 5,
2287 "candidatesTokenCount": 2,
2288 "totalTokenCount": 7,
2289 },
2290 "responseId": "resp-1",
2291 "modelVersion": "gemini-2.5-pro",
2292 }))
2293 }
2294
2295 fn interleaved_thought_fixture() -> InterleavedReasoningFixture {
2298 InterleavedReasoningFixture {
2299 frames: vec![
2300 chunk(&json!([{"text": "before tool", "thought": true}])),
2301 chunk(&json!([{
2302 "functionCall": {"name": "get_weather", "args": {"city": "Tokyo"}},
2303 }])),
2304 chunk(&json!([{"text": "after tool", "thought": true}])),
2305 terminal_frame(),
2306 ],
2307 first_reasoning: "before tool",
2308 tool_name: "get_weather",
2309 second_reasoning: "after tool",
2310 }
2311 }
2312
2313 pub fn interleaved_signed_thought_frames()
2316 -> (Vec<WireInput>, &'static str, &'static str, &'static str) {
2317 let frames = vec![
2318 chunk(&json!([{"text": "before tool", "thought": true}])),
2319 chunk(&json!([{
2320 "functionCall": {"name": "get_weather", "args": {"city": "Tokyo"}},
2321 }])),
2322 chunk(&json!([{
2323 "text": "signed conclusion",
2324 "thought": true,
2325 "thoughtSignature": "sig-1",
2326 }])),
2327 terminal_frame(),
2328 ];
2329 (frames, "before tool", "get_weather", "signed conclusion")
2330 }
2331 }
2332
2333 pub mod interactions {
2335 use super::*;
2336
2337 fn driver() -> WireDriver {
2338 byte_driver("gemini", |transport| {
2339 crate::driver::Model::new(
2340 crate::providers::gemini::GeminiConfig::new("test-key")
2341 .interactions("gemini-2.5-pro"),
2342 transport,
2343 )
2344 })
2345 }
2346
2347 fn completed(usage: Option<serde_json::Value>) -> WireInput {
2348 let mut interaction = json!({
2349 "id": "int-1",
2350 "model": "gemini-2.5-pro",
2351 "status": "completed",
2352 });
2353 if let (Some(usage), Some(object)) = (usage, interaction.as_object_mut()) {
2354 object.insert("usage".to_string(), usage);
2355 }
2356 sse(&json!({
2357 "event_type": "interaction.completed",
2358 "interaction": interaction,
2359 }))
2360 }
2361
2362 pub fn fixture() -> ProviderWireFixture {
2364 ProviderWireFixture {
2365 driver: driver(),
2366 text_frames: vec![sse(&json!({
2367 "event_type": "step.delta",
2368 "index": 0,
2369 "delta": {"type": "text", "text": "hi"},
2370 }))],
2371 expected_texts: vec!["hi"],
2372 tool_call_frames: vec![sse(&json!({
2373 "event_type": "step.delta",
2374 "index": 0,
2375 "delta": {
2376 "type": "function_call",
2377 "name": "get_weather",
2378 "arguments": {"city": "Tokyo"},
2379 "id": "call-1",
2380 },
2381 }))],
2382 expected_tool_name: "get_weather",
2383 partial_tool_call_frames: None,
2386 terminal_frames: vec![completed(Some(json!({
2387 "total_input_tokens": 5,
2388 "total_output_tokens": 2,
2389 "total_tokens": 7,
2390 })))],
2391 expected_usage_total: 7,
2392 expected_finish_reason: Some(FinishReason::Stop),
2393 zero_usage_terminal_frames: Some(vec![completed(None)]),
2394 bare_terminal_frames: None,
2395 malformed_frame: Some(sse_raw("{not json")),
2396 unknown_event_frame: Some(sse(&json!({
2397 "event_type": "future.event",
2398 "index": 0,
2399 }))),
2400 defective_known_frame: Some(sse_raw(
2403 r#"{"event_type":"step.delta","index":0,"delta":42}"#,
2404 )),
2405 delta_less_prelude_frame: None,
2406 refusal: None,
2407 interleaved_reasoning: Some(interleaved_thought_fixture()),
2408 }
2409 }
2410
2411 fn interleaved_thought_fixture() -> InterleavedReasoningFixture {
2415 let frames = vec![
2416 sse(&json!({
2417 "event_type": "step.delta",
2418 "index": 0,
2419 "delta": {
2420 "type": "thought_summary",
2421 "content": {"type": "text", "text": "before tool"},
2422 },
2423 })),
2424 sse(&json!({
2425 "event_type": "step.delta",
2426 "index": 0,
2427 "delta": {
2428 "type": "function_call",
2429 "name": "get_weather",
2430 "arguments": {"city": "Tokyo"},
2431 "id": "call-1",
2432 },
2433 })),
2434 sse(&json!({
2435 "event_type": "step.delta",
2436 "index": 0,
2437 "delta": {
2438 "type": "thought_summary",
2439 "content": {"type": "text", "text": "after tool"},
2440 },
2441 })),
2442 completed(Some(json!({
2443 "total_input_tokens": 5,
2444 "total_output_tokens": 2,
2445 "total_tokens": 7,
2446 }))),
2447 ];
2448 InterleavedReasoningFixture {
2449 frames,
2450 first_reasoning: "before tool",
2451 tool_name: "get_weather",
2452 second_reasoning: "after tool",
2453 }
2454 }
2455 }
2456
2457 pub mod anthropic {
2459 use super::*;
2460
2461 fn driver() -> WireDriver {
2462 byte_driver("anthropic", |transport| {
2463 crate::driver::Model::new(
2464 crate::providers::anthropic::wire::AnthropicConfig::new("test-key")
2465 .completion(crate::providers::anthropic::completion::CLAUDE_SONNET_4_6),
2466 transport,
2467 )
2468 })
2469 }
2470
2471 fn message_start() -> WireInput {
2472 sse(&json!({
2473 "type": "message_start",
2474 "message": {
2475 "id": "msg_1",
2476 "role": "assistant",
2477 "content": [],
2478 "model": "claude-sonnet-4-6",
2479 "stop_reason": null,
2480 "stop_sequence": null,
2481 "usage": {"input_tokens": 5, "output_tokens": 0},
2482 },
2483 }))
2484 }
2485
2486 pub fn fixture() -> ProviderWireFixture {
2488 ProviderWireFixture {
2489 driver: driver(),
2490 text_frames: vec![
2491 message_start(),
2492 sse(&json!({
2493 "type": "content_block_start",
2494 "index": 0,
2495 "content_block": {"type": "text", "text": ""},
2496 })),
2497 sse(&json!({
2498 "type": "content_block_delta",
2499 "index": 0,
2500 "delta": {"type": "text_delta", "text": "hi"},
2501 })),
2502 ],
2503 expected_texts: vec!["hi"],
2504 tool_call_frames: vec![
2505 sse(&json!({
2506 "type": "content_block_start",
2507 "index": 0,
2508 "content_block": {
2509 "type": "tool_use",
2510 "id": "toolu_1",
2511 "name": "get_weather",
2512 "input": {},
2513 },
2514 })),
2515 sse(&json!({
2516 "type": "content_block_delta",
2517 "index": 0,
2518 "delta": {"type": "input_json_delta", "partial_json": "{\"city\":\"Tokyo\"}"},
2519 })),
2520 sse(&json!({"type": "content_block_stop", "index": 0})),
2523 ],
2524 expected_tool_name: "get_weather",
2525 partial_tool_call_frames: Some(vec![
2526 sse(&json!({
2527 "type": "content_block_start",
2528 "index": 0,
2529 "content_block": {
2530 "type": "tool_use",
2531 "id": "toolu_1",
2532 "name": "get_weather",
2533 "input": {},
2534 },
2535 })),
2536 sse(&json!({
2537 "type": "content_block_delta",
2538 "index": 0,
2539 "delta": {"type": "input_json_delta", "partial_json": "{\"cit"},
2540 })),
2541 ]),
2542 terminal_frames: vec![sse(&json!({
2543 "type": "message_delta",
2544 "delta": {"stop_reason": "end_turn", "stop_sequence": null},
2545 "usage": {"output_tokens": 4},
2546 }))],
2547 expected_usage_total: 9,
2549 expected_finish_reason: Some(FinishReason::Stop),
2550 zero_usage_terminal_frames: None,
2553 bare_terminal_frames: Some(vec![sse(&json!({"type": "message_stop"}))]),
2556 malformed_frame: Some(sse_raw("{not json")),
2557 unknown_event_frame: Some(sse(&json!({
2558 "type": "content_block_heartbeat",
2559 "index": 0,
2560 }))),
2561 defective_known_frame: Some(sse_raw(
2564 r#"{"type":"content_block_delta","index":0,"delta":42}"#,
2565 )),
2566 delta_less_prelude_frame: None,
2567 refusal: None,
2568 interleaved_reasoning: None,
2569 }
2570 }
2571 }
2572
2573 pub mod cohere {
2575 use super::*;
2576
2577 fn driver() -> WireDriver {
2578 byte_driver("cohere", |transport| {
2579 crate::driver::Model::new(
2580 crate::providers::cohere::wire::CohereConfig::new("test-key")
2581 .completion(crate::providers::cohere::COMMAND_R_08_2024),
2582 transport,
2583 )
2584 })
2585 }
2586
2587 pub fn fixture() -> ProviderWireFixture {
2589 ProviderWireFixture {
2590 driver: driver(),
2591 text_frames: vec![
2592 sse(&json!({"type": "message-start", "id": "msg_1"})),
2593 sse(&json!({
2594 "type": "content-delta",
2595 "delta": {"message": {"content": {"text": "hi"}}},
2596 })),
2597 ],
2598 expected_texts: vec!["hi"],
2599 tool_call_frames: vec![
2600 sse(&json!({
2601 "type": "tool-call-start",
2602 "delta": {"message": {"tool_calls": {
2603 "id": "call_1",
2604 "function": {"name": "get_weather", "arguments": ""},
2605 }}},
2606 })),
2607 sse(&json!({
2608 "type": "tool-call-delta",
2609 "delta": {"message": {"tool_calls": {
2610 "function": {"arguments": "{\"city\":\"Tokyo\"}"},
2611 }}},
2612 })),
2613 sse(&json!({"type": "tool-call-end"})),
2614 ],
2615 expected_tool_name: "get_weather",
2616 partial_tool_call_frames: Some(vec![sse(&json!({
2617 "type": "tool-call-start",
2618 "delta": {"message": {"tool_calls": {
2619 "id": "call_1",
2620 "function": {"name": "get_weather", "arguments": "{\"cit"},
2621 }}},
2622 }))]),
2623 terminal_frames: vec![sse(&json!({
2624 "type": "message-end",
2625 "delta": {
2626 "finish_reason": "COMPLETE",
2627 "usage": {"tokens": {"input_tokens": 10, "output_tokens": 4}},
2628 },
2629 }))],
2630 expected_usage_total: 14,
2631 expected_finish_reason: Some(FinishReason::Stop),
2632 zero_usage_terminal_frames: Some(vec![sse(&json!({"type": "message-end"}))]),
2633 bare_terminal_frames: None,
2634 malformed_frame: Some(sse_raw("{not json")),
2635 unknown_event_frame: Some(sse(&json!({
2636 "type": "citation-start",
2637 "delta": {"message": {"citations": {}}},
2638 }))),
2639 defective_known_frame: Some(sse_raw(r#"{"type":"content-delta","delta":42}"#)),
2640 delta_less_prelude_frame: None,
2641 refusal: None,
2642 interleaved_reasoning: Some(interleaved_thinking_fixture()),
2643 }
2644 }
2645
2646 fn interleaved_thinking_fixture() -> InterleavedReasoningFixture {
2650 let frames = vec![
2651 sse(&json!({"type": "message-start", "id": "msg_1"})),
2652 sse(&json!({
2653 "type": "content-delta",
2654 "delta": {"message": {"content": {"thinking": "before tool"}}},
2655 })),
2656 sse(&json!({
2657 "type": "tool-call-start",
2658 "delta": {"message": {"tool_calls": {
2659 "id": "call_1",
2660 "function": {"name": "get_weather", "arguments": "{\"city\":\"Tokyo\"}"},
2661 }}},
2662 })),
2663 sse(&json!({"type": "tool-call-end"})),
2664 sse(&json!({
2665 "type": "content-delta",
2666 "delta": {"message": {"content": {"thinking": "after tool"}}},
2667 })),
2668 sse(&json!({
2669 "type": "message-end",
2670 "delta": {
2671 "finish_reason": "COMPLETE",
2672 "usage": {"tokens": {"input_tokens": 10, "output_tokens": 4}},
2673 },
2674 })),
2675 ];
2676 InterleavedReasoningFixture {
2677 frames,
2678 first_reasoning: "before tool",
2679 tool_name: "get_weather",
2680 second_reasoning: "after tool",
2681 }
2682 }
2683 }
2684
2685 pub mod ollama {
2687 use super::*;
2688
2689 fn driver() -> WireDriver {
2690 byte_driver("ollama", |transport| {
2691 crate::driver::Model::new(
2692 crate::providers::ollama::wire::OllamaConfig::new().completion("llama3.2"),
2693 transport,
2694 )
2695 })
2696 }
2697
2698 pub fn fixture() -> ProviderWireFixture {
2700 ProviderWireFixture {
2701 driver: driver(),
2702 text_frames: vec![ndjson(&json!({
2703 "model": "llama3.2",
2704 "created_at": "2023-08-04T19:22:45.499127Z",
2705 "message": {"role": "assistant", "content": "hi"},
2706 "done": false,
2707 }))],
2708 expected_texts: vec!["hi"],
2709 tool_call_frames: vec![ndjson(&json!({
2710 "model": "llama3.2",
2711 "created_at": "2023-08-04T19:22:45.499127Z",
2712 "message": {"role": "assistant", "content": "", "tool_calls": [{
2713 "function": {"name": "get_weather", "arguments": {"city": "Tokyo"}},
2714 }]},
2715 "done": false,
2716 }))],
2717 expected_tool_name: "get_weather",
2718 partial_tool_call_frames: None,
2720 terminal_frames: vec![ndjson(&json!({
2721 "model": "llama3.2",
2722 "created_at": "2023-08-04T19:22:47.499127Z",
2723 "message": {"role": "assistant", "content": ""},
2724 "done": true,
2725 "done_reason": "stop",
2726 "prompt_eval_count": 10,
2727 "eval_count": 4,
2728 }))],
2729 expected_usage_total: 14,
2730 expected_finish_reason: Some(FinishReason::Stop),
2731 zero_usage_terminal_frames: Some(vec![ndjson(&json!({
2732 "model": "llama3.2",
2733 "created_at": "2023-08-04T19:22:47.499127Z",
2734 "message": {"role": "assistant", "content": ""},
2735 "done": true,
2736 "done_reason": "stop",
2737 }))]),
2738 bare_terminal_frames: None,
2739 malformed_frame: Some(WireInput::Bytes(Bytes::from_static(b"{not json\n"))),
2740 unknown_event_frame: None,
2741 defective_known_frame: Some(ndjson(&json!({
2742 "model": "llama3.2",
2743 "created_at": "2023-08-04T19:22:46.499127Z",
2744 "message": {"role": "assistant", "content": 42},
2745 "done": false,
2746 }))),
2747 delta_less_prelude_frame: None,
2748 refusal: None,
2749 interleaved_reasoning: Some(interleaved_thinking_fixture()),
2750 }
2751 }
2752
2753 fn interleaved_thinking_fixture() -> InterleavedReasoningFixture {
2756 let frames = vec![
2757 ndjson(&json!({
2758 "model": "llama3.2",
2759 "created_at": "2023-08-04T19:22:45.499127Z",
2760 "message": {"role": "assistant", "content": "", "thinking": "before tool"},
2761 "done": false,
2762 })),
2763 ndjson(&json!({
2764 "model": "llama3.2",
2765 "created_at": "2023-08-04T19:22:45.599127Z",
2766 "message": {"role": "assistant", "content": "", "tool_calls": [{
2767 "function": {"name": "get_weather", "arguments": {"city": "Tokyo"}},
2768 }]},
2769 "done": false,
2770 })),
2771 ndjson(&json!({
2772 "model": "llama3.2",
2773 "created_at": "2023-08-04T19:22:45.699127Z",
2774 "message": {"role": "assistant", "content": "", "thinking": "after tool"},
2775 "done": false,
2776 })),
2777 ndjson(&json!({
2778 "model": "llama3.2",
2779 "created_at": "2023-08-04T19:22:47.499127Z",
2780 "message": {"role": "assistant", "content": ""},
2781 "done": true,
2782 "done_reason": "stop",
2783 "prompt_eval_count": 10,
2784 "eval_count": 4,
2785 })),
2786 ];
2787 InterleavedReasoningFixture {
2788 frames,
2789 first_reasoning: "before tool",
2790 tool_name: "get_weather",
2791 second_reasoning: "after tool",
2792 }
2793 }
2794 }
2795}