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
4 changes: 4 additions & 0 deletions codex-rs/core/src/config/network_proxy_spec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,10 @@ impl NetworkProxySpec {
self.config.enabled
}

pub(crate) fn credential_broker_enabled(&self) -> bool {
self.config.credential_broker && self.constraints.enabled != Some(false)
}

pub fn proxy_host_and_port(&self) -> String {
host_and_port_from_network_addr(&self.config.proxy_url, /*default_port*/ 3128)
}
Expand Down
109 changes: 99 additions & 10 deletions codex-rs/core/src/environment_selection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ use codex_exec_server::ExecutorFileSystem;
use codex_exec_server::SelectedCapabilityRootsStatus;
use codex_protocol::capabilities::CapabilityRootLocation;
use codex_protocol::capabilities::SelectedCapabilityRoot;
use codex_protocol::config_types::ShellEnvironmentPolicy;
use codex_protocol::models::PermissionProfile;
use codex_protocol::protocol::AskForApproval;
use codex_protocol::protocol::EnvironmentConfig;
Expand All @@ -37,6 +38,7 @@ use crate::session::turn_context::ShellSnapshotTask;
use crate::session::turn_context::TurnEnvironment;
use crate::shell::Shell;
use crate::shell_snapshot::ShellSnapshot;
use crate::shell_snapshot::SnapshotCredentialBrokerState;

/// Records whether a normalized config should follow later thread setting updates.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
Expand Down Expand Up @@ -142,6 +144,7 @@ struct ResolvedEnvironment {
executor_platform_os: Option<String>,
temporary_directories: Option<Vec<PathUri>>,
shell_snapshot: ShellSnapshotTask,
shell_snapshot_builder: ShellSnapshot,
shell_snapshot_v2_supported: bool,
installed_config: Option<EnvironmentConfig>,
}
Expand Down Expand Up @@ -256,15 +259,25 @@ impl ThreadEnvironments {
let selection = environment.selection;
let config_origin = environment.config_origin;
let selected_environment = Arc::clone(&environment.environment);
let inherited_snapshot = if !selected_environment.is_remote()
&& shell_snapshot.should_rebuild_inherited()
{
futures::future::ready(None).boxed().shared()
} else {
environment.shell_snapshot
};
let resolution: TurnEnvironmentResolution =
futures::future::ready(Ok(ResolvedEnvironment {
environment: environment.environment,
shell: environment.shell,
user_home_dir: environment.user_home_dir,
executor_platform_os: environment.executor_platform_os,
temporary_directories: environment.temporary_directories,
shell_snapshot: environment.shell_snapshot,
shell_snapshot_v2_supported: environment.shell_snapshot_v2_supported,
shell_snapshot: inherited_snapshot,
shell_snapshot_builder: shell_snapshot.clone(),
shell_snapshot_v2_supported: environment.shell_snapshot_v2_supported
&& (selected_environment.is_remote()
|| !shell_snapshot.should_rebuild_inherited()),
installed_config: None,
}))
.boxed()
Expand Down Expand Up @@ -292,6 +305,36 @@ impl ThreadEnvironments {
}
}

fn start_shell_snapshot_task(
shell_snapshot: ShellSnapshot,
environment: Arc<Environment>,
cwd: PathUri,
shell: Option<Shell>,
) -> ShellSnapshotTask {
// Protected snapshots are captured by the command, using its sandbox.
if shell_snapshot.should_rebuild_inherited() {
return futures::future::ready(None).boxed().shared();
}
let shell_snapshot = shell_snapshot
.build(
environment,
cwd,
shell,
/*allow_login_shell*/ true,
ShellEnvironmentPolicy::default(),
/*sandbox*/ None,
)
.boxed()
.shared();
drop(tokio::spawn(
shell_snapshot
.clone()
.in_current_span()
.with_current_subscriber(),
));
shell_snapshot
}

pub(crate) fn update_selections(
&self,
environments: &[TurnEnvironmentSelection],
Expand Down Expand Up @@ -482,6 +525,47 @@ impl ThreadEnvironments {
self.environments.store(Arc::new(environments));
}

pub(crate) fn set_snapshot_credential_broker(&self, state: SnapshotCredentialBrokerState) {
if !self.shell_snapshot.set_credential_broker(state) {
return;
}

let mut environments = Vec::clone(&self.environments.load());
let mut changed = false;
for selected in &mut environments {
if !selected.environment.is_remote()
&& let Some(Ok(resolved)) = selected.resolution.clone().now_or_never()
{
self.restart_shell_snapshot(selected, resolved);
changed = true;
}
}
if changed {
self.environments.store(Arc::new(environments));
}
}

fn restart_shell_snapshot(
&self,
selected: &mut SelectedTurnEnvironment,
resolved: ResolvedEnvironment,
) {
let shell_snapshot = Self::start_shell_snapshot_task(
self.shell_snapshot.clone(),
Arc::clone(&resolved.environment),
selected.selection.cwd.clone(),
resolved.shell.clone(),
);
selected.resolution = futures::future::ready(Ok(ResolvedEnvironment {
shell_snapshot,
shell_snapshot_v2_supported: cfg!(unix)
&& !self.shell_snapshot.should_rebuild_inherited(),
..resolved
}))
.boxed()
.shared();
}

/// Combines persisted thread roots with installed attachment roots, keeping
/// thread roots first and hiding attachments that are not ready yet.
pub(crate) fn inspect_selected_capability_roots(
Expand Down Expand Up @@ -667,21 +751,24 @@ impl ThreadEnvironments {
cfg!(unix),
)
};
let task = shell_snapshot
.build(Arc::clone(&environment), selection.cwd, shell.clone())
.boxed()
.shared();
drop(tokio::spawn(
task.clone().in_current_span().with_current_subscriber(),
));
let shell_snapshot_builder = shell_snapshot.clone();
let task = Self::start_shell_snapshot_task(
shell_snapshot,
Arc::clone(&environment),
selection.cwd,
shell.clone(),
);
let shell_snapshot_v2_supported = snapshot_v2
&& (environment.is_remote() || !shell_snapshot_builder.should_rebuild_inherited());
Ok(ResolvedEnvironment {
environment,
shell,
user_home_dir,
executor_platform_os,
temporary_directories: temporary_dirs,
shell_snapshot: task,
shell_snapshot_v2_supported: snapshot_v2,
shell_snapshot_builder,
shell_snapshot_v2_supported,
installed_config,
})
}
Expand Down Expand Up @@ -762,6 +849,8 @@ impl TurnEnvironmentState {
);
turn_environment.executor_platform_os = environment.executor_platform_os;
turn_environment.shell_snapshot = environment.shell_snapshot;
turn_environment.shell_snapshot_builder =
Some(Box::new(environment.shell_snapshot_builder));
turn_environment.shell_snapshot_v2_supported =
environment.shell_snapshot_v2_supported;
turn_environment.user_home_dir = environment.user_home_dir;
Expand Down
29 changes: 29 additions & 0 deletions codex-rs/core/src/session/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ use crate::session::step_settings::ResolvedStepSettings;
use crate::session::step_settings::StepSettings;
use crate::session::turn_context::TurnEnvironment;
use crate::session_prefix::format_inter_agent_completion_message;
use crate::shell_snapshot::SnapshotCredentialBrokerState;
use crate::skills_load_input_from_config;
use crate::stream_events_utils::mark_thread_memory_mode_polluted_if_external_context;
use crate::turn_metadata::TurnMetadataState;
Expand Down Expand Up @@ -1165,6 +1166,9 @@ impl Session {
.cloned()
else {
self.services.network_proxy.store(None);
self.services
.turn_environments
.set_snapshot_credential_broker(SnapshotCredentialBrokerState::Inactive);
return;
};

Expand All @@ -1191,11 +1195,22 @@ impl Session {
// listeners and must not be exposed as active managed proxy runtimes.
if !spec.enabled() {
self.services.network_proxy.store(None);
self.services
.turn_environments
.set_snapshot_credential_broker(SnapshotCredentialBrokerState::Inactive);
return;
}
if let Some(started_proxy) = self.services.network_proxy.load_full() {
if let Err(err) = spec.apply_to_started_proxy(started_proxy.as_ref()).await {
warn!("failed to refresh managed network proxy for sandbox change: {err}");
} else {
self.services
.turn_environments
.set_snapshot_credential_broker(if spec.credential_broker_enabled() {
SnapshotCredentialBrokerState::Ready(started_proxy.proxy())
} else {
SnapshotCredentialBrokerState::Inactive
});
}
return;
}
Expand All @@ -1214,11 +1229,25 @@ impl Session {
.await
{
Ok((started_proxy, _session_network_proxy)) => {
if spec.credential_broker_enabled() {
self.services
.turn_environments
.set_snapshot_credential_broker(SnapshotCredentialBrokerState::Ready(
started_proxy.proxy(),
));
}
self.services
.network_proxy
.store(Some(Arc::new(started_proxy)));
}
Err(err) => {
self.services
.turn_environments
.set_snapshot_credential_broker(if spec.credential_broker_enabled() {
SnapshotCredentialBrokerState::Unavailable
} else {
SnapshotCredentialBrokerState::Inactive
});
warn!("failed to start managed network proxy for sandbox change: {err}");
}
}
Expand Down
55 changes: 48 additions & 7 deletions codex-rs/core/src/session/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ use crate::hook_mcp_executor::CoreHookMcpExecutor;
use crate::responses_metadata::CodexResponsesMetadata;
use crate::responses_metadata::CodexResponsesRequestKind;
use crate::shell_snapshot::ShellSnapshot;
use crate::shell_snapshot::SnapshotCredentialBrokerState;
use crate::state::ActiveTurn;
use codex_extension_api::ExtensionDataInit;
use codex_http_client::ClientRouteClass;
Expand Down Expand Up @@ -1091,13 +1092,10 @@ impl Session {
}),
});
}
let effective_config = config.config_layer_stack.effective_config();
let config_path = config.codex_home.join(CONFIG_TOML_FILE);
if let Some(event) = unstable_features_warning_event(
config
.config_layer_stack
.effective_config()
.get("features")
.and_then(TomlValue::as_table),
effective_config.get("features").and_then(TomlValue::as_table),
config.suppress_unstable_features_warning,
&config.features,
&config_path.display().to_string(),
Expand Down Expand Up @@ -1217,7 +1215,27 @@ impl Session {
} else {
shell::default_user_shell()
};
let use_executor_shell_snapshots = config.features.enabled(Feature::ShellSnapshotV2)
let credential_broker_available = config.features.enabled(Feature::NetworkProxy)
&& config
.config_layer_stack
.requirements()
.network
.as_ref()
.is_none_or(|network| network.value.enabled != Some(false));
let credential_broker_configured = credential_broker_available
&& effective_config
.get("features")
.and_then(|features| features.get("network_proxy"))
.and_then(|network_proxy| network_proxy.get("credential_broker"))
.and_then(TomlValue::as_bool)
.unwrap_or(false);
let credential_broker_active = credential_broker_configured
&& config
.permissions
.network
.as_ref()
.is_some_and(crate::config::NetworkProxySpec::credential_broker_enabled);
let prefer_executor_shell_snapshots = config.features.enabled(Feature::ShellSnapshotV2)
&& config.features.enabled(Feature::ShellTool)
&& config.features.enabled(Feature::UnifiedExec)
&& matches!(
Expand All @@ -1229,14 +1247,26 @@ impl Session {
),
codex_tools::UnifiedExecShellMode::Direct
);
let use_executor_shell_snapshots =
prefer_executor_shell_snapshots && !credential_broker_active;
let shell_snapshot = if config.features.enabled(Feature::ShellSnapshot)
&& !use_executor_shell_snapshots
&& (!use_executor_shell_snapshots || credential_broker_available)
{
let snapshot_credential_broker = credential_broker_available.then(|| {
let state = if credential_broker_active {
SnapshotCredentialBrokerState::Starting
} else {
SnapshotCredentialBrokerState::Inactive
};
watch::channel(state).0
});
ShellSnapshot::new(
config.codex_home.clone(),
thread_id,
session_telemetry.clone(),
state_db_ctx.clone(),
snapshot_credential_broker,
prefer_executor_shell_snapshots,
)
} else {
ShellSnapshot::disabled()
Expand Down Expand Up @@ -1355,6 +1385,17 @@ impl Session {
} else {
(None, None)
};
if let Some(network_proxy) = network_proxy.as_ref()
&& config
.permissions
.network
.as_ref()
.is_some_and(crate::config::NetworkProxySpec::credential_broker_enabled)
{
turn_environments.set_snapshot_credential_broker(
SnapshotCredentialBrokerState::Ready(network_proxy.proxy()),
);
}

// Hooks and extensions share one stable thread-owned MCP runtime handle.
let mcp_runtime = Arc::new(McpRuntime::empty(
Expand Down
16 changes: 14 additions & 2 deletions codex-rs/core/src/session/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5644,8 +5644,12 @@ async fn session_configuration_apply_permission_profile_accepts_direct_write_roo
);
}

#[test_case::test_case(false; "ordinary_proxy")]
#[test_case::test_case(true; "credential_broker")]
#[tokio::test]
async fn active_profile_update_rebuilds_network_proxy_config() -> std::io::Result<()> {
async fn active_profile_update_rebuilds_network_proxy_config(
credential_broker: bool,
) -> std::io::Result<()> {
let codex_home = tempfile::tempdir().expect("create codex home");
let cwd = tempfile::tempdir().expect("create cwd");
let permissions = PermissionsToml {
Expand Down Expand Up @@ -5690,7 +5694,14 @@ async fn active_profile_update_rebuilds_network_proxy_config() -> std::io::Resul
]),
};
let base_config = ConfigToml {
features: Some(toml::from_str("network_proxy = true").expect("valid features")),
features: Some(
toml::from_str(if credential_broker {
"network_proxy = { enabled = true, credential_broker = true }"
} else {
"network_proxy = true"
})
.expect("valid features"),
),
default_permissions: Some("locked-down".to_string()),
permissions: Some(permissions),
..Default::default()
Expand Down Expand Up @@ -5752,6 +5763,7 @@ async fn active_profile_update_rebuilds_network_proxy_config() -> std::io::Resul
.expect("selected profile proxy should become the session proxy config");
assert_eq!(network.proxy_host_and_port(), "127.0.0.1:43128");
assert!(!network.socks_enabled());
assert_eq!(network.credential_broker_enabled(), credential_broker);
Ok(())
}

Expand Down
Loading
Loading