From bb33b3c6dd7c99b4938824f75ae089a6302ef245 Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Tue, 16 Jun 2026 17:42:25 -0700 Subject: [PATCH 1/3] Surface tool calls and outputs in OpenAI Responses API MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The non-streaming Responses API (`POST /v1/responses` create-and-wait and `GET /v1/responses/{id}`) previously returned only the assistant's final text message. Tool activity in the run was invisible to clients. This projects each run's tool calls and their raw outputs into the response `output` array: a `function_call` item (call_id / name / arguments) paired with its `function_call_output` (the raw `model_observation` tool result), followed by the assistant `Message`. Implementation lives in the composition projection reader (`OpenAiResponsesThreadProjectionReader`). Neither thread-service read alone carries everything — the history projection keeps `turn_run_id` plus the tool-result envelope content but strips `tool_result_provider_call`, while the context projection preserves the provider call but drops run attribution — so the reader joins `list_thread_history` and `load_context_messages` by `message_id`. Tool items are always emitted before the assistant message: a run's single assistant draft reserves its sequence at turn start, which can predate the tool results it produces, so transcript order alone would mis-order them. The wait poll loop checks a cheap completion gate (`finalized_assistant_message_by_run`) and builds the full projection once, rather than re-reading the whole transcript every tick. Notes: - Only the non-streaming Responses path is covered; `stream: true` still emits text only (separate projection-streamer mechanism). - This deliberately exposes the raw `model_observation` and provider tool name/arguments through the API surface, a documented exception to the crate's narrow-DTO policy. --- .../src/openai_compat_serve.rs | 571 ++++++++++++++++-- 1 file changed, 515 insertions(+), 56 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs index b6a6082165b..01401c25b4e 100644 --- a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs +++ b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs @@ -4,6 +4,7 @@ //! authority-bearing wiring: authenticated callers, ProductWorkflow, //! conversation binding, durable idempotency/ref stores, and projection reads. +use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; @@ -11,7 +12,7 @@ use async_trait::async_trait; use ironclaw_filesystem::{RootFilesystem, ScopedFilesystem}; use ironclaw_host_api::{ AgentId, InvocationId, MountAlias, MountGrant, MountPermissions, MountView, ProjectId, - ResourceScope, TenantId, UserId, VirtualPath, + ResourceScope, TenantId, ThreadId, UserId, VirtualPath, }; use ironclaw_product_adapters::{ AdapterInstallationId, ProductAdapterId, ProductInboundAck, ProductOutboundEnvelope, @@ -29,7 +30,7 @@ use ironclaw_reborn_openai_compat::{ OpenAiChatCompletionProjectionRequest, OpenAiChatCompletionsWorkflow, OpenAiChatProjectionStreamRequest, OpenAiCompatErrorKind, OpenAiCompatHttpError, OpenAiCompatProjectionStreamer, OpenAiCompatRefStore, OpenAiCompatResourceBinding, - OpenAiCompatRouterState, OpenAiResponseObject, OpenAiResponseOutputItem, + OpenAiCompatRouterState, OpenAiResponseId, OpenAiResponseObject, OpenAiResponseOutputItem, OpenAiResponseOutputItemStatus, OpenAiResponseProjection, OpenAiResponseProjectionStreamRequest, OpenAiResponseReadRequest, OpenAiResponseStatus, OpenAiResponseWaitRequest, OpenAiResponsesMessageRole, OpenAiResponsesProjectionReader, @@ -37,7 +38,10 @@ use ironclaw_reborn_openai_compat::{ }; use ironclaw_reborn_openai_compat_storage::FilesystemOpenAiCompatRefStore; use ironclaw_threads::{ - FinalizedAssistantMessageByRunRequest, SessionThreadError, SessionThreadService, ThreadScope, + FinalizedAssistantMessageByRunRequest, LoadContextMessagesRequest, MessageKind, MessageStatus, + ProviderToolCallReferenceEnvelope, SessionThreadError, SessionThreadService, + ThreadHistoryRequest, ThreadMessageId, ThreadMessageRecord, ThreadScope, + ToolResultReferenceEnvelope, }; use crate::RebornBuildError; @@ -290,13 +294,130 @@ impl OpenAiResponsesThreadProjectionReader { } } - async fn read_finalized_response_message( + /// Project a single response run into OpenAI Responses `output` items: every + /// tool call paired with its raw output, then the run's finalized assistant + /// message. Tool items always precede the assistant message regardless of + /// transcript sequence — the assistant draft's sequence is reserved when the + /// turn starts, which can predate the tool results it later produces. + /// + /// The data is joined from two thread-service reads because neither one + /// alone carries everything: the history projection keeps `turn_run_id` + /// (run attribution) plus the tool-result envelope content (the raw + /// `model_observation`) but strips `tool_result_provider_call`, while the + /// context projection preserves the provider call (function name, arguments, + /// call id) but drops `turn_run_id`. Joining by `message_id` recovers both. + async fn read_run_output( &self, request: &ProjectionReadRequest, turn_run_id: String, - ) -> Result, OpenAiCompatHttpError> { + public_id: &OpenAiResponseId, + ) -> Result { let thread_scope = thread_scope_from_projection_read(request)?; - match self + let thread_id = request.scope.thread_id.clone(); + + let history = self + .thread_service + .list_thread_history(ThreadHistoryRequest { + scope: thread_scope.clone(), + thread_id: thread_id.clone(), + }) + .await + .map_err(map_thread_read_error)?; + + let run_messages: Vec<&ThreadMessageRecord> = history + .messages + .iter() + .filter(|message| message.turn_run_id.as_deref() == Some(turn_run_id.as_str())) + .collect(); + + // Tool calls/outputs, in transcript order. + let mut tool_results: Vec<&ThreadMessageRecord> = run_messages + .iter() + .copied() + .filter(|message| message.kind == MessageKind::ToolResultReference) + .collect(); + tool_results.sort_by_key(|message| message.sequence); + // The run's single finalized assistant reply (drafts dedup by run). + let assistant = run_messages + .iter() + .copied() + .filter(|message| { + message.kind == MessageKind::Assistant && message.status == MessageStatus::Finalized + }) + .max_by_key(|message| message.sequence); + + let provider_calls = self + .load_provider_calls( + &thread_scope, + &thread_id, + tool_results + .iter() + .map(|message| message.message_id) + .collect(), + ) + .await?; + + let mut items = Vec::with_capacity(tool_results.len() * 2 + 1); + for message in tool_results { + push_tool_output_items(&mut items, message, &provider_calls); + } + if let Some(message) = assistant { + items.push(OpenAiResponseOutputItem::Message { + id: format!("msg_{}", public_id.as_str()), + status: Some(OpenAiResponseOutputItemStatus::Completed), + role: OpenAiResponsesMessageRole::Assistant, + content: serde_json::json!([{ + "type": "output_text", + "text": message.content.clone().unwrap_or_default(), + }]), + }); + } + + Ok(RunResponseProjection { + items, + assistant_finalized: assistant.is_some(), + }) + } + + /// Re-read the run's tool-result messages through the context projection to + /// recover `tool_result_provider_call`, which the history projection strips. + async fn load_provider_calls( + &self, + scope: &ThreadScope, + thread_id: &ThreadId, + message_ids: Vec, + ) -> Result, OpenAiCompatHttpError> + { + if message_ids.is_empty() { + return Ok(HashMap::new()); + } + let context = self + .thread_service + .load_context_messages(LoadContextMessagesRequest { + scope: scope.clone(), + thread_id: thread_id.clone(), + message_ids, + }) + .await + .map_err(map_thread_read_error)?; + Ok(context + .messages + .into_iter() + .filter_map(|message| Some((message.message_id?, message.tool_result_provider_call?))) + .collect()) + } + + /// Cheap completion gate for the wait poll loop: the run has produced its + /// final reply once a finalized assistant message exists for it. Kept + /// separate from [`read_run_output`] so polling does not repeatedly read the + /// full transcript while waiting. + async fn run_completed( + &self, + request: &ProjectionReadRequest, + turn_run_id: String, + ) -> Result { + let thread_scope = thread_scope_from_projection_read(request)?; + let message = self .thread_service .finalized_assistant_message_by_run(FinalizedAssistantMessageByRunRequest { scope: thread_scope, @@ -304,27 +425,94 @@ impl OpenAiResponsesThreadProjectionReader { turn_run_id, }) .await - { - Ok(message) => Ok(message.map(|message| message.content.unwrap_or_default())), - Err( - SessionThreadError::UnknownThread { .. } - | SessionThreadError::ThreadScopeMismatch { .. }, - ) => Err(OpenAiCompatHttpError::not_found(Some( - "response_id".to_string(), - ))), - Err(error) => { - tracing::warn!( - target = "ironclaw::reborn::openai_compat", - error = %error, - "failed to read finalized assistant message for OpenAI-compatible response" - ); - Err(OpenAiCompatHttpError::from_kind( - 503, - true, - OpenAiCompatErrorKind::ServiceUnavailable, - None, - )) - } + .map_err(map_thread_read_error)?; + Ok(message.is_some()) + } +} + +/// A response run's projected `output` items plus whether the run's final +/// assistant message has been finalized (i.e. the run is complete). +struct RunResponseProjection { + items: Vec, + assistant_finalized: bool, +} + +/// Append a tool call and its raw output to `items`. When the provider-call side +/// channel is available, both a `function_call` (name/arguments) and the paired +/// `function_call_output` are emitted; otherwise only the output is emitted, +/// keyed by the opaque tool-result ref so the call id still correlates. +fn push_tool_output_items( + items: &mut Vec, + message: &ThreadMessageRecord, + provider_calls: &HashMap, +) { + let output = tool_result_output(message); + match provider_calls.get(&message.message_id) { + Some(provider_call) => { + let call_id = provider_call.provider_call_id.clone(); + let arguments = serde_json::to_string(&provider_call.arguments) + .unwrap_or_else(|_| "{}".to_string()); + items.push(OpenAiResponseOutputItem::FunctionCall { + id: format!("fc_{call_id}"), + status: Some(OpenAiResponseOutputItemStatus::Completed), + call_id: call_id.clone(), + name: provider_call.provider_tool_name.clone(), + arguments, + }); + items.push(OpenAiResponseOutputItem::FunctionCallOutput { + id: format!("fco_{call_id}"), + status: Some(OpenAiResponseOutputItemStatus::Completed), + call_id, + output, + }); + } + None => { + let call_id = message + .tool_result_ref + .clone() + .unwrap_or_else(|| format!("call_{}", message.sequence)); + items.push(OpenAiResponseOutputItem::FunctionCallOutput { + id: format!("fco_{call_id}"), + status: Some(OpenAiResponseOutputItemStatus::Completed), + call_id, + output, + }); + } + } +} + +/// The raw tool output for a tool-result message: the model-visible +/// `model_observation` JSON when present, falling back to the safe summary. +fn tool_result_output(message: &ThreadMessageRecord) -> serde_json::Value { + let Some(content) = message.content.as_deref() else { + return serde_json::Value::Null; + }; + match ToolResultReferenceEnvelope::from_json_str(content) { + Ok(envelope) => envelope.model_observation.unwrap_or_else(|| { + serde_json::Value::String(envelope.safe_summary.as_str().to_string()) + }), + Err(_) => serde_json::Value::String(content.to_string()), + } +} + +fn map_thread_read_error(error: SessionThreadError) -> OpenAiCompatHttpError { + match error { + SessionThreadError::UnknownThread { .. } + | SessionThreadError::ThreadScopeMismatch { .. } => { + OpenAiCompatHttpError::not_found(Some("response_id".to_string())) + } + error => { + tracing::warn!( + target = "ironclaw::reborn::openai_compat", + error = %error, + "failed to read thread projection for OpenAI-compatible response" + ); + OpenAiCompatHttpError::from_kind( + 503, + true, + OpenAiCompatErrorKind::ServiceUnavailable, + None, + ) } } } @@ -341,21 +529,26 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { } => submitted_run_id.to_string(), _ => return Err(OpenAiCompatHttpError::internal()), }; - loop { - if let Some(content) = self - .read_finalized_response_message(&request.projection_read, submitted_run_id.clone()) - .await? - { - return Ok(OpenAiResponseProjection::new(response_object( - request.public_id, - request.mapping.created_at, - request.requested_model, - OpenAiResponseStatus::Completed, - Some(content), - ))); - } + while !self + .run_completed(&request.projection_read, submitted_run_id.clone()) + .await? + { tokio::time::sleep(self.poll_interval).await; } + let projection = self + .read_run_output( + &request.projection_read, + submitted_run_id, + &request.public_id, + ) + .await?; + Ok(OpenAiResponseProjection::new(response_object( + request.public_id, + request.mapping.created_at, + request.requested_model, + OpenAiResponseStatus::Completed, + projection.items, + ))) } async fn read_response( @@ -363,10 +556,14 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { request: OpenAiResponseReadRequest, ) -> Result { let submitted_run_id = response_turn_run_ref_from_mapping(&request)?; - let content = self - .read_finalized_response_message(&request.projection_read, submitted_run_id) + let projection = self + .read_run_output( + &request.projection_read, + submitted_run_id, + &request.public_id, + ) .await?; - let status = if content.is_some() { + let status = if projection.assistant_finalized { OpenAiResponseStatus::Completed } else { OpenAiResponseStatus::InProgress @@ -378,7 +575,7 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { .requested_model .unwrap_or_else(|| "reborn".to_string()), status, - content, + projection.items, )) } } @@ -400,22 +597,12 @@ fn response_turn_run_ref_from_mapping( } fn response_object( - id: ironclaw_reborn_openai_compat::OpenAiResponseId, + id: OpenAiResponseId, created_at: u64, model: String, status: OpenAiResponseStatus, - content: Option, + output: Vec, ) -> OpenAiResponseObject { - let output = content - .map(|text| { - vec![OpenAiResponseOutputItem::Message { - id: format!("msg_{}", id.as_str()), - status: Some(OpenAiResponseOutputItemStatus::Completed), - role: OpenAiResponsesMessageRole::Assistant, - content: serde_json::json!([{"type": "output_text", "text": text}]), - }] - }) - .unwrap_or_default(); OpenAiResponseObject { id, object: "response".to_string(), @@ -495,3 +682,275 @@ fn invalid_openai_compat_config( reason: format!("invalid OpenAI-compatible {field}: {error}"), } } + +#[cfg(test)] +mod tests { + use super::*; + + use ironclaw_host_api::CapabilityId; + use ironclaw_threads::{ + AppendAssistantDraftRequest, AppendToolResultReferenceRequest, EnsureThreadRequest, + InMemorySessionThreadService, MessageContent, ToolResultSafeSummary, + }; + use ironclaw_turns::{TurnActor, TurnScope}; + + const RUN_ID: &str = "turn-run-1"; + + fn test_scope() -> ThreadScope { + ThreadScope { + tenant_id: TenantId::new("tenant-1").expect("tenant"), + agent_id: AgentId::new("agent-1").expect("agent"), + project_id: Some(ProjectId::new("project-1").expect("project")), + owner_user_id: Some(UserId::new("user-1").expect("user")), + mission_id: None, + } + } + + /// A schema-valid `model_observation` so the thread service preserves it as + /// the raw, model-visible tool output rather than dropping it. + fn model_observation() -> serde_json::Value { + serde_json::json!({ + "schema_version": 1, + "status": "error", + "summary": "search failed", + "detail": { + "kind": "invalid_input", + "issues": [{ "path": "query", "code": "invalid_value" }] + }, + "trust": "untrusted_tool_output" + }) + } + + fn provider_call() -> ProviderToolCallReferenceEnvelope { + ProviderToolCallReferenceEnvelope { + provider_id: "openai".to_string(), + provider_model_id: "gpt-test".to_string(), + provider_turn_id: "turn-1".to_string(), + provider_call_id: "call_abc".to_string(), + provider_tool_name: "web_search".to_string(), + capability_id: CapabilityId::new("web.search").expect("capability id"), + arguments: serde_json::json!({ "query": "rust" }), + response_reasoning: None, + reasoning: None, + signature: None, + } + } + + fn projection_read(scope: &ThreadScope, thread_id: &ThreadId) -> ProjectionReadRequest { + ProjectionReadRequest { + actor: TurnActor::new(scope.owner_user_id.clone().expect("owner")), + scope: TurnScope::new_with_owner( + scope.tenant_id.clone(), + Some(scope.agent_id.clone()), + scope.project_id.clone(), + thread_id.clone(), + scope.owner_user_id.clone(), + ), + after_cursor: None, + limit: None, + } + } + + async fn ensure_thread( + service: &InMemorySessionThreadService, + scope: &ThreadScope, + thread_id: &ThreadId, + ) { + service + .ensure_thread(EnsureThreadRequest { + scope: scope.clone(), + thread_id: Some(thread_id.clone()), + created_by_actor_id: "actor".to_string(), + title: None, + metadata_json: None, + }) + .await + .expect("ensure thread"); + } + + async fn append_tool_result( + service: &InMemorySessionThreadService, + scope: &ThreadScope, + thread_id: &ThreadId, + ) { + service + .append_tool_result_reference(AppendToolResultReferenceRequest { + scope: scope.clone(), + thread_id: thread_id.clone(), + turn_run_id: RUN_ID.to_string(), + result_ref: "result:tool-1".to_string(), + safe_summary: ToolResultSafeSummary::new("search failed").expect("summary"), + provider_call: Some(provider_call()), + model_observation: Some(model_observation()), + }) + .await + .expect("append tool result"); + } + + #[tokio::test] + async fn read_run_output_emits_paired_function_call_and_raw_output() { + let service = Arc::new(InMemorySessionThreadService::default()); + let scope = test_scope(); + let thread_id = ThreadId::new("thread-1").expect("thread"); + ensure_thread(&service, &scope, &thread_id).await; + append_tool_result(&service, &scope, &thread_id).await; + let draft = service + .append_assistant_draft(AppendAssistantDraftRequest { + scope: scope.clone(), + thread_id: thread_id.clone(), + turn_run_id: RUN_ID.to_string(), + content: MessageContent::text("here is the answer"), + }) + .await + .expect("append draft"); + service + .finalize_assistant_message( + &scope, + &thread_id, + draft.message_id, + MessageContent::text("here is the answer"), + ) + .await + .expect("finalize assistant message"); + + let reader = OpenAiResponsesThreadProjectionReader::new(service); + let public_id = OpenAiResponseId::generate(); + let projection = reader + .read_run_output( + &projection_read(&scope, &thread_id), + RUN_ID.to_string(), + &public_id, + ) + .await + .expect("read run output"); + + assert!(projection.assistant_finalized); + assert_eq!(projection.items.len(), 3); + + match &projection.items[0] { + OpenAiResponseOutputItem::FunctionCall { + call_id, + name, + arguments, + .. + } => { + assert_eq!(call_id, "call_abc"); + assert_eq!(name, "web_search"); + let parsed: serde_json::Value = + serde_json::from_str(arguments).expect("arguments json"); + assert_eq!(parsed, serde_json::json!({ "query": "rust" })); + } + other => panic!("expected function_call, got {other:?}"), + } + match &projection.items[1] { + OpenAiResponseOutputItem::FunctionCallOutput { + call_id, output, .. + } => { + assert_eq!(call_id, "call_abc"); + // The raw model_observation flows through verbatim, not the summary. + assert_eq!(output, &model_observation()); + } + other => panic!("expected function_call_output, got {other:?}"), + } + match &projection.items[2] { + OpenAiResponseOutputItem::Message { + id, role, content, .. + } => { + assert!(matches!(role, OpenAiResponsesMessageRole::Assistant)); + // Message ids are sequence-keyed so multi-step runs stay unique. + assert!(id.starts_with("msg_")); + assert_eq!( + content, + &serde_json::json!([{ "type": "output_text", "text": "here is the answer" }]) + ); + } + other => panic!("expected message, got {other:?}"), + } + } + + #[tokio::test] + async fn read_run_output_orders_tool_items_before_assistant_message() { + // The assistant draft is created (sequence reserved) BEFORE the tool runs, + // so a naive sequence sort would place the final message first. Tool items + // must still come before the assistant message. + let service = Arc::new(InMemorySessionThreadService::default()); + let scope = test_scope(); + let thread_id = ThreadId::new("thread-3").expect("thread"); + ensure_thread(&service, &scope, &thread_id).await; + let draft = service + .append_assistant_draft(AppendAssistantDraftRequest { + scope: scope.clone(), + thread_id: thread_id.clone(), + turn_run_id: RUN_ID.to_string(), + content: MessageContent::text("draft"), + }) + .await + .expect("append draft"); + append_tool_result(&service, &scope, &thread_id).await; + service + .finalize_assistant_message( + &scope, + &thread_id, + draft.message_id, + MessageContent::text("here is the answer"), + ) + .await + .expect("finalize assistant message"); + + let reader = OpenAiResponsesThreadProjectionReader::new(service); + let public_id = OpenAiResponseId::generate(); + let projection = reader + .read_run_output( + &projection_read(&scope, &thread_id), + RUN_ID.to_string(), + &public_id, + ) + .await + .expect("read run output"); + + assert_eq!(projection.items.len(), 3); + assert!(matches!( + projection.items[0], + OpenAiResponseOutputItem::FunctionCall { .. } + )); + assert!(matches!( + projection.items[1], + OpenAiResponseOutputItem::FunctionCallOutput { .. } + )); + assert!(matches!( + projection.items[2], + OpenAiResponseOutputItem::Message { .. } + )); + } + + #[tokio::test] + async fn read_run_output_in_progress_surfaces_tool_output_without_final_message() { + let service = Arc::new(InMemorySessionThreadService::default()); + let scope = test_scope(); + let thread_id = ThreadId::new("thread-2").expect("thread"); + ensure_thread(&service, &scope, &thread_id).await; + append_tool_result(&service, &scope, &thread_id).await; + + let reader = OpenAiResponsesThreadProjectionReader::new(service); + let public_id = OpenAiResponseId::generate(); + let projection = reader + .read_run_output( + &projection_read(&scope, &thread_id), + RUN_ID.to_string(), + &public_id, + ) + .await + .expect("read run output"); + + assert!(!projection.assistant_finalized); + assert_eq!(projection.items.len(), 2); + assert!(matches!( + projection.items[0], + OpenAiResponseOutputItem::FunctionCall { .. } + )); + assert!(matches!( + projection.items[1], + OpenAiResponseOutputItem::FunctionCallOutput { .. } + )); + } +} From 71eb66e993aeeb14f683a70fff4ac93e38fe31cc Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Wed, 17 Jun 2026 10:43:52 -0700 Subject: [PATCH 2/3] Address PR #5010 review: harden tool-output projection - tool_result_output: best-effort envelope parse instead of strict from_json_str, so a version/model_observation-shape mismatch falls back to safe_summary rather than leaking the raw envelope string - read_run_output: filter tool results to MessageStatus::Finalized to avoid surfacing redacted/deleted tool activity - fix misleading test comment (message ids are response-id-keyed) --- .../src/openai_compat_serve.rs | 22 +++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs index 01401c25b4e..b92b0e672ae 100644 --- a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs +++ b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs @@ -330,11 +330,17 @@ impl OpenAiResponsesThreadProjectionReader { .filter(|message| message.turn_run_id.as_deref() == Some(turn_run_id.as_str())) .collect(); - // Tool calls/outputs, in transcript order. + // Tool calls/outputs, in transcript order. Finalized only: history can + // surface redacted/deleted messages, and matching the finalized status + // these rows are appended with avoids leaking redacted tool activity + // (including a `tool_result_ref` fallback call id). let mut tool_results: Vec<&ThreadMessageRecord> = run_messages .iter() .copied() - .filter(|message| message.kind == MessageKind::ToolResultReference) + .filter(|message| { + message.kind == MessageKind::ToolResultReference + && message.status == MessageStatus::Finalized + }) .collect(); tool_results.sort_by_key(|message| message.sequence); // The run's single finalized assistant reply (drafts dedup by run). @@ -483,11 +489,19 @@ fn push_tool_output_items( /// The raw tool output for a tool-result message: the model-visible /// `model_observation` JSON when present, falling back to the safe summary. +/// +/// Parsing is best-effort: the envelope shape is deserialized directly without +/// the strict `from_json_str` validation, so a version mismatch or a legacy +/// `model_observation`/`result_ref` shape still yields the intended +/// `model_observation`/`safe_summary` rather than leaking the entire raw +/// envelope. Only content that does not deserialize as an envelope at all +/// (`safe_summary` itself is still validated on deserialize) falls back to the +/// raw string. fn tool_result_output(message: &ThreadMessageRecord) -> serde_json::Value { let Some(content) = message.content.as_deref() else { return serde_json::Value::Null; }; - match ToolResultReferenceEnvelope::from_json_str(content) { + match serde_json::from_str::(content) { Ok(envelope) => envelope.model_observation.unwrap_or_else(|| { serde_json::Value::String(envelope.safe_summary.as_str().to_string()) }), @@ -857,7 +871,7 @@ mod tests { id, role, content, .. } => { assert!(matches!(role, OpenAiResponsesMessageRole::Assistant)); - // Message ids are sequence-keyed so multi-step runs stay unique. + // Message ids are response-id-keyed (`msg_{response_id}`). assert!(id.starts_with("msg_")); assert_eq!( content, From b033765b1b4da493bc0c5544beb056609b3e432e Mon Sep 17 00:00:00 2001 From: Illia Polosukhin Date: Wed, 17 Jun 2026 11:09:06 -0700 Subject: [PATCH 3/3] Reconcile openai_compat_serve after origin/main merge The merge left two duplicate `tests` modules (inline + file-based) and my read_run_output tests called the old single-arg reader constructor. Combine the wait/read paths so projected run status (failed/cancelled + error object, from main) and the full tool-call/output projection (from this branch) both flow through, and fold the three read_run_output tests into the file-based tests module against the 2-arg constructor. All 10 openai_compat_serve tests pass; clippy clean with --features openai-compat-beta. --- .../src/openai_compat_serve.rs | 272 ------------------ .../src/openai_compat_serve/tests.rs | 251 +++++++++++++++- 2 files changed, 249 insertions(+), 274 deletions(-) diff --git a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs index b6a996f85a9..d54cb231406 100644 --- a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs +++ b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs @@ -887,275 +887,3 @@ fn invalid_openai_compat_config( reason: format!("invalid OpenAI-compatible {field}: {error}"), } } - -#[cfg(test)] -mod tests { - use super::*; - - use ironclaw_host_api::CapabilityId; - use ironclaw_threads::{ - AppendAssistantDraftRequest, AppendToolResultReferenceRequest, EnsureThreadRequest, - InMemorySessionThreadService, MessageContent, ToolResultSafeSummary, - }; - use ironclaw_turns::{TurnActor, TurnScope}; - - const RUN_ID: &str = "turn-run-1"; - - fn test_scope() -> ThreadScope { - ThreadScope { - tenant_id: TenantId::new("tenant-1").expect("tenant"), - agent_id: AgentId::new("agent-1").expect("agent"), - project_id: Some(ProjectId::new("project-1").expect("project")), - owner_user_id: Some(UserId::new("user-1").expect("user")), - mission_id: None, - } - } - - /// A schema-valid `model_observation` so the thread service preserves it as - /// the raw, model-visible tool output rather than dropping it. - fn model_observation() -> serde_json::Value { - serde_json::json!({ - "schema_version": 1, - "status": "error", - "summary": "search failed", - "detail": { - "kind": "invalid_input", - "issues": [{ "path": "query", "code": "invalid_value" }] - }, - "trust": "untrusted_tool_output" - }) - } - - fn provider_call() -> ProviderToolCallReferenceEnvelope { - ProviderToolCallReferenceEnvelope { - provider_id: "openai".to_string(), - provider_model_id: "gpt-test".to_string(), - provider_turn_id: "turn-1".to_string(), - provider_call_id: "call_abc".to_string(), - provider_tool_name: "web_search".to_string(), - capability_id: CapabilityId::new("web.search").expect("capability id"), - arguments: serde_json::json!({ "query": "rust" }), - response_reasoning: None, - reasoning: None, - signature: None, - } - } - - fn projection_read(scope: &ThreadScope, thread_id: &ThreadId) -> ProjectionReadRequest { - ProjectionReadRequest { - actor: TurnActor::new(scope.owner_user_id.clone().expect("owner")), - scope: TurnScope::new_with_owner( - scope.tenant_id.clone(), - Some(scope.agent_id.clone()), - scope.project_id.clone(), - thread_id.clone(), - scope.owner_user_id.clone(), - ), - after_cursor: None, - limit: None, - } - } - - async fn ensure_thread( - service: &InMemorySessionThreadService, - scope: &ThreadScope, - thread_id: &ThreadId, - ) { - service - .ensure_thread(EnsureThreadRequest { - scope: scope.clone(), - thread_id: Some(thread_id.clone()), - created_by_actor_id: "actor".to_string(), - title: None, - metadata_json: None, - }) - .await - .expect("ensure thread"); - } - - async fn append_tool_result( - service: &InMemorySessionThreadService, - scope: &ThreadScope, - thread_id: &ThreadId, - ) { - service - .append_tool_result_reference(AppendToolResultReferenceRequest { - scope: scope.clone(), - thread_id: thread_id.clone(), - turn_run_id: RUN_ID.to_string(), - result_ref: "result:tool-1".to_string(), - safe_summary: ToolResultSafeSummary::new("search failed").expect("summary"), - provider_call: Some(provider_call()), - model_observation: Some(model_observation()), - }) - .await - .expect("append tool result"); - } - - #[tokio::test] - async fn read_run_output_emits_paired_function_call_and_raw_output() { - let service = Arc::new(InMemorySessionThreadService::default()); - let scope = test_scope(); - let thread_id = ThreadId::new("thread-1").expect("thread"); - ensure_thread(&service, &scope, &thread_id).await; - append_tool_result(&service, &scope, &thread_id).await; - let draft = service - .append_assistant_draft(AppendAssistantDraftRequest { - scope: scope.clone(), - thread_id: thread_id.clone(), - turn_run_id: RUN_ID.to_string(), - content: MessageContent::text("here is the answer"), - }) - .await - .expect("append draft"); - service - .finalize_assistant_message( - &scope, - &thread_id, - draft.message_id, - MessageContent::text("here is the answer"), - ) - .await - .expect("finalize assistant message"); - - let reader = OpenAiResponsesThreadProjectionReader::new(service); - let public_id = OpenAiResponseId::generate(); - let projection = reader - .read_run_output( - &projection_read(&scope, &thread_id), - RUN_ID.to_string(), - &public_id, - ) - .await - .expect("read run output"); - - assert!(projection.assistant_finalized); - assert_eq!(projection.items.len(), 3); - - match &projection.items[0] { - OpenAiResponseOutputItem::FunctionCall { - call_id, - name, - arguments, - .. - } => { - assert_eq!(call_id, "call_abc"); - assert_eq!(name, "web_search"); - let parsed: serde_json::Value = - serde_json::from_str(arguments).expect("arguments json"); - assert_eq!(parsed, serde_json::json!({ "query": "rust" })); - } - other => panic!("expected function_call, got {other:?}"), - } - match &projection.items[1] { - OpenAiResponseOutputItem::FunctionCallOutput { - call_id, output, .. - } => { - assert_eq!(call_id, "call_abc"); - // The raw model_observation flows through verbatim, not the summary. - assert_eq!(output, &model_observation()); - } - other => panic!("expected function_call_output, got {other:?}"), - } - match &projection.items[2] { - OpenAiResponseOutputItem::Message { - id, role, content, .. - } => { - assert!(matches!(role, OpenAiResponsesMessageRole::Assistant)); - // Message ids are response-id-keyed (`msg_{response_id}`). - assert!(id.starts_with("msg_")); - assert_eq!( - content, - &serde_json::json!([{ "type": "output_text", "text": "here is the answer" }]) - ); - } - other => panic!("expected message, got {other:?}"), - } - } - - #[tokio::test] - async fn read_run_output_orders_tool_items_before_assistant_message() { - // The assistant draft is created (sequence reserved) BEFORE the tool runs, - // so a naive sequence sort would place the final message first. Tool items - // must still come before the assistant message. - let service = Arc::new(InMemorySessionThreadService::default()); - let scope = test_scope(); - let thread_id = ThreadId::new("thread-3").expect("thread"); - ensure_thread(&service, &scope, &thread_id).await; - let draft = service - .append_assistant_draft(AppendAssistantDraftRequest { - scope: scope.clone(), - thread_id: thread_id.clone(), - turn_run_id: RUN_ID.to_string(), - content: MessageContent::text("draft"), - }) - .await - .expect("append draft"); - append_tool_result(&service, &scope, &thread_id).await; - service - .finalize_assistant_message( - &scope, - &thread_id, - draft.message_id, - MessageContent::text("here is the answer"), - ) - .await - .expect("finalize assistant message"); - - let reader = OpenAiResponsesThreadProjectionReader::new(service); - let public_id = OpenAiResponseId::generate(); - let projection = reader - .read_run_output( - &projection_read(&scope, &thread_id), - RUN_ID.to_string(), - &public_id, - ) - .await - .expect("read run output"); - - assert_eq!(projection.items.len(), 3); - assert!(matches!( - projection.items[0], - OpenAiResponseOutputItem::FunctionCall { .. } - )); - assert!(matches!( - projection.items[1], - OpenAiResponseOutputItem::FunctionCallOutput { .. } - )); - assert!(matches!( - projection.items[2], - OpenAiResponseOutputItem::Message { .. } - )); - } - - #[tokio::test] - async fn read_run_output_in_progress_surfaces_tool_output_without_final_message() { - let service = Arc::new(InMemorySessionThreadService::default()); - let scope = test_scope(); - let thread_id = ThreadId::new("thread-2").expect("thread"); - ensure_thread(&service, &scope, &thread_id).await; - append_tool_result(&service, &scope, &thread_id).await; - - let reader = OpenAiResponsesThreadProjectionReader::new(service); - let public_id = OpenAiResponseId::generate(); - let projection = reader - .read_run_output( - &projection_read(&scope, &thread_id), - RUN_ID.to_string(), - &public_id, - ) - .await - .expect("read run output"); - - assert!(!projection.assistant_finalized); - assert_eq!(projection.items.len(), 2); - assert!(matches!( - projection.items[0], - OpenAiResponseOutputItem::FunctionCall { .. } - )); - assert!(matches!( - projection.items[1], - OpenAiResponseOutputItem::FunctionCallOutput { .. } - )); - } -} diff --git a/crates/ironclaw_reborn_composition/src/openai_compat_serve/tests.rs b/crates/ironclaw_reborn_composition/src/openai_compat_serve/tests.rs index 937bfa033d1..16cb7c16664 100644 --- a/crates/ironclaw_reborn_composition/src/openai_compat_serve/tests.rs +++ b/crates/ironclaw_reborn_composition/src/openai_compat_serve/tests.rs @@ -3,7 +3,7 @@ use super::*; use std::collections::VecDeque; use std::sync::Mutex; -use ironclaw_host_api::ThreadId; +use ironclaw_host_api::{CapabilityId, ThreadId}; use ironclaw_product_adapters::{ AdapterInstallationId, ExternalConversationRef, ProductAdapterError, ProductAdapterId, ProductOutboundTarget, ProjectionCursor, @@ -14,7 +14,8 @@ use ironclaw_reborn_openai_compat::{ OpenAiCompatRouteSurface, OpenAiCompatTurnRunRef, OpenAiResponseId, }; use ironclaw_threads::{ - AppendAssistantDraftRequest, EnsureThreadRequest, InMemorySessionThreadService, MessageContent, + AppendAssistantDraftRequest, AppendToolResultReferenceRequest, EnsureThreadRequest, + InMemorySessionThreadService, MessageContent, ToolResultSafeSummary, }; use ironclaw_turns::{ReplyTargetBindingRef, TurnActor, TurnRunId, TurnScope}; @@ -412,3 +413,249 @@ fn run_status_envelope( }, ) } + +const RUN_ID: &str = "turn-run-1"; + +fn run_output_scope() -> ThreadScope { + ThreadScope { + tenant_id: TenantId::new("tenant-1").expect("tenant"), + agent_id: AgentId::new("agent-1").expect("agent"), + project_id: Some(ProjectId::new("project-1").expect("project")), + owner_user_id: Some(UserId::new("user-1").expect("user")), + mission_id: None, + } +} + +/// A schema-valid `model_observation` so the thread service preserves it as +/// the raw, model-visible tool output rather than dropping it. +fn model_observation() -> serde_json::Value { + serde_json::json!({ + "schema_version": 1, + "status": "error", + "summary": "search failed", + "detail": { + "kind": "invalid_input", + "issues": [{ "path": "query", "code": "invalid_value" }] + }, + "trust": "untrusted_tool_output" + }) +} + +fn run_output_provider_call() -> ProviderToolCallReferenceEnvelope { + ProviderToolCallReferenceEnvelope { + provider_id: "openai".to_string(), + provider_model_id: "gpt-test".to_string(), + provider_turn_id: "turn-1".to_string(), + provider_call_id: "call_abc".to_string(), + provider_tool_name: "web_search".to_string(), + capability_id: CapabilityId::new("web.search").expect("capability id"), + arguments: serde_json::json!({ "query": "rust" }), + response_reasoning: None, + reasoning: None, + signature: None, + } +} + +fn run_output_projection_read(scope: &ThreadScope, thread_id: &ThreadId) -> ProjectionReadRequest { + ProjectionReadRequest { + actor: TurnActor::new(scope.owner_user_id.clone().expect("owner")), + scope: TurnScope::new_with_owner( + scope.tenant_id.clone(), + Some(scope.agent_id.clone()), + scope.project_id.clone(), + thread_id.clone(), + scope.owner_user_id.clone(), + ), + after_cursor: None, + limit: None, + } +} + +async fn ensure_run_output_thread( + service: &InMemorySessionThreadService, + scope: &ThreadScope, + thread_id: &ThreadId, +) { + service + .ensure_thread(EnsureThreadRequest { + scope: scope.clone(), + thread_id: Some(thread_id.clone()), + created_by_actor_id: "actor".to_string(), + title: None, + metadata_json: None, + }) + .await + .expect("ensure thread"); +} + +async fn append_run_output_tool_result( + service: &InMemorySessionThreadService, + scope: &ThreadScope, + thread_id: &ThreadId, +) { + service + .append_tool_result_reference(AppendToolResultReferenceRequest { + scope: scope.clone(), + thread_id: thread_id.clone(), + turn_run_id: RUN_ID.to_string(), + result_ref: "result:tool-1".to_string(), + safe_summary: ToolResultSafeSummary::new("search failed").expect("summary"), + provider_call: Some(run_output_provider_call()), + model_observation: Some(model_observation()), + }) + .await + .expect("append tool result"); +} + +fn run_output_reader( + service: Arc, +) -> OpenAiResponsesThreadProjectionReader { + OpenAiResponsesThreadProjectionReader::new( + service, + Arc::new(StaticProjectionStream::new(vec![])), + ) +} + +#[tokio::test] +async fn read_run_output_emits_paired_function_call_and_raw_output() { + let service = Arc::new(InMemorySessionThreadService::default()); + let scope = run_output_scope(); + let thread_id = ThreadId::new("thread-1").expect("thread"); + ensure_run_output_thread(&service, &scope, &thread_id).await; + append_run_output_tool_result(&service, &scope, &thread_id).await; + + let reader = run_output_reader(service); + let public_id = OpenAiResponseId::new("resp_test").expect("response id"); + let projection = reader + .read_run_output( + &run_output_projection_read(&scope, &thread_id), + RUN_ID.to_string(), + &public_id, + ) + .await + .expect("read run output"); + + assert!(!projection.assistant_finalized); + match &projection.items[0] { + OpenAiResponseOutputItem::FunctionCall { + call_id, + name, + arguments, + .. + } => { + assert_eq!(call_id, "call_abc"); + assert_eq!(name, "web_search"); + let parsed: serde_json::Value = + serde_json::from_str(arguments).expect("arguments json"); + assert_eq!(parsed, serde_json::json!({ "query": "rust" })); + } + other => panic!("expected function_call, got {other:?}"), + } + match &projection.items[1] { + OpenAiResponseOutputItem::FunctionCallOutput { + call_id, output, .. + } => { + assert_eq!(call_id, "call_abc"); + // The raw model_observation flows through verbatim, not the summary. + assert_eq!(output, &model_observation()); + } + other => panic!("expected function_call_output, got {other:?}"), + } +} + +#[tokio::test] +async fn read_run_output_orders_tool_items_before_assistant_message() { + // The assistant draft is created (sequence reserved) BEFORE the tool runs, + // so a naive sequence sort would place the final message first. Tool items + // must still come before the assistant message. + let service = Arc::new(InMemorySessionThreadService::default()); + let scope = run_output_scope(); + let thread_id = ThreadId::new("thread-3").expect("thread"); + ensure_run_output_thread(&service, &scope, &thread_id).await; + + let draft = service + .append_assistant_draft(AppendAssistantDraftRequest { + scope: scope.clone(), + thread_id: thread_id.clone(), + turn_run_id: RUN_ID.to_string(), + content: MessageContent::text("here is the answer"), + }) + .await + .expect("append assistant draft"); + append_run_output_tool_result(&service, &scope, &thread_id).await; + service + .finalize_assistant_message( + &scope, + &thread_id, + draft.message_id, + MessageContent::text("here is the answer"), + ) + .await + .expect("finalize assistant message"); + + let reader = run_output_reader(service); + let public_id = OpenAiResponseId::new("resp_test").expect("response id"); + let projection = reader + .read_run_output( + &run_output_projection_read(&scope, &thread_id), + RUN_ID.to_string(), + &public_id, + ) + .await + .expect("read run output"); + + assert!(projection.assistant_finalized); + assert!(matches!( + projection.items[0], + OpenAiResponseOutputItem::FunctionCall { .. } + )); + assert!(matches!( + projection.items[1], + OpenAiResponseOutputItem::FunctionCallOutput { .. } + )); + match &projection.items[2] { + OpenAiResponseOutputItem::Message { + id, role, content, .. + } => { + assert!(matches!(role, OpenAiResponsesMessageRole::Assistant)); + // Message ids are response-id-keyed (`msg_{response_id}`). + assert!(id.starts_with("msg_")); + assert_eq!( + content, + &serde_json::json!([{ "type": "output_text", "text": "here is the answer" }]) + ); + } + other => panic!("expected message, got {other:?}"), + } +} + +#[tokio::test] +async fn read_run_output_in_progress_surfaces_tool_output_without_final_message() { + let service = Arc::new(InMemorySessionThreadService::default()); + let scope = run_output_scope(); + let thread_id = ThreadId::new("thread-2").expect("thread"); + ensure_run_output_thread(&service, &scope, &thread_id).await; + append_run_output_tool_result(&service, &scope, &thread_id).await; + + let reader = run_output_reader(service); + let public_id = OpenAiResponseId::new("resp_test").expect("response id"); + let projection = reader + .read_run_output( + &run_output_projection_read(&scope, &thread_id), + RUN_ID.to_string(), + &public_id, + ) + .await + .expect("read run output"); + + assert!(!projection.assistant_finalized); + assert_eq!(projection.items.len(), 2); + assert!(matches!( + projection.items[0], + OpenAiResponseOutputItem::FunctionCall { .. } + )); + assert!(matches!( + projection.items[1], + OpenAiResponseOutputItem::FunctionCallOutput { .. } + )); +}