| use crate::context::GuardianContextMode; |
| use std::sync::Arc; |
| use std::time::Instant; |
|
|
| use crate::Prompt; |
| use crate::client::ModelClientSession; |
| use crate::client_common::ResponseEvent; |
| use crate::context::CompactionSummary; |
| use crate::context::ContextualUserFragment; |
| use crate::context::world_state::WorldState; |
| use crate::hook_runtime::PostCompactHookOutcome; |
| use crate::hook_runtime::PreCompactHookOutcome; |
| use crate::hook_runtime::run_post_compact_hooks; |
| use crate::hook_runtime::run_pre_compact_hooks; |
| use crate::responses_metadata::CodexResponsesMetadata; |
| use crate::responses_metadata::CompactionTurnMetadata; |
| use crate::session::RequestEffortUsage; |
| use crate::session::session::Session; |
| use crate::session::step_context::StepContext; |
| use crate::session::turn::get_last_assistant_message_from_turn; |
| use crate::session::turn_context::TurnContext; |
| use crate::state::AutoCompactWindowIds; |
| use crate::util::backoff; |
| use codex_analytics::CodexCompactionEvent; |
| use codex_analytics::CompactionImplementation; |
| use codex_analytics::CompactionPhase; |
| use codex_analytics::CompactionReason; |
| use codex_analytics::CompactionStatus; |
| use codex_analytics::CompactionStrategy; |
| use codex_analytics::CompactionTrigger; |
| use codex_analytics::now_unix_seconds; |
| use codex_context_fragments::AnnotatedContent; |
| use codex_context_fragments::set_annotated_content; |
| use codex_history::CodexHarnessMetadata; |
| use codex_history::ResponseItemEnvelope; |
| use codex_protocol::ResponseItemId; |
| use codex_protocol::error::CodexErr; |
| use codex_protocol::error::CodexErrorDetails; |
| use codex_protocol::error::Result as CodexResult; |
| use codex_protocol::items::ContextCompactionItem; |
| use codex_protocol::items::TurnItem; |
| use codex_protocol::models::AgentMessageInputContent; |
| use codex_protocol::models::ContentItem; |
| use codex_protocol::models::ContentItemKind; |
| use codex_protocol::models::InternalChatMessageMetadataPassthrough; |
| use codex_protocol::models::ResponseInputItem; |
| use codex_protocol::models::ResponseItem; |
| use codex_protocol::protocol::EventMsg; |
| use codex_protocol::protocol::WarningEvent; |
| use codex_protocol::user_input::UserInput; |
| use codex_rollout_trace::InferenceTraceContext; |
| use codex_utils_output_truncation::TruncationPolicy; |
| use codex_utils_output_truncation::approx_token_count; |
| use codex_utils_output_truncation::truncate_text; |
| use futures::prelude::*; |
| use tracing::error; |
|
|
| pub use codex_prompts::SUMMARIZATION_PROMPT; |
| pub use codex_prompts::SUMMARY_PREFIX; |
| const COMPACT_USER_MESSAGE_MAX_TOKENS: usize = 20_000; |
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| pub(crate) enum InitialContextInjection { |
| BeforeLastUserMessage { |
| world_state: Arc<WorldState>, |
| step_context: Arc<StepContext>, |
| }, |
| DoNotInject, |
| } |
|
|
| |
| |
| |
| |
| pub(crate) struct CompactedHistoryMetadata { |
| pub(crate) message: String, |
| pub(crate) window_number: u64, |
| pub(crate) window_ids: AutoCompactWindowIds, |
| pub(crate) compaction_response_id: Option<String>, |
| pub(crate) compaction_model_hash: Option<String>, |
| pub(crate) reviewer_compaction_hash: Option<String>, |
| } |
|
|
| pub(crate) async fn build_compaction_initial_context( |
| sess: &Session, |
| initial_context_injection: &InitialContextInjection, |
| ) -> (Vec<ResponseItemEnvelope>, Option<Arc<WorldState>>) { |
| |
| match initial_context_injection { |
| InitialContextInjection::BeforeLastUserMessage { |
| world_state, |
| step_context, |
| } => { |
| let items = sess |
| .build_initial_context_with_world_state(step_context, world_state.as_ref()) |
| .await; |
| ( |
| items.into_iter().map(ResponseItemEnvelope::new).collect(), |
| Some(Arc::clone(world_state)), |
| ) |
| } |
| InitialContextInjection::DoNotInject => (Vec::new(), None), |
| } |
| } |
|
|
| pub(crate) async fn run_inline_auto_compact_task( |
| sess: Arc<Session>, |
| turn_context: Arc<TurnContext>, |
| initial_context_injection: InitialContextInjection, |
| reason: CompactionReason, |
| phase: CompactionPhase, |
| ) -> CodexResult<()> { |
| let prompt = turn_context |
| .config |
| .compact_prompt |
| .as_deref() |
| .unwrap_or(SUMMARIZATION_PROMPT) |
| .to_string(); |
| let input = vec![UserInput::Text { |
| text: prompt, |
| |
| text_elements: Vec::new(), |
| }]; |
|
|
| run_compact_task_inner( |
| sess, |
| turn_context, |
| input, |
| initial_context_injection, |
| CompactionTrigger::Auto, |
| reason, |
| phase, |
| ) |
| .await?; |
| Ok(()) |
| } |
|
|
| pub(crate) async fn run_compact_task( |
| sess: Arc<Session>, |
| turn_context: Arc<TurnContext>, |
| input: Vec<UserInput>, |
| ) -> CodexResult<()> { |
| sess.emit_turn_started(&turn_context).await; |
| run_compact_task_inner( |
| sess.clone(), |
| turn_context, |
| input, |
| InitialContextInjection::DoNotInject, |
| CompactionTrigger::Manual, |
| CompactionReason::UserRequested, |
| CompactionPhase::StandaloneTurn, |
| ) |
| .await?; |
| Ok(()) |
| } |
|
|
| async fn run_compact_task_inner( |
| sess: Arc<Session>, |
| turn_context: Arc<TurnContext>, |
| input: Vec<UserInput>, |
| initial_context_injection: InitialContextInjection, |
| trigger: CompactionTrigger, |
| reason: CompactionReason, |
| phase: CompactionPhase, |
| ) -> CodexResult<()> { |
| let compaction_metadata = |
| CompactionTurnMetadata::new(trigger, reason, CompactionImplementation::Responses, phase); |
| let attempt = CompactionAnalyticsAttempt::begin( |
| sess.as_ref(), |
| turn_context.as_ref(), |
| trigger, |
| reason, |
| CompactionImplementation::Responses, |
| phase, |
| ) |
| .await; |
| let pre_compact_outcome = run_pre_compact_hooks(&sess, &turn_context, trigger).await; |
| match pre_compact_outcome { |
| PreCompactHookOutcome::Continue => {} |
| PreCompactHookOutcome::Stopped => { |
| let error = CodexErr::TurnAborted; |
| attempt |
| .track( |
| sess.as_ref(), |
| CompactionStatus::Interrupted, |
| Some(&error), |
| CompactionAnalyticsDetails::default(), |
| ) |
| .await; |
| return Err(error); |
| } |
| } |
| let result = run_compact_task_inner_impl( |
| Arc::clone(&sess), |
| Arc::clone(&turn_context), |
| input, |
| initial_context_injection, |
| compaction_metadata, |
| ) |
| .await; |
| let status = compaction_status_from_result(&result); |
| let codex_error = result.as_ref().err(); |
| if result.is_ok() { |
| let post_compact_outcome = run_post_compact_hooks(&sess, &turn_context, trigger).await; |
| if let PostCompactHookOutcome::Stopped = post_compact_outcome { |
| attempt |
| .track( |
| sess.as_ref(), |
| status, |
| codex_error, |
| CompactionAnalyticsDetails::default(), |
| ) |
| .await; |
| return Err(CodexErr::TurnAborted); |
| } |
| } |
| attempt |
| .track( |
| sess.as_ref(), |
| status, |
| codex_error, |
| CompactionAnalyticsDetails::default(), |
| ) |
| .await; |
| result.map(|_| ()) |
| } |
|
|
| async fn run_compact_task_inner_impl( |
| sess: Arc<Session>, |
| turn_context: Arc<TurnContext>, |
| input: Vec<UserInput>, |
| initial_context_injection: InitialContextInjection, |
| compaction_metadata: CompactionTurnMetadata, |
| ) -> CodexResult<String> { |
| let compaction_item = TurnItem::ContextCompaction(ContextCompactionItem::new()); |
| sess.emit_turn_item_started(&turn_context, &compaction_item) |
| .await; |
| let initial_input_for_turn: ResponseInputItem = ResponseInputItem::from(input); |
|
|
| let mut history = sess.clone_history().await; |
| history.record_items( |
| &[initial_input_for_turn.into()], |
| turn_context.model_info().truncation_policy.into(), |
| ); |
|
|
| let max_retries = turn_context.provider.info().stream_max_retries(); |
| let mut retries = 0; |
| let mut client_session = sess.services.model_client.new_session(); |
| |
| |
| |
| let responses_metadata = sess |
| .compaction_responses_metadata(turn_context.as_ref(), compaction_metadata) |
| .await; |
|
|
| let compaction_response_id = loop { |
| |
| let mut turn_input = history |
| .clone() |
| .for_prompt(&turn_context.model_info().input_modalities); |
| sess.services |
| .executed_tool_calls |
| .strip_disabled_direct_metadata(&mut turn_input); |
| let turn_input_len = turn_input.len(); |
| let prompt = Prompt { |
| input: turn_input, |
| base_instructions: sess.get_prompt_base_instructions().await, |
| ..Default::default() |
| }; |
| let attempt_result = drain_to_completed( |
| &sess, |
| turn_context.as_ref(), |
| &mut client_session, |
| &responses_metadata, |
| &prompt, |
| ) |
| .await; |
|
|
| match attempt_result { |
| Ok(response_id) => { |
| break response_id; |
| } |
| Err(err) |
| if matches!( |
| err.details(), |
| CodexErrorDetails::Interrupted | CodexErrorDetails::TurnAborted |
| ) => |
| { |
| return Err(err); |
| } |
| Err(e) if matches!(e.details(), CodexErrorDetails::SessionBudgetExceeded) => { |
| sess.track_turn_codex_error(turn_context.as_ref(), &e); |
| |
| if !matches!(compaction_metadata.phase(), CompactionPhase::PreTurn) { |
| let event = EventMsg::Error(e.to_error_event( None)); |
| sess.send_event(&turn_context, event).await; |
| } |
| return Err(e); |
| } |
| Err(e) if matches!(e.details(), CodexErrorDetails::ContextWindowExceeded) => { |
| if turn_input_len > 1 { |
| |
| error!( |
| "Context window exceeded while compacting; removing oldest history item. Error: {e}" |
| ); |
| history.remove_first_item(); |
| retries = 0; |
| continue; |
| } |
| sess.set_total_tokens_full(turn_context.as_ref()).await; |
| sess.track_turn_codex_error(turn_context.as_ref(), &e); |
| if !matches!(compaction_metadata.phase(), CompactionPhase::PreTurn) { |
| let event = EventMsg::Error(e.to_error_event( None)); |
| sess.send_event(&turn_context, event).await; |
| } |
| return Err(e); |
| } |
| Err(e) => { |
| if retries < max_retries { |
| retries += 1; |
| let delay = backoff(retries); |
| sess.notify_stream_error( |
| turn_context.as_ref(), |
| format!("Reconnecting... {retries}/{max_retries}"), |
| e, |
| ) |
| .await; |
| tokio::time::sleep(delay).await; |
| continue; |
| } else { |
| sess.track_turn_codex_error(turn_context.as_ref(), &e); |
| if !matches!(compaction_metadata.phase(), CompactionPhase::PreTurn) { |
| let event = EventMsg::Error(e.to_error_event( None)); |
| sess.send_event(&turn_context, event).await; |
| } |
| return Err(e); |
| } |
| } |
| } |
| }; |
|
|
| let history_snapshot = sess.clone_history().await; |
| let history_items = history_snapshot.annotated_items(); |
| let summary_suffix = |
| get_last_assistant_message_from_turn(history_snapshot.raw_items()).unwrap_or_default(); |
| let summary_text = format!("{SUMMARY_PREFIX}\n{summary_suffix}"); |
| let identity = if sess.guardian_context_mode == GuardianContextMode::ThreadOwned { |
| CompactedMessageIdentity::Preserve |
| } else { |
| CompactedMessageIdentity::Regenerate |
| }; |
| let user_messages = collect_annotated_user_messages(history_items, identity); |
|
|
| let mut new_history = build_compacted_history(Vec::new(), &user_messages, &summary_text); |
| if let Some(summary_item) = new_history.last_mut() { |
| |
| |
| summary_item.set_turn_id_if_missing(&turn_context.sub_id); |
| } |
| let (window_number, window_ids) = sess.advance_auto_compact_window().await; |
|
|
| let (initial_context, world_state_baseline) = |
| build_compaction_initial_context(sess.as_ref(), &initial_context_injection).await; |
| if !initial_context.is_empty() { |
| new_history = |
| insert_initial_context_before_last_real_user_or_summary(new_history, initial_context); |
| } |
| let reference_context_item = match initial_context_injection { |
| InitialContextInjection::DoNotInject => None, |
| InitialContextInjection::BeforeLastUserMessage { step_context, .. } => { |
| Some(step_context.to_turn_context_item()) |
| } |
| }; |
| sess.replace_compacted_history( |
| new_history, |
| reference_context_item, |
| world_state_baseline, |
| CompactedHistoryMetadata { |
| message: summary_text, |
| window_number, |
| window_ids, |
| compaction_response_id: Some(compaction_response_id), |
| compaction_model_hash: turn_context.model_info().comp_hash.clone(), |
| reviewer_compaction_hash: None, |
| }, |
| ) |
| .await; |
| sess.recompute_token_usage(&turn_context).await; |
|
|
| sess.emit_turn_item_completed(&turn_context, compaction_item) |
| .await; |
| let warning = EventMsg::Warning(WarningEvent { |
| message: "Heads up: Long threads and multiple compactions can cause the model to be less accurate. Start a new thread when possible to keep threads small and targeted.".to_string(), |
| }); |
| sess.send_event(&turn_context, warning).await; |
| Ok(summary_suffix) |
| } |
|
|
| pub(crate) struct CompactionAnalyticsAttempt { |
| thread_id: String, |
| turn_id: String, |
| trigger: CompactionTrigger, |
| reason: CompactionReason, |
| implementation: CompactionImplementation, |
| phase: CompactionPhase, |
| active_context_tokens_before: i64, |
| started_at: u64, |
| start_instant: Instant, |
| } |
|
|
| #[derive(Clone, Copy, Default)] |
| pub(crate) struct CompactionAnalyticsDetails { |
| pub(crate) active_context_tokens_before: Option<i64>, |
| pub(crate) retained_image_count: Option<usize>, |
| pub(crate) compaction_summary_tokens: Option<i64>, |
| pub(crate) cached_input_tokens: Option<i64>, |
| pub(crate) cache_write_input_tokens: Option<i64>, |
| } |
|
|
| impl CompactionAnalyticsAttempt { |
| pub(crate) async fn begin( |
| sess: &Session, |
| turn_context: &TurnContext, |
| trigger: CompactionTrigger, |
| reason: CompactionReason, |
| implementation: CompactionImplementation, |
| phase: CompactionPhase, |
| ) -> Self { |
| let active_context_tokens_before = sess.get_total_token_usage().await; |
| Self { |
| thread_id: sess.thread_id.to_string(), |
| turn_id: turn_context.sub_id.clone(), |
| trigger, |
| reason, |
| implementation, |
| phase, |
| active_context_tokens_before, |
| started_at: now_unix_seconds(), |
| start_instant: Instant::now(), |
| } |
| } |
|
|
| pub(crate) async fn track( |
| self, |
| sess: &Session, |
| status: CompactionStatus, |
| codex_error: Option<&CodexErr>, |
| details: CompactionAnalyticsDetails, |
| ) { |
| let CompactionAnalyticsDetails { |
| active_context_tokens_before, |
| retained_image_count, |
| compaction_summary_tokens, |
| cached_input_tokens, |
| cache_write_input_tokens, |
| } = details; |
| let active_context_tokens_before = |
| active_context_tokens_before.unwrap_or(self.active_context_tokens_before); |
| let active_context_tokens_after = sess.get_total_token_usage().await; |
| sess.services |
| .analytics_events_client |
| .track_compaction(CodexCompactionEvent { |
| thread_id: self.thread_id, |
| turn_id: self.turn_id, |
| trigger: self.trigger, |
| reason: self.reason, |
| implementation: self.implementation, |
| phase: self.phase, |
| strategy: CompactionStrategy::Memento, |
| status, |
| codex_error_kind: codex_error.map(Into::into), |
| codex_error_http_status_code: codex_error |
| .and_then(CodexErr::http_status_code_value), |
| active_context_tokens_before, |
| active_context_tokens_after, |
| retained_image_count, |
| compaction_summary_tokens, |
| cached_input_tokens, |
| cache_write_input_tokens, |
| started_at: self.started_at, |
| completed_at: now_unix_seconds(), |
| duration_ms: Some( |
| u64::try_from(self.start_instant.elapsed().as_millis()).unwrap_or(u64::MAX), |
| ), |
| }); |
| } |
| } |
|
|
| pub(crate) fn compaction_status_from_result<T>(result: &CodexResult<T>) -> CompactionStatus { |
| match result { |
| Ok(_) => CompactionStatus::Completed, |
| Err(err) |
| if matches!( |
| err.details(), |
| CodexErrorDetails::Interrupted | CodexErrorDetails::TurnAborted |
| ) => |
| { |
| CompactionStatus::Interrupted |
| } |
| Err(_) => CompactionStatus::Failed, |
| } |
| } |
|
|
| pub fn content_items_to_text(content: &[ContentItem]) -> Option<String> { |
| let mut pieces = Vec::new(); |
| for item in content { |
| match item { |
| ContentItem::InputText { text } | ContentItem::OutputText { text } => { |
| if !text.is_empty() { |
| pieces.push(text.as_str()); |
| } |
| } |
| ContentItem::InputImage { .. } | ContentItem::InputAudio { .. } => {} |
| } |
| } |
| if pieces.is_empty() { |
| None |
| } else { |
| Some(pieces.join("\n")) |
| } |
| } |
|
|
| #[derive(Clone, Debug, PartialEq)] |
| pub(crate) struct CompactedUserMessage { |
| |
| |
| id: Option<ResponseItemId>, |
| message: String, |
| internal_chat_message_metadata_passthrough: Option<InternalChatMessageMetadataPassthrough>, |
| harness_metadata: Option<CodexHarnessMetadata>, |
| } |
|
|
| #[cfg(test)] |
| pub(crate) fn collect_user_messages(items: &[ResponseItem]) -> Vec<CompactedUserMessage> { |
| items |
| .iter() |
| .filter_map(|item| compacted_user_message(item, None)) |
| .collect() |
| } |
|
|
| pub(crate) enum CompactedMessageIdentity { |
| Preserve, |
| Regenerate, |
| } |
|
|
| pub(crate) fn collect_annotated_user_messages( |
| items: &[ResponseItemEnvelope], |
| identity: CompactedMessageIdentity, |
| ) -> Vec<CompactedUserMessage> { |
| items |
| .iter() |
| .filter_map(|envelope| compacted_user_message(&envelope.item, envelope.metadata.clone())) |
| .map(|mut message| { |
| if matches!(identity, CompactedMessageIdentity::Regenerate) { |
| message.id = None; |
| } |
| message |
| }) |
| .collect() |
| } |
|
|
| fn compacted_user_message( |
| item: &ResponseItem, |
| harness_metadata: Option<CodexHarnessMetadata>, |
| ) -> Option<CompactedUserMessage> { |
| let Some(TurnItem::UserMessage(user)) = crate::event_mapping::parse_turn_item(item) else { |
| return None; |
| }; |
| if is_summary_message(&user.message()) { |
| return None; |
| } |
| Some(CompactedUserMessage { |
| id: item.id().cloned(), |
| message: user.message(), |
| internal_chat_message_metadata_passthrough: match item { |
| ResponseItem::Message { |
| internal_chat_message_metadata_passthrough, |
| .. |
| } => internal_chat_message_metadata_passthrough.clone(), |
| _ => None, |
| }, |
| harness_metadata, |
| }) |
| } |
|
|
| pub(crate) fn is_summary_message(message: &str) -> bool { |
| message.starts_with(format!("{SUMMARY_PREFIX}\n").as_str()) |
| } |
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| pub(crate) fn insert_initial_context_before_last_real_user_or_summary( |
| mut compacted_history: Vec<ResponseItemEnvelope>, |
| initial_context: Vec<ResponseItemEnvelope>, |
| ) -> Vec<ResponseItemEnvelope> { |
| let mut last_user_or_summary_index = None; |
| let mut last_real_user_index = None; |
| for (i, item) in compacted_history.iter().enumerate().rev() { |
| if let ResponseItem::AgentMessage { content, .. } = &item.item |
| && !matches!( |
| content.first(), |
| Some(AgentMessageInputContent::InputText { text }) |
| if text.starts_with("Message Type: FINAL_ANSWER\n") |
| ) |
| { |
| last_real_user_index = Some(i); |
| break; |
| } |
| let Some(TurnItem::UserMessage(user)) = crate::event_mapping::parse_turn_item(&item.item) |
| else { |
| continue; |
| }; |
| |
| |
| |
| last_user_or_summary_index.get_or_insert(i); |
| if !is_summary_message(&user.message()) { |
| last_real_user_index = Some(i); |
| break; |
| } |
| } |
| let last_compaction_index = compacted_history |
| .iter() |
| .enumerate() |
| .rev() |
| .find_map(|(i, item)| { |
| matches!( |
| &item.item, |
| ResponseItem::Compaction { .. } | ResponseItem::ContextCompaction { .. } |
| ) |
| .then_some(i) |
| }); |
| let insertion_index = last_real_user_index |
| .or(last_user_or_summary_index) |
| .or(last_compaction_index); |
|
|
| |
| |
| |
| |
| if let Some(insertion_index) = insertion_index { |
| compacted_history.splice(insertion_index..insertion_index, initial_context); |
| } else { |
| compacted_history.extend(initial_context); |
| } |
|
|
| compacted_history |
| } |
|
|
| pub(crate) fn build_compacted_history( |
| initial_context: Vec<ResponseItemEnvelope>, |
| user_messages: &[CompactedUserMessage], |
| summary_text: &str, |
| ) -> Vec<ResponseItemEnvelope> { |
| build_compacted_history_with_limit( |
| initial_context, |
| user_messages, |
| summary_text, |
| COMPACT_USER_MESSAGE_MAX_TOKENS, |
| ) |
| } |
|
|
| fn build_compacted_history_with_limit( |
| mut history: Vec<ResponseItemEnvelope>, |
| user_messages: &[CompactedUserMessage], |
| summary_text: &str, |
| max_tokens: usize, |
| ) -> Vec<ResponseItemEnvelope> { |
| let mut selected_messages: Vec<CompactedUserMessage> = Vec::new(); |
| if max_tokens > 0 { |
| let mut remaining = max_tokens; |
| for message in user_messages.iter().rev() { |
| if remaining == 0 { |
| break; |
| } |
| let tokens = approx_token_count(&message.message); |
| if tokens <= remaining { |
| selected_messages.push(message.clone()); |
| remaining = remaining.saturating_sub(tokens); |
| } else { |
| let truncated = |
| truncate_text(&message.message, TruncationPolicy::Tokens(remaining)); |
| selected_messages.push(CompactedUserMessage { |
| id: message.id.clone(), |
| message: truncated, |
| internal_chat_message_metadata_passthrough: message |
| .internal_chat_message_metadata_passthrough |
| .clone(), |
| harness_metadata: message.harness_metadata.clone(), |
| }); |
| break; |
| } |
| } |
| selected_messages.reverse(); |
| } |
|
|
| for message in &selected_messages { |
| let mut item = ResponseItem::Message { |
| id: message.id.clone(), |
| role: "user".to_string(), |
| content: vec![ContentItem::InputText { |
| text: message.message.clone(), |
| }], |
| phase: None, |
| internal_chat_message_metadata_passthrough: message |
| .internal_chat_message_metadata_passthrough |
| .clone(), |
| }; |
| if message |
| .internal_chat_message_metadata_passthrough |
| .as_ref() |
| .and_then(|metadata| metadata.content_item_kinds.as_ref()) |
| .is_some() |
| { |
| let _ = set_annotated_content( |
| &mut item, |
| vec![AnnotatedContent::input_text( |
| &message.message, |
| ContentItemKind("user.text".to_string()), |
| )], |
| ); |
| } |
| history.push(ResponseItemEnvelope { |
| item, |
| metadata: message.harness_metadata.clone(), |
| }); |
| } |
|
|
| let summary_text = if summary_text.is_empty() { |
| "(no summary available)".to_string() |
| } else { |
| summary_text.to_string() |
| }; |
|
|
| history.push(ResponseItemEnvelope::new(ContextualUserFragment::into( |
| CompactionSummary::new(summary_text), |
| ))); |
|
|
| history |
| } |
|
|
| async fn drain_to_completed( |
| sess: &Session, |
| turn_context: &TurnContext, |
| client_session: &mut ModelClientSession, |
| responses_metadata: &CodexResponsesMetadata, |
| prompt: &Prompt, |
| ) -> CodexResult<String> { |
| let mut stream = client_session |
| .stream( |
| prompt, |
| turn_context.model_info(), |
| &turn_context.session_telemetry, |
| sess.reasoning_effort_for_request( |
| &turn_context.initial_settings, |
| RequestEffortUsage::Compaction, |
| ) |
| .await, |
| turn_context.reasoning_summary(), |
| turn_context.config.service_tier.clone(), |
| responses_metadata, |
| |
| |
| &InferenceTraceContext::disabled(), |
| ) |
| .await?; |
| loop { |
| let maybe_event = stream.next().await; |
| let Some(event) = maybe_event else { |
| return Err(CodexErr::Stream( |
| "stream closed before response.completed".into(), |
| )); |
| }; |
| match event { |
| Ok(ResponseEvent::OutputItemDone(item)) => { |
| sess.record_conversation_items( |
| turn_context, |
| turn_context.model_info(), |
| std::slice::from_ref(&item), |
| ) |
| .await; |
| } |
| Ok(ResponseEvent::ServerReasoningIncluded(included)) => { |
| sess.set_server_reasoning_included(included).await; |
| } |
| Ok(ResponseEvent::RateLimits(snapshot)) => { |
| sess.update_rate_limits(turn_context, snapshot).await; |
| } |
| Ok(ResponseEvent::Completed { |
| response_id, |
| token_usage, |
| usage_metadata, |
| .. |
| }) => { |
| sess.record_observed_response_completed( |
| turn_context, |
| &response_id, |
| token_usage.as_ref(), |
| usage_metadata.as_ref(), |
| ) |
| .await; |
| sess.update_token_usage_info(turn_context, token_usage.as_ref()) |
| .await?; |
| return Ok(response_id); |
| } |
| Ok(_) => continue, |
| Err(e) => return Err(e), |
| } |
| } |
| } |
|
|
| #[cfg(test)] |
| #[path = "compact_tests.rs"] |
| mod tests; |
|
|