use std::borrow::Cow;
use std::error::Error as StdError;
use std::sync::Arc;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use crate::runtime::config::AgentLoopConfig;
use crate::runtime::runner::prompt_context::{
append_core_agent_directives, strip_existing_core_directives, strip_existing_external_memory,
strip_existing_goal, strip_existing_plan_mode_instructions,
strip_existing_plan_runtime_context, strip_existing_task_list,
};
use crate::runtime::runner::session_setup::prompt_envelope::{
build_active_workflow_context_block, build_agent_hook_context_block,
build_conversation_summary_context_block, build_external_memory_context_block,
build_goal_context_block, build_plan_mode_context_block, build_plan_runtime_context_block,
build_project_resources_context_block, build_task_list_context_block,
};
use crate::runtime::runner::session_setup::prompt_setup::{
build_stable_prompt_frame_with_sections, StablePrefixSection,
};
use bamboo_agent_core::agent::events::TokenBudgetUsage;
use bamboo_agent_core::tools::ToolSchema;
use bamboo_agent_core::{
AgentError, AgentEvent, ContextBlock, ContextBlockPriority, ContextBlockStability,
ContextBlockType, Message, MessagePhase, Role, Session,
};
use bamboo_compression::{PreparedContext, TiktokenTokenCounter, TokenCounter};
use bamboo_domain::ReasoningEffort;
use bamboo_domain::MAX_MODEL_CONTEXT_RENDERED_BYTES;
use bamboo_llm::provider::ResponsesRequestOptions;
use bamboo_llm::{
CacheTtl, Continuation, LLMProvider, LLMRequestOptions, PromptCachePlan, PromptIR, Segment,
SegmentRole,
};
use bamboo_tools::exposure::activated_discoverable_tools;
use sha2::{Digest, Sha256};
pub(in crate::runtime::runner) struct LlmStreamFrame<'a> {
pub event_tx: &'a mpsc::Sender<AgentEvent>,
pub cancel_token: &'a CancellationToken,
pub session_id: &'a str,
pub model: &'a str,
pub provider_name: Option<&'a str>,
pub provider_type: Option<&'a str>,
pub reasoning_effort: Option<ReasoningEffort>,
pub max_context_tokens: u32,
pub max_output_tokens: u32,
}
const SESSION_RESPONSES_PREVIOUS_RESPONSE_ID_KEY: &str = "responses.previous_response_id";
const CONVERSATION_SUMMARY_START_MARKER: &str = "<!-- CONVERSATION_SUMMARY_START -->";
const INTERRUPTED_ASSISTANT_OUTPUT_KIND: &str = "interrupted_assistant_output";
const AGENT_LOOP_REQUEST_PURPOSE: &str = "agent_loop";
const PROMPT_CACHE_KEY_DOMAIN: &[u8] = b"bamboo/openai/responses/prompt-cache-key/v1\0";
fn interruption_kind(error: &AgentError) -> &'static str {
match error {
AgentError::Cancelled => "cancelled",
AgentError::StreamTimeout(_) => "stream_timeout",
AgentError::LLMOverflow(_) => "llm_overflow",
AgentError::EmptyAssistantResponse { .. } => "empty_assistant_response",
AgentError::LLM(_) => "llm_error",
_ => "execution_error",
}
}
fn append_interrupted_assistant_output(
session: &mut Session,
partial: crate::runtime::stream::handler::InterruptedStreamOutput,
error: &AgentError,
) -> bool {
if partial.content.is_empty()
&& partial.reasoning_content.is_empty()
&& partial.partial_tool_calls.is_empty()
{
return false;
}
let partial_tool_calls = (!partial.partial_tool_calls.is_empty()).then(|| {
serde_json::to_value(&partial.partial_tool_calls).unwrap_or(serde_json::Value::Null)
});
let mut message = Message::assistant_with_reasoning(
partial.content,
None,
(!partial.reasoning_content.is_empty()).then_some(partial.reasoning_content),
);
message.phase = Some(MessagePhase::Commentary);
message.reasoning_signature = None;
message.metadata = Some(serde_json::json!({
"runtime_kind": INTERRUPTED_ASSISTANT_OUTPUT_KIND,
"interrupted": true,
"interruption_kind": interruption_kind(error),
"partial_tool_calls": partial_tool_calls,
}));
session.add_message(message);
true
}
pub(crate) fn discard_latest_interrupted_assistant_output(
session: &mut Session,
attempt_tail_message_id: Option<&str>,
) -> bool {
let interrupted = session.messages.last().is_some_and(|message| {
if Some(message.id.as_str()) == attempt_tail_message_id {
return false;
}
message
.metadata
.as_ref()
.and_then(|metadata| metadata.get("runtime_kind"))
.and_then(serde_json::Value::as_str)
== Some(INTERRUPTED_ASSISTANT_OUTPUT_KIND)
});
if interrupted {
session.messages.pop();
session.updated_at = chrono::Utc::now();
}
interrupted
}
fn session_previous_response_id(session: &Session) -> Option<&str> {
session
.metadata
.get(SESSION_RESPONSES_PREVIOUS_RESPONSE_ID_KEY)
.map(String::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
}
fn provider_supports_previous_response_id(provider_type: Option<&str>) -> bool {
!matches!(provider_type.map(str::trim), Some("copilot"))
}
fn engine_responses_policy() -> ResponsesRequestOptions {
ResponsesRequestOptions {
store: Some(false),
text_verbosity: Some("high".to_string()),
reasoning_summary: Some("auto".to_string()),
include: Some(vec!["reasoning.encrypted_content".to_string()]),
..Default::default()
}
}
fn agent_loop_prompt_cache_key(
session_id: Option<&str>,
request_purpose: Option<&str>,
) -> Option<String> {
let purpose = request_purpose.filter(|purpose| *purpose == AGENT_LOOP_REQUEST_PURPOSE)?;
let session_id = session_id.filter(|session_id| !session_id.trim().is_empty())?;
let mut hasher = Sha256::new();
hasher.update(PROMPT_CACHE_KEY_DOMAIN);
for component in [purpose.as_bytes(), session_id.as_bytes()] {
hasher.update((component.len() as u64).to_be_bytes());
hasher.update(component);
}
Some(hex::encode(hasher.finalize()))
}
fn responses_continuation_enabled(
policy: &ResponsesRequestOptions,
provider_type: Option<&str>,
) -> bool {
policy.store == Some(true) && provider_supports_previous_response_id(provider_type)
}
fn format_reqwest_transport_error(error: &reqwest::Error) -> String {
let mut kinds = Vec::new();
if error.is_timeout() {
kinds.push("timeout");
}
if error.is_connect() {
kinds.push("connect");
}
if error.is_request() {
kinds.push("request");
}
if error.is_body() {
kinds.push("body");
}
if error.is_decode() {
kinds.push("decode");
}
if error.is_redirect() {
kinds.push("redirect");
}
if error.is_builder() {
kinds.push("builder");
}
if error.is_status() {
kinds.push("status");
}
let kind = if kinds.is_empty() {
"unknown".to_string()
} else {
kinds.join("+")
};
let url = error
.url()
.map(ToString::to_string)
.unwrap_or_else(|| "<unknown>".to_string());
let mut causes = Vec::new();
let mut source = StdError::source(error);
while let Some(cause) = source {
causes.push(cause.to_string());
source = cause.source();
if causes.len() >= 4 {
break;
}
}
if causes.is_empty() {
format!(
"HTTP transport error [{}] for url ({}): {}",
kind, url, error
)
} else {
format!(
"HTTP transport error [{}] for url ({}): {} | causes: {}",
kind,
url,
error,
causes.join(" | ")
)
}
}
fn format_provider_error(error: bamboo_llm::provider::LLMError) -> String {
match error {
bamboo_llm::provider::LLMError::Http(http) => format_reqwest_transport_error(&http),
other => other.to_string(),
}
}
fn is_llm_overflow_error(message: &str) -> bool {
let normalized = message.trim().to_ascii_lowercase();
if normalized.is_empty() {
return false;
}
let overflow_patterns = [
"prompt too long",
"context too long",
"maximum context length",
"maximum context size",
"context length exceeded",
"context window exceeded",
"request too large",
"too many tokens",
"input is too long",
"input too long",
"token limit exceeded",
];
overflow_patterns
.iter()
.any(|pattern| normalized.contains(pattern))
}
fn is_conversation_summary_message(message: &Message) -> bool {
matches!(message.role, Role::System)
&& message.content.contains(CONVERSATION_SUMMARY_START_MARKER)
}
fn derive_system_remainder_message(
message: &Message,
stable_instructions: &str,
) -> Option<Message> {
if !matches!(message.role, Role::System) || is_conversation_summary_message(message) {
return None;
}
let without_external_memory = strip_existing_external_memory(&message.content);
let without_task_list = strip_existing_task_list(&without_external_memory);
let without_plan_mode = strip_existing_plan_mode_instructions(&without_task_list);
let without_plan_runtime = strip_existing_plan_runtime_context(&without_plan_mode);
let without_goal = strip_existing_goal(&without_plan_runtime);
let without_directives = strip_existing_core_directives(&without_goal);
let trimmed = without_directives.trim();
if trimmed.is_empty() {
return None;
}
let stable_without_directives = strip_existing_core_directives(stable_instructions);
let stable_trimmed = stable_without_directives.trim();
if stable_trimmed.is_empty() {
return Some(Message::system(trimmed.to_string()));
}
if trimmed == stable_trimmed {
return None;
}
if let Some(remainder) = trimmed.strip_prefix(stable_trimmed) {
let remainder = remainder.trim();
return (!remainder.is_empty()).then(|| Message::system(remainder.to_string()));
}
if stable_trimmed.starts_with(&format!("{trimmed}\n")) {
return None;
}
Some(Message::system(trimmed.to_string()))
}
struct PreparedRequestEnvelope {
ir: PromptIR,
stable_prefix_sections: Vec<StablePrefixSection>,
ledger_changed: bool,
prefix_epoch: u64,
prefix_reset_reason: Option<bamboo_domain::ModelContextResetReason>,
}
#[derive(Debug, Clone, Copy)]
pub(super) struct ProjectedRequestUsage {
pub input_tokens: u32,
pub ledger_rendered_bytes: usize,
}
fn measure_request_usage(
session: &Session,
envelope: &PreparedRequestEnvelope,
) -> ProjectedRequestUsage {
let input_tokens = TiktokenTokenCounter::default().count_messages(&envelope.ir.flatten());
let ledger_rendered_bytes = session
.model_context_state
.as_ref()
.map(|state| {
state.events.iter().fold(0usize, |total, event| {
total.saturating_add(event.rendered_text.len())
})
})
.unwrap_or(0);
ProjectedRequestUsage {
input_tokens,
ledger_rendered_bytes,
}
}
fn required_tool_for_session(session: &Session) -> Option<&'static str> {
crate::runtime::runner::session_setup::skill_context::explicit_activation_pending(session)
.then_some("load_skill")
}
fn effective_tool_schemas<'a>(
session: &Session,
tool_schemas: &'a [ToolSchema],
) -> Cow<'a, [ToolSchema]> {
let Some(required_tool) = required_tool_for_session(session) else {
return Cow::Borrowed(tool_schemas);
};
Cow::Owned(
tool_schemas
.iter()
.filter(|schema| schema.function.name == required_tool)
.cloned()
.collect(),
)
}
pub(super) fn project_request_usage(
session: &Session,
prepared_context: &PreparedContext,
config: &AgentLoopConfig,
tool_schemas: &[ToolSchema],
model: &str,
) -> ProjectedRequestUsage {
let mut shadow = session.clone();
let effective_tool_schemas = effective_tool_schemas(&shadow, tool_schemas);
let envelope = build_request_envelope_reconciled(
&mut shadow,
prepared_context,
config,
effective_tool_schemas.as_ref(),
model,
);
measure_request_usage(&shadow, &envelope)
}
fn build_request_envelope_reconciled(
session: &mut Session,
prepared_context: &PreparedContext,
config: &AgentLoopConfig,
tool_schemas: &[ToolSchema],
model: &str,
) -> PreparedRequestEnvelope {
let activated = activated_discoverable_tools(session);
let (stable_frame, stable_prefix_sections) =
build_stable_prompt_frame_with_sections(session, config, tool_schemas, &activated);
let stable_instructions = stable_frame.stable_instructions.clone();
let mut context_blocks = Vec::new();
if let Some(block) = build_active_workflow_context_block(session) {
context_blocks.push(block);
}
if let Some(block) = build_external_memory_context_block(session) {
context_blocks.push(block);
}
if let Some(block) = build_project_resources_context_block(session) {
context_blocks.push(block);
}
if let Some(block) = build_task_list_context_block(session) {
context_blocks.push(block);
}
if let Some(block) = build_goal_context_block(config.active_goal()) {
context_blocks.push(block);
}
if let Some(block) = build_agent_hook_context_block(session) {
context_blocks.push(block);
}
if let Some(block) = build_plan_runtime_context_block(session, config.app_data_dir.as_deref()) {
context_blocks.push(block);
}
if let Some(block) = build_plan_mode_context_block(session) {
context_blocks.push(block);
}
if let Some(block) = build_conversation_summary_context_block(session) {
context_blocks.push(block);
}
let mut system_remainder_messages = Vec::new();
let mut conversation_messages = Vec::new();
for message in &prepared_context.messages {
if matches!(message.role, Role::System) {
if let Some(remainder_message) =
derive_system_remainder_message(message, &stable_instructions)
{
system_remainder_messages.push(remainder_message);
}
} else {
conversation_messages.push(message.clone());
}
}
let section = |name: &str| -> String {
stable_prefix_sections
.iter()
.find(|s| s.name == name)
.map(|s| s.content.clone())
.unwrap_or_default()
};
let tool_guide = section("tool_guide");
let relocate_tool_guide = !tool_guide.trim().is_empty();
let mut system_blocks: Vec<bamboo_domain::PromptBlock> = [
("base", bamboo_domain::ContextBlockType::Base),
(
"core_directives",
bamboo_domain::ContextBlockType::CoreDirectives,
),
]
.into_iter()
.filter_map(|(name, kind)| {
let text = section(name);
let text = if name == "core_directives" {
append_core_agent_directives("", &text)
} else {
text
};
(!text.trim().is_empty()).then(|| {
bamboo_domain::PromptBlock::new(name, kind, text)
.with_stability(bamboo_domain::ContextBlockStability::Stable)
})
})
.collect();
if let Some(last) = system_blocks.last_mut() {
last.cache_anchor = true;
}
let lane_system = system_blocks
.iter()
.map(|b| b.text.as_str())
.collect::<Vec<_>>()
.join("\n\n");
let tool_guide_message = relocate_tool_guide.then(|| {
ContextBlock::new(
ContextBlockType::ToolGuide,
ContextBlockPriority::High,
ContextBlockStability::SessionStable,
"Tool & Connected-Server Guide",
tool_guide,
)
.render_runtime_context_message()
});
let tool_guide_breakpoint_id = tool_guide_message.as_ref().map(|m| m.id.clone());
if !relocate_tool_guide {
system_blocks.clear();
}
[
(
ContextBlockType::Workspace,
ContextBlockPriority::High,
"Project & Workspace",
[section("project"), section("workspace")]
.into_iter()
.filter(|value| !value.trim().is_empty())
.collect::<Vec<_>>()
.join("\n\n"),
),
(
ContextBlockType::InstructionOverlay,
ContextBlockPriority::Critical,
"Project Instructions",
section("instruction"),
),
(
ContextBlockType::SkillContext,
ContextBlockPriority::High,
"Loaded Skills",
section("skill"),
),
(
ContextBlockType::EnvSnapshot,
ContextBlockPriority::High,
"Environment Snapshot",
section("env"),
),
]
.into_iter()
.filter(|(_, _, _, content)| !content.trim().is_empty())
.map(|(block_type, priority, title, content)| {
ContextBlock::new(
block_type,
priority,
ContextBlockStability::SessionStable,
title,
content,
)
})
.for_each(|block| context_blocks.push(block));
let mut stable_prefix_messages = stable_frame.stable_prefix_messages.clone();
if let Some(message) = tool_guide_message {
stable_prefix_messages.push(message);
}
let mut breakpoint_message_ids = Vec::new();
if let Some(id) = tool_guide_breakpoint_id {
breakpoint_message_ids.push(id);
}
let cache_plan = PromptCachePlan {
cache_tools: true,
cache_system: true,
breakpoint_message_ids,
ttl: CacheTtl::Extended,
};
let mut real_transcript = system_remainder_messages;
real_transcript.extend(conversation_messages);
let cache_scope = serde_json::json!({
"model": model,
"system": &lane_system,
"stable_prefix": stable_prefix_messages
.iter()
.map(|message| (&message.role, message.content.as_str()))
.collect::<Vec<_>>(),
"tools": tool_schemas,
});
let cache_scope_sha256 = hex::encode(Sha256::digest(
serde_json::to_vec(&cache_scope).expect("cache scope is serializable"),
));
let ledger = super::context_ledger::reconcile_model_context(
session,
context_blocks,
&real_transcript,
cache_scope_sha256,
prepared_context.truncation_occurred,
);
let ir = PromptIR {
system_text: lane_system,
system_blocks,
segments: vec![
Segment::new(SegmentRole::StablePrefix, stable_prefix_messages),
Segment::new(SegmentRole::ModelTranscript, ledger.transcript),
],
cache: cache_plan,
continuation: None,
};
PreparedRequestEnvelope {
ir,
stable_prefix_sections,
ledger_changed: ledger.changed,
prefix_epoch: ledger.prefix_epoch,
prefix_reset_reason: ledger.reset_reason,
}
}
#[cfg(test)]
fn build_request_envelope(
session: &Session,
prepared_context: &PreparedContext,
config: &AgentLoopConfig,
tool_schemas: &[ToolSchema],
) -> PreparedRequestEnvelope {
let mut shadow = session.clone();
let model = config.model_name.as_deref().unwrap_or(&session.model);
build_request_envelope_reconciled(&mut shadow, prepared_context, config, tool_schemas, model)
}
const SESSION_REQUEST_RENDER_KEY: &str = "llm_request_render";
#[derive(Debug, Clone, serde::Serialize)]
pub(crate) struct RequestRenderObservability {
pub wire: &'static str,
pub system_block_count: usize,
pub system_chars: usize,
pub stable_prefix_messages: usize,
pub model_transcript_messages: usize,
pub dynamic_context_messages: usize,
pub conversation_messages: usize,
pub volatile_context_messages: usize,
pub request_message_count: usize,
pub tool_count: usize,
pub cache_system: bool,
pub cache_tools: bool,
pub cache_breakpoints: usize,
pub cache_ttl: &'static str,
pub prefix_epoch: u64,
pub prefix_reset_reason: Option<&'static str>,
}
impl RequestRenderObservability {
fn log(&self, session_id: &str) {
tracing::info!(
"[{}] LLM request render: wire={} system_blocks={} system_chars={} stable_prefix_msgs={} model_transcript_msgs={} dynamic_ctx_msgs={} conversation_msgs={} volatile_ctx_msgs={} request_msgs={} tools={} cache(system={}, tools={}, breakpoints={}, ttl={}) prefix_epoch={} prefix_reset_reason={}",
session_id,
self.wire,
self.system_block_count,
self.system_chars,
self.stable_prefix_messages,
self.model_transcript_messages,
self.dynamic_context_messages,
self.conversation_messages,
self.volatile_context_messages,
self.request_message_count,
self.tool_count,
self.cache_system,
self.cache_tools,
self.cache_breakpoints,
self.cache_ttl,
self.prefix_epoch,
self.prefix_reset_reason.unwrap_or("none"),
);
}
}
struct LlmRequestPlan {
request_options: LLMRequestOptions,
render: RequestRenderObservability,
}
fn cache_ttl_label(ttl: CacheTtl) -> &'static str {
if ttl == CacheTtl::Extended {
"1h"
} else {
"5m"
}
}
fn plan_llm_request(
envelope: &PreparedRequestEnvelope,
session_id: &str,
reasoning_effort: Option<ReasoningEffort>,
tool_count: usize,
required_tool: Option<&str>,
) -> LlmRequestPlan {
let is_continuation = envelope.ir.continuation.is_some();
let mut responses_options = engine_responses_policy();
responses_options.prompt_cache_key =
agent_loop_prompt_cache_key(Some(session_id), Some(AGENT_LOOP_REQUEST_PURPOSE));
responses_options.prefix_epoch = Some(envelope.prefix_epoch);
responses_options.prefix_reset_reason = envelope.prefix_reset_reason;
let request_options = LLMRequestOptions {
session_id: Some(session_id.to_string()),
reasoning_effort,
parallel_tool_calls: Some(required_tool.is_none()),
required_tool: required_tool.map(str::to_string),
responses: Some(responses_options),
request_purpose: Some(AGENT_LOOP_REQUEST_PURPOSE.to_string()),
cache: Some(envelope.ir.cache.clone()),
};
let render = RequestRenderObservability {
wire: if is_continuation {
"responses_continuation"
} else {
"model_transcript"
},
system_block_count: envelope.ir.system_blocks.len(),
system_chars: envelope.ir.system_field().len(),
stable_prefix_messages: envelope.ir.run(SegmentRole::StablePrefix).len(),
model_transcript_messages: envelope.ir.run(SegmentRole::ModelTranscript).len(),
dynamic_context_messages: envelope.ir.run(SegmentRole::DynamicContext).len(),
conversation_messages: envelope.ir.run(SegmentRole::Conversation).len(),
volatile_context_messages: envelope.ir.run(SegmentRole::VolatileTail).len(),
request_message_count: if is_continuation {
envelope.ir.continuation_delta().len()
} else {
envelope.ir.flatten().len()
},
tool_count,
cache_system: envelope.ir.cache.cache_system,
cache_tools: envelope.ir.cache.cache_tools,
cache_breakpoints: envelope.ir.cache.breakpoint_message_ids.len(),
cache_ttl: cache_ttl_label(envelope.ir.cache.ttl),
prefix_epoch: envelope.prefix_epoch,
prefix_reset_reason: envelope
.prefix_reset_reason
.map(bamboo_domain::ModelContextResetReason::as_str),
};
LlmRequestPlan {
request_options,
render,
}
}
fn persist_request_render_metadata(session: &mut Session, render: &RequestRenderObservability) {
if let Ok(value) = serde_json::to_string(render) {
session
.metadata
.insert(SESSION_REQUEST_RENDER_KEY.to_string(), value);
}
}
pub(super) async fn execute_llm_stream(
session: &mut Session,
config: &AgentLoopConfig,
llm: &Arc<dyn LLMProvider>,
prepared_context: &PreparedContext,
tool_schemas: &[ToolSchema],
frame: &LlmStreamFrame<'_>,
) -> Result<(crate::runtime::stream::handler::StreamHandlingOutput, u128), AgentError> {
let event_tx = frame.event_tx;
let cancel_token = frame.cancel_token;
let max_context_tokens = frame.max_context_tokens;
let max_output_tokens = frame.max_output_tokens;
let model = frame.model;
let provider_name = frame.provider_name;
let provider_type = frame.provider_type;
let reasoning_effort = frame.reasoning_effort;
let session_id = frame.session_id;
let llm_started_at = std::time::Instant::now();
let required_tool = required_tool_for_session(session);
let effective_tool_schemas = effective_tool_schemas(session, tool_schemas);
let tool_schemas = effective_tool_schemas.as_ref();
let responses_policy = engine_responses_policy();
let continuation_enabled = responses_continuation_enabled(&responses_policy, provider_type);
let previous_response_id = if continuation_enabled {
session_previous_response_id(session).map(str::to_string)
} else {
None
};
let previous_model_context_state = session.model_context_state.clone();
let mut prepared_envelope =
build_request_envelope_reconciled(session, prepared_context, config, tool_schemas, model);
let final_usage = measure_request_usage(session, &prepared_envelope);
let request_input_limit = max_context_tokens.saturating_sub(max_output_tokens);
if final_usage.input_tokens > request_input_limit
|| final_usage.ledger_rendered_bytes > MAX_MODEL_CONTEXT_RENDERED_BYTES
{
session.model_context_state = previous_model_context_state;
return Err(AgentError::Budget(format!(
"final PromptIR message input exceeds ledger-safe limits: input_tokens={}, input_limit={request_input_limit}, ledger_bytes={}, ledger_byte_limit={MAX_MODEL_CONTEXT_RENDERED_BYTES}",
final_usage.input_tokens, final_usage.ledger_rendered_bytes,
)));
}
if prepared_envelope.ledger_changed {
if let Some(persistence) = config.persistence.as_ref() {
if let Err(error) = persistence.checkpoint_runtime_session(session).await {
session.model_context_state = previous_model_context_state;
return Err(AgentError::LLM(format!(
"model-context ledger checkpoint failed before provider dispatch: {error}"
)));
}
}
}
super::prefix_drift::record_prefix_drift(
session,
config.app_data_dir.as_deref(),
&prepared_envelope.stable_prefix_sections,
);
if let Some(response_id) = previous_response_id.as_deref() {
let last_committed_assistant_id = prepared_envelope
.ir
.run(SegmentRole::ModelTranscript)
.iter()
.rev()
.find(|message| matches!(message.role, Role::Assistant))
.map(|message| message.id.clone());
prepared_envelope.ir.continuation = Some(Continuation {
previous_response_id: response_id.to_string(),
last_committed_assistant_id,
});
}
let planned = plan_llm_request(
&prepared_envelope,
session_id,
reasoning_effort,
tool_schemas.len(),
required_tool,
);
if !continuation_enabled {
tracing::debug!(
"[{}] Responses API previous_response_id disabled (store={:?}, provider={})",
session_id,
responses_policy.store,
provider_name.unwrap_or("unknown")
);
}
planned.render.log(session_id);
persist_request_render_metadata(session, &planned.render);
let timeout_context = crate::runtime::stream::handler::StreamTimeoutContext::new(
config.stream_timeout,
provider_name,
Some(model),
)
.allow_turn_retry_before_semantic_output()
.begin_request();
let stream = crate::runtime::stream::handler::await_stream_bootstrap(
llm.chat_stream_ir(
&prepared_envelope.ir,
tool_schemas,
Some(max_output_tokens),
model,
Some(&planned.request_options),
),
cancel_token,
session_id,
&timeout_context,
)
.await?
.map_err(|error| {
let message = format_provider_error(error);
if is_llm_overflow_error(&message) {
AgentError::LLMOverflow(message)
} else {
AgentError::LLM(message)
}
})?;
let usage = TokenBudgetUsage {
system_tokens: prepared_context.token_usage.system_tokens,
summary_tokens: prepared_context.token_usage.summary_tokens,
window_tokens: prepared_context.token_usage.window_tokens,
total_tokens: prepared_context.token_usage.total_tokens,
max_context_tokens,
budget_limit: prepared_context.token_usage.budget_limit,
truncation_occurred: prepared_context.truncation_occurred,
segments_removed: prepared_context.segments_removed,
prompt_cached_tool_outputs: prepared_context.prompt_cached_tool_outputs,
prompt_cached_tool_tokens_saved: prepared_context.prompt_cached_tool_tokens_saved,
thinking_tokens: 0,
cache_read_input_tokens: 0,
};
session.token_usage = Some(usage.clone());
let budget_event = AgentEvent::TokenBudgetUpdated { usage };
if let Err(error) = event_tx.send(budget_event).await {
tracing::warn!(
"[{}] Failed to send token budget event: {}",
session_id,
error
);
}
let stream_output_result =
if crate::runtime::runner::session_setup::skill_context::explicit_activation_pending(
session,
) {
crate::runtime::stream::handler::consume_llm_stream_silent_with_context_and_partial(
stream,
cancel_token,
session_id,
&timeout_context,
)
.await
} else {
crate::runtime::stream::handler::consume_llm_stream_with_context_and_partial(
stream,
event_tx,
cancel_token,
session_id,
&timeout_context,
)
.await
};
let stream_output = match stream_output_result {
Ok(output) => output,
Err(failure) => {
let appended = append_interrupted_assistant_output(
session,
*failure.partial_output,
&failure.error,
);
if appended {
tracing::warn!(
"[{}] Preserved interrupted partial assistant output in the live transcript",
session_id
);
}
return Err(failure.error);
}
};
if let Some(ref mut usage) = session.token_usage {
usage.thinking_tokens = stream_output.thinking_tokens as u32;
usage.cache_read_input_tokens = stream_output.cache_read_input_tokens as u32;
}
if let Some(usage) = session.token_usage.clone() {
let final_budget_event = AgentEvent::TokenBudgetUpdated { usage };
if let Err(error) = event_tx.send(final_budget_event).await {
tracing::warn!(
"[{}] Failed to send final token budget event: {}",
session_id,
error
);
}
}
let cache_write_input_tokens = stream_output
.provider_usage
.and_then(|usage| usage.cache_write_input_tokens)
.unwrap_or(0);
if stream_output.cache_creation_input_tokens > 0
|| stream_output.cache_read_input_tokens > 0
|| cache_write_input_tokens > 0
{
tracing::info!(
"[{}] Provider prompt cache: creation={}, read={}, write={}, output={}, thinking={}",
session_id,
stream_output.cache_creation_input_tokens,
stream_output.cache_read_input_tokens,
cache_write_input_tokens,
stream_output.output_tokens,
stream_output.thinking_tokens,
);
}
if cfg!(any(debug_assertions, feature = "token-usage-log")) {
if let Some(persistence) = config.persistence.as_ref() {
let record = crate::token_usage_log::TokenUsageRecord::new(
chrono::Utc::now().to_rfc3339(),
session_id,
model,
provider_name.unwrap_or(""),
session.messages.len(),
session.token_usage.as_ref(),
stream_output.cache_creation_input_tokens,
stream_output.cache_read_input_tokens,
cache_write_input_tokens,
stream_output.input_tokens,
stream_output.output_tokens,
stream_output.thinking_tokens,
);
match record.to_json_line() {
Ok(line) => {
if let Err(error) = persistence
.append_token_usage_record(session_id, &line)
.await
{
tracing::warn!(
"[{}] Failed to append token-usage record: {}",
session_id,
error
);
}
}
Err(error) => {
tracing::warn!(
"[{}] Failed to serialize token-usage record: {}",
session_id,
error
);
}
}
}
}
if continuation_enabled {
if let Some(response_id) = stream_output
.response_id
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
{
session.metadata.insert(
SESSION_RESPONSES_PREVIOUS_RESPONSE_ID_KEY.to_string(),
response_id.to_string(),
);
} else {
session
.metadata
.remove(SESSION_RESPONSES_PREVIOUS_RESPONSE_ID_KEY);
}
} else {
session
.metadata
.remove(SESSION_RESPONSES_PREVIOUS_RESPONSE_ID_KEY);
}
let llm_duration = llm_started_at.elapsed().as_millis();
Ok((stream_output, llm_duration))
}
#[cfg(test)]
mod tests;