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
1 change: 1 addition & 0 deletions codex-rs/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

8 changes: 8 additions & 0 deletions codex-rs/app-server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,14 @@ and migration remain supported so resuming, reading, and forking those threads
preserves the surviving history. This disk compatibility does not require restoring
support for new `thread/rollback` requests.

# MCP configuration reload

`config/mcpServer/reload` returns an error when a loaded thread rejects the
refreshed enterprise policy. The rejected thread retains its previous configuration
layers with enterprise MCP disabled. Other planned thread refreshes are processed
before the rejection is reported, so an error does not imply that no changes were
applied. Correct the policy before retrying the reload.

# Selected workspace routing

The experimental `account/read.workspaceRouting` response field returns the selected ChatGPT workspace's `chatgptAccountId`, resolved HTTPS `backendOrigin`, and backend-provided `accountRoutingOverride`. The routing value is `us`, `us_cr`, or the explicit `NO_CONSTRAINT` value. API-only and signed-out accounts return `null` and do not need `accounts/check`.
Expand Down
247 changes: 203 additions & 44 deletions codex-rs/app-server/src/mcp_refresh.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
use crate::config_manager::ConfigManager;
use codex_core::CodexThread;
use codex_core::ConfigRefreshOutcome;
use codex_core::ThreadManager;
use codex_core::config::Config;
use std::io;
Expand All @@ -10,56 +11,153 @@ pub(crate) async fn reload_mcp_config(
thread_manager: &Arc<ThreadManager>,
config_manager: &ConfigManager,
) -> io::Result<()> {
config_manager
.load_latest_config(/*fallback_cwd*/ None)
.await?;
let mut refreshes = Vec::new();
for thread_id in thread_manager.list_thread_ids().await {
let thread = thread_manager
.get_thread(thread_id)
.await
.map_err(|err| io::Error::other(format!("failed to load thread {thread_id}: {err}")))?;
let config = load_refresh_config(thread.as_ref(), config_manager).await?;
refreshes.push((thread, config));
}
for (thread, config) in refreshes {
thread.refresh_mcp_config(config).await;
}
Ok(())
// Keep the reload state machine out of the request dispatcher's stack frame.
Box::pin(reload_mcp_config_with_policy(
thread_manager,
config_manager,
ReloadPolicy::Strict,
))
.await
}

pub(crate) async fn reload_mcp_config_best_effort(
thread_manager: &Arc<ThreadManager>,
config_manager: &ConfigManager,
) {
for thread_id in thread_manager.list_thread_ids().await {
let thread = match thread_manager.get_thread(thread_id).await {
Ok(thread) => thread,
Err(err) => {
warn!(%thread_id, %err, "failed to load thread for MCP configuration refresh");
continue;
}
if let Err(err) = Box::pin(reload_mcp_config_with_policy(
thread_manager,
config_manager,
ReloadPolicy::BestEffort,
))
.await
{
warn!(%err, "failed to reload MCP configuration");
}
}

#[derive(Clone, Copy)]
enum ReloadPolicy {
Strict,
BestEffort,
}

async fn reload_mcp_config_with_policy(
thread_manager: &ThreadManager,
config_manager: &ConfigManager,
policy: ReloadPolicy,
) -> io::Result<()> {
let mut rejected = false;
let mut load_error = None;
let mut completed: Vec<Arc<CodexThread>> = Vec::new();
// Bound both discovery and stale-owner retries so concurrent churn cannot
// keep the reload request alive indefinitely. Completed publications remain valid.
for pass in 0..3 {
let host_loaded = match policy {
ReloadPolicy::BestEffort => true,
ReloadPolicy::Strict => match config_manager
.load_latest_config(/*fallback_cwd*/ None)
.await
{
Ok(_) => true,
Err(err) => {
load_error.get_or_insert(err);
false
}
},
};
let config = match load_refresh_config(thread.as_ref(), config_manager).await {
Ok(config) => config,
Err(err) => {
warn!(%thread_id, %err, "failed to load thread MCP configuration");
continue;
let mut targets = Vec::new();
for thread_id in thread_manager.list_thread_ids().await {
match thread_manager.get_thread(thread_id).await {
Ok(thread) => {
if !completed.iter().any(|done| Arc::ptr_eq(done, &thread)) {
targets.push((thread_id, thread));
}
}
Err(err) => {
warn!(%thread_id, %err, "failed to load thread for MCP configuration refresh");
}
}
};
thread.refresh_mcp_config(config).await;
}
if targets.is_empty() {
break;
}
for (thread_id, target) in targets {
for attempt in 0..3 {
let current_config = target.config().await;
let mut target_load_error = None;
let next_config = if host_loaded {
match load_refresh_config(&current_config, config_manager).await {
Ok(next_config) => next_config,
Err(err) => {
if matches!(policy, ReloadPolicy::BestEffort) {
warn!(%thread_id, %err, "failed to load session configuration");
}
target_load_error = Some(err);
disabled_enterprise_config(&current_config)
}
}
} else {
disabled_enterprise_config(&current_config)
};
match target.refresh_mcp_config(current_config, next_config).await {
ConfigRefreshOutcome::Published => {}
ConfigRefreshOutcome::Rejected => rejected = true,
ConfigRefreshOutcome::Stale => {
if attempt == 2 {
target.disable_mcp_enterprise_auth().await;
rejected = true;
}
continue;
}
}
if matches!(policy, ReloadPolicy::Strict)
&& let Some(err) = target_load_error
{
load_error.get_or_insert(err);
}
break;
}
completed.push(target);
}
if pass == 2 {
for thread_id in thread_manager.list_thread_ids().await {
if let Ok(thread) = thread_manager.get_thread(thread_id).await
&& !completed.iter().any(|done| Arc::ptr_eq(done, &thread))
{
thread.disable_mcp_enterprise_auth().await;
rejected = true;
}
}
}
}
// Finish publishing every session's fail-closed configuration before reporting failure.
if matches!(policy, ReloadPolicy::Strict)
&& let Some(err) = load_error
{
return Err(err);
}
if rejected && matches!(policy, ReloadPolicy::Strict) {
return Err(io::Error::other(
"enterprise MCP configuration was rejected; affected sessions retained their previous configuration with enterprise MCP disabled",
));
}
Ok(())
}

fn disabled_enterprise_config(current_config: &Config) -> Config {
let mut config = current_config.clone();
config.disable_mcp_enterprise_auth();
config
}

async fn load_refresh_config(
thread: &CodexThread,
current_config: &Config,
config_manager: &ConfigManager,
) -> io::Result<Config> {
let thread_config = thread.config().await;
config_manager
.load_latest_config_with_session_layers(
&thread_config.config_layer_stack,
&thread_config.cwd,
&current_config.config_layer_stack,
&current_config.cwd,
)
.await
}
Expand Down Expand Up @@ -96,6 +194,27 @@ mod tests {
use std::sync::atomic::Ordering;
use tempfile::TempDir;

#[tokio::test]
async fn stale_refresh_stops_after_bounded_attempts() -> anyhow::Result<()> {
let (_temp_dir, thread_manager, config_manager, loader) = refresh_test_state().await?;
for thread_id in thread_manager.list_thread_ids().await {
let thread = thread_manager.get_thread(thread_id).await?;
if thread.config().await.cwd.ends_with("good") {
*loader.churn_thread.lock().unwrap() = Some(thread);
}
}
tokio::time::timeout(
std::time::Duration::from_secs(10),
reload_mcp_config(&thread_manager, &config_manager),
)
.await
.expect("reload must terminate under owner churn")
.expect_err("reload failed closed");
assert_eq!(loader.good_loads.load(Ordering::Relaxed), 3);
assert_eq!(loader.bad_loads.load(Ordering::Relaxed), 1);
Ok(())
}

#[tokio::test]
async fn strict_refresh_reports_thread_planning_failures() -> anyhow::Result<()> {
let (temp_dir, thread_manager, config_manager, _loader) = refresh_test_state().await?;
Expand All @@ -110,14 +229,15 @@ mod tests {

assert_eq!(err.to_string(), "failed to load refresh config");
for thread_id in thread_manager.list_thread_ids().await {
let config = thread_manager.get_thread(thread_id).await?.config().await;
let failed = config.cwd.ends_with("bad");
assert_eq!(
thread_manager
.get_thread(thread_id)
.await?
.config()
.await
.auth_keyring_backend_kind(),
AuthKeyringBackendKind::Direct
config.auth_keyring_backend_kind(),
if failed {
AuthKeyringBackendKind::Direct
} else {
AuthKeyringBackendKind::Secrets
}
);
}
Ok(())
Expand Down Expand Up @@ -148,6 +268,30 @@ mod tests {
Ok(())
}

#[tokio::test]
async fn strict_refresh_preserves_config_when_host_config_cannot_load() -> anyhow::Result<()> {
let (temp_dir, thread_manager, config_manager, _loader) = refresh_test_state().await?;
std::fs::write(
temp_dir.path().join(codex_config::CONFIG_TOML_FILE),
"[features",
)?;

assert!(
reload_mcp_config(&thread_manager, &config_manager)
.await
.is_err()
);

for thread_id in thread_manager.list_thread_ids().await {
let config = thread_manager.get_thread(thread_id).await?.config().await;
assert_eq!(
config.auth_keyring_backend_kind(),
AuthKeyringBackendKind::Direct
);
}
Ok(())
}

#[tokio::test]
async fn invalidation_does_not_reload_thread_config() -> anyhow::Result<()> {
let (_temp_dir, thread_manager, _config_manager, loader) = refresh_test_state().await?;
Expand Down Expand Up @@ -179,8 +323,11 @@ mod tests {
let thread = good_thread.expect("good test thread should exist");
let original_model = thread.config().await.model.clone();

let refresh_config = load_refresh_config(thread.as_ref(), &config_manager).await?;
thread.refresh_mcp_config(refresh_config).await;
let current_config = thread.config().await;
let refresh_config = load_refresh_config(&current_config, &config_manager).await?;
let _ = thread
.refresh_mcp_config(current_config, refresh_config)
.await;

assert_eq!(
thread.config().await.auth_keyring_backend_kind(),
Expand Down Expand Up @@ -221,7 +368,8 @@ enabled = false
"#,
)?;

let refresh_config = load_refresh_config(thread.as_ref(), &config_manager).await?;
let current_config = thread.config().await;
let refresh_config = load_refresh_config(&current_config, &config_manager).await?;
let mut actual = refresh_config.mcp_servers.get().clone();
actual.remove(codex_mcp::CODEX_APPS_MCP_SERVER_NAME);
let expected = serde_json::from_value::<HashMap<String, McpServerConfig>>(json!({
Expand Down Expand Up @@ -366,6 +514,7 @@ enabled = false
bad_cwd: AbsolutePathBuf::try_from(bad_cwd)?,
good_loads: AtomicUsize::new(0),
bad_loads: AtomicUsize::new(0),
churn_thread: std::sync::Mutex::new(None),
});
let config_manager = ConfigManager::new(
temp_dir.path().to_path_buf(),
Expand All @@ -385,6 +534,7 @@ enabled = false
bad_cwd: AbsolutePathBuf,
good_loads: AtomicUsize,
bad_loads: AtomicUsize,
churn_thread: std::sync::Mutex<Option<Arc<CodexThread>>>,
}

impl CountingThreadConfigLoader {
Expand All @@ -394,6 +544,15 @@ enabled = false
) -> Result<Vec<ThreadConfigSource>, ThreadConfigLoadError> {
if context.cwd.as_ref() == Some(&self.good_cwd) {
self.good_loads.fetch_add(1, Ordering::Relaxed);
let thread = self.churn_thread.lock().unwrap().clone();
if let Some(thread) = thread {
let current = thread.config().await;
let next = current.as_ref().clone();
assert_eq!(
thread.refresh_runtime_config(current, next).await,
ConfigRefreshOutcome::Published
);
}
}
if context.cwd.as_ref() == Some(&self.bad_cwd) {
self.bad_loads.fetch_add(1, Ordering::Relaxed);
Expand Down
Loading
Loading