mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-06 02:48:25 +00:00
Merge pull request #626 from fabro-sh/feat/passive-reasoning-capture
feat(reasoning): passive reasoning capture in agent.message
This commit is contained in:
commit
ae3b6702e2
35 changed files with 1763 additions and 46 deletions
|
|
@ -10045,6 +10045,46 @@ 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.
|
||||
oneOf:
|
||||
- $ref: "#/components/schemas/ReasoningOutputWithSummary"
|
||||
- $ref: "#/components/schemas/ReasoningOutputTraceOnly"
|
||||
|
||||
ReasoningOutputWithSummary:
|
||||
type: object
|
||||
required:
|
||||
- summary
|
||||
properties:
|
||||
summary:
|
||||
type: string
|
||||
description: Model-authored summary of its reasoning.
|
||||
trace:
|
||||
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.
|
||||
|
|
|
|||
|
|
@ -534,6 +534,7 @@ mod tests {
|
|||
cost_source: None,
|
||||
tool_call_count: 0,
|
||||
context_window: None,
|
||||
reasoning: None,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}),
|
||||
),
|
||||
]
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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::new(
|
||||
"inspect the sink first",
|
||||
"read events.rs, then attach",
|
||||
)),
|
||||
},
|
||||
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::<serde_json::Value>(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();
|
||||
|
|
|
|||
|
|
@ -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,111 @@ 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<SessionEvent>,
|
||||
) -> Vec<Option<ReasoningOutput>> {
|
||||
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::new(
|
||||
"the user wants 2+2",
|
||||
"2+2 is 4",
|
||||
))]);
|
||||
}
|
||||
|
||||
#[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::from_summary("call the tool")),
|
||||
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::new(
|
||||
"final summary",
|
||||
"final trace"
|
||||
))]);
|
||||
}
|
||||
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn stream_quota_error_does_not_replay() {
|
||||
let quota_error = LlmError::Provider {
|
||||
|
|
|
|||
|
|
@ -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<StageContextWindowProjection>,
|
||||
/// 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<ReasoningOutput>,
|
||||
},
|
||||
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 {
|
||||
|
|
|
|||
|
|
@ -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::{
|
||||
|
|
@ -13,18 +13,24 @@ pub(super) fn decode_response(
|
|||
ctx: &CodecCtx<'_>,
|
||||
rate_limit: Option<RateLimitInfo>,
|
||||
) -> Result<Response, Error> {
|
||||
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.take() {
|
||||
content_parts.extend(ReasoningDetails::from_complete_payload(payload).into_content_part());
|
||||
}
|
||||
if let Some(reasoning) = choice.message.reasoning() {
|
||||
if !reasoning.is_empty() {
|
||||
content_parts.push(ContentPart::Thinking(ThinkingData {
|
||||
|
|
|
|||
|
|
@ -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<AccumulatedToolCall>,
|
||||
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,
|
||||
|
|
@ -54,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<Vec<StreamEvent>> {
|
||||
fn process_chunk(&mut self, mut chunk: StreamChunk) -> Option<Vec<StreamEvent>> {
|
||||
// Capture response metadata from the first chunk.
|
||||
if let Some(id) = &chunk.id {
|
||||
if self.response_id.is_empty() {
|
||||
|
|
@ -74,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();
|
||||
|
||||
|
|
@ -84,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() {
|
||||
|
|
@ -93,6 +95,11 @@ impl StreamState {
|
|||
}
|
||||
}
|
||||
|
||||
// Accumulate structured reasoning detail fragments in wire order.
|
||||
if let Some(payload) = delta.reasoning_details.take() {
|
||||
self.reasoning_details.push_stream_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 {
|
||||
|
|
@ -248,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<StreamEvent> {
|
||||
|
|
@ -359,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 { .. }));
|
||||
|
|
@ -367,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 { .. }));
|
||||
|
||||
|
|
@ -381,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 { .. }));
|
||||
|
||||
|
|
|
|||
|
|
@ -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<String>,
|
||||
/// OpenRouter's normalized spelling for reasoning text.
|
||||
pub reasoning: Option<String>,
|
||||
/// 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<serde_json::Value>,
|
||||
pub tool_calls: Option<Vec<ApiToolCall>>,
|
||||
}
|
||||
|
||||
|
|
@ -140,6 +145,125 @@ 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<serde_json::Value>,
|
||||
}
|
||||
|
||||
impl ReasoningDetails {
|
||||
/// Preserve a complete-response `reasoning_details` payload.
|
||||
///
|
||||
/// Providers document an array of detail objects; a lone object is
|
||||
/// 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. 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,
|
||||
payload @ serde_json::Value::Object(_) => vec![payload],
|
||||
_ => Vec::new(),
|
||||
};
|
||||
for entry in incoming {
|
||||
if !entry.is_object() {
|
||||
continue;
|
||||
}
|
||||
match self
|
||||
.entries
|
||||
.iter_mut()
|
||||
.rev()
|
||||
.find(|existing| continues_detail(existing, &entry))
|
||||
{
|
||||
Some(existing) => merge_detail_fragment(existing, 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<ContentPart> {
|
||||
(!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 {
|
||||
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;
|
||||
};
|
||||
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`
|
||||
/// has not seen yet.
|
||||
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) {
|
||||
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, value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(serde::Deserialize)]
|
||||
pub(super) struct ApiToolCall {
|
||||
pub id: String,
|
||||
|
|
@ -245,6 +369,10 @@ pub(super) struct StreamDelta {
|
|||
pub reasoning_content: Option<String>,
|
||||
/// OpenRouter's normalized spelling for reasoning text.
|
||||
pub reasoning: Option<String>,
|
||||
/// Structured reasoning channel, streamed as fragments of the entries
|
||||
/// the non-streaming response returns whole.
|
||||
#[serde(default)]
|
||||
pub reasoning_details: Option<serde_json::Value>,
|
||||
pub tool_calls: Option<Vec<StreamToolCall>>,
|
||||
}
|
||||
|
||||
|
|
@ -280,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() {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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(), Some("weighed both"));
|
||||
assert_eq!(reasoning.trace(), Some("step one"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn complete_returns_error_on_502() {
|
||||
let server = MockServer::start();
|
||||
|
|
|
|||
342
lib/components/fabro-llm/src/reasoning.rs
Normal file
342
lib/components/fabro-llm/src/reasoning.rs
Normal file
|
|
@ -0,0 +1,342 @@
|
|||
//! 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";
|
||||
|
||||
/// 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 trace
|
||||
/// that no explicit trace produced.
|
||||
#[derive(Default)]
|
||||
struct Blocks<'a> {
|
||||
explicit_summary: Vec<&'a str>,
|
||||
explicit_trace: Vec<&'a str>,
|
||||
fallback_trace: Vec<&'a str>,
|
||||
}
|
||||
|
||||
impl Blocks<'_> {
|
||||
fn into_output(self) -> Option<ReasoningOutput> {
|
||||
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 retained complete blocks in provider order. Text is never trimmed or
|
||||
/// rewritten.
|
||||
fn join_blocks(blocks: &[&str]) -> Option<String> {
|
||||
(!blocks.is_empty()).then(|| blocks.join(BLOCK_SEPARATOR))
|
||||
}
|
||||
|
||||
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.
|
||||
///
|
||||
/// `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<'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() {
|
||||
push_block(&mut blocks.explicit_summary, text);
|
||||
} else if let Some(text) = entry.get("text").and_then(serde_json::Value::as_str) {
|
||||
push_block(&mut blocks.explicit_summary, text);
|
||||
}
|
||||
}
|
||||
}
|
||||
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" {
|
||||
push_block(&mut blocks.explicit_trace, text);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Extract readable text from OpenAI-compatible `reasoning_details` entries.
|
||||
fn collect_reasoning_details<'a>(details: &'a serde_json::Value, blocks: &mut Blocks<'a>) {
|
||||
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();
|
||||
match detail_type {
|
||||
"reasoning.text" => {
|
||||
if let Some(text) = readable_member(entry, "text") {
|
||||
push_block(&mut blocks.explicit_trace, text);
|
||||
}
|
||||
}
|
||||
"reasoning.summary" => {
|
||||
if let Some(text) = readable_member(entry, "summary") {
|
||||
push_block(&mut blocks.explicit_summary, 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<ReasoningOutput> {
|
||||
let mut blocks = Blocks::default();
|
||||
for part in content {
|
||||
match part {
|
||||
ContentPart::Thinking(thinking) if !thinking.redacted => {
|
||||
push_block(&mut blocks.fallback_trace, &thinking.text);
|
||||
}
|
||||
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_trace() {
|
||||
let output = normalize(&[thinking("weighing the options")]).unwrap();
|
||||
assert!(output.summary().is_none());
|
||||
assert_eq!(output.trace(), Some("weighing the options"));
|
||||
}
|
||||
|
||||
#[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(), Some("inspect first"));
|
||||
assert_eq!(output.trace(), 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(), 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]
|
||||
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(), Some("checked the parser"));
|
||||
assert_eq!(output.trace(), 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(), 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_remain_opaque() {
|
||||
assert!(
|
||||
normalize(&[reasoning_details(json!([
|
||||
{"type": "reasoning.future", "text": "new channel"},
|
||||
]))])
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[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(), Some("checked the parser"));
|
||||
assert!(output.trace().is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn structured_trace_takes_precedence_over_flattened_trace() {
|
||||
let output = normalize(&[
|
||||
reasoning_details(json!([{"type": "reasoning.text", "text": "verbatim"}])),
|
||||
thinking("flattened"),
|
||||
])
|
||||
.unwrap();
|
||||
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]
|
||||
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.trace(), 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());
|
||||
}
|
||||
}
|
||||
|
|
@ -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<ReasoningOutput> {
|
||||
reasoning::normalize(&self.message.content)
|
||||
}
|
||||
}
|
||||
|
||||
// --- 3.13 StreamEvent ---
|
||||
|
|
|
|||
|
|
@ -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,6 +443,28 @@ async fn decode_response(body: serde_json::Value) -> fabro_llm::types::Response
|
|||
response
|
||||
}
|
||||
|
||||
/// 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 accumulator = StreamAccumulator::new();
|
||||
while let Some(item) = stream.next().await {
|
||||
accumulator.process(&item.expect("stream event should decode"));
|
||||
}
|
||||
mock.assert();
|
||||
accumulator
|
||||
.response()
|
||||
.cloned()
|
||||
.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 +513,230 @@ 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(), 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
|
||||
/// 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(), 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
|
||||
/// 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.");
|
||||
assert!(response.reasoning_output().is_none());
|
||||
}
|
||||
|
||||
/// 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(), Some("the user wants 2+2"));
|
||||
assert!(reasoning.trace().is_none());
|
||||
}
|
||||
|
||||
/// A structured trace takes precedence over the flattened trace channel.
|
||||
#[tokio::test]
|
||||
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 trace",
|
||||
"reasoning_details": [{"type": "reasoning.text", "text": "verbatim trace", "index": 0}]
|
||||
})))
|
||||
.await;
|
||||
|
||||
let reasoning = response.reasoning_output().expect("reasoning present");
|
||||
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
|
||||
/// 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},{"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]",
|
||||
]);
|
||||
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());
|
||||
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"));
|
||||
}
|
||||
|
||||
/// 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]
|
||||
|
|
|
|||
|
|
@ -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(), 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"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn decode_incomplete_status_maps_to_length() {
|
||||
let response = decode_response(serde_json::json!({
|
||||
|
|
|
|||
|
|
@ -4032,6 +4032,7 @@ mod tests {
|
|||
visit: 1,
|
||||
message: None,
|
||||
context_window: None,
|
||||
reasoning: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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,52 @@ 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::new(
|
||||
"inspect the conversion first",
|
||||
"read convert.rs, then the sink",
|
||||
)),
|
||||
},
|
||||
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(), 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");
|
||||
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 {
|
||||
|
|
|
|||
|
|
@ -30,7 +30,10 @@ pub fn event_payload_from_redacted_json(line: &str, run_id: &RunId) -> Result<Ev
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use ::fabro_types::{fixtures, run_event as fabro_types};
|
||||
use ::fabro_types::{ReasoningOutput, fixtures, run_event as fabro_types};
|
||||
use fabro_agent::AgentEvent;
|
||||
use fabro_llm::types::TokenCounts as LlmTokenCounts;
|
||||
use fabro_model::{ModelRef, ProviderId};
|
||||
|
||||
use super::*;
|
||||
use crate::event::{Event, to_run_event};
|
||||
|
|
@ -72,4 +75,45 @@ mod tests {
|
|||
"plain stderr"
|
||||
);
|
||||
}
|
||||
|
||||
/// Reasoning is model-authored text like any other, so it goes through
|
||||
/// the same canonical redaction pass as assistant output.
|
||||
#[test]
|
||||
fn build_redacted_event_payload_redacts_secrets_inside_reasoning() {
|
||||
let secret = "sk-ant-api03-xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6pA";
|
||||
let stored = to_run_event(&fixtures::RUN_8, &Event::Agent {
|
||||
stage: "code".to_string(),
|
||||
visit: 1,
|
||||
event: AgentEvent::AssistantMessage {
|
||||
text: "done".to_string(),
|
||||
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: 0,
|
||||
context_window: None,
|
||||
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,
|
||||
tool_call_id: None,
|
||||
});
|
||||
|
||||
let payload = build_redacted_event_payload(&stored, &fixtures::RUN_8).unwrap();
|
||||
let reasoning = &payload.as_value()["properties"]["reasoning"];
|
||||
let summary = reasoning["summary"].as_str().unwrap();
|
||||
let trace = reasoning["trace"].as_str().unwrap();
|
||||
|
||||
assert!(!summary.contains(secret));
|
||||
assert!(!trace.contains(secret));
|
||||
assert!(summary.contains("REDACTED"));
|
||||
assert!(trace.contains("REDACTED"));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -309,6 +309,55 @@ mod tests {
|
|||
assert_eq!(payload.as_value()["properties"]["action"], "pause");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_event_sink_json_lines_carries_agent_message_reasoning() {
|
||||
use tokio::io::{AsyncBufReadExt, BufReader};
|
||||
|
||||
let (writer, reader) = tokio::io::duplex(4096);
|
||||
let sink = RunEventSink::json_lines(writer);
|
||||
let event = to_run_event(&fixtures::RUN_7, &Event::Agent {
|
||||
stage: "code".to_string(),
|
||||
visit: 1,
|
||||
event: fabro_agent::AgentEvent::AssistantMessage {
|
||||
text: String::new(),
|
||||
model: fabro_model::ModelRef {
|
||||
provider: fabro_model::ProviderId::openai(),
|
||||
model_id: "gpt-5.4".into(),
|
||||
speed: None,
|
||||
},
|
||||
usage: fabro_llm::types::TokenCounts::default(),
|
||||
cost_usd: None,
|
||||
cost_source: None,
|
||||
tool_call_count: 1,
|
||||
context_window: None,
|
||||
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,
|
||||
tool_call_id: None,
|
||||
});
|
||||
|
||||
sink.write_run_event(&event).await.unwrap();
|
||||
|
||||
let mut reader = BufReader::new(reader);
|
||||
let mut line = String::new();
|
||||
reader.read_line(&mut line).await.unwrap();
|
||||
|
||||
let payload = event_payload_from_redacted_json(line.trim_end(), &fixtures::RUN_7).unwrap();
|
||||
assert_eq!(payload.as_value()["event"], "agent.message");
|
||||
assert_eq!(
|
||||
payload.as_value()["properties"]["reasoning"]["summary"],
|
||||
"inspect the sink first"
|
||||
);
|
||||
assert_eq!(
|
||||
payload.as_value()["properties"]["reasoning"]["trace"],
|
||||
"write the line, then read it back"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn run_event_sink_map_applies_transform_before_fanout() {
|
||||
let first = Arc::new(AsyncMutex::new(Vec::new()));
|
||||
|
|
|
|||
|
|
@ -680,6 +680,7 @@ fn main() {
|
|||
("SessionRecord", "fabro_types::SessionRecord", &[]),
|
||||
("SessionSummary", "fabro_types::SessionSummary", &[]),
|
||||
("SessionDetail", "fabro_types::SessionDetail", &[]),
|
||||
("ReasoningOutput", "fabro_types::ReasoningOutput", &[]),
|
||||
("CompletionMessage", "fabro_types::Message", &[]),
|
||||
("CompletionMessageRole", "fabro_types::Role", &[]),
|
||||
("CompletionContentPart", "fabro_types::ContentPart", &[]),
|
||||
|
|
|
|||
|
|
@ -55,23 +55,23 @@ pub mod types {
|
|||
PairTranscriptResponse, ParallelBranchResult, PendingInterviewRecord, PermissionLevel,
|
||||
PreRunPushOutcome, Principal, PullRequest, PullRequestDetails, PullRequestDetailsStatus,
|
||||
PullRequestDetailsUnavailableReason, PullRequestLink, PullRequestMeta, PullRequestResponse,
|
||||
QuestionType, RepositoryRef, Role, Run, RunApproval, RunApprovalState, RunClientProvenance,
|
||||
RunEvent, RunEventDetailContentKind, RunEventDetailResponse, RunFailure,
|
||||
RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource, RunSandbox,
|
||||
RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan, RunSandboxRuntime,
|
||||
RunServerProvenance, RunSize, SandboxDetails, SandboxInfo, SandboxListMeta,
|
||||
SandboxListResponse, SandboxNetwork, SandboxNetworkPolicy, SandboxNetworkPolicyMode,
|
||||
SandboxProviderKind, SandboxProviderLookupError, SandboxResources, SandboxService,
|
||||
SandboxServiceListResponse, SandboxState, SandboxTimestamps, SecretMetadata, SecretType,
|
||||
ServerSettings, SessionDetail, SessionId, SessionMessage, SessionRecord, SessionStatus,
|
||||
SessionSummary, SessionTurn, SkillsProjection, StageCompletion, StageContextWindow,
|
||||
StageContextWindowBreakdownItem, StageContextWindowCategory, StageContextWindowCountMethod,
|
||||
StageContextWindowProjection, StageContextWindowStaleness,
|
||||
StageContextWindowUnavailableReason, StageContextWindowWarning, StageHandler, StageId,
|
||||
StageModelUsage, StageOutcome, StageProjection, StageState, SubAgentProjection,
|
||||
SubAgentStatus, SystemActorKind, SystemIntegrationStatus, SystemIntegrationsResponse,
|
||||
TodoListProjection, TurnId, UpdateVariableRequest, UserPrincipal, Variable,
|
||||
VariableListResponse, WorkflowSettings,
|
||||
QuestionType, ReasoningOutput, RepositoryRef, Role, Run, RunApproval, RunApprovalState,
|
||||
RunClientProvenance, RunEvent, RunEventDetailContentKind, RunEventDetailResponse,
|
||||
RunFailure, RunPairStatusResponse, RunProjection, RunProvenance, RunRunnableSource,
|
||||
RunSandbox, RunSandboxFailure, RunSandboxInstance, RunSandboxKind, RunSandboxPlan,
|
||||
RunSandboxRuntime, RunServerProvenance, RunSize, SandboxDetails, SandboxInfo,
|
||||
SandboxListMeta, SandboxListResponse, SandboxNetwork, SandboxNetworkPolicy,
|
||||
SandboxNetworkPolicyMode, SandboxProviderKind, SandboxProviderLookupError,
|
||||
SandboxResources, SandboxService, SandboxServiceListResponse, SandboxState,
|
||||
SandboxTimestamps, SecretMetadata, SecretType, ServerSettings, SessionDetail, SessionId,
|
||||
SessionMessage, SessionRecord, SessionStatus, SessionSummary, SessionTurn,
|
||||
SkillsProjection, StageCompletion, StageContextWindow, StageContextWindowBreakdownItem,
|
||||
StageContextWindowCategory, StageContextWindowCountMethod, StageContextWindowProjection,
|
||||
StageContextWindowStaleness, StageContextWindowUnavailableReason,
|
||||
StageContextWindowWarning, StageHandler, StageId, StageModelUsage, StageOutcome,
|
||||
StageProjection, StageState, SubAgentProjection, SubAgentStatus, SystemActorKind,
|
||||
SystemIntegrationStatus, SystemIntegrationsResponse, TodoListProjection, TurnId,
|
||||
UpdateVariableRequest, UserPrincipal, Variable, VariableListResponse, WorkflowSettings,
|
||||
};
|
||||
|
||||
pub use crate::generated::types::*;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,94 @@
|
|||
use std::any::{TypeId, type_name};
|
||||
|
||||
use fabro_api::types::{
|
||||
AgentMessageProps as ApiAgentMessageProps, ReasoningOutput as ApiReasoningOutput,
|
||||
};
|
||||
use fabro_types::ReasoningOutput;
|
||||
use serde_json::json;
|
||||
|
||||
#[test]
|
||||
fn reasoning_output_reuses_canonical_type() {
|
||||
assert_same_type::<ApiReasoningOutput, ReasoningOutput>();
|
||||
}
|
||||
|
||||
#[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(), 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();
|
||||
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 reasoning_output_rejects_an_empty_object() {
|
||||
let error = serde_json::from_value::<ApiReasoningOutput>(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!({
|
||||
"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(), 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"})
|
||||
);
|
||||
}
|
||||
|
||||
fn assert_same_type<T: 'static, U: 'static>() {
|
||||
assert_eq!(
|
||||
TypeId::of::<T>(),
|
||||
TypeId::of::<U>(),
|
||||
"{} should be the same type as {}",
|
||||
type_name::<T>(),
|
||||
type_name::<U>()
|
||||
);
|
||||
}
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
142
lib/foundation/fabro-types/src/reasoning.rs
Normal file
142
lib/foundation/fabro-types/src/reasoning.rs
Normal file
|
|
@ -0,0 +1,142 @@
|
|||
use serde::{Deserialize, Serialize, de};
|
||||
|
||||
/// 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, 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")]
|
||||
summary: Option<String>,
|
||||
/// 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")]
|
||||
trace: Option<String>,
|
||||
}
|
||||
|
||||
impl ReasoningOutput {
|
||||
/// Creates reasoning output with both a model-authored summary and a
|
||||
/// verbatim trace.
|
||||
#[must_use]
|
||||
pub fn new(summary: impl Into<String>, trace: impl Into<String>) -> 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<String>) -> 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<String>) -> 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<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
|
||||
#[derive(Deserialize)]
|
||||
struct Fields {
|
||||
#[serde(default)]
|
||||
summary: Option<String>,
|
||||
#[serde(default)]
|
||||
trace: Option<String>,
|
||||
}
|
||||
|
||||
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",
|
||||
)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use serde_json::json;
|
||||
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn summary_only_round_trips_without_trace_member() {
|
||||
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!(
|
||||
serde_json::from_value::<ReasoningOutput>(v).unwrap(),
|
||||
output
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn trace_only_round_trips_without_summary_member() {
|
||||
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!(
|
||||
serde_json::from_value::<ReasoningOutput>(v).unwrap(),
|
||||
output
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn both_fields_round_trip() {
|
||||
let output = ReasoningOutput::new("summary", "trace");
|
||||
let v = serde_json::to_value(&output).unwrap();
|
||||
assert_eq!(v, json!({"summary": "summary", "trace": "trace"}));
|
||||
assert_eq!(
|
||||
serde_json::from_value::<ReasoningOutput>(v).unwrap(),
|
||||
output
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_object_is_rejected() {
|
||||
let error = serde_json::from_value::<ReasoningOutput>(json!({})).unwrap_err();
|
||||
assert!(error.to_string().contains("requires a summary or trace"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn explicit_nulls_are_rejected() {
|
||||
let error =
|
||||
serde_json::from_value::<ReasoningOutput>(json!({"summary": null, "trace": null}))
|
||||
.unwrap_err();
|
||||
assert!(error.to_string().contains("requires a summary or trace"));
|
||||
}
|
||||
}
|
||||
|
|
@ -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<StageContextWindowProjection>,
|
||||
/// 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<ReasoningOutput>,
|
||||
}
|
||||
|
||||
#[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::new(
|
||||
"inspect the implementation first",
|
||||
"read convert.rs, then the sink",
|
||||
)),
|
||||
};
|
||||
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");
|
||||
|
|
|
|||
|
|
@ -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::new(
|
||||
"inspect the implementation first",
|
||||
"read convert.rs, then the sink",
|
||||
)),
|
||||
});
|
||||
|
||||
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();
|
||||
|
|
|
|||
|
|
@ -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<String>) -> Self {
|
||||
Self::Text(text.into())
|
||||
|
|
|
|||
|
|
@ -321,6 +321,9 @@ 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
|
||||
models/render-workflow-graph-format.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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -291,6 +291,9 @@ export * from './pull-request-user';
|
|||
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';
|
||||
|
|
|
|||
22
lib/packages/fabro-api-client/src/models/reasoning-output-trace-only.ts
generated
Normal file
22
lib/packages/fabro-api-client/src/models/reasoning-output-trace-only.ts
generated
Normal file
|
|
@ -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;
|
||||
}
|
||||
26
lib/packages/fabro-api-client/src/models/reasoning-output-with-summary.ts
generated
Normal file
26
lib/packages/fabro-api-client/src/models/reasoning-output-with-summary.ts
generated
Normal file
|
|
@ -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;
|
||||
}
|
||||
27
lib/packages/fabro-api-client/src/models/reasoning-output.ts
generated
Normal file
27
lib/packages/fabro-api-client/src/models/reasoning-output.ts
generated
Normal file
|
|
@ -0,0 +1,27 @@
|
|||
/* 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.
|
||||
*/
|
||||
|
||||
|
||||
// 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 type ReasoningOutput = ReasoningOutputTraceOnly | ReasoningOutputWithSummary;
|
||||
|
|
@ -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;
|
||||
Loading…
Add table
Reference in a new issue