diff --git a/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs b/crates/ironclaw_reborn_composition/src/openai_compat_serve.rs index 79b0aacbf25..d54cb231406 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; @@ -12,7 +13,7 @@ use ironclaw_attachments::InboundAttachment; 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, ProductAdapterError, ProductAdapterId, ProductInboundAck, @@ -34,15 +35,18 @@ use ironclaw_reborn_openai_compat::{ OpenAiChatProjectionStreamRequest, OpenAiCompatErrorKind, OpenAiCompatHttpError, OpenAiCompatInboundAttachmentSubmit, OpenAiCompatProjectionStreamer, OpenAiCompatRefStore, OpenAiCompatResourceBinding, OpenAiCompatRouterState, OpenAiResponseErrorObject, - OpenAiResponseObject, OpenAiResponseOutputItem, OpenAiResponseOutputItemStatus, - OpenAiResponseProjection, OpenAiResponseProjectionStreamRequest, OpenAiResponseReadRequest, - OpenAiResponseStatus, OpenAiResponseWaitRequest, OpenAiResponsesMessageRole, - OpenAiResponsesProjectionReader, OpenAiResponsesWorkflow, openai_compat_router_with_state, - openai_compat_routes, + OpenAiResponseId, OpenAiResponseObject, OpenAiResponseOutputItem, + OpenAiResponseOutputItemStatus, OpenAiResponseProjection, + OpenAiResponseProjectionStreamRequest, OpenAiResponseReadRequest, OpenAiResponseStatus, + OpenAiResponseWaitRequest, OpenAiResponsesMessageRole, OpenAiResponsesProjectionReader, + OpenAiResponsesWorkflow, openai_compat_router_with_state, openai_compat_routes, }; 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; @@ -352,13 +356,136 @@ 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. 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 + && message.status == MessageStatus::Finalized + }) + .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, @@ -366,30 +493,13 @@ 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()) } + /// Drain the projection event stream for the run's latest status. Used while + /// polling so a failed/cancelled run surfaces even before (or without) a + /// finalized assistant message. async fn read_projected_response_status( &self, request: &ProjectionReadRequest, @@ -412,6 +522,101 @@ impl OpenAiResponsesThreadProjectionReader { } } +/// 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. +/// +/// 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 serde_json::from_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, + ) + } + } +} + #[async_trait] impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { async fn wait_for_response_completion( @@ -426,16 +631,26 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { }; let mut projection_after_cursor = request.projection_read.after_cursor.clone(); loop { - if let Some(content) = self - .read_finalized_response_message(&request.projection_read, submitted_run_id.clone()) + // Cheap completion gate first; only read the full transcript + // projection (tool calls/outputs + assistant message) once the run + // has produced its finalized reply. + if self + .run_completed(&request.projection_read, submitted_run_id.clone()) .await? { + let projection = self + .read_run_output( + &request.projection_read, + submitted_run_id, + &request.public_id, + ) + .await?; return Ok(OpenAiResponseProjection::new(response_object( request.public_id, request.mapping.created_at, request.requested_model, OpenAiResponseStatus::Completed, - Some(content), + projection.items, ))); } let projected = self @@ -459,7 +674,7 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { request.mapping.created_at, request.requested_model, status, - None, + Vec::new(), ))); } tokio::time::sleep(self.poll_interval).await; @@ -471,10 +686,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.clone()) + let projection = self + .read_run_output( + &request.projection_read, + submitted_run_id.clone(), + &request.public_id, + ) .await?; - let projected_status = if content.is_some() { + let projected_status = if projection.assistant_finalized { None } else { self.read_projected_response_status( @@ -485,7 +704,7 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { .await? .status }; - let status = match (content.is_some(), projected_status) { + let status = match (projection.assistant_finalized, projected_status) { (true, _) => OpenAiResponseStatus::Completed, (false, Some(OpenAiResponseStatus::Completed | OpenAiResponseStatus::InProgress)) => { OpenAiResponseStatus::InProgress @@ -500,7 +719,7 @@ impl OpenAiResponsesProjectionReader for OpenAiResponsesThreadProjectionReader { .requested_model .unwrap_or_else(|| "reborn".to_string()), status, - content, + projection.items, )) } } @@ -576,22 +795,12 @@ fn response_status_from_projection_run_status(status: &str) -> 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(); let error = if matches!(status, OpenAiResponseStatus::Failed) { Some(OpenAiResponseErrorObject::from_kind( OpenAiCompatErrorKind::Internal, 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 { .. } + )); +}