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::{
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::{
assemble_prompt_envelope, build_active_workflow_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_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, Role, Session,
};
use bamboo_compression::PreparedContext;
use bamboo_domain::ReasoningEffort;
use bamboo_llm::provider::ResponsesRequestOptions;
use bamboo_llm::{
CacheTtl, Continuation, LLMProvider, LLMRequestOptions, PromptCachePlan, PromptIR, Segment,
SegmentRole,
};
use bamboo_tools::exposure::activated_discoverable_tools;
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 -->";
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 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>,
}
fn build_request_envelope(
session: &Session,
prepared_context: &PreparedContext,
config: &AgentLoopConfig,
tool_schemas: &[ToolSchema],
) -> 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 front_blocks = Vec::new();
let mut volatile_blocks = Vec::new();
if let Some(block) = build_active_workflow_context_block(session) {
front_blocks.push(block);
}
if let Some(block) = build_external_memory_context_block(session) {
volatile_blocks.push(block);
}
if let Some(block) = build_task_list_context_block(session) {
volatile_blocks.push(block);
}
if let Some(block) = build_goal_context_block(config.active_goal()) {
volatile_blocks.push(block);
}
if let Some(block) = build_plan_runtime_context_block(session, config.app_data_dir.as_deref()) {
volatile_blocks.push(block);
}
if let Some(block) = build_plan_mode_context_block(session) {
volatile_blocks.push(block);
}
if let Some(block) = build_conversation_summary_context_block(session) {
front_blocks.push(block);
}
let volatile_context_messages: Vec<Message> = volatile_blocks
.iter()
.map(|block| block.render_runtime_context_message())
.collect();
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 conversation_breakpoint_id = conversation_messages
.last()
.or_else(|| system_remainder_messages.last())
.map(|message| message.id.clone());
let envelope = assemble_prompt_envelope(stable_frame, front_blocks);
let summary_breakpoint_id = envelope
.dynamic_context_messages
.last()
.map(|message| message.id.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,
),
("env", bamboo_domain::ContextBlockType::EnvSnapshot),
]
.into_iter()
.filter_map(|(name, kind)| {
let text = section(name);
(!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 = if relocate_tool_guide {
system_blocks
.iter()
.map(|b| b.text.as_str())
.collect::<Vec<_>>()
.join("\n\n")
} else {
system_blocks.clear();
envelope.stable_instructions.clone()
};
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());
let session_context_messages: Vec<Message> = if relocate_tool_guide {
[
(
ContextBlockType::Workspace,
ContextBlockPriority::High,
"Workspace",
section("workspace"),
),
(
ContextBlockType::InstructionOverlay,
ContextBlockPriority::Critical,
"Project Instructions",
section("instruction"),
),
(
ContextBlockType::SkillContext,
ContextBlockPriority::High,
"Loaded Skills",
section("skill"),
),
]
.into_iter()
.filter(|(_, _, _, content)| !content.trim().is_empty())
.map(|(block_type, priority, title, content)| {
ContextBlock::new(
block_type,
priority,
ContextBlockStability::SessionStable,
title,
content,
)
.render_runtime_context_message()
})
.collect()
} else {
Vec::new()
};
let mut stable_prefix_messages = envelope.stable_prefix_messages.clone();
if let Some(message) = tool_guide_message {
stable_prefix_messages.push(message);
}
stable_prefix_messages.extend(session_context_messages);
let mut breakpoint_message_ids = Vec::new();
if let Some(id) = tool_guide_breakpoint_id {
breakpoint_message_ids.push(id);
} else if let Some(id) = summary_breakpoint_id {
breakpoint_message_ids.push(id);
}
if let Some(id) = conversation_breakpoint_id {
breakpoint_message_ids.push(id);
}
let cache_plan = PromptCachePlan {
cache_tools: true,
cache_system: true,
breakpoint_message_ids,
ttl: CacheTtl::Extended,
};
let ir = PromptIR {
system_text: lane_system,
system_blocks,
segments: vec![
Segment::new(SegmentRole::StablePrefix, stable_prefix_messages),
Segment::new(
SegmentRole::DynamicContext,
envelope.dynamic_context_messages.clone(),
),
Segment::new(
SegmentRole::SystemRemainder,
system_remainder_messages.clone(),
),
Segment::new(SegmentRole::Conversation, conversation_messages.clone()),
Segment::new(SegmentRole::VolatileTail, volatile_context_messages.clone()),
],
cache: cache_plan,
continuation: None,
};
PreparedRequestEnvelope {
ir,
stable_prefix_sections,
}
}
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 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,
}
impl RequestRenderObservability {
fn log(&self, session_id: &str) {
tracing::info!(
"[{}] LLM request render: wire={} system_blocks={} system_chars={} stable_prefix_msgs={} dynamic_ctx_msgs={} conversation_msgs={} volatile_ctx_msgs={} request_msgs={} tools={} cache(system={}, tools={}, breakpoints={}, ttl={})",
session_id,
self.wire,
self.system_block_count,
self.system_chars,
self.stable_prefix_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,
);
}
}
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 responses_options = engine_responses_policy();
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".to_string()),
cache: Some(envelope.ir.cache.clone()),
};
let render = RequestRenderObservability {
wire: if is_continuation {
"responses_continuation"
} else {
"lanes"
},
system_block_count: envelope.ir.system_blocks.len(),
system_chars: envelope.ir.system_field().len(),
stable_prefix_messages: envelope.ir.run(SegmentRole::StablePrefix).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),
};
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 =
crate::runtime::runner::session_setup::skill_context::explicit_activation_pending(session)
.then_some("load_skill");
let restricted_tool_schemas;
let tool_schemas = if let Some(required_tool) = required_tool {
restricted_tool_schemas = tool_schemas
.iter()
.filter(|schema| schema.function.name == required_tool)
.cloned()
.collect::<Vec<_>>();
restricted_tool_schemas.as_slice()
} else {
tool_schemas
};
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 mut prepared_envelope =
build_request_envelope(session, prepared_context, config, tool_schemas);
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::Conversation)
.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 stream = llm
.chat_stream_ir(
&prepared_envelope.ir,
tool_schemas,
Some(max_output_tokens),
model,
Some(&planned.request_options),
)
.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 timeout_context = crate::runtime::stream::handler::StreamTimeoutContext::new(
config.stream_timeout,
provider_name,
Some(model),
);
let stream_output =
if crate::runtime::runner::session_setup::skill_context::explicit_activation_pending(
session,
) {
crate::runtime::stream::handler::consume_llm_stream_silent_with_context(
stream,
cancel_token,
session_id,
&timeout_context,
)
.await?
} else {
crate::runtime::stream::handler::consume_llm_stream_with_context(
stream,
event_tx,
cancel_token,
session_id,
&timeout_context,
)
.await?
};
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
);
}
}
if stream_output.cache_creation_input_tokens > 0 || stream_output.cache_read_input_tokens > 0 {
tracing::info!(
"[{}] Anthropic prompt cache: creation={}, read={}, output={}, thinking={}",
session_id,
stream_output.cache_creation_input_tokens,
stream_output.cache_read_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,
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;