From e7740b4acb0a55444f414fdf8c4d9b9afd91e66e Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 14:21:15 -0400 Subject: [PATCH 1/3] feat(reasoning): capture provider reasoning in agent.message Normalize the readable reasoning providers already return into a canonical `ReasoningOutput` and carry it through the `agent.message` run event to storage, SSE, and JSONL. The shape is derived from the final response's canonical message content rather than stored a second time, so there is no duplicate source of truth and retried or replaced streaming buffers never become durable reasoning. OpenAI-compatible `reasoning_details` are now preserved verbatim as an opaque content part; only known readable members are normalized out of them, leaving encrypted entries for a later provider-aware replay phase. This phase is passive: no request parameters change, no capability guessing, and no newly observed provider field is replayed. Co-Authored-By: Claude Opus 5 (1M context) --- docs/public/api-reference/fabro-api.yaml | 21 ++ .../src/commands/run/run_progress/mod.rs | 1 + lib/apps/fabro-server/src/demo/mod.rs | 2 + .../fabro-server/src/server/handler/pair.rs | 2 + lib/apps/fabro-server/src/server/tests.rs | 74 ++++ lib/components/fabro-agent/src/session.rs | 113 +++++- lib/components/fabro-agent/src/types.rs | 8 +- .../src/codec/openai_compatible/response.rs | 7 +- .../src/codec/openai_compatible/stream.rs | 12 +- .../src/codec/openai_compatible/wire.rs | 93 ++++- lib/components/fabro-llm/src/lib.rs | 1 + .../fabro-llm/src/providers/fabro_server.rs | 47 +++ lib/components/fabro-llm/src/reasoning.rs | 337 ++++++++++++++++++ lib/components/fabro-llm/src/types.rs | 18 +- .../tests/it/wire/openai_compatible.rs | 194 ++++++++++ .../tests/it/wire/openai_responses.rs | 36 ++ lib/components/fabro-store/src/run_state.rs | 1 + .../fabro-workflow/src/event/convert.rs | 58 +++ .../fabro-workflow/src/event/redaction.rs | 46 ++- .../fabro-workflow/src/event/sink.rs | 49 +++ lib/foundation/fabro-api/build.rs | 1 + lib/foundation/fabro-api/src/lib.rs | 34 +- .../tests/reasoning_output_round_trip.rs | 94 +++++ lib/foundation/fabro-types/src/lib.rs | 2 + lib/foundation/fabro-types/src/reasoning.rs | 96 +++++ .../fabro-types/src/run_event/agent.rs | 35 +- .../fabro-types/src/run_event/mod.rs | 64 ++++ lib/foundation/fabro-types/src/transcript.rs | 5 + .../src/.openapi-generator/FILES | 1 + .../src/models/agent-message-props.ts | 4 + .../fabro-api-client/src/models/index.ts | 1 + .../src/models/reasoning-output.ts | 29 ++ 32 files changed, 1460 insertions(+), 26 deletions(-) create mode 100644 lib/components/fabro-llm/src/reasoning.rs create mode 100644 lib/foundation/fabro-api/tests/reasoning_output_round_trip.rs create mode 100644 lib/foundation/fabro-types/src/reasoning.rs create mode 100644 lib/packages/fabro-api-client/src/models/reasoning-output.ts diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 81794e101..9759233bf 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -10045,6 +10045,27 @@ components: - $ref: "#/components/schemas/StageContextWindowProjection" - type: "null" description: Latest content-free context-window projection for this agent stage. + reasoning: + oneOf: + - $ref: "#/components/schemas/ReasoningOutput" + - type: "null" + description: Readable reasoning the provider returned with this response, if any. + + ReasoningOutput: + description: >- + Readable model reasoning normalized into a provider-neutral shape. + Both members may be present for the same response, and at least one is + present whenever the object is emitted. Opaque provider material + (signatures, item IDs, encrypted or redacted payloads) never appears + here. + type: object + properties: + summary: + type: string + description: Model-authored summary of its reasoning. + trace: + type: string + description: Verbatim readable reasoning text, when the provider returns it. AgentToolsAvailableProps: description: Properties for the `agent.tools.available` event. diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index c5ef114d3..97799bfc1 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -534,6 +534,7 @@ mod tests { cost_source: None, tool_call_count: 0, context_window: None, + reasoning: None, }) } diff --git a/lib/apps/fabro-server/src/demo/mod.rs b/lib/apps/fabro-server/src/demo/mod.rs index 9fa69379b..aeb913826 100644 --- a/lib/apps/fabro-server/src/demo/mod.rs +++ b/lib/apps/fabro-server/src/demo/mod.rs @@ -1502,6 +1502,7 @@ mod runs { visit: 1, message: None, context_window: None, + reasoning: None, }), ), make_envelope( @@ -1572,6 +1573,7 @@ mod runs { visit: 1, message: None, context_window: None, + reasoning: None, }), ), ] diff --git a/lib/apps/fabro-server/src/server/handler/pair.rs b/lib/apps/fabro-server/src/server/handler/pair.rs index 67bc1ad05..82208b592 100644 --- a/lib/apps/fabro-server/src/server/handler/pair.rs +++ b/lib/apps/fabro-server/src/server/handler/pair.rs @@ -892,6 +892,7 @@ mod tests { visit: 1, message: None, context_window: None, + reasoning: None, }), ), ) @@ -925,6 +926,7 @@ mod tests { visit: 1, message: None, context_window: None, + reasoning: None, }), ), ) diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 1a86ef59f..64fe5f345 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -4171,6 +4171,7 @@ fn context_window_event( cost_source: None, tool_call_count: 0, context_window: Some(context_window), + reasoning: None, }, session_id: Some("session-1".to_string()), parent_session_id: None, @@ -15382,6 +15383,79 @@ async fn cancel_before_run_transitions_to_running_returns_empty_attach_stream() assert!(body.is_empty(), "expected an empty attach stream"); } +/// Reasoning has to survive the whole durable path, not just the local +/// struct conversion: emitted event → run store → attach SSE JSON. +#[tokio::test] +async fn attach_stream_replays_agent_message_reasoning() { + let state = test_app_state(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = fixtures::RUN_1; + + create_durable_run_with_events(&state, run_id, &[ + stage_started_event("code", "agent"), + workflow_event::Event::Agent { + stage: "code".to_string(), + visit: 1, + event: fabro_agent::AgentEvent::AssistantMessage { + text: String::new(), + model: ModelRef { + provider: ProviderId::openai(), + model_id: "gpt-5.4".into(), + speed: None, + }, + usage: TokenCounts::default(), + cost_usd: None, + cost_source: None, + tool_call_count: 1, + context_window: None, + reasoning: Some(fabro_types::ReasoningOutput { + summary: Some("inspect the sink first".to_string()), + trace: Some("read events.rs, then attach".to_string()), + }), + }, + session_id: Some("session-1".to_string()), + parent_session_id: None, + tool_call_id: None, + }, + workflow_event::Event::WorkflowRunCompleted { + timing: fabro_types::RunTiming::wall_only(1000), + artifact_count: 0, + status: "succeeded".to_string(), + reason: SuccessReason::Completed, + total_usd_micros: None, + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + billing: None, + }, + ]) + .await; + + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/attach?since_seq=1"))) + .body(Body::empty()) + .unwrap(); + let response = app.oneshot(req).await.unwrap(); + let body = response_bytes!(response, StatusCode::OK).await; + let text = String::from_utf8(body.clone()).unwrap(); + + let message = text + .lines() + .filter_map(|line| line.strip_prefix("data: ")) + .filter_map(|data| serde_json::from_str::(data).ok()) + .find(|value| value["event"] == "agent.message") + .expect("attach stream should replay the agent message"); + assert_eq!( + message["properties"]["reasoning"]["summary"], + "inspect the sink first" + ); + assert_eq!( + message["properties"]["reasoning"]["trace"], + "read events.rs, then attach" + ); +} + #[tokio::test] async fn queue_position_reported_for_runnable_runs() { let state = test_app_state(); diff --git a/lib/components/fabro-agent/src/session.rs b/lib/components/fabro-agent/src/session.rs index b36c8a639..a7048d833 100644 --- a/lib/components/fabro-agent/src/session.rs +++ b/lib/components/fabro-agent/src/session.rs @@ -1762,6 +1762,8 @@ impl Session { // Record assistant turn let text = response.text(); let tool_calls = response.tool_calls(); + // Normalize before the response's content moves into history. + let reasoning = response.reasoning_output(); let provider_parts: Vec<_> = response .message .content @@ -1805,6 +1807,7 @@ impl Session { cost_source: response.cost_source, tool_call_count: tool_calls.len(), context_window, + reasoning, }); // Post-response compaction: trim context after appending assistant turn @@ -2113,7 +2116,7 @@ mod tests { ContentPart, ReasoningEffort, Request, Response, Role, StreamEvent, TokenCounts, ToolCall, ToolDefinition, }; - use fabro_types::StageContextWindowCountMethod; + use fabro_types::{ReasoningOutput, StageContextWindowCountMethod}; use futures::stream; use tokio::time::{sleep, timeout}; @@ -3924,6 +3927,114 @@ mod tests { ]); } + /// Builds a response whose provider parts carry both reasoning channels. + fn reasoning_response(text: &str, summary: &str, trace: &str) -> Response { + let mut response = text_response(text); + let mut content = vec![ContentPart::Other { + kind: ContentPart::OPENAI_COMPAT_REASONING_DETAILS.to_string(), + data: serde_json::json!([ + {"type": "reasoning.summary", "summary": summary}, + {"type": "reasoning.text", "text": trace}, + ]), + }]; + content.extend(response.message.content); + response.message.content = content; + response + } + + fn collect_message_reasoning( + rx: &mut broadcast::Receiver, + ) -> Vec> { + let mut collected = Vec::new(); + while let Ok(event) = rx.try_recv() { + if let AgentEvent::AssistantMessage { reasoning, .. } = event.event { + collected.push(reasoning); + } + } + collected + } + + #[tokio::test] + async fn completed_response_emits_normalized_reasoning_once() { + let mut session = make_session(vec![reasoning_response( + "4.", + "the user wants 2+2", + "2+2 is 4", + )]) + .await; + let mut rx = session.subscribe(); + + session.process_input("What is 2+2?").await.unwrap(); + + let reasoning = collect_message_reasoning(&mut rx); + assert_eq!(reasoning, vec![Some(ReasoningOutput { + summary: Some("the user wants 2+2".to_string()), + trace: Some("2+2 is 4".to_string()), + })]); + } + + #[tokio::test] + async fn tool_call_response_with_no_visible_text_still_carries_reasoning() { + let mut tool_call = tool_call_response("nonexistent_tool", "call_1", serde_json::json!({})); + // Drop the visible text so only the tool call and reasoning remain. + tool_call.message.content = vec![ + ContentPart::Other { + kind: ContentPart::OPENAI_COMPAT_REASONING_DETAILS.to_string(), + data: serde_json::json!([{"type": "reasoning.summary", "summary": "call the tool"}]), + }, + ContentPart::ToolCall(ToolCall::new( + "call_1", + "nonexistent_tool", + serde_json::json!({}), + )), + ]; + + let mut session = make_session(vec![tool_call, text_response("OK")]).await; + let mut rx = session.subscribe(); + + session.process_input("Do something").await.unwrap(); + + let reasoning = collect_message_reasoning(&mut rx); + assert_eq!(reasoning, vec![ + Some(ReasoningOutput { + summary: Some("call the tool".to_string()), + trace: None, + }), + None, + ]); + } + + #[tokio::test(start_paused = true)] + async fn only_the_final_response_contributes_reasoning_after_a_retry() { + let provider = Arc::new(ScriptedStreamProvider::new(vec![ + ScriptedStreamCall::Events(vec![ + Ok(StreamEvent::ReasoningDelta { + delta: "discarded thinking".to_string(), + }), + Err(LlmError::Stream { + message: "connection reset".into(), + source: None, + }), + ]), + ScriptedStreamCall::Response(Box::new(reasoning_response( + "Recovered", + "final summary", + "final trace", + ))), + ])); + let mut session = make_session_with_provider(provider.clone()).await; + let mut rx = session.subscribe(); + + session.process_input("Hello").await.unwrap(); + + assert_eq!(provider.call_index.load(Ordering::SeqCst), 2); + let reasoning = collect_message_reasoning(&mut rx); + assert_eq!(reasoning, vec![Some(ReasoningOutput { + summary: Some("final summary".to_string()), + trace: Some("final trace".to_string()), + })]); + } + #[tokio::test(start_paused = true)] async fn stream_quota_error_does_not_replay() { let quota_error = LlmError::Provider { diff --git a/lib/components/fabro-agent/src/types.rs b/lib/components/fabro-agent/src/types.rs index a39600549..fc3581c57 100644 --- a/lib/components/fabro-agent/src/types.rs +++ b/lib/components/fabro-agent/src/types.rs @@ -4,7 +4,7 @@ use chrono::{DateTime, Utc}; use fabro_llm::Error as LlmError; use fabro_llm::types::{ContentPart, ThinkingData, TokenCounts, ToolCall, ToolResult}; use fabro_model::{CostSource, ModelRef}; -use fabro_types::{SessionMessage, StageContextWindowProjection}; +use fabro_types::{ReasoningOutput, SessionMessage, StageContextWindowProjection}; use serde::de::DeserializeOwned; use serde::{Deserialize, Serialize}; @@ -254,6 +254,11 @@ pub enum AgentEvent { tool_call_count: usize, #[serde(default, skip_serializing_if = "Option::is_none")] context_window: Option, + /// Readable reasoning normalized from the final response. Derived + /// once the response is complete, so retried or replaced streaming + /// buffers never become durable reasoning. + #[serde(default, skip_serializing_if = "Option::is_none")] + reasoning: Option, }, TextDelta { delta: String, @@ -892,6 +897,7 @@ mod tests { cost_source: Some(CostSource::Authoritative), tool_call_count: 2, context_window: None, + reasoning: None, }; match &event { AgentEvent::AssistantMessage { diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/response.rs b/lib/components/fabro-llm/src/codec/openai_compatible/response.rs index 8b853acb9..6f5e6d40e 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/response.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/response.rs @@ -1,7 +1,7 @@ //! Response decoding: Chat Completions body → canonical `Response`. use super::translate::{self, map_finish_reason}; -use super::wire::{ApiResponse, ApiUsage}; +use super::wire::{ApiResponse, ApiUsage, ReasoningDetails}; use crate::codec::CodecCtx; use crate::error::{Error, ProviderErrorDetail, ProviderErrorKind}; use crate::types::{ @@ -25,6 +25,11 @@ pub(super) fn decode_response( })?; let mut content_parts = Vec::new(); + if let Some(payload) = &choice.message.reasoning_details { + let mut details = ReasoningDetails::default(); + details.push_payload(payload); + content_parts.extend(details.into_content_part()); + } if let Some(reasoning) = choice.message.reasoning() { if !reasoning.is_empty() { content_parts.push(ContentPart::Thinking(ThinkingData { diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs b/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs index e1c811ee4..39a35e9e7 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs @@ -4,7 +4,7 @@ //! already-stripped payloads (including the `[DONE]` sentinel) via `on_event`. use super::translate::{map_finish_reason, parse_tool_arguments}; -use super::wire::{AccumulatedToolCall, StreamChunk}; +use super::wire::{AccumulatedToolCall, ReasoningDetails, StreamChunk}; use crate::codec::{CodecCtx, RawEvent, StreamDecoder}; use crate::error::Error; use crate::types::{ @@ -20,6 +20,7 @@ pub(super) struct StreamState { response_model: String, accumulated_text: String, accumulated_reasoning: String, + reasoning_details: ReasoningDetails, tool_calls: Vec, usage: TokenCounts, finish_reason: FinishReason, @@ -42,6 +43,7 @@ impl StreamState { response_model: String::new(), accumulated_text: String::new(), accumulated_reasoning: String::new(), + reasoning_details: ReasoningDetails::default(), tool_calls: Vec::new(), usage: TokenCounts::default(), finish_reason: FinishReason::Stop, @@ -93,6 +95,11 @@ impl StreamState { } } + // Accumulate structured reasoning detail fragments in wire order. + if let Some(payload) = &delta.reasoning_details { + self.reasoning_details.push_payload(payload); + } + // Handle text content delta. if let Some(content) = &delta.content { if !content.is_empty() { @@ -170,6 +177,9 @@ impl StreamState { let mut content_parts = Vec::new(); + // Preserve the structured reasoning channel verbatim. + content_parts.extend(std::mem::take(&mut self.reasoning_details).into_content_part()); + // Include reasoning/thinking content if present (Kimi, etc.). if !self.accumulated_reasoning.is_empty() { content_parts.push(ContentPart::Thinking(ThinkingData { diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs index 5bc0a3b93..01e88d500 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs @@ -2,7 +2,7 @@ use crate::codec::cache::CacheControl; use crate::codec::split_inclusive_token_total; -use crate::types::{ReasoningEffort, TokenCounts}; +use crate::types::{ContentPart, ReasoningEffort, TokenCounts}; #[derive(serde::Serialize)] pub(super) struct ApiRequest { @@ -129,6 +129,11 @@ pub(super) struct ApiChoiceMessage { pub reasoning_content: Option, /// OpenRouter's normalized spelling for reasoning text. pub reasoning: Option, + /// Structured reasoning channel (OpenRouter and compatible + /// aggregators). Kept as an untyped value so unknown detail variants + /// cannot fail an otherwise valid completion. + #[serde(default)] + pub reasoning_details: Option, pub tool_calls: Option>, } @@ -140,6 +145,88 @@ impl ApiChoiceMessage { } } +/// Structured `reasoning_details` entries accumulated in wire order. +/// +/// The entries are preserved verbatim as an opaque content part so +/// encrypted material survives for future provider-aware replay; only known +/// readable members are ever normalized out of them. +#[derive(Default)] +pub(super) struct ReasoningDetails { + entries: Vec, +} + +impl ReasoningDetails { + /// Absorb one `reasoning_details` payload. + /// + /// Providers document an array of detail objects; a lone object is + /// accepted as a single entry. Scalars carry nothing replayable and are + /// dropped. Streaming deltas repeat the same logical detail across + /// chunks, so a fragment that continues the previous entry is coalesced + /// into it rather than becoming a separate block. + pub(super) fn push_payload(&mut self, payload: &serde_json::Value) { + let incoming = match payload { + serde_json::Value::Array(entries) => entries.clone(), + serde_json::Value::Object(_) => vec![payload.clone()], + _ => Vec::new(), + }; + for entry in incoming { + if !entry.is_object() { + continue; + } + match self.entries.last_mut() { + Some(last) if continues_detail(last, &entry) => merge_detail_fragment(last, &entry), + _ => self.entries.push(entry), + } + } + } + + /// Opaque content part holding the accumulated entries, or `None` when + /// nothing usable arrived. + pub(super) fn into_content_part(self) -> Option { + (!self.entries.is_empty()).then(|| ContentPart::Other { + kind: ContentPart::OPENAI_COMPAT_REASONING_DETAILS.to_string(), + data: serde_json::Value::Array(self.entries), + }) + } +} + +/// Text-bearing members whose fragments concatenate across stream chunks. +const DETAIL_TEXT_MEMBERS: [&str; 3] = ["text", "summary", "data"]; + +/// Whether `entry` continues the logical detail already in `last`. +/// +/// Aggregators tag each logical detail with a stable `type` and `index`; +/// fragment streams that omit `index` are matched on `type` alone. +fn continues_detail(last: &serde_json::Value, entry: &serde_json::Value) -> bool { + last.get("type") == entry.get("type") && last.get("index") == entry.get("index") +} + +/// Append `entry`'s text fragments onto `last` and fill in members `last` +/// has not seen yet. +fn merge_detail_fragment(last: &mut serde_json::Value, entry: &serde_json::Value) { + let Some(entry_members) = entry.as_object() else { + return; + }; + let Some(last_members) = last.as_object_mut() else { + return; + }; + for (key, value) in entry_members { + match last_members.get_mut(key) { + Some(serde_json::Value::String(existing)) + if DETAIL_TEXT_MEMBERS.contains(&key.as_str()) => + { + if let Some(fragment) = value.as_str() { + existing.push_str(fragment); + } + } + Some(_) => {} + None => { + last_members.insert(key.clone(), value.clone()); + } + } + } +} + #[derive(serde::Deserialize)] pub(super) struct ApiToolCall { pub id: String, @@ -245,6 +332,10 @@ pub(super) struct StreamDelta { pub reasoning_content: Option, /// OpenRouter's normalized spelling for reasoning text. pub reasoning: Option, + /// Structured reasoning channel, streamed as fragments of the entries + /// the non-streaming response returns whole. + #[serde(default)] + pub reasoning_details: Option, pub tool_calls: Option>, } diff --git a/lib/components/fabro-llm/src/lib.rs b/lib/components/fabro-llm/src/lib.rs index 2ad5d2d3f..0b33bfece 100644 --- a/lib/components/fabro-llm/src/lib.rs +++ b/lib/components/fabro-llm/src/lib.rs @@ -9,6 +9,7 @@ pub mod middleware; pub mod model_test; pub mod provider; pub mod providers; +mod reasoning; pub mod retry; pub mod token_count; pub mod tools; diff --git a/lib/components/fabro-llm/src/providers/fabro_server.rs b/lib/components/fabro-llm/src/providers/fabro_server.rs index 7a5819199..bd5c42ea6 100644 --- a/lib/components/fabro-llm/src/providers/fabro_server.rs +++ b/lib/components/fabro-llm/src/providers/fabro_server.rs @@ -347,6 +347,53 @@ data: {\"type\":\"text_delta\",\"delta\":\" world\",\"text_id\":null}\n\ assert_eq!(response.cost_source, Some(CostSource::Estimated)); } + /// Reasoning needs no dedicated wire field on this hop: the canonical + /// message already transports the provider parts it is derived from. + #[tokio::test] + async fn complete_normalizes_reasoning_from_the_transported_message() { + let server = MockServer::start(); + + server.mock(|when, then| { + when.method(POST).path("/completions"); + then.status(200) + .header("content-type", "application/json") + .json_body(serde_json::json!({ + "id": "resp-123", + "model": "test-model", + "message": { + "role": "assistant", + "content": [ + { + "kind": "openai_compat_reasoning_details", + "data": [ + {"type": "reasoning.summary", "summary": "weighed both"}, + {"type": "reasoning.text", "text": "step one"}, + ] + }, + {"kind": "text", "data": "Hello there!"}, + ], + "name": null, + "tool_call_id": null + }, + "stop_reason": "end_turn", + "usage": {"input_tokens": 10, "output_tokens": 5} + })); + }); + + let adapter = Adapter::new( + fabro_test::test_http_client(), + server.base_url(), + "test-provider", + ); + + let response = adapter.complete(&make_request()).await.unwrap(); + + assert_eq!(response.text(), "Hello there!"); + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("weighed both")); + assert_eq!(reasoning.trace.as_deref(), Some("step one")); + } + #[tokio::test] async fn complete_returns_error_on_502() { let server = MockServer::start(); diff --git a/lib/components/fabro-llm/src/reasoning.rs b/lib/components/fabro-llm/src/reasoning.rs new file mode 100644 index 000000000..a4909ec1b --- /dev/null +++ b/lib/components/fabro-llm/src/reasoning.rs @@ -0,0 +1,337 @@ +//! Normalization of provider reasoning material into [`ReasoningOutput`]. +//! +//! Every provider that returns readable reasoning does it differently, and +//! several return more than one channel at once. This module reduces the +//! final response's content parts to the two normalized fields without +//! reaching into opaque material (signatures, item IDs, encrypted payloads) +//! and without failing a completion it cannot classify. +//! +//! Parsing is deliberately tolerant: provider payloads are read as +//! `serde_json::Value` with optional lookups, so unknown detail variants, +//! missing members, extra members, and unexpected member types are ignored +//! rather than surfaced as errors. + +use fabro_types::{ContentPart, ReasoningOutput}; + +/// Separator between distinct complete reasoning blocks. Fragments of one +/// logical block are coalesced by the streaming decoders before they reach +/// this module. +const BLOCK_SEPARATOR: &str = "\n\n"; + +/// Detail types whose payload is opaque and must never reach the public +/// normalized fields. +fn is_opaque_detail_type(detail_type: &str) -> bool { + detail_type.contains("encrypted") || detail_type.contains("redacted") +} + +/// Readable blocks collected per normalized field. +/// +/// Explicit blocks come from a channel with documented reasoning semantics. +/// Fallback blocks come from flattened provider strings, which aggregators +/// commonly duplicate alongside a structured channel; they only fill a +/// summary that no explicit block produced. +#[derive(Default)] +struct Blocks { + explicit_summary: Vec, + explicit_trace: Vec, + fallback_summary: Vec, +} + +impl Blocks { + fn into_output(self) -> Option { + let summary = + join_blocks(&self.explicit_summary).or_else(|| join_blocks(&self.fallback_summary)); + let trace = join_blocks(&self.explicit_trace); + let output = ReasoningOutput { summary, trace }; + (!output.is_empty()).then_some(output) + } +} + +/// Join complete blocks in provider order, dropping empty and +/// whitespace-only fragments. Retained text is never trimmed or rewritten. +fn join_blocks(blocks: &[String]) -> Option { + let joined = blocks + .iter() + .filter(|block| !block.trim().is_empty()) + .map(String::as_str) + .collect::>() + .join(BLOCK_SEPARATOR); + (!joined.is_empty()).then_some(joined) +} + +/// Read a text-bearing member, preferring the member the provider's +/// documented semantics name and accepting the other readable spelling. +fn readable_member(entry: &serde_json::Value, preferred: &str) -> Option { + let other = if preferred == "text" { + "summary" + } else { + "text" + }; + entry + .get(preferred) + .and_then(serde_json::Value::as_str) + .or_else(|| entry.get(other).and_then(serde_json::Value::as_str)) + .map(str::to_string) +} + +/// Extract readable text from an OpenAI Responses `reasoning` output item. +/// +/// `summary[].text` is the model-authored summary; `content[]` entries typed +/// `reasoning_text` are the verbatim trace. `encrypted_content`, `id`, and +/// `status` are opaque and ignored. +fn collect_openai_reasoning_item(item: &serde_json::Value, blocks: &mut Blocks) { + if let Some(entries) = item.get("summary").and_then(serde_json::Value::as_array) { + for entry in entries { + if let Some(text) = entry.as_str() { + blocks.explicit_summary.push(text.to_string()); + } else if let Some(text) = entry.get("text").and_then(serde_json::Value::as_str) { + blocks.explicit_summary.push(text.to_string()); + } + } + } + if let Some(entries) = item.get("content").and_then(serde_json::Value::as_array) { + for entry in entries { + let Some(text) = entry.get("text").and_then(serde_json::Value::as_str) else { + continue; + }; + let entry_type = entry + .get("type") + .and_then(serde_json::Value::as_str) + .unwrap_or_default(); + if entry_type == "reasoning_text" { + blocks.explicit_trace.push(text.to_string()); + } else if !is_opaque_detail_type(entry_type) { + // Readable text from the reasoning channel that carries no + // recognized classification is treated as a summary. + blocks.explicit_summary.push(text.to_string()); + } + } + } +} + +/// Extract readable text from OpenAI-compatible `reasoning_details` entries. +fn collect_reasoning_details(details: &serde_json::Value, blocks: &mut Blocks) { + let Some(entries) = details.as_array() else { + return; + }; + for entry in entries { + let detail_type = entry + .get("type") + .and_then(serde_json::Value::as_str) + .unwrap_or_default(); + if is_opaque_detail_type(detail_type) { + continue; + } + match detail_type { + "reasoning.text" => { + if let Some(text) = readable_member(entry, "text") { + blocks.explicit_trace.push(text); + } + } + // Documented summaries and unknown detail variants both come + // from an established reasoning channel, so a readable member + // is kept as a summary either way. + _ => { + if let Some(text) = readable_member(entry, "summary") { + blocks.explicit_summary.push(text); + } + } + } + } +} + +/// Normalize the content parts of a final response into readable reasoning. +/// +/// Returns `None` when the response carries no readable reasoning, so an +/// event without reasoning keeps its previous serialized shape. +pub(crate) fn normalize(content: &[ContentPart]) -> Option { + let mut blocks = Blocks::default(); + for part in content { + match part { + ContentPart::Thinking(thinking) if !thinking.redacted => { + blocks.fallback_summary.push(thinking.text.clone()); + } + ContentPart::Other { kind, data } if kind == ContentPart::OPENAI_REASONING => { + collect_openai_reasoning_item(data, &mut blocks); + } + ContentPart::Other { kind, data } + if kind == ContentPart::OPENAI_COMPAT_REASONING_DETAILS => + { + collect_reasoning_details(data, &mut blocks); + } + _ => {} + } + } + blocks.into_output() +} + +#[cfg(test)] +mod tests { + use fabro_types::ThinkingData; + use serde_json::json; + + use super::*; + + fn thinking(text: &str) -> ContentPart { + ContentPart::Thinking(ThinkingData { + text: text.to_string(), + signature: None, + redacted: false, + }) + } + + fn openai_reasoning(item: serde_json::Value) -> ContentPart { + ContentPart::Other { + kind: ContentPart::OPENAI_REASONING.to_string(), + data: item, + } + } + + fn reasoning_details(details: serde_json::Value) -> ContentPart { + ContentPart::Other { + kind: ContentPart::OPENAI_COMPAT_REASONING_DETAILS.to_string(), + data: details, + } + } + + #[test] + fn non_redacted_thinking_becomes_a_summary() { + let output = normalize(&[thinking("weighing the options")]).unwrap(); + assert_eq!(output.summary.as_deref(), Some("weighing the options")); + assert!(output.trace.is_none()); + } + + #[test] + fn redacted_thinking_yields_no_readable_reasoning() { + let redacted = ContentPart::Thinking(ThinkingData { + text: "AAAAopaque".to_string(), + signature: Some("sig".to_string()), + redacted: true, + }); + assert!(normalize(&[redacted]).is_none()); + } + + #[test] + fn responses_item_with_summary_and_reasoning_text_produces_both_fields() { + let output = normalize(&[openai_reasoning(json!({ + "type": "reasoning", + "id": "rs_1", + "encrypted_content": "gAAAAA", + "summary": [{"type": "summary_text", "text": "inspect first"}], + "content": [{"type": "reasoning_text", "text": "step one"}], + }))]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("inspect first")); + assert_eq!(output.trace.as_deref(), Some("step one")); + } + + #[test] + fn responses_blocks_join_in_provider_order() { + let output = normalize(&[openai_reasoning(json!({ + "summary": [ + {"type": "summary_text", "text": "first"}, + {"type": "summary_text", "text": "second"}, + ], + }))]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("first\n\nsecond")); + } + + #[test] + fn structured_details_produce_summary_and_trace() { + let output = normalize(&[reasoning_details(json!([ + {"type": "reasoning.summary", "summary": "checked the parser"}, + {"type": "reasoning.text", "text": "read convert.rs", "signature": "sig"}, + ]))]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("checked the parser")); + assert_eq!(output.trace.as_deref(), Some("read convert.rs")); + } + + #[test] + fn encrypted_details_are_excluded() { + let output = normalize(&[reasoning_details(json!([ + {"type": "reasoning.encrypted", "data": "gAAAAAsecret", "format": "openai-responses-v1"}, + {"type": "reasoning.summary", "summary": "visible"}, + ])),]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("visible")); + assert!(output.trace.is_none()); + } + + #[test] + fn encrypted_only_details_produce_no_reasoning() { + assert!( + normalize(&[reasoning_details(json!([ + {"type": "reasoning.encrypted", "data": "gAAAAAsecret"}, + ]))]) + .is_none() + ); + } + + #[test] + fn unknown_detail_variants_are_treated_as_summary() { + let output = normalize(&[reasoning_details(json!([ + {"type": "reasoning.future", "text": "new channel"}, + ]))]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("new channel")); + } + + #[test] + fn malformed_details_are_ignored_without_failing() { + assert!(normalize(&[reasoning_details(json!("not-an-array"))]).is_none()); + assert!( + normalize(&[reasoning_details(json!([ + 42, + {"type": "reasoning.summary", "summary": 7}, + {"no_type": true}, + ]))]) + .is_none() + ); + } + + #[test] + fn structured_details_suppress_a_duplicate_flattened_value() { + let output = normalize(&[ + reasoning_details(json!([ + {"type": "reasoning.summary", "summary": "checked the parser"}, + ])), + thinking("checked the parser"), + ]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("checked the parser")); + assert!(output.trace.is_none()); + } + + #[test] + fn trace_only_detail_lets_the_fallback_fill_the_summary() { + let output = normalize(&[ + reasoning_details(json!([{"type": "reasoning.text", "text": "verbatim"}])), + thinking("flattened"), + ]) + .unwrap(); + assert_eq!(output.summary.as_deref(), Some("flattened")); + assert_eq!(output.trace.as_deref(), Some("verbatim")); + } + + #[test] + fn whitespace_only_fragments_do_not_create_reasoning() { + assert!(normalize(&[thinking(" \n ")]).is_none()); + } + + #[test] + fn non_empty_text_is_preserved_verbatim() { + let output = normalize(&[thinking(" indented thought\n")]).unwrap(); + assert_eq!(output.summary.as_deref(), Some(" indented thought\n")); + } + + #[test] + fn unrelated_content_parts_are_ignored() { + let parts = vec![ContentPart::text("answer"), ContentPart::Other { + kind: ContentPart::OPENAI_MESSAGE.to_string(), + data: json!({"type": "message", "content": [{"text": "answer"}]}), + }]; + assert!(normalize(&parts).is_none()); + } +} diff --git a/lib/components/fabro-llm/src/types.rs b/lib/components/fabro-llm/src/types.rs index bd5212f1b..9c41263fd 100644 --- a/lib/components/fabro-llm/src/types.rs +++ b/lib/components/fabro-llm/src/types.rs @@ -10,13 +10,14 @@ use std::sync::Arc; // model. They are re-exported here so existing `fabro_llm::types::*` // imports keep working. pub use fabro_types::{ - AudioData, ContentPart, DocumentData, ImageData, Message, Role, ThinkingData, ToolCall, - ToolResult, + AudioData, ContentPart, DocumentData, ImageData, Message, ReasoningOutput, Role, ThinkingData, + ToolCall, ToolResult, }; use fabro_util::backoff::BackoffPolicy; use serde::{Deserialize, Serialize}; use crate::error::Error; +use crate::reasoning; // --- 3.8 FinishReason --- @@ -281,6 +282,19 @@ impl Response { Some(reasoning) } } + + /// Readable reasoning normalized from this response's canonical message + /// content, or `None` when the provider returned none. + /// + /// The message is the single source of truth: opaque provider reasoning + /// items are already preserved there, so normalization needs no second + /// stored field and cannot drift from what will be replayed. Deriving it + /// from the final response also keeps retried or replaced streaming + /// buffers out of the durable result. + #[must_use] + pub fn reasoning_output(&self) -> Option { + reasoning::normalize(&self.message.content) + } } // --- 3.13 StreamEvent --- diff --git a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs index c8facd920..4ba98dfa2 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs @@ -437,6 +437,39 @@ async fn decode_response(body: serde_json::Value) -> fabro_llm::types::Response response } +/// Wraps an assistant message in a complete Chat Completions body. +fn body_with_message(message: &serde_json::Value) -> serde_json::Value { + serde_json::json!({ + "id": "chatcmpl_test", + "object": "chat.completion", + "created": CREATED_TS, + "model": MODEL, + "choices": [{"index": 0, "message": message, "finish_reason": "stop"}], + "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15} + }) +} + +/// Streams an SSE transcript and returns the final accumulated response. +async fn stream_final_response(sse_body: &str) -> fabro_llm::types::Response { + use futures::StreamExt; + + let server = MockServer::start(); + let (mock, _slot) = mount_capture_sse(&server, "/chat/completions", sse_body); + let adapter = adapter(&server); + let mut stream = adapter + .stream(&base_request(MODEL)) + .await + .expect("stream should start"); + let mut final_response = None; + while let Some(item) = stream.next().await { + if let Ok(fabro_llm::types::StreamEvent::Finish { response, .. }) = item { + final_response = Some(*response); + } + } + mock.assert(); + final_response.expect("stream should emit a finish event") +} + #[tokio::test] async fn decode_tool_calls_with_string_arguments() { let response = decode_response(serde_json::json!({ @@ -485,6 +518,167 @@ async fn decode_reasoning_content_as_thinking() { fabro_test::fabro_json_snapshot!(response); } +// --------------------------------------------------------------------------- +// Structured reasoning details +// --------------------------------------------------------------------------- + +/// The structured channel classifies summary and trace independently. +#[tokio::test] +async fn decode_reasoning_details_normalize_summary_and_trace() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning_details": [ + {"type": "reasoning.summary", "summary": "the user wants 2+2", "index": 0}, + {"type": "reasoning.text", "text": "2 plus 2 is 4", "index": 1}, + ] + }))) + .await; + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("the user wants 2+2")); + assert_eq!(reasoning.trace.as_deref(), Some("2 plus 2 is 4")); +} + +/// Encrypted entries stay in the opaque provider part for future replay but +/// never reach the normalized output. +#[tokio::test] +async fn decode_reasoning_details_preserve_encrypted_entries_opaquely() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning_details": [ + {"type": "reasoning.encrypted", "data": "gAAAAAopaque", "index": 0}, + {"type": "reasoning.summary", "summary": "visible", "index": 1}, + ] + }))) + .await; + + let opaque = response + .message + .content + .iter() + .find_map(|part| match part { + fabro_llm::types::ContentPart::Other { kind, data } + if kind == fabro_llm::types::ContentPart::OPENAI_COMPAT_REASONING_DETAILS => + { + Some(data) + } + _ => None, + }) + .expect("opaque reasoning details preserved"); + assert_eq!(opaque[0]["data"], "gAAAAAopaque"); + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("visible")); + assert!(reasoning.trace.is_none()); +} + +/// Unknown and malformed detail entries must not fail an otherwise valid +/// completion. +#[tokio::test] +async fn decode_tolerates_unknown_and_malformed_reasoning_details() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning_details": [ + {"type": "reasoning.future", "text": "new channel", "extra": {"nested": true}}, + {"type": "reasoning.summary", "summary": 7}, + "not-an-object", + 42, + ] + }))) + .await; + + assert_eq!(response.text(), "4."); + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("new channel")); +} + +/// A scalar `reasoning_details` carries nothing replayable and is dropped +/// without disturbing the rest of the response. +#[tokio::test] +async fn decode_ignores_scalar_reasoning_details() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning_details": "unexpected" + }))) + .await; + + assert_eq!(response.text(), "4."); + assert!(response.reasoning_output().is_none()); +} + +/// OpenRouter returns both the structured channel and a flattened copy of +/// the same material; the summary must not appear twice. +#[tokio::test] +async fn decode_structured_details_suppress_the_duplicate_flattened_value() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning": "the user wants 2+2", + "reasoning_details": [ + {"type": "reasoning.summary", "summary": "the user wants 2+2", "index": 0}, + ] + }))) + .await; + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("the user wants 2+2")); + assert!(reasoning.trace.is_none()); +} + +/// A structured trace with no structured summary leaves room for the +/// flattened value to fill the summary. +#[tokio::test] +async fn decode_trace_only_details_let_the_flattened_value_fill_the_summary() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning": "flattened summary", + "reasoning_details": [{"type": "reasoning.text", "text": "verbatim trace", "index": 0}] + }))) + .await; + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("flattened summary")); + assert_eq!(reasoning.trace.as_deref(), Some("verbatim trace")); +} + +/// Streamed detail fragments coalesce back into the same normalized output +/// the non-streaming body produces. +#[tokio::test] +async fn stream_reasoning_details_normalize_like_the_non_streaming_body() { + let sse = support::sse_data_transcript(&[ + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"role":"assistant","reasoning_details":[{"type":"reasoning.summary","summary":"the user ","index":0}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"reasoning.summary","summary":"wants 2+2","index":0}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"reasoning.text","text":"2 plus 2 is 4","index":1}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"content":"4."},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}"#, + "[DONE]", + ]); + let streamed = stream_final_response(&sse).await; + + let non_streamed = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning_details": [ + {"type": "reasoning.summary", "summary": "the user wants 2+2", "index": 0}, + {"type": "reasoning.text", "text": "2 plus 2 is 4", "index": 1}, + ] + }))) + .await; + + assert_eq!(streamed.reasoning_output(), non_streamed.reasoning_output()); + assert_eq!( + streamed + .reasoning_output() + .and_then(|reasoning| reasoning.summary), + Some("the user wants 2+2".to_string()) + ); +} + /// Cached and reasoning detail tokens are split into their own disjoint /// buckets and subtracted out of input/output. #[tokio::test] diff --git a/lib/components/fabro-llm/tests/it/wire/openai_responses.rs b/lib/components/fabro-llm/tests/it/wire/openai_responses.rs index 4f9d8e077..fd281ce62 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_responses.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_responses.rs @@ -461,6 +461,42 @@ async fn decode_reasoning_and_function_call_items() { fabro_test::fabro_json_snapshot!(response); } +/// A reasoning item carrying both readable channels normalizes into both +/// fields while its encrypted payload stays opaque. +#[tokio::test] +async fn decode_reasoning_item_normalizes_summary_and_trace() { + let response = decode_response(serde_json::json!({ + "id": "resp_test", + "object": "response", + "model": MODEL, + "status": "completed", + "output": [ + { + "type": "reasoning", + "id": "rs_1", + "encrypted_content": "gAAAAAopaque", + "summary": [{"type": "summary_text", "text": "Adding two numbers."}], + "content": [{"type": "reasoning_text", "text": "2 plus 2 is 4."}] + }, + { + "type": "message", + "role": "assistant", + "id": "msg_out", + "content": [{"type": "output_text", "text": "4."}] + } + ], + "usage": {"input_tokens": 30, "output_tokens": 12} + })) + .await; + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary.as_deref(), Some("Adding two numbers.")); + assert_eq!(reasoning.trace.as_deref(), Some("2 plus 2 is 4.")); + + let normalized = serde_json::to_string(&reasoning).unwrap(); + assert!(!normalized.contains("gAAAAAopaque")); +} + #[tokio::test] async fn decode_incomplete_status_maps_to_length() { let response = decode_response(serde_json::json!({ diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index 60f3eb877..5d1f83861 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -4032,6 +4032,7 @@ mod tests { visit: 1, message: None, context_window: None, + reasoning: None, } } diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 2cbc34d07..4beb3b888 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -612,6 +612,7 @@ fn event_body_from_event(event: &Event) -> EventBody { cost_source, tool_call_count, context_window, + reasoning, } => { let billing = billed_token_counts_from_llm(usage) .with_reported_cost(cost_usd.map(UsdMicros::from_usd)); @@ -624,6 +625,7 @@ fn event_body_from_event(event: &Event) -> EventBody { visit: *visit, message: None, context_window: context_window.clone(), + reasoning: reasoning.clone(), }) } AgentEvent::ToolCallStarted { @@ -2230,6 +2232,7 @@ mod tests { cost_source: None, tool_call_count: 0, context_window: None, + reasoning: None, }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2264,6 +2267,7 @@ mod tests { cost_source: None, tool_call_count: 0, context_window: None, + reasoning: None, }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2301,6 +2305,7 @@ mod tests { cost_source: Some(fabro_model::CostSource::Authoritative), tool_call_count: 0, context_window: None, + reasoning: None, }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2351,6 +2356,7 @@ mod tests { cost_source: None, tool_call_count: 0, context_window: Some(context_window), + reasoning: None, }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2368,6 +2374,58 @@ mod tests { ); } + #[test] + fn agent_assistant_message_copies_reasoning_into_canonical_event() { + let stored = to_run_event(&fixtures::RUN_1, &Event::Agent { + stage: "code".to_string(), + visit: 1, + event: AgentEvent::AssistantMessage { + text: String::new(), + model: ModelRef { + provider: ProviderId::openai(), + model_id: "gpt-5.4".into(), + speed: None, + }, + usage: LlmTokenCounts::default(), + cost_usd: None, + cost_source: None, + tool_call_count: 1, + context_window: None, + reasoning: Some(::fabro_types::ReasoningOutput { + summary: Some("inspect the conversion first".to_string()), + trace: Some("read convert.rs, then the sink".to_string()), + }), + }, + session_id: Some("ses_agent".to_string()), + parent_session_id: None, + tool_call_id: None, + }); + + let EventBody::AgentMessage(message) = &stored.body else { + panic!("expected agent message body"); + }; + let reasoning = message.reasoning.as_ref().expect("reasoning copied"); + assert_eq!( + reasoning.summary.as_deref(), + Some("inspect the conversion first") + ); + assert_eq!( + reasoning.trace.as_deref(), + Some("read convert.rs, then the sink") + ); + + let value = stored.to_value().unwrap(); + assert_eq!(value["event"], "agent.message"); + assert_eq!( + value["properties"]["reasoning"]["summary"], + "inspect the conversion first" + ); + assert_eq!( + value["properties"]["reasoning"]["trace"], + "read convert.rs, then the sink" + ); + } + #[test] fn agent_acp_events_map_to_event_bodies_with_stage_scope() { let scope = StageScope { diff --git a/lib/components/fabro-workflow/src/event/redaction.rs b/lib/components/fabro-workflow/src/event/redaction.rs index 5ae44d111..d0d7d0974 100644 --- a/lib/components/fabro-workflow/src/event/redaction.rs +++ b/lib/components/fabro-workflow/src/event/redaction.rs @@ -30,7 +30,10 @@ pub fn event_payload_from_redacted_json(line: &str, run_id: &RunId) -> Result(); +} + +#[test] +fn reasoning_output_matches_openapi_json_shape() { + let value = json!({ + "summary": "inspect the conversion first", + "trace": "read convert.rs, then the sink", + }); + + let output: ReasoningOutput = serde_json::from_value(value.clone()).unwrap(); + assert_eq!( + output.summary.as_deref(), + Some("inspect the conversion first") + ); + assert_eq!( + output.trace.as_deref(), + Some("read convert.rs, then the sink") + ); + assert_eq!(serde_json::to_value(&output).unwrap(), value); + + let api_output: ApiReasoningOutput = serde_json::from_value(value).unwrap(); + assert_eq!(api_output, output); +} + +#[test] +fn reasoning_output_members_are_individually_optional() { + let summary_only: ReasoningOutput = + serde_json::from_value(json!({"summary": "only a summary"})).unwrap(); + assert!(summary_only.trace.is_none()); + assert_eq!( + serde_json::to_value(&summary_only).unwrap(), + json!({"summary": "only a summary"}) + ); + + let trace_only: ReasoningOutput = + serde_json::from_value(json!({"trace": "only a trace"})).unwrap(); + assert!(trace_only.summary.is_none()); + assert_eq!( + serde_json::to_value(&trace_only).unwrap(), + json!({"trace": "only a trace"}) + ); +} + +#[test] +fn agent_message_props_reasoning_is_optional_on_the_wire() { + let without = json!({ + "text": "ok", + "model": {"provider": "openai", "model_id": "gpt-5.4"}, + "billing": {"input_tokens": 1, "output_tokens": 1, "total_tokens": 2}, + "tool_call_count": 0, + "visit": 1, + }); + let props: ApiAgentMessageProps = serde_json::from_value(without.clone()).unwrap(); + assert!(props.reasoning.is_none()); + assert!( + !serde_json::to_value(&props) + .unwrap() + .as_object() + .unwrap() + .contains_key("reasoning") + ); + + let mut with = without; + with["reasoning"] = json!({"summary": "checked the parser", "trace": "step one"}); + let props: ApiAgentMessageProps = serde_json::from_value(with).unwrap(); + let reasoning = props.reasoning.as_ref().unwrap(); + assert_eq!(reasoning.summary.as_deref(), Some("checked the parser")); + assert_eq!(reasoning.trace.as_deref(), Some("step one")); + assert_eq!( + serde_json::to_value(&props).unwrap()["reasoning"], + json!({"summary": "checked the parser", "trace": "step one"}) + ); +} + +fn assert_same_type() { + assert_eq!( + TypeId::of::(), + TypeId::of::(), + "{} should be the same type as {}", + type_name::(), + type_name::() + ); +} diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index f8418b955..1b19c9248 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -22,6 +22,7 @@ pub mod pair; pub mod parallel; pub mod principal; pub mod pull_request; +pub mod reasoning; pub mod repository; pub mod run; pub mod run_blob_id; @@ -101,6 +102,7 @@ pub use pull_request::{ PullRequestDetailsUnavailableReason, PullRequestGithubDetail, PullRequestLink, PullRequestMeta, PullRequestRef, PullRequestResponse, PullRequestTimestamps, PullRequestUser, }; +pub use reasoning::ReasoningOutput; pub use repository::{RepositoryProvider, RepositoryRef}; pub use run::{ DirtyStatus, ForkSourceRef, GitContext, PreRunPushOutcome, RunClientProvenance, RunProvenance, diff --git a/lib/foundation/fabro-types/src/reasoning.rs b/lib/foundation/fabro-types/src/reasoning.rs new file mode 100644 index 000000000..0363c6deb --- /dev/null +++ b/lib/foundation/fabro-types/src/reasoning.rs @@ -0,0 +1,96 @@ +use serde::{Deserialize, Serialize}; + +/// Readable model reasoning normalized into a provider-neutral shape. +/// +/// Providers expose reasoning through several unrelated channels: OpenAI +/// Responses reasoning items, OpenAI-compatible `reasoning_details`, and +/// flattened `reasoning`/`reasoning_content`/`thinking` strings. This type +/// reduces all of them to the two capabilities consumers actually care +/// about, so the durable event contract does not change shape when a +/// provider dialect does. +/// +/// Both fields may be populated for the same response. An emitted object +/// always carries at least one of them; opaque provider material +/// (signatures, IDs, encrypted or redacted payloads) never appears here. +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct ReasoningOutput { + /// Model-authored summary of its reasoning, safe to show to users. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub summary: Option, + /// Verbatim readable reasoning text, when the provider returns it in + /// addition to (or instead of) a summary. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub trace: Option, +} + +impl ReasoningOutput { + /// Returns `true` when neither readable field is present, meaning the + /// object carries nothing worth emitting. + #[must_use] + pub fn is_empty(&self) -> bool { + self.summary.is_none() && self.trace.is_none() + } +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::*; + + #[test] + fn summary_only_round_trips_without_trace_member() { + let output = ReasoningOutput { + summary: Some("checked the parser first".to_string()), + trace: None, + }; + let v = serde_json::to_value(&output).unwrap(); + assert_eq!(v, json!({"summary": "checked the parser first"})); + assert_eq!( + serde_json::from_value::(v).unwrap(), + output + ); + } + + #[test] + fn trace_only_round_trips_without_summary_member() { + let output = ReasoningOutput { + summary: None, + trace: Some("step one, step two".to_string()), + }; + let v = serde_json::to_value(&output).unwrap(); + assert_eq!(v, json!({"trace": "step one, step two"})); + assert_eq!( + serde_json::from_value::(v).unwrap(), + output + ); + } + + #[test] + fn both_fields_round_trip() { + let output = ReasoningOutput { + summary: Some("summary".to_string()), + trace: Some("trace".to_string()), + }; + let v = serde_json::to_value(&output).unwrap(); + assert_eq!(v, json!({"summary": "summary", "trace": "trace"})); + assert_eq!( + serde_json::from_value::(v).unwrap(), + output + ); + } + + #[test] + fn absent_members_are_omitted_rather_than_null() { + let v = serde_json::to_value(ReasoningOutput::default()).unwrap(); + assert_eq!(v, json!({})); + assert!(ReasoningOutput::default().is_empty()); + } + + #[test] + fn explicit_nulls_deserialize_as_absent() { + let output: ReasoningOutput = + serde_json::from_value(json!({"summary": null, "trace": null})).unwrap(); + assert!(output.is_empty()); + } +} diff --git a/lib/foundation/fabro-types/src/run_event/agent.rs b/lib/foundation/fabro-types/src/run_event/agent.rs index ea7476a3c..9fc70b821 100644 --- a/lib/foundation/fabro-types/src/run_event/agent.rs +++ b/lib/foundation/fabro-types/src/run_event/agent.rs @@ -7,7 +7,7 @@ use super::BilledTokenCounts; use crate::transcript::{ToolCall, ToolResult, TranscriptMessage}; use crate::{ MessageId, ModelRef, PairId, PairMessageId, PairSystemMessageKind, PermissionLevel, - StageContextWindowProjection, TurnId, + ReasoningOutput, StageContextWindowProjection, TurnId, }; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -135,6 +135,11 @@ pub struct AgentMessageProps { /// computed from the request that produced this assistant response. #[serde(default, skip_serializing_if = "Option::is_none")] pub context_window: Option, + /// Readable reasoning the provider returned with this response, if any. + /// Absent when the provider returned none or returned only opaque + /// material, so events without reasoning keep their previous shape. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub reasoning: Option, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -409,6 +414,33 @@ mod tests { assert!(props.cost_source.is_none()); assert!(props.message.is_none()); assert!(props.context_window.is_none()); + assert!(props.reasoning.is_none()); + } + + #[test] + fn agent_message_props_round_trips_reasoning_with_both_fields() { + let props = AgentMessageProps { + text: String::new(), + model: sample_model_ref(), + billing: BilledTokenCounts::default(), + cost_source: None, + tool_call_count: 1, + visit: 1, + message: None, + context_window: None, + reasoning: Some(ReasoningOutput { + summary: Some("inspect the implementation first".to_string()), + trace: Some("read convert.rs, then the sink".to_string()), + }), + }; + let v = serde_json::to_value(&props).unwrap(); + assert_eq!( + v["reasoning"]["summary"], + "inspect the implementation first" + ); + assert_eq!(v["reasoning"]["trace"], "read convert.rs, then the sink"); + let back: AgentMessageProps = serde_json::from_value(v).unwrap(); + assert_eq!(back, props); } #[test] @@ -425,6 +457,7 @@ mod tests { visit: 1, message: Some(msg.clone()), context_window: None, + reasoning: None, }; let v = serde_json::to_value(&props).unwrap(); assert_eq!(v["message"]["kind"], "agent"); diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index a1c864f7d..7a528e3d7 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -2155,6 +2155,7 @@ mod tests { visit: 1, message: None, context_window: None, + reasoning: None, }); let value = serde_json::to_value(&body).unwrap(); @@ -2170,6 +2171,68 @@ mod tests { assert_eq!(parsed.event_name(), "agent.message"); } + #[test] + fn agent_message_omits_reasoning_when_absent() { + let body = EventBody::AgentMessage(AgentMessageProps { + text: "ok".to_string(), + model: crate::ModelRef { + provider: fabro_model::ProviderId::openai(), + model_id: "gpt-5.4".into(), + speed: None, + }, + billing: BilledTokenCounts::default(), + cost_source: None, + tool_call_count: 0, + visit: 1, + message: None, + context_window: None, + reasoning: None, + }); + + let value = serde_json::to_value(&body).unwrap(); + assert!( + value["properties"] + .as_object() + .unwrap() + .get("reasoning") + .is_none() + ); + } + + #[test] + fn agent_message_carries_reasoning_through_canonical_json() { + let body = EventBody::AgentMessage(AgentMessageProps { + text: String::new(), + model: crate::ModelRef { + provider: fabro_model::ProviderId::openai(), + model_id: "gpt-5.4".into(), + speed: None, + }, + billing: BilledTokenCounts::default(), + cost_source: None, + tool_call_count: 1, + visit: 1, + message: None, + context_window: None, + reasoning: Some(crate::ReasoningOutput { + summary: Some("inspect the implementation first".to_string()), + trace: Some("read convert.rs, then the sink".to_string()), + }), + }); + + let value = serde_json::to_value(&body).unwrap(); + assert_eq!(value["event"], "agent.message"); + assert_eq!( + value["properties"]["reasoning"], + serde_json::json!({ + "summary": "inspect the implementation first", + "trace": "read convert.rs, then the sink", + }) + ); + let parsed: EventBody = serde_json::from_value(value).unwrap(); + assert_eq!(parsed, body); + } + #[test] fn agent_message_round_trips_optional_context_window() { let context_window = crate::StageContextWindowProjection { @@ -2208,6 +2271,7 @@ mod tests { visit: 1, message: None, context_window: Some(context_window), + reasoning: None, }); let value = serde_json::to_value(&body).unwrap(); diff --git a/lib/foundation/fabro-types/src/transcript.rs b/lib/foundation/fabro-types/src/transcript.rs index a9969e34e..fce181dba 100644 --- a/lib/foundation/fabro-types/src/transcript.rs +++ b/lib/foundation/fabro-types/src/transcript.rs @@ -237,6 +237,11 @@ impl ContentPart { pub const OPENAI_REASONING: &str = "openai_reasoning"; /// Kind string for opaque OpenAI message output items. pub const OPENAI_MESSAGE: &str = "openai_message"; + /// Kind string for opaque OpenAI-compatible `reasoning_details` entries. + /// The data is the received array of detail objects, preserved verbatim + /// so encrypted entries survive for future provider-aware replay. Only + /// known readable members are ever normalized out of it. + pub const OPENAI_COMPAT_REASONING_DETAILS: &str = "openai_compat_reasoning_details"; pub fn text(text: impl Into) -> Self { Self::Text(text.into()) diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 9ea572eb3..6d7c3803e 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -321,6 +321,7 @@ models/pull-request.ts models/question-type.ts models/reasoning-effort-feature.ts models/reasoning-effort.ts +models/reasoning-output.ts models/related-workflow-diagnostic.ts models/render-workflow-graph-direction.ts models/render-workflow-graph-format.ts diff --git a/lib/packages/fabro-api-client/src/models/agent-message-props.ts b/lib/packages/fabro-api-client/src/models/agent-message-props.ts index a0c268de2..bb98aeeeb 100644 --- a/lib/packages/fabro-api-client/src/models/agent-message-props.ts +++ b/lib/packages/fabro-api-client/src/models/agent-message-props.ts @@ -21,6 +21,9 @@ import type { BilledTokenCounts } from './billed-token-counts'; import type { BillingModelRef } from './billing-model-ref'; // May contain unused imports in some cases // @ts-ignore +import type { ReasoningOutput } from './reasoning-output'; +// May contain unused imports in some cases +// @ts-ignore import type { StageContextWindowProjection } from './stage-context-window-projection'; /** @@ -34,4 +37,5 @@ export interface AgentMessageProps { 'visit': number; 'message'?: { [key: string]: any; } | null; 'context_window'?: StageContextWindowProjection | null; + 'reasoning'?: ReasoningOutput | null; } diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index d80dd8a4c..5d9d663ef 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -291,6 +291,7 @@ export * from './pull-request-user'; export * from './question-type'; export * from './reasoning-effort'; export * from './reasoning-effort-feature'; +export * from './reasoning-output'; export * from './related-workflow-diagnostic'; export * from './render-workflow-graph-direction'; export * from './render-workflow-graph-format'; diff --git a/lib/packages/fabro-api-client/src/models/reasoning-output.ts b/lib/packages/fabro-api-client/src/models/reasoning-output.ts new file mode 100644 index 000000000..ab970e292 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/reasoning-output.ts @@ -0,0 +1,29 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +/** + * Readable model reasoning normalized into a provider-neutral shape. Both members may be present for the same response, and at least one is present whenever the object is emitted. Opaque provider material (signatures, item IDs, encrypted or redacted payloads) never appears here. + */ +export interface ReasoningOutput { + /** + * Model-authored summary of its reasoning. + */ + 'summary'?: string; + /** + * Verbatim readable reasoning text, when the provider returns it. + */ + 'trace'?: string; +} From 4d3de5f5640f55907e3e1d38ad9a85cf985d9356 Mon Sep 17 00:00:00 2001 From: Release Repro Date: Fri, 24 Jul 2026 17:31:14 -0400 Subject: [PATCH 2/3] fix(reasoning): tighten capture normalization --- docs/public/api-reference/fabro-api.yaml | 5 + lib/apps/fabro-server/src/server/tests.rs | 8 +- lib/components/fabro-agent/src/session.rs | 21 +-- .../src/codec/openai_compatible/response.rs | 25 +-- .../src/codec/openai_compatible/stream.rs | 22 +-- .../src/codec/openai_compatible/wire.rs | 52 ++++-- .../fabro-llm/src/providers/fabro_server.rs | 4 +- lib/components/fabro-llm/src/reasoning.rs | 165 +++++++++--------- .../tests/it/wire/openai_compatible.rs | 123 ++++++++----- .../tests/it/wire/openai_responses.rs | 4 +- .../fabro-workflow/src/event/convert.rs | 18 +- .../fabro-workflow/src/event/redaction.rs | 8 +- .../fabro-workflow/src/event/sink.rs | 8 +- .../tests/reasoning_output_round_trip.rs | 24 +-- lib/foundation/fabro-types/src/reasoning.rs | 102 ++++++++--- .../fabro-types/src/run_event/agent.rs | 8 +- .../fabro-types/src/run_event/mod.rs | 8 +- 17 files changed, 358 insertions(+), 247 deletions(-) diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 9759233bf..6c38d71f7 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -10059,6 +10059,11 @@ components: (signatures, item IDs, encrypted or redacted payloads) never appears here. type: object + anyOf: + - required: + - summary + - required: + - trace properties: summary: type: string diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 64fe5f345..0ab365454 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -15408,10 +15408,10 @@ async fn attach_stream_replays_agent_message_reasoning() { cost_source: None, tool_call_count: 1, context_window: None, - reasoning: Some(fabro_types::ReasoningOutput { - summary: Some("inspect the sink first".to_string()), - trace: Some("read events.rs, then attach".to_string()), - }), + reasoning: Some(fabro_types::ReasoningOutput::new( + "inspect the sink first", + "read events.rs, then attach", + )), }, session_id: Some("session-1".to_string()), parent_session_id: None, diff --git a/lib/components/fabro-agent/src/session.rs b/lib/components/fabro-agent/src/session.rs index a7048d833..afb8a0c63 100644 --- a/lib/components/fabro-agent/src/session.rs +++ b/lib/components/fabro-agent/src/session.rs @@ -3967,10 +3967,10 @@ mod tests { session.process_input("What is 2+2?").await.unwrap(); let reasoning = collect_message_reasoning(&mut rx); - assert_eq!(reasoning, vec![Some(ReasoningOutput { - summary: Some("the user wants 2+2".to_string()), - trace: Some("2+2 is 4".to_string()), - })]); + assert_eq!(reasoning, vec![Some(ReasoningOutput::new( + "the user wants 2+2", + "2+2 is 4", + ))]); } #[tokio::test] @@ -3996,10 +3996,7 @@ mod tests { let reasoning = collect_message_reasoning(&mut rx); assert_eq!(reasoning, vec![ - Some(ReasoningOutput { - summary: Some("call the tool".to_string()), - trace: None, - }), + Some(ReasoningOutput::from_summary("call the tool")), None, ]); } @@ -4029,10 +4026,10 @@ mod tests { assert_eq!(provider.call_index.load(Ordering::SeqCst), 2); let reasoning = collect_message_reasoning(&mut rx); - assert_eq!(reasoning, vec![Some(ReasoningOutput { - summary: Some("final summary".to_string()), - trace: Some("final trace".to_string()), - })]); + assert_eq!(reasoning, vec![Some(ReasoningOutput::new( + "final summary", + "final trace" + ))]); } #[tokio::test(start_paused = true)] diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/response.rs b/lib/components/fabro-llm/src/codec/openai_compatible/response.rs index 6f5e6d40e..f7df78729 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/response.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/response.rs @@ -13,22 +13,23 @@ pub(super) fn decode_response( ctx: &CodecCtx<'_>, rate_limit: Option, ) -> Result { - let api_resp: ApiResponse = serde_json::from_str(body) + let mut api_resp: ApiResponse = serde_json::from_str(body) .map_err(|e| Error::network(format!("failed to parse response: {e}"), e))?; - let choice = api_resp.choices.first().ok_or_else(|| Error::Provider { - kind: ProviderErrorKind::Server, - detail: Box::new(ProviderErrorDetail::new( - "no choices in response", - ctx.provider_name, - )), - })?; + let choice = api_resp + .choices + .first_mut() + .ok_or_else(|| Error::Provider { + kind: ProviderErrorKind::Server, + detail: Box::new(ProviderErrorDetail::new( + "no choices in response", + ctx.provider_name, + )), + })?; let mut content_parts = Vec::new(); - if let Some(payload) = &choice.message.reasoning_details { - let mut details = ReasoningDetails::default(); - details.push_payload(payload); - content_parts.extend(details.into_content_part()); + if let Some(payload) = choice.message.reasoning_details.take() { + content_parts.extend(ReasoningDetails::from_complete_payload(payload).into_content_part()); } if let Some(reasoning) = choice.message.reasoning() { if !reasoning.is_empty() { diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs b/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs index 39a35e9e7..f389ee086 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/stream.rs @@ -56,7 +56,7 @@ impl StreamState { } /// Process a parsed SSE chunk and return events to emit, if any. - fn process_chunk(&mut self, chunk: &StreamChunk) -> Option> { + fn process_chunk(&mut self, mut chunk: StreamChunk) -> Option> { // Capture response metadata from the first chunk. if let Some(id) = &chunk.id { if self.response_id.is_empty() { @@ -76,8 +76,8 @@ impl StreamState { self.cost_usd = usage.cost.or(self.cost_usd); } - let choices = chunk.choices.as_ref()?; - let choice = choices.first()?; + let choices = chunk.choices.as_mut()?; + let choice = choices.first_mut()?; let mut events = Vec::new(); @@ -86,7 +86,7 @@ impl StreamState { self.finish_reason = map_finish_reason(Some(reason.as_str())); } - let delta = choice.delta.as_ref()?; + let delta = choice.delta.as_mut()?; // Accumulate reasoning/thinking content (Kimi, etc.). if let Some(reasoning) = delta.reasoning() { @@ -96,8 +96,8 @@ impl StreamState { } // Accumulate structured reasoning detail fragments in wire order. - if let Some(payload) = &delta.reasoning_details { - self.reasoning_details.push_payload(payload); + if let Some(payload) = delta.reasoning_details.take() { + self.reasoning_details.push_stream_payload(payload); } // Handle text content delta. @@ -258,7 +258,7 @@ impl StreamDecoder for StreamState { let chunk: StreamChunk = serde_json::from_str(ev.data) .map_err(|e| Error::stream_error(format!("failed to parse SSE chunk: {e}"), e))?; - Ok(self.process_chunk(&chunk).unwrap_or_default()) + Ok(self.process_chunk(chunk).unwrap_or_default()) } fn finish(&mut self) -> Vec { @@ -369,7 +369,7 @@ mod tests { let chunk1: StreamChunk = serde_json::from_str( r#"{"id":"c1","model":"m1","choices":[{"delta":{"content":"Hello"},"finish_reason":null}]}"#, ).unwrap(); - let events1 = state.process_chunk(&chunk1).unwrap(); + let events1 = state.process_chunk(chunk1).unwrap(); assert_eq!(events1.len(), 2); assert!(matches!(events1[0], StreamEvent::TextStart { .. })); assert!(matches!(events1[1], StreamEvent::TextDelta { .. })); @@ -377,7 +377,7 @@ mod tests { let chunk2: StreamChunk = serde_json::from_str( r#"{"id":"c1","model":"m1","choices":[{"delta":{"content":" world"},"finish_reason":null}]}"#, ).unwrap(); - let events2 = state.process_chunk(&chunk2).unwrap(); + let events2 = state.process_chunk(chunk2).unwrap(); assert_eq!(events2.len(), 1); assert!(matches!(events2[0], StreamEvent::TextDelta { .. })); @@ -391,14 +391,14 @@ mod tests { let chunk1: StreamChunk = serde_json::from_str( r#"{"id":"c1","model":"m1","choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"fn1","arguments":"{\"k"}}]},"finish_reason":null}]}"#, ).unwrap(); - let events1 = state.process_chunk(&chunk1).unwrap(); + let events1 = state.process_chunk(chunk1).unwrap(); assert_eq!(events1.len(), 1); assert!(matches!(events1[0], StreamEvent::ToolCallStart { .. })); let chunk2: StreamChunk = serde_json::from_str( r#"{"id":"c1","model":"m1","choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"ey\"}"}}]},"finish_reason":null}]}"#, ).unwrap(); - let events2 = state.process_chunk(&chunk2).unwrap(); + let events2 = state.process_chunk(chunk2).unwrap(); assert_eq!(events2.len(), 1); assert!(matches!(events2[0], StreamEvent::ToolCallDelta { .. })); diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs index 01e88d500..25668a023 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs @@ -156,25 +156,44 @@ pub(super) struct ReasoningDetails { } impl ReasoningDetails { - /// Absorb one `reasoning_details` payload. + /// Preserve a complete-response `reasoning_details` payload. /// /// Providers document an array of detail objects; a lone object is - /// accepted as a single entry. Scalars carry nothing replayable and are - /// dropped. Streaming deltas repeat the same logical detail across - /// chunks, so a fragment that continues the previous entry is coalesced - /// into it rather than becoming a separate block. - pub(super) fn push_payload(&mut self, payload: &serde_json::Value) { + /// accepted as a single entry. Complete entries retain their received + /// order and shape; scalars carry nothing replayable and are dropped. + pub(super) fn from_complete_payload(payload: serde_json::Value) -> Self { + let entries = match payload { + serde_json::Value::Array(entries) => entries + .into_iter() + .filter(serde_json::Value::is_object) + .collect(), + payload @ serde_json::Value::Object(_) => vec![payload], + _ => Vec::new(), + }; + Self { entries } + } + + /// Absorb one streamed `reasoning_details` payload. + /// + /// Fragments carrying the same `type` and `index` are coalesced even when + /// other logical details appear between them. First-seen detail order is + /// retained. + pub(super) fn push_stream_payload(&mut self, payload: serde_json::Value) { let incoming = match payload { - serde_json::Value::Array(entries) => entries.clone(), - serde_json::Value::Object(_) => vec![payload.clone()], + serde_json::Value::Array(entries) => entries, + payload @ serde_json::Value::Object(_) => vec![payload], _ => Vec::new(), }; for entry in incoming { if !entry.is_object() { continue; } - match self.entries.last_mut() { - Some(last) if continues_detail(last, &entry) => merge_detail_fragment(last, &entry), + match self + .entries + .iter_mut() + .find(|existing| continues_detail(existing, &entry)) + { + Some(existing) => merge_detail_fragment(existing, entry), _ => self.entries.push(entry), } } @@ -198,20 +217,23 @@ const DETAIL_TEXT_MEMBERS: [&str; 3] = ["text", "summary", "data"]; /// Aggregators tag each logical detail with a stable `type` and `index`; /// fragment streams that omit `index` are matched on `type` alone. fn continues_detail(last: &serde_json::Value, entry: &serde_json::Value) -> bool { - last.get("type") == entry.get("type") && last.get("index") == entry.get("index") + let Some(entry_type) = entry.get("type") else { + return false; + }; + last.get("type") == Some(entry_type) && last.get("index") == entry.get("index") } /// Append `entry`'s text fragments onto `last` and fill in members `last` /// has not seen yet. -fn merge_detail_fragment(last: &mut serde_json::Value, entry: &serde_json::Value) { - let Some(entry_members) = entry.as_object() else { +fn merge_detail_fragment(last: &mut serde_json::Value, entry: serde_json::Value) { + let serde_json::Value::Object(entry_members) = entry else { return; }; let Some(last_members) = last.as_object_mut() else { return; }; for (key, value) in entry_members { - match last_members.get_mut(key) { + match last_members.get_mut(&key) { Some(serde_json::Value::String(existing)) if DETAIL_TEXT_MEMBERS.contains(&key.as_str()) => { @@ -221,7 +243,7 @@ fn merge_detail_fragment(last: &mut serde_json::Value, entry: &serde_json::Value } Some(_) => {} None => { - last_members.insert(key.clone(), value.clone()); + last_members.insert(key, value); } } } diff --git a/lib/components/fabro-llm/src/providers/fabro_server.rs b/lib/components/fabro-llm/src/providers/fabro_server.rs index bd5c42ea6..5e4d6d990 100644 --- a/lib/components/fabro-llm/src/providers/fabro_server.rs +++ b/lib/components/fabro-llm/src/providers/fabro_server.rs @@ -390,8 +390,8 @@ data: {\"type\":\"text_delta\",\"delta\":\" world\",\"text_id\":null}\n\ assert_eq!(response.text(), "Hello there!"); let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("weighed both")); - assert_eq!(reasoning.trace.as_deref(), Some("step one")); + assert_eq!(reasoning.summary(), Some("weighed both")); + assert_eq!(reasoning.trace(), Some("step one")); } #[tokio::test] diff --git a/lib/components/fabro-llm/src/reasoning.rs b/lib/components/fabro-llm/src/reasoning.rs index a4909ec1b..cead5d180 100644 --- a/lib/components/fabro-llm/src/reasoning.rs +++ b/lib/components/fabro-llm/src/reasoning.rs @@ -18,60 +18,50 @@ use fabro_types::{ContentPart, ReasoningOutput}; /// this module. const BLOCK_SEPARATOR: &str = "\n\n"; -/// Detail types whose payload is opaque and must never reach the public -/// normalized fields. -fn is_opaque_detail_type(detail_type: &str) -> bool { - detail_type.contains("encrypted") || detail_type.contains("redacted") -} - /// Readable blocks collected per normalized field. /// /// Explicit blocks come from a channel with documented reasoning semantics. /// Fallback blocks come from flattened provider strings, which aggregators -/// commonly duplicate alongside a structured channel; they only fill a -/// summary that no explicit block produced. +/// commonly duplicate alongside a structured channel. They only fill a trace +/// that no explicit trace produced. #[derive(Default)] -struct Blocks { - explicit_summary: Vec, - explicit_trace: Vec, - fallback_summary: Vec, +struct Blocks<'a> { + explicit_summary: Vec<&'a str>, + explicit_trace: Vec<&'a str>, + fallback_trace: Vec<&'a str>, } -impl Blocks { +impl Blocks<'_> { fn into_output(self) -> Option { - let summary = - join_blocks(&self.explicit_summary).or_else(|| join_blocks(&self.fallback_summary)); - let trace = join_blocks(&self.explicit_trace); - let output = ReasoningOutput { summary, trace }; - (!output.is_empty()).then_some(output) + let summary = join_blocks(&self.explicit_summary); + let trace = join_blocks(&self.explicit_trace) + .or_else(|| join_blocks(&self.fallback_trace)) + .filter(|trace| summary.as_ref() != Some(trace)); + + match (summary, trace) { + (Some(summary), Some(trace)) => Some(ReasoningOutput::new(summary, trace)), + (Some(summary), None) => Some(ReasoningOutput::from_summary(summary)), + (None, Some(trace)) => Some(ReasoningOutput::from_trace(trace)), + (None, None) => None, + } } } /// Join complete blocks in provider order, dropping empty and /// whitespace-only fragments. Retained text is never trimmed or rewritten. -fn join_blocks(blocks: &[String]) -> Option { - let joined = blocks - .iter() - .filter(|block| !block.trim().is_empty()) - .map(String::as_str) - .collect::>() - .join(BLOCK_SEPARATOR); - (!joined.is_empty()).then_some(joined) +fn join_blocks(blocks: &[&str]) -> Option { + (!blocks.is_empty()).then(|| blocks.join(BLOCK_SEPARATOR)) } -/// Read a text-bearing member, preferring the member the provider's -/// documented semantics name and accepting the other readable spelling. -fn readable_member(entry: &serde_json::Value, preferred: &str) -> Option { - let other = if preferred == "text" { - "summary" - } else { - "text" - }; - entry - .get(preferred) - .and_then(serde_json::Value::as_str) - .or_else(|| entry.get(other).and_then(serde_json::Value::as_str)) - .map(str::to_string) +fn push_block<'a>(blocks: &mut Vec<&'a str>, block: &'a str) { + if !block.trim().is_empty() { + blocks.push(block); + } +} + +/// Read a text-bearing member with the provider's documented semantics. +fn readable_member<'a>(entry: &'a serde_json::Value, member: &str) -> Option<&'a str> { + entry.get(member).and_then(serde_json::Value::as_str) } /// Extract readable text from an OpenAI Responses `reasoning` output item. @@ -79,13 +69,13 @@ fn readable_member(entry: &serde_json::Value, preferred: &str) -> Option /// `summary[].text` is the model-authored summary; `content[]` entries typed /// `reasoning_text` are the verbatim trace. `encrypted_content`, `id`, and /// `status` are opaque and ignored. -fn collect_openai_reasoning_item(item: &serde_json::Value, blocks: &mut Blocks) { +fn collect_openai_reasoning_item<'a>(item: &'a serde_json::Value, blocks: &mut Blocks<'a>) { if let Some(entries) = item.get("summary").and_then(serde_json::Value::as_array) { for entry in entries { if let Some(text) = entry.as_str() { - blocks.explicit_summary.push(text.to_string()); + push_block(&mut blocks.explicit_summary, text); } else if let Some(text) = entry.get("text").and_then(serde_json::Value::as_str) { - blocks.explicit_summary.push(text.to_string()); + push_block(&mut blocks.explicit_summary, text); } } } @@ -99,18 +89,14 @@ fn collect_openai_reasoning_item(item: &serde_json::Value, blocks: &mut Blocks) .and_then(serde_json::Value::as_str) .unwrap_or_default(); if entry_type == "reasoning_text" { - blocks.explicit_trace.push(text.to_string()); - } else if !is_opaque_detail_type(entry_type) { - // Readable text from the reasoning channel that carries no - // recognized classification is treated as a summary. - blocks.explicit_summary.push(text.to_string()); + push_block(&mut blocks.explicit_trace, text); } } } } /// Extract readable text from OpenAI-compatible `reasoning_details` entries. -fn collect_reasoning_details(details: &serde_json::Value, blocks: &mut Blocks) { +fn collect_reasoning_details<'a>(details: &'a serde_json::Value, blocks: &mut Blocks<'a>) { let Some(entries) = details.as_array() else { return; }; @@ -119,23 +105,18 @@ fn collect_reasoning_details(details: &serde_json::Value, blocks: &mut Blocks) { .get("type") .and_then(serde_json::Value::as_str) .unwrap_or_default(); - if is_opaque_detail_type(detail_type) { - continue; - } match detail_type { "reasoning.text" => { if let Some(text) = readable_member(entry, "text") { - blocks.explicit_trace.push(text); + push_block(&mut blocks.explicit_trace, text); } } - // Documented summaries and unknown detail variants both come - // from an established reasoning channel, so a readable member - // is kept as a summary either way. - _ => { + "reasoning.summary" => { if let Some(text) = readable_member(entry, "summary") { - blocks.explicit_summary.push(text); + push_block(&mut blocks.explicit_summary, text); } } + _ => {} } } } @@ -149,7 +130,7 @@ pub(crate) fn normalize(content: &[ContentPart]) -> Option { for part in content { match part { ContentPart::Thinking(thinking) if !thinking.redacted => { - blocks.fallback_summary.push(thinking.text.clone()); + push_block(&mut blocks.fallback_trace, &thinking.text); } ContentPart::Other { kind, data } if kind == ContentPart::OPENAI_REASONING => { collect_openai_reasoning_item(data, &mut blocks); @@ -195,10 +176,10 @@ mod tests { } #[test] - fn non_redacted_thinking_becomes_a_summary() { + fn non_redacted_thinking_becomes_a_trace() { let output = normalize(&[thinking("weighing the options")]).unwrap(); - assert_eq!(output.summary.as_deref(), Some("weighing the options")); - assert!(output.trace.is_none()); + assert!(output.summary().is_none()); + assert_eq!(output.trace(), Some("weighing the options")); } #[test] @@ -221,8 +202,8 @@ mod tests { "content": [{"type": "reasoning_text", "text": "step one"}], }))]) .unwrap(); - assert_eq!(output.summary.as_deref(), Some("inspect first")); - assert_eq!(output.trace.as_deref(), Some("step one")); + assert_eq!(output.summary(), Some("inspect first")); + assert_eq!(output.trace(), Some("step one")); } #[test] @@ -234,7 +215,17 @@ mod tests { ], }))]) .unwrap(); - assert_eq!(output.summary.as_deref(), Some("first\n\nsecond")); + assert_eq!(output.summary(), Some("first\n\nsecond")); + } + + #[test] + fn unknown_responses_content_types_remain_opaque() { + assert!( + normalize(&[openai_reasoning(json!({ + "content": [{"type": "reasoning_future", "text": "not classified"}], + }))]) + .is_none() + ); } #[test] @@ -244,8 +235,8 @@ mod tests { {"type": "reasoning.text", "text": "read convert.rs", "signature": "sig"}, ]))]) .unwrap(); - assert_eq!(output.summary.as_deref(), Some("checked the parser")); - assert_eq!(output.trace.as_deref(), Some("read convert.rs")); + assert_eq!(output.summary(), Some("checked the parser")); + assert_eq!(output.trace(), Some("read convert.rs")); } #[test] @@ -255,8 +246,8 @@ mod tests { {"type": "reasoning.summary", "summary": "visible"}, ])),]) .unwrap(); - assert_eq!(output.summary.as_deref(), Some("visible")); - assert!(output.trace.is_none()); + assert_eq!(output.summary(), Some("visible")); + assert!(output.trace().is_none()); } #[test] @@ -270,12 +261,13 @@ mod tests { } #[test] - fn unknown_detail_variants_are_treated_as_summary() { - let output = normalize(&[reasoning_details(json!([ - {"type": "reasoning.future", "text": "new channel"}, - ]))]) - .unwrap(); - assert_eq!(output.summary.as_deref(), Some("new channel")); + fn unknown_detail_variants_remain_opaque() { + assert!( + normalize(&[reasoning_details(json!([ + {"type": "reasoning.future", "text": "new channel"}, + ]))]) + .is_none() + ); } #[test] @@ -300,19 +292,32 @@ mod tests { thinking("checked the parser"), ]) .unwrap(); - assert_eq!(output.summary.as_deref(), Some("checked the parser")); - assert!(output.trace.is_none()); + assert_eq!(output.summary(), Some("checked the parser")); + assert!(output.trace().is_none()); } #[test] - fn trace_only_detail_lets_the_fallback_fill_the_summary() { + fn structured_trace_takes_precedence_over_flattened_trace() { let output = normalize(&[ reasoning_details(json!([{"type": "reasoning.text", "text": "verbatim"}])), thinking("flattened"), ]) .unwrap(); - assert_eq!(output.summary.as_deref(), Some("flattened")); - assert_eq!(output.trace.as_deref(), Some("verbatim")); + assert!(output.summary().is_none()); + assert_eq!(output.trace(), Some("verbatim")); + } + + #[test] + fn structured_summary_keeps_a_distinct_flattened_trace() { + let output = normalize(&[ + reasoning_details(json!([ + {"type": "reasoning.summary", "summary": "short summary"}, + ])), + thinking("full verbatim trace"), + ]) + .unwrap(); + assert_eq!(output.summary(), Some("short summary")); + assert_eq!(output.trace(), Some("full verbatim trace")); } #[test] @@ -323,7 +328,7 @@ mod tests { #[test] fn non_empty_text_is_preserved_verbatim() { let output = normalize(&[thinking(" indented thought\n")]).unwrap(); - assert_eq!(output.summary.as_deref(), Some(" indented thought\n")); + assert_eq!(output.trace(), Some(" indented thought\n")); } #[test] diff --git a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs index 4ba98dfa2..a35ba6427 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs @@ -4,6 +4,7 @@ use std::sync::Arc; +use fabro_llm::generate::StreamAccumulator; use fabro_llm::provider::ProviderAdapter; use fabro_llm::providers::OpenAiCompatibleAdapter; use fabro_llm::types::{ @@ -29,6 +30,11 @@ const CREATED_TS: i64 = 1_700_000_000; /// Minimal valid Chat Completions body for encode-side tests. fn minimal_body() -> serde_json::Value { + body_with_message(&serde_json::json!({"role": "assistant", "content": "ok"})) +} + +/// Wraps an assistant message in a complete Chat Completions body. +fn body_with_message(message: &serde_json::Value) -> serde_json::Value { serde_json::json!({ "id": "chatcmpl_test", "object": "chat.completion", @@ -36,7 +42,7 @@ fn minimal_body() -> serde_json::Value { "model": MODEL, "choices": [{ "index": 0, - "message": {"role": "assistant", "content": "ok"}, + "message": message, "finish_reason": "stop" }], "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2} @@ -437,18 +443,6 @@ async fn decode_response(body: serde_json::Value) -> fabro_llm::types::Response response } -/// Wraps an assistant message in a complete Chat Completions body. -fn body_with_message(message: &serde_json::Value) -> serde_json::Value { - serde_json::json!({ - "id": "chatcmpl_test", - "object": "chat.completion", - "created": CREATED_TS, - "model": MODEL, - "choices": [{"index": 0, "message": message, "finish_reason": "stop"}], - "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15} - }) -} - /// Streams an SSE transcript and returns the final accumulated response. async fn stream_final_response(sse_body: &str) -> fabro_llm::types::Response { use futures::StreamExt; @@ -460,14 +454,15 @@ async fn stream_final_response(sse_body: &str) -> fabro_llm::types::Response { .stream(&base_request(MODEL)) .await .expect("stream should start"); - let mut final_response = None; + let mut accumulator = StreamAccumulator::new(); while let Some(item) = stream.next().await { - if let Ok(fabro_llm::types::StreamEvent::Finish { response, .. }) = item { - final_response = Some(*response); - } + accumulator.process(&item.expect("stream event should decode")); } mock.assert(); - final_response.expect("stream should emit a finish event") + accumulator + .response() + .cloned() + .expect("stream should emit a finish event") } #[tokio::test] @@ -536,8 +531,8 @@ async fn decode_reasoning_details_normalize_summary_and_trace() { .await; let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("the user wants 2+2")); - assert_eq!(reasoning.trace.as_deref(), Some("2 plus 2 is 4")); + assert_eq!(reasoning.summary(), Some("the user wants 2+2")); + assert_eq!(reasoning.trace(), Some("2 plus 2 is 4")); } /// Encrypted entries stay in the opaque provider part for future replay but @@ -570,8 +565,42 @@ async fn decode_reasoning_details_preserve_encrypted_entries_opaquely() { assert_eq!(opaque[0]["data"], "gAAAAAopaque"); let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("visible")); - assert!(reasoning.trace.is_none()); + assert_eq!(reasoning.summary(), Some("visible")); + assert!(reasoning.trace().is_none()); +} + +/// Complete-response details are already assembled and must retain their +/// received block boundaries. +#[tokio::test] +async fn decode_reasoning_details_preserves_complete_entries_verbatim() { + let details = serde_json::json!([ + {"type": "reasoning.summary", "summary": "first"}, + {"type": "reasoning.summary", "summary": "second"}, + ]); + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning_details": details, + }))) + .await; + + let opaque = response + .message + .content + .iter() + .find_map(|part| match part { + fabro_llm::types::ContentPart::Other { kind, data } + if kind == fabro_llm::types::ContentPart::OPENAI_COMPAT_REASONING_DETAILS => + { + Some(data) + } + _ => None, + }) + .expect("opaque reasoning details preserved"); + assert_eq!(opaque, &details); + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary(), Some("first\n\nsecond")); } /// Unknown and malformed detail entries must not fail an otherwise valid @@ -591,8 +620,7 @@ async fn decode_tolerates_unknown_and_malformed_reasoning_details() { .await; assert_eq!(response.text(), "4."); - let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("new channel")); + assert!(response.reasoning_output().is_none()); } /// A scalar `reasoning_details` carries nothing replayable and is dropped @@ -625,25 +653,42 @@ async fn decode_structured_details_suppress_the_duplicate_flattened_value() { .await; let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("the user wants 2+2")); - assert!(reasoning.trace.is_none()); + assert_eq!(reasoning.summary(), Some("the user wants 2+2")); + assert!(reasoning.trace().is_none()); } -/// A structured trace with no structured summary leaves room for the -/// flattened value to fill the summary. +/// A structured trace takes precedence over the flattened trace channel. #[tokio::test] -async fn decode_trace_only_details_let_the_flattened_value_fill_the_summary() { +async fn decode_structured_trace_takes_precedence_over_flattened_trace() { let response = decode_response(body_with_message(&serde_json::json!({ "role": "assistant", "content": "4.", - "reasoning": "flattened summary", + "reasoning": "flattened trace", "reasoning_details": [{"type": "reasoning.text", "text": "verbatim trace", "index": 0}] }))) .await; let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("flattened summary")); - assert_eq!(reasoning.trace.as_deref(), Some("verbatim trace")); + assert!(reasoning.summary().is_none()); + assert_eq!(reasoning.trace(), Some("verbatim trace")); +} + +/// A structured summary and distinct flattened trace are both retained. +#[tokio::test] +async fn decode_structured_summary_keeps_distinct_flattened_trace() { + let response = decode_response(body_with_message(&serde_json::json!({ + "role": "assistant", + "content": "4.", + "reasoning": "full verbatim trace", + "reasoning_details": [ + {"type": "reasoning.summary", "summary": "short summary", "index": 0}, + ] + }))) + .await; + + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary(), Some("short summary")); + assert_eq!(reasoning.trace(), Some("full verbatim trace")); } /// Streamed detail fragments coalesce back into the same normalized output @@ -651,9 +696,8 @@ async fn decode_trace_only_details_let_the_flattened_value_fill_the_summary() { #[tokio::test] async fn stream_reasoning_details_normalize_like_the_non_streaming_body() { let sse = support::sse_data_transcript(&[ - r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"role":"assistant","reasoning_details":[{"type":"reasoning.summary","summary":"the user ","index":0}]},"finish_reason":null}]}"#, - r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"reasoning.summary","summary":"wants 2+2","index":0}]},"finish_reason":null}]}"#, - r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"reasoning.text","text":"2 plus 2 is 4","index":1}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"role":"assistant","reasoning_details":[{"type":"reasoning.summary","summary":"the user ","index":0},{"type":"reasoning.text","text":"2 plus ","index":1}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"reasoning.summary","summary":"wants 2+2","index":0},{"type":"reasoning.text","text":"2 is 4","index":1}]},"finish_reason":null}]}"#, r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"content":"4."},"finish_reason":null}]}"#, r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}"#, "[DONE]", @@ -671,12 +715,9 @@ async fn stream_reasoning_details_normalize_like_the_non_streaming_body() { .await; assert_eq!(streamed.reasoning_output(), non_streamed.reasoning_output()); - assert_eq!( - streamed - .reasoning_output() - .and_then(|reasoning| reasoning.summary), - Some("the user wants 2+2".to_string()) - ); + let reasoning = streamed.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.summary(), Some("the user wants 2+2")); + assert_eq!(reasoning.trace(), Some("2 plus 2 is 4")); } /// Cached and reasoning detail tokens are split into their own disjoint diff --git a/lib/components/fabro-llm/tests/it/wire/openai_responses.rs b/lib/components/fabro-llm/tests/it/wire/openai_responses.rs index fd281ce62..1a3f2fb6a 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_responses.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_responses.rs @@ -490,8 +490,8 @@ async fn decode_reasoning_item_normalizes_summary_and_trace() { .await; let reasoning = response.reasoning_output().expect("reasoning present"); - assert_eq!(reasoning.summary.as_deref(), Some("Adding two numbers.")); - assert_eq!(reasoning.trace.as_deref(), Some("2 plus 2 is 4.")); + assert_eq!(reasoning.summary(), Some("Adding two numbers.")); + assert_eq!(reasoning.trace(), Some("2 plus 2 is 4.")); let normalized = serde_json::to_string(&reasoning).unwrap(); assert!(!normalized.contains("gAAAAAopaque")); diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 4beb3b888..d86c17424 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -2391,10 +2391,10 @@ mod tests { cost_source: None, tool_call_count: 1, context_window: None, - reasoning: Some(::fabro_types::ReasoningOutput { - summary: Some("inspect the conversion first".to_string()), - trace: Some("read convert.rs, then the sink".to_string()), - }), + reasoning: Some(::fabro_types::ReasoningOutput::new( + "inspect the conversion first", + "read convert.rs, then the sink", + )), }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2405,14 +2405,8 @@ mod tests { panic!("expected agent message body"); }; let reasoning = message.reasoning.as_ref().expect("reasoning copied"); - assert_eq!( - reasoning.summary.as_deref(), - Some("inspect the conversion first") - ); - assert_eq!( - reasoning.trace.as_deref(), - Some("read convert.rs, then the sink") - ); + assert_eq!(reasoning.summary(), Some("inspect the conversion first")); + assert_eq!(reasoning.trace(), Some("read convert.rs, then the sink")); let value = stored.to_value().unwrap(); assert_eq!(value["event"], "agent.message"); diff --git a/lib/components/fabro-workflow/src/event/redaction.rs b/lib/components/fabro-workflow/src/event/redaction.rs index d0d7d0974..16aabccab 100644 --- a/lib/components/fabro-workflow/src/event/redaction.rs +++ b/lib/components/fabro-workflow/src/event/redaction.rs @@ -96,10 +96,10 @@ mod tests { cost_source: None, tool_call_count: 0, context_window: None, - reasoning: Some(ReasoningOutput { - summary: Some(format!("the key is {secret}")), - trace: Some(format!("reading {secret} from the env")), - }), + reasoning: Some(ReasoningOutput::new( + format!("the key is {secret}"), + format!("reading {secret} from the env"), + )), }, session_id: Some("ses_agent".to_string()), parent_session_id: None, diff --git a/lib/components/fabro-workflow/src/event/sink.rs b/lib/components/fabro-workflow/src/event/sink.rs index 0b9404b31..eb987f620 100644 --- a/lib/components/fabro-workflow/src/event/sink.rs +++ b/lib/components/fabro-workflow/src/event/sink.rs @@ -330,10 +330,10 @@ mod tests { cost_source: None, tool_call_count: 1, context_window: None, - reasoning: Some(::fabro_types::ReasoningOutput { - summary: Some("inspect the sink first".to_string()), - trace: Some("write the line, then read it back".to_string()), - }), + reasoning: Some(::fabro_types::ReasoningOutput::new( + "inspect the sink first", + "write the line, then read it back", + )), }, session_id: Some("ses_agent".to_string()), parent_session_id: None, diff --git a/lib/foundation/fabro-api/tests/reasoning_output_round_trip.rs b/lib/foundation/fabro-api/tests/reasoning_output_round_trip.rs index e4f02f719..035e66c23 100644 --- a/lib/foundation/fabro-api/tests/reasoning_output_round_trip.rs +++ b/lib/foundation/fabro-api/tests/reasoning_output_round_trip.rs @@ -19,14 +19,8 @@ fn reasoning_output_matches_openapi_json_shape() { }); let output: ReasoningOutput = serde_json::from_value(value.clone()).unwrap(); - assert_eq!( - output.summary.as_deref(), - Some("inspect the conversion first") - ); - assert_eq!( - output.trace.as_deref(), - Some("read convert.rs, then the sink") - ); + assert_eq!(output.summary(), Some("inspect the conversion first")); + assert_eq!(output.trace(), Some("read convert.rs, then the sink")); assert_eq!(serde_json::to_value(&output).unwrap(), value); let api_output: ApiReasoningOutput = serde_json::from_value(value).unwrap(); @@ -37,7 +31,7 @@ fn reasoning_output_matches_openapi_json_shape() { fn reasoning_output_members_are_individually_optional() { let summary_only: ReasoningOutput = serde_json::from_value(json!({"summary": "only a summary"})).unwrap(); - assert!(summary_only.trace.is_none()); + assert!(summary_only.trace().is_none()); assert_eq!( serde_json::to_value(&summary_only).unwrap(), json!({"summary": "only a summary"}) @@ -45,13 +39,19 @@ fn reasoning_output_members_are_individually_optional() { let trace_only: ReasoningOutput = serde_json::from_value(json!({"trace": "only a trace"})).unwrap(); - assert!(trace_only.summary.is_none()); + assert!(trace_only.summary().is_none()); assert_eq!( serde_json::to_value(&trace_only).unwrap(), json!({"trace": "only a trace"}) ); } +#[test] +fn reasoning_output_rejects_an_empty_object() { + let error = serde_json::from_value::(json!({})).unwrap_err(); + assert!(error.to_string().contains("requires a summary or trace")); +} + #[test] fn agent_message_props_reasoning_is_optional_on_the_wire() { let without = json!({ @@ -75,8 +75,8 @@ fn agent_message_props_reasoning_is_optional_on_the_wire() { with["reasoning"] = json!({"summary": "checked the parser", "trace": "step one"}); let props: ApiAgentMessageProps = serde_json::from_value(with).unwrap(); let reasoning = props.reasoning.as_ref().unwrap(); - assert_eq!(reasoning.summary.as_deref(), Some("checked the parser")); - assert_eq!(reasoning.trace.as_deref(), Some("step one")); + assert_eq!(reasoning.summary(), Some("checked the parser")); + assert_eq!(reasoning.trace(), Some("step one")); assert_eq!( serde_json::to_value(&props).unwrap()["reasoning"], json!({"summary": "checked the parser", "trace": "step one"}) diff --git a/lib/foundation/fabro-types/src/reasoning.rs b/lib/foundation/fabro-types/src/reasoning.rs index 0363c6deb..89d521e9f 100644 --- a/lib/foundation/fabro-types/src/reasoning.rs +++ b/lib/foundation/fabro-types/src/reasoning.rs @@ -1,4 +1,4 @@ -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Serialize, de}; /// Readable model reasoning normalized into a provider-neutral shape. /// @@ -12,23 +12,78 @@ use serde::{Deserialize, Serialize}; /// Both fields may be populated for the same response. An emitted object /// always carries at least one of them; opaque provider material /// (signatures, IDs, encrypted or redacted payloads) never appears here. -#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +#[derive(Debug, Clone, PartialEq, Eq, Serialize)] pub struct ReasoningOutput { /// Model-authored summary of its reasoning, safe to show to users. #[serde(default, skip_serializing_if = "Option::is_none")] - pub summary: Option, + summary: Option, /// Verbatim readable reasoning text, when the provider returns it in /// addition to (or instead of) a summary. #[serde(default, skip_serializing_if = "Option::is_none")] - pub trace: Option, + trace: Option, } impl ReasoningOutput { - /// Returns `true` when neither readable field is present, meaning the - /// object carries nothing worth emitting. + /// Creates reasoning output with both a model-authored summary and a + /// verbatim trace. #[must_use] - pub fn is_empty(&self) -> bool { - self.summary.is_none() && self.trace.is_none() + pub fn new(summary: impl Into, trace: impl Into) -> Self { + Self { + summary: Some(summary.into()), + trace: Some(trace.into()), + } + } + + /// Creates reasoning output containing only a model-authored summary. + #[must_use] + pub fn from_summary(summary: impl Into) -> Self { + Self { + summary: Some(summary.into()), + trace: None, + } + } + + /// Creates reasoning output containing only a verbatim trace. + #[must_use] + pub fn from_trace(trace: impl Into) -> Self { + Self { + summary: None, + trace: Some(trace.into()), + } + } + + /// Returns the model-authored summary, when present. + #[must_use] + pub fn summary(&self) -> Option<&str> { + self.summary.as_deref() + } + + /// Returns the verbatim readable reasoning trace, when present. + #[must_use] + pub fn trace(&self) -> Option<&str> { + self.trace.as_deref() + } +} + +impl<'de> Deserialize<'de> for ReasoningOutput { + fn deserialize>(deserializer: D) -> Result { + #[derive(Deserialize)] + struct Fields { + #[serde(default)] + summary: Option, + #[serde(default)] + trace: Option, + } + + let Fields { summary, trace } = Fields::deserialize(deserializer)?; + match (summary, trace) { + (Some(summary), Some(trace)) => Ok(Self::new(summary, trace)), + (Some(summary), None) => Ok(Self::from_summary(summary)), + (None, Some(trace)) => Ok(Self::from_trace(trace)), + (None, None) => Err(de::Error::custom( + "reasoning output requires a summary or trace", + )), + } } } @@ -40,10 +95,7 @@ mod tests { #[test] fn summary_only_round_trips_without_trace_member() { - let output = ReasoningOutput { - summary: Some("checked the parser first".to_string()), - trace: None, - }; + let output = ReasoningOutput::from_summary("checked the parser first"); let v = serde_json::to_value(&output).unwrap(); assert_eq!(v, json!({"summary": "checked the parser first"})); assert_eq!( @@ -54,10 +106,7 @@ mod tests { #[test] fn trace_only_round_trips_without_summary_member() { - let output = ReasoningOutput { - summary: None, - trace: Some("step one, step two".to_string()), - }; + let output = ReasoningOutput::from_trace("step one, step two"); let v = serde_json::to_value(&output).unwrap(); assert_eq!(v, json!({"trace": "step one, step two"})); assert_eq!( @@ -68,10 +117,7 @@ mod tests { #[test] fn both_fields_round_trip() { - let output = ReasoningOutput { - summary: Some("summary".to_string()), - trace: Some("trace".to_string()), - }; + let output = ReasoningOutput::new("summary", "trace"); let v = serde_json::to_value(&output).unwrap(); assert_eq!(v, json!({"summary": "summary", "trace": "trace"})); assert_eq!( @@ -81,16 +127,16 @@ mod tests { } #[test] - fn absent_members_are_omitted_rather_than_null() { - let v = serde_json::to_value(ReasoningOutput::default()).unwrap(); - assert_eq!(v, json!({})); - assert!(ReasoningOutput::default().is_empty()); + fn empty_object_is_rejected() { + let error = serde_json::from_value::(json!({})).unwrap_err(); + assert!(error.to_string().contains("requires a summary or trace")); } #[test] - fn explicit_nulls_deserialize_as_absent() { - let output: ReasoningOutput = - serde_json::from_value(json!({"summary": null, "trace": null})).unwrap(); - assert!(output.is_empty()); + fn explicit_nulls_are_rejected() { + let error = + serde_json::from_value::(json!({"summary": null, "trace": null})) + .unwrap_err(); + assert!(error.to_string().contains("requires a summary or trace")); } } diff --git a/lib/foundation/fabro-types/src/run_event/agent.rs b/lib/foundation/fabro-types/src/run_event/agent.rs index 9fc70b821..9cd852407 100644 --- a/lib/foundation/fabro-types/src/run_event/agent.rs +++ b/lib/foundation/fabro-types/src/run_event/agent.rs @@ -428,10 +428,10 @@ mod tests { visit: 1, message: None, context_window: None, - reasoning: Some(ReasoningOutput { - summary: Some("inspect the implementation first".to_string()), - trace: Some("read convert.rs, then the sink".to_string()), - }), + reasoning: Some(ReasoningOutput::new( + "inspect the implementation first", + "read convert.rs, then the sink", + )), }; let v = serde_json::to_value(&props).unwrap(); assert_eq!( diff --git a/lib/foundation/fabro-types/src/run_event/mod.rs b/lib/foundation/fabro-types/src/run_event/mod.rs index 7a528e3d7..7c8d1621a 100644 --- a/lib/foundation/fabro-types/src/run_event/mod.rs +++ b/lib/foundation/fabro-types/src/run_event/mod.rs @@ -2214,10 +2214,10 @@ mod tests { visit: 1, message: None, context_window: None, - reasoning: Some(crate::ReasoningOutput { - summary: Some("inspect the implementation first".to_string()), - trace: Some("read convert.rs, then the sink".to_string()), - }), + reasoning: Some(crate::ReasoningOutput::new( + "inspect the implementation first", + "read convert.rs, then the sink", + )), }); let value = serde_json::to_value(&body).unwrap(); From 512ab50f9c6001171766ee5d1f097840249bf35a Mon Sep 17 00:00:00 2001 From: Release Repro Date: Fri, 24 Jul 2026 17:42:44 -0400 Subject: [PATCH 3/3] fix(reasoning): align stream and client invariants --- docs/public/api-reference/fabro-api.yaml | 24 ++++-- .../src/codec/openai_compatible/wire.rs | 83 +++++++++++++++++-- lib/components/fabro-llm/src/reasoning.rs | 4 +- .../tests/it/wire/openai_compatible.rs | 17 ++++ .../src/.openapi-generator/FILES | 2 + .../fabro-api-client/src/models/index.ts | 2 + .../src/models/reasoning-output-trace-only.ts | 22 +++++ .../models/reasoning-output-with-summary.ts | 26 ++++++ .../src/models/reasoning-output.ts | 18 ++-- .../tests/reasoning-output-invariant.ts | 20 +++++ 10 files changed, 195 insertions(+), 23 deletions(-) create mode 100644 lib/packages/fabro-api-client/src/models/reasoning-output-trace-only.ts create mode 100644 lib/packages/fabro-api-client/src/models/reasoning-output-with-summary.ts create mode 100644 lib/packages/fabro-api-client/tests/reasoning-output-invariant.ts diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 6c38d71f7..8b011aef7 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -10058,12 +10058,14 @@ components: present whenever the object is emitted. Opaque provider material (signatures, item IDs, encrypted or redacted payloads) never appears here. + oneOf: + - $ref: "#/components/schemas/ReasoningOutputWithSummary" + - $ref: "#/components/schemas/ReasoningOutputTraceOnly" + + ReasoningOutputWithSummary: type: object - anyOf: - - required: - - summary - - required: - - trace + required: + - summary properties: summary: type: string @@ -10072,6 +10074,18 @@ components: type: string description: Verbatim readable reasoning text, when the provider returns it. + ReasoningOutputTraceOnly: + type: object + required: + - trace + not: + required: + - summary + properties: + trace: + type: string + description: Verbatim readable reasoning text, when the provider returns it. + AgentToolsAvailableProps: description: Properties for the `agent.tools.available` event. type: object diff --git a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs index 25668a023..1c4f70b27 100644 --- a/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs +++ b/lib/components/fabro-llm/src/codec/openai_compatible/wire.rs @@ -176,8 +176,9 @@ impl ReasoningDetails { /// Absorb one streamed `reasoning_details` payload. /// /// Fragments carrying the same `type` and `index` are coalesced even when - /// other logical details appear between them. First-seen detail order is - /// retained. + /// other logical details appear between them. Without an index, a fragment + /// continues the most recently seen detail of the same type. First-seen + /// detail order is retained. pub(super) fn push_stream_payload(&mut self, payload: serde_json::Value) { let incoming = match payload { serde_json::Value::Array(entries) => entries, @@ -191,6 +192,7 @@ impl ReasoningDetails { match self .entries .iter_mut() + .rev() .find(|existing| continues_detail(existing, &entry)) { Some(existing) => merge_detail_fragment(existing, entry), @@ -217,10 +219,23 @@ const DETAIL_TEXT_MEMBERS: [&str; 3] = ["text", "summary", "data"]; /// Aggregators tag each logical detail with a stable `type` and `index`; /// fragment streams that omit `index` are matched on `type` alone. fn continues_detail(last: &serde_json::Value, entry: &serde_json::Value) -> bool { - let Some(entry_type) = entry.get("type") else { + let (Some(last_type), Some(entry_type)) = ( + last.get("type").and_then(serde_json::Value::as_str), + entry.get("type").and_then(serde_json::Value::as_str), + ) else { return false; }; - last.get("type") == Some(entry_type) && last.get("index") == entry.get("index") + if last_type != entry_type { + return false; + } + + match ( + last.get("index").and_then(serde_json::Value::as_u64), + entry.get("index").and_then(serde_json::Value::as_u64), + ) { + (Some(last_index), Some(entry_index)) => last_index == entry_index, + _ => true, + } } /// Append `entry`'s text fragments onto `last` and fill in members `last` @@ -393,9 +408,65 @@ pub(super) struct AccumulatedToolCall { #[cfg(test)] mod tests { - use super::{ApiResponse, ApiUsage, ChatContent, ChatTextPart, StreamChunk}; + use super::{ + ApiResponse, ApiUsage, ChatContent, ChatTextPart, ReasoningDetails, StreamChunk, + continues_detail, + }; use crate::codec::cache::CacheControl; - use crate::types::TokenCounts; + use crate::types::{ContentPart, TokenCounts}; + + #[test] + fn reasoning_detail_continuation_uses_type_when_either_index_is_missing() { + let indexed = serde_json::json!({"type": "reasoning.text", "index": 0}); + let unindexed = serde_json::json!({"type": "reasoning.text"}); + + assert!(continues_detail(&indexed, &unindexed)); + assert!(continues_detail(&unindexed, &indexed)); + assert!(continues_detail(&unindexed, &unindexed)); + } + + #[test] + fn reasoning_detail_continuation_requires_a_matching_string_type() { + let detail = serde_json::json!({"type": "reasoning.text", "index": 0}); + + assert!(!continues_detail( + &detail, + &serde_json::json!({"type": "reasoning.summary", "index": 0}) + )); + assert!(!continues_detail( + &serde_json::json!({"index": 0}), + &serde_json::json!({"index": 0}) + )); + assert!(!continues_detail( + &serde_json::json!({"type": 7, "index": 0}), + &serde_json::json!({"type": 7, "index": 0}) + )); + } + + #[test] + fn unindexed_reasoning_fragment_continues_the_latest_matching_type() { + let mut details = ReasoningDetails::default(); + details.push_stream_payload(serde_json::json!([ + {"type": "reasoning.text", "text": "first", "index": 0}, + {"type": "reasoning.text", "text": "second", "index": 1}, + ])); + details.push_stream_payload(serde_json::json!([ + {"type": "reasoning.text", "text": " continued"}, + ])); + + let ContentPart::Other { data, .. } = + details.into_content_part().expect("reasoning detail part") + else { + panic!("expected opaque reasoning detail part"); + }; + assert_eq!( + data, + serde_json::json!([ + {"type": "reasoning.text", "text": "first", "index": 0}, + {"type": "reasoning.text", "text": "second continued", "index": 1}, + ]) + ); + } #[test] fn chat_content_text_serializes_as_plain_string() { diff --git a/lib/components/fabro-llm/src/reasoning.rs b/lib/components/fabro-llm/src/reasoning.rs index cead5d180..6c12ee762 100644 --- a/lib/components/fabro-llm/src/reasoning.rs +++ b/lib/components/fabro-llm/src/reasoning.rs @@ -47,8 +47,8 @@ impl Blocks<'_> { } } -/// Join complete blocks in provider order, dropping empty and -/// whitespace-only fragments. Retained text is never trimmed or rewritten. +/// Join retained complete blocks in provider order. Text is never trimmed or +/// rewritten. fn join_blocks(blocks: &[&str]) -> Option { (!blocks.is_empty()).then(|| blocks.join(BLOCK_SEPARATOR)) } diff --git a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs index a35ba6427..79d61223c 100644 --- a/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs +++ b/lib/components/fabro-llm/tests/it/wire/openai_compatible.rs @@ -720,6 +720,23 @@ async fn stream_reasoning_details_normalize_like_the_non_streaming_body() { assert_eq!(reasoning.trace(), Some("2 plus 2 is 4")); } +/// Providers may omit the optional index after the first fragment; the type +/// still identifies the logical detail being continued. +#[tokio::test] +async fn stream_reasoning_details_coalesce_when_a_later_fragment_omits_index() { + let sse = support::sse_data_transcript(&[ + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"role":"assistant","reasoning_details":[{"type":"reasoning.text","text":"first ","index":0}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"reasoning.text","text":"second"}]},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{"content":"done"},"finish_reason":null}]}"#, + r#"{"id":"chatcmpl_stream","object":"chat.completion.chunk","created":1700000000,"model":"test-model","choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}"#, + "[DONE]", + ]); + + let response = stream_final_response(&sse).await; + let reasoning = response.reasoning_output().expect("reasoning present"); + assert_eq!(reasoning.trace(), Some("first second")); +} + /// Cached and reasoning detail tokens are split into their own disjoint /// buckets and subtracted out of input/output. #[tokio::test] diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 6d7c3803e..02add4f1d 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -321,6 +321,8 @@ models/pull-request.ts models/question-type.ts models/reasoning-effort-feature.ts models/reasoning-effort.ts +models/reasoning-output-trace-only.ts +models/reasoning-output-with-summary.ts models/reasoning-output.ts models/related-workflow-diagnostic.ts models/render-workflow-graph-direction.ts diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 5d9d663ef..474bf0b2b 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -292,6 +292,8 @@ export * from './question-type'; export * from './reasoning-effort'; export * from './reasoning-effort-feature'; export * from './reasoning-output'; +export * from './reasoning-output-trace-only'; +export * from './reasoning-output-with-summary'; export * from './related-workflow-diagnostic'; export * from './render-workflow-graph-direction'; export * from './render-workflow-graph-format'; diff --git a/lib/packages/fabro-api-client/src/models/reasoning-output-trace-only.ts b/lib/packages/fabro-api-client/src/models/reasoning-output-trace-only.ts new file mode 100644 index 000000000..8f757fddf --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/reasoning-output-trace-only.ts @@ -0,0 +1,22 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +export interface ReasoningOutputTraceOnly { + /** + * Verbatim readable reasoning text, when the provider returns it. + */ + 'trace': string; +} diff --git a/lib/packages/fabro-api-client/src/models/reasoning-output-with-summary.ts b/lib/packages/fabro-api-client/src/models/reasoning-output-with-summary.ts new file mode 100644 index 000000000..c93d38405 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/reasoning-output-with-summary.ts @@ -0,0 +1,26 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + + +export interface ReasoningOutputWithSummary { + /** + * Model-authored summary of its reasoning. + */ + 'summary': string; + /** + * Verbatim readable reasoning text, when the provider returns it. + */ + 'trace'?: string; +} diff --git a/lib/packages/fabro-api-client/src/models/reasoning-output.ts b/lib/packages/fabro-api-client/src/models/reasoning-output.ts index ab970e292..b2c60fb2c 100644 --- a/lib/packages/fabro-api-client/src/models/reasoning-output.ts +++ b/lib/packages/fabro-api-client/src/models/reasoning-output.ts @@ -13,17 +13,15 @@ */ +// May contain unused imports in some cases +// @ts-ignore +import type { ReasoningOutputTraceOnly } from './reasoning-output-trace-only'; +// May contain unused imports in some cases +// @ts-ignore +import type { ReasoningOutputWithSummary } from './reasoning-output-with-summary'; /** + * @type ReasoningOutput * Readable model reasoning normalized into a provider-neutral shape. Both members may be present for the same response, and at least one is present whenever the object is emitted. Opaque provider material (signatures, item IDs, encrypted or redacted payloads) never appears here. */ -export interface ReasoningOutput { - /** - * Model-authored summary of its reasoning. - */ - 'summary'?: string; - /** - * Verbatim readable reasoning text, when the provider returns it. - */ - 'trace'?: string; -} +export type ReasoningOutput = ReasoningOutputTraceOnly | ReasoningOutputWithSummary; diff --git a/lib/packages/fabro-api-client/tests/reasoning-output-invariant.ts b/lib/packages/fabro-api-client/tests/reasoning-output-invariant.ts new file mode 100644 index 000000000..1a5f165ac --- /dev/null +++ b/lib/packages/fabro-api-client/tests/reasoning-output-invariant.ts @@ -0,0 +1,20 @@ +import type { ReasoningOutput } from "../src"; + +const summaryOnly: ReasoningOutput = { + summary: "short summary", +}; +const traceOnly: ReasoningOutput = { + trace: "full trace", +}; +const summaryAndTrace: ReasoningOutput = { + summary: "short summary", + trace: "full trace", +}; + +// @ts-expect-error ReasoningOutput requires at least one readable member. +const empty: ReasoningOutput = {}; + +void summaryOnly; +void traceOnly; +void summaryAndTrace; +void empty;