Skip to content
Prev Previous commit
Next Next commit
codex-mcp: reject discarded runtime work
  • Loading branch information
mzeng-openai committed Jul 8, 2026
commit a0a02e540c2ae0e73b9d514a40e02431375de875
2 changes: 1 addition & 1 deletion codex-rs/codex-mcp/src/connection_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -557,7 +557,7 @@ impl McpConnectionManager {
fetch_ticket,
) {
(Some(cache_context), Some(fetch_ticket)) => cache_context
.publish_if_newest_accepted(fetch_ticket, &managed_client.server_info, tools),
.publish_if_newest_accepted(fetch_ticket, &managed_client.server_info, tools)?,
(None, None) => tools,
_ => unreachable!("Codex Apps fetch ticket requires cache context"),
};
Expand Down
79 changes: 79 additions & 0 deletions codex-rs/codex-mcp/src/connection_manager_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1186,6 +1186,85 @@ async fn list_all_tools_uses_shared_codex_apps_cache_when_client_startup_fails()
);
}

#[tokio::test]
async fn context_discard_while_checking_failed_startup_does_not_reconnect() {
let codex_home = tempdir().expect("tempdir");
let cache = CodexAppsToolsCache::default();
let context = cache.context(
codex_home.path().to_path_buf(),
CodexAppsToolsCacheKey {
account_id: Some("account-one".to_string()),
chatgpt_user_id: Some("user-one".to_string()),
is_workspace_account: false,
},
);
let (failure_started_tx, failure_started_rx) = tokio::sync::oneshot::channel();
let release_failure = Arc::new(tokio::sync::Notify::new());
let release_failure_for_client = Arc::clone(&release_failure);
let client: ManagedClientFuture = async move {
let _ = failure_started_tx.send(());
release_failure_for_client.notified().await;
Err(StartupOutcomeError::Failed {
error: "startup failed".to_string(),
is_authentication_required: false,
})
}
.boxed()
.shared();
let reconnect_attempts = Arc::new(AtomicUsize::new(0));
let reconnect_attempts_for_factory = Arc::clone(&reconnect_attempts);
let reconnect_factory = Arc::new(move || {
reconnect_attempts_for_factory.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
futures::future::ready(Err(StartupOutcomeError::Failed {
error: "reconnect should not start".to_string(),
is_authentication_required: false,
}))
.boxed()
.shared()
});
let client = AsyncManagedClient {
client,
is_codex_apps_mcp_server: true,
cached_server_info: None,
codex_apps_tools_cache_context: Some(context),
tool_filter: ToolFilter::default(),
startup_complete: Arc::new(std::sync::atomic::AtomicBool::new(true)),
startup_reconnect: Some(Arc::new(CodexAppsStartupReconnect::new(reconnect_factory))),
tool_plugin_provenance: Arc::new(ToolPluginProvenance::default()),
cancel_token: CancellationToken::new(),
};
let reconnect_task = tokio::spawn({
let client = client.clone();
async move { client.reconnect_failed_startup().await }
});

failure_started_rx
.await
.expect("failed startup check should begin");
let _new_context = cache.context(
codex_home.path().to_path_buf(),
CodexAppsToolsCacheKey {
account_id: Some("account-two".to_string()),
chatgpt_user_id: Some("user-two".to_string()),
is_workspace_account: false,
},
);
release_failure.notify_one();
reconnect_task.await.expect("reconnect check should finish");
tokio::task::yield_now().await;
assert_eq!(
reconnect_attempts.load(std::sync::atomic::Ordering::SeqCst),
0
);

client.reconnect_failed_startup().await;
tokio::task::yield_now().await;
assert_eq!(
reconnect_attempts.load(std::sync::atomic::Ordering::SeqCst),
0
);
}

#[tokio::test]
async fn list_all_tools_reconnects_failed_codex_apps_startup_and_reuses_client() {
let recovered_client = create_test_managed_client(vec![create_test_tool(
Expand Down
3 changes: 1 addition & 2 deletions codex-rs/codex-mcp/src/connector_runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -274,10 +274,9 @@ impl ConnectorRuntimeContext {
ticket: CodexAppsToolsFetchTicket,
server_info: &McpServerInfo,
tools: Vec<ToolInfo>,
) -> Vec<ToolInfo> {
) -> Result<Vec<ToolInfo>, ConnectorRuntimeContextDiscarded> {
self.publish_runtime_if_newest_accepted(ticket, server_info, tools)
.map(|snapshot| snapshot.tools.clone())
.unwrap_or_default()
}

#[cfg(test)]
Expand Down
54 changes: 45 additions & 9 deletions codex-rs/codex-mcp/src/connector_runtime_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -382,8 +382,9 @@ fn codex_apps_tools_cache_publishes_newest_shared_snapshot() {
let newer_tools = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "newer")];
let older_tools = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "older")];

let published_tools =
cache_context_2.publish_if_newest_accepted(newer_ticket, &server_info, newer_tools);
let published_tools = cache_context_2
.publish_if_newest_accepted(newer_ticket, &server_info, newer_tools)
.expect("publish newer tools");
assert_eq!(
model_tool_names(&published_tools),
model_tool_names(
Expand All @@ -392,8 +393,9 @@ fn codex_apps_tools_cache_publishes_newest_shared_snapshot() {
.expect("new snapshot should publish")
)
);
let current_tools =
cache_context_1.publish_if_newest_accepted(older_ticket, &server_info, older_tools);
let current_tools = cache_context_1
.publish_if_newest_accepted(older_ticket, &server_info, older_tools)
.expect("stale publish should return current tools");

assert_eq!(current_tools[0].callable_name, "newer");
assert_eq!(
Expand Down Expand Up @@ -421,11 +423,13 @@ fn codex_apps_tools_cache_keeps_live_publish_when_disk_persistence_fails() {
},
);
let tools = vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "live")];
let published_tools = cache_context.publish_if_newest_accepted(
cache_context.begin_fetch(CodexAppsToolsFetchSource::HardRefresh),
&create_test_server_info("Codex Apps"),
tools.clone(),
);
let published_tools = cache_context
.publish_if_newest_accepted(
cache_context.begin_fetch(CodexAppsToolsFetchSource::HardRefresh),
&create_test_server_info("Codex Apps"),
tools.clone(),
)
.expect("live publish should survive persistence failure");

assert_eq!(model_tool_names(&published_tools), model_tool_names(&tools));
assert_eq!(
Expand Down Expand Up @@ -533,6 +537,38 @@ fn activating_new_context_discards_old_context_without_state_bleed() {
));
}

#[test]
fn discarded_compatibility_publish_reports_an_error() {
let codex_home = tempdir().expect("tempdir");
let manager = ConnectorRuntimeManager::default();
let context_a = manager.context(
codex_home.path().to_path_buf(),
ConnectorRuntimeContextKey {
account_id: Some("account-a".to_string()),
chatgpt_user_id: Some("user-a".to_string()),
is_workspace_account: false,
},
);
let ticket = context_a.begin_fetch(CodexAppsToolsFetchSource::HardRefresh);
let _context_b = manager.context(
codex_home.path().to_path_buf(),
ConnectorRuntimeContextKey {
account_id: Some("account-b".to_string()),
chatgpt_user_id: Some("user-b".to_string()),
is_workspace_account: false,
},
);

assert!(matches!(
context_a.publish_if_newest_accepted(
ticket,
&create_test_server_info("Codex Apps"),
vec![create_test_tool(CODEX_APPS_MCP_SERVER_NAME, "stale-a")],
),
Err(ConnectorRuntimeContextDiscarded)
));
}

#[test]
fn oversized_tools_cache_is_ignored_during_initial_load() {
let codex_home = tempdir().expect("tempdir");
Expand Down
13 changes: 9 additions & 4 deletions codex-rs/codex-mcp/src/rmcp_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -481,13 +481,18 @@ impl AsyncManagedClient {
}

pub(crate) async fn reconnect_failed_startup(&self) {
if !self.connector_runtime_context_is_active() {
return;
}
let Some(startup_reconnect) = self.startup_reconnect.as_ref() else {
return;
};
if !self.startup_complete.load(Ordering::Acquire) {
return;
}
if matches!(self.client().await, Err(StartupOutcomeError::Failed { .. })) {
if matches!(self.client().await, Err(StartupOutcomeError::Failed { .. }))
&& self.connector_runtime_context_is_active()
{
startup_reconnect.reconnect_in_background();
}
}
Expand Down Expand Up @@ -875,9 +880,9 @@ async fn start_server_task(
.map_err(StartupOutcomeError::from)?;
let server_info = mcp_server_info_from_implementation(initialize_result.server_info);
let tools = match (codex_apps_tools_cache_context.as_ref(), fetch_ticket) {
(Some(cache_context), Some(fetch_ticket)) => {
cache_context.publish_if_newest_accepted(fetch_ticket, &server_info, tools)
}
(Some(cache_context), Some(fetch_ticket)) => cache_context
.publish_if_newest_accepted(fetch_ticket, &server_info, tools)
.map_err(|error| StartupOutcomeError::from(anyhow!(error)))?,
(None, None) => tools,
_ => unreachable!("Codex Apps fetch ticket requires cache context"),
};
Expand Down
Loading