Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 12 additions & 7 deletions codex-rs/core/src/compact_remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ use crate::hook_runtime::run_pre_compact_hooks;
use crate::responses_metadata::CodexResponsesRequestKind;
use crate::responses_metadata::CompactionTurnMetadata;
use crate::session::session::Session;
use crate::session::step_context::StepContext;
use crate::session::turn::built_tools;
use crate::session::turn_context::TurnContext;
use codex_analytics::CompactionImplementation;
Expand Down Expand Up @@ -45,15 +46,15 @@ const CONTEXT_WINDOW_TRUNCATED_OUTPUT_MESSAGE: &str =

pub(crate) async fn run_inline_remote_auto_compact_task(
sess: Arc<Session>,
turn_context: Arc<TurnContext>,
step_context: Arc<StepContext>,
turn_state: Arc<OnceLock<String>>,
initial_context_injection: InitialContextInjection,
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
run_remote_compact_task_inner(
&sess,
&turn_context,
&step_context,
Some(turn_state),
initial_context_injection,
CompactionTrigger::Auto,
Expand All @@ -68,6 +69,8 @@ pub(crate) async fn run_remote_compact_task(
sess: Arc<Session>,
turn_context: Arc<TurnContext>,
) -> CodexResult<()> {
// Standalone compaction is its own request boundary, so it captures a fresh step.
let step_context = sess.capture_step_context(Arc::clone(&turn_context)).await;
let start_event = EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_context.sub_id.clone(),
trace_id: turn_context.trace_id.clone(),
Expand All @@ -79,7 +82,7 @@ pub(crate) async fn run_remote_compact_task(

run_remote_compact_task_inner(
&sess,
&turn_context,
&step_context,
/*turn_state*/ None,
InitialContextInjection::DoNotInject,
CompactionTrigger::Manual,
Expand All @@ -92,13 +95,14 @@ pub(crate) async fn run_remote_compact_task(

async fn run_remote_compact_task_inner(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
step_context: &Arc<StepContext>,
turn_state: Option<Arc<OnceLock<String>>>,
initial_context_injection: InitialContextInjection,
trigger: CompactionTrigger,
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
let turn_context = &step_context.turn;
let compaction_metadata = CompactionTurnMetadata::new(
trigger,
reason,
Expand Down Expand Up @@ -136,7 +140,7 @@ async fn run_remote_compact_task_inner(
}
let result = run_remote_compact_task_inner_impl(
sess,
turn_context,
step_context,
turn_state,
initial_context_injection,
compaction_metadata,
Expand Down Expand Up @@ -170,12 +174,13 @@ async fn run_remote_compact_task_inner(

async fn run_remote_compact_task_inner_impl(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
step_context: &Arc<StepContext>,
turn_state: Option<Arc<OnceLock<String>>>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
analytics_details: &mut CompactionAnalyticsDetails,
) -> CodexResult<()> {
let turn_context = &step_context.turn;
let context_compaction_item = ContextCompactionItem::new();
// Use the UI compaction item ID as the trace compaction ID so protocol lifecycle events,
// endpoint attempts, and the installed history checkpoint all have one join key.
Expand Down Expand Up @@ -221,7 +226,7 @@ async fn run_remote_compact_task_inner_impl(
let prompt_input = history.for_prompt(&turn_context.model_info.input_modalities);
let tool_router = built_tools(
sess.as_ref(),
turn_context.as_ref(),
step_context.as_ref(),
&CancellationToken::new(),
)
.await?;
Expand Down
19 changes: 12 additions & 7 deletions codex-rs/core/src/compact_remote_v2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ use crate::responses_metadata::CompactionTurnMetadata;
use crate::responses_retry::ResponsesStreamRequest;
use crate::responses_retry::handle_retryable_response_stream_error;
use crate::session::session::Session;
use crate::session::step_context::StepContext;
use crate::session::turn::built_tools;
use crate::session::turn_context::TurnContext;
use codex_analytics::CompactionImplementation;
Expand Down Expand Up @@ -55,15 +56,15 @@ const MAX_REMOTE_COMPACTION_V2_STREAM_RETRIES: u64 = 2;

pub(crate) async fn run_inline_remote_auto_compact_task(
sess: Arc<Session>,
turn_context: Arc<TurnContext>,
step_context: Arc<StepContext>,
client_session: &mut ModelClientSession,
initial_context_injection: InitialContextInjection,
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
run_remote_compact_task_inner(
&sess,
&turn_context,
&step_context,
Some(client_session),
initial_context_injection,
CompactionTrigger::Auto,
Expand All @@ -77,6 +78,8 @@ pub(crate) async fn run_remote_compact_task(
sess: Arc<Session>,
turn_context: Arc<TurnContext>,
) -> CodexResult<()> {
// Standalone compaction is its own request boundary, so it captures a fresh step.
let step_context = sess.capture_step_context(Arc::clone(&turn_context)).await;
let start_event = EventMsg::TurnStarted(TurnStartedEvent {
turn_id: turn_context.sub_id.clone(),
trace_id: turn_context.trace_id.clone(),
Expand All @@ -88,7 +91,7 @@ pub(crate) async fn run_remote_compact_task(

run_remote_compact_task_inner(
&sess,
&turn_context,
&step_context,
/*client_session*/ None,
InitialContextInjection::DoNotInject,
CompactionTrigger::Manual,
Expand All @@ -100,13 +103,14 @@ pub(crate) async fn run_remote_compact_task(

async fn run_remote_compact_task_inner(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
step_context: &Arc<StepContext>,
client_session: Option<&mut ModelClientSession>,
initial_context_injection: InitialContextInjection,
trigger: CompactionTrigger,
reason: CompactionReason,
phase: CompactionPhase,
) -> CodexResult<()> {
let turn_context = &step_context.turn;
let compaction_metadata = CompactionTurnMetadata::new(
trigger,
reason,
Expand Down Expand Up @@ -144,7 +148,7 @@ async fn run_remote_compact_task_inner(
}
let result = run_remote_compact_task_inner_impl(
sess,
turn_context,
step_context,
client_session,
initial_context_injection,
compaction_metadata,
Expand Down Expand Up @@ -181,12 +185,13 @@ async fn run_remote_compact_task_inner(

async fn run_remote_compact_task_inner_impl(
sess: &Arc<Session>,
turn_context: &Arc<TurnContext>,
step_context: &Arc<StepContext>,
client_session: Option<&mut ModelClientSession>,
initial_context_injection: InitialContextInjection,
compaction_metadata: CompactionTurnMetadata,
analytics_details: &mut CompactionAnalyticsDetails,
) -> CodexResult<()> {
let turn_context = &step_context.turn;
let context_compaction_item = ContextCompactionItem::new();
let compaction_trace = sess.services.rollout_thread_trace.compaction_trace_context(
turn_context.sub_id.as_str(),
Expand Down Expand Up @@ -229,7 +234,7 @@ async fn run_remote_compact_task_inner_impl(
let prompt_input = history.for_prompt(&turn_context.model_info.input_modalities);
let tool_router = built_tools(
sess.as_ref(),
turn_context.as_ref(),
step_context.as_ref(),
&CancellationToken::new(),
)
.await?;
Expand Down
16 changes: 14 additions & 2 deletions codex-rs/core/src/prompt_debug.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ use std::sync::Arc;
use codex_exec_server::EnvironmentManager;
use codex_exec_server::ExecServerRuntimePaths;
use codex_extension_api::UserInstructionsProvider;
use codex_features::Feature;
use codex_login::AuthManager;
use codex_protocol::error::CodexErr;
use codex_protocol::error::Result as CodexResult;
Expand Down Expand Up @@ -77,7 +78,8 @@ pub(crate) async fn build_prompt_input_from_session(
input: Vec<UserInput>,
) -> CodexResult<Vec<ResponseItem>> {
let turn_context = sess.new_default_turn().await;
sess.record_context_updates_and_set_reference_context_item(turn_context.as_ref())
let world_state = sess
.record_context_updates_and_set_reference_context_item(turn_context.as_ref())
.await;

if !input.is_empty() {
Expand All @@ -86,11 +88,21 @@ pub(crate) async fn build_prompt_input_from_session(
.await;
}

let step_context = sess.capture_step_context(Arc::clone(&turn_context)).await;
if turn_context
.config
.features
.enabled(Feature::DeferredExecutor)
{
sess.record_step_environment_context_if_changed(&world_state, step_context.as_ref())
.await;
}

let prompt_input = sess
.clone_history()
.await
.for_prompt(&turn_context.model_info.input_modalities);
let router = built_tools(sess, turn_context.as_ref(), &CancellationToken::new()).await?;
let router = built_tools(sess, step_context.as_ref(), &CancellationToken::new()).await?;
let base_instructions = sess.get_base_instructions().await;
let prompt = build_prompt(
prompt_input,
Expand Down
22 changes: 20 additions & 2 deletions codex-rs/core/src/session/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ use crate::exec_policy::ExecPolicyManager;
use crate::image_preparation::prepare_response_items;
use crate::parse_turn_item;
use crate::realtime_conversation::RealtimeConversationManager;
use crate::session::step_context::StepContext;
use crate::session::turn_context::TurnEnvironment;
use crate::session_prefix::format_inter_agent_completion_message;
use crate::skills::SkillRenderSideEffects;
Expand Down Expand Up @@ -388,7 +389,6 @@ use codex_protocol::protocol::TokenUsageInfo;
use codex_protocol::protocol::TurnModerationMetadataEvent;
use codex_protocol::protocol::WarningEvent;
use codex_protocol::user_input::UserInput;
use codex_tools::ToolEnvironmentMode;
use codex_tools::UnifiedExecShellMode;
use codex_utils_absolute_path::AbsolutePathBuf;
#[cfg(test)]
Expand Down Expand Up @@ -2774,10 +2774,11 @@ impl Session {

pub(crate) async fn record_step_environment_context_if_changed(
&self,
turn_context: &TurnContext,
previous_world_state: &Arc<WorldState>,
step_context: &step_context::StepContext,
) -> Arc<WorldState> {
let turn_context = step_context.turn.as_ref();
// Render model-visible state from the same step used to build and run tools.
let world_state = Arc::new(
self.build_world_state_for_environments(turn_context, &step_context.environments)
.await,
Expand All @@ -2798,6 +2799,23 @@ impl Session {
world_state
}

pub(crate) async fn capture_step_context(
Comment thread
sayan-oai marked this conversation as resolved.
&self,
turn_context: Arc<TurnContext>,
) -> Arc<StepContext> {
// Keep the old turn-frozen view unless deferred executors are explicitly enabled.
let environments = if turn_context
.config
.features
.enabled(Feature::DeferredExecutor)
{
self.services.turn_environments.snapshot().await
} else {
turn_context.environments.clone()
};
Arc::new(StepContext::new(turn_context, environments))
}

pub(crate) async fn record_inter_agent_communication(
&self,
turn_context: &TurnContext,
Expand Down
10 changes: 10 additions & 0 deletions codex-rs/core/src/session/step_context.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,17 @@
use std::sync::Arc;

use crate::environment_selection::TurnEnvironmentSnapshot;
use crate::session::turn_context::TurnContext;

/// Request-scoped state that may change between model sampling requests.
#[derive(Debug)]
pub(crate) struct StepContext {
pub(crate) turn: Arc<TurnContext>,
pub(crate) environments: TurnEnvironmentSnapshot,
}

impl StepContext {
pub(crate) fn new(turn: Arc<TurnContext>, environments: TurnEnvironmentSnapshot) -> Self {
Self { turn, environments }
}
}
Loading
Loading