Skip to content
Closed
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
22 changes: 15 additions & 7 deletions codex-rs/rmcp-client/src/oauth.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
use anyhow::Context;
use anyhow::Error;
use anyhow::Result;
use anyhow::anyhow;
use codex_config::types::OAuthCredentialsStoreMode;
use oauth2::AccessToken;
use oauth2::RefreshToken;
Expand Down Expand Up @@ -52,6 +53,8 @@ use codex_utils_home_dir::find_codex_home;

const KEYRING_SERVICE: &str = "Codex MCP Credentials";
const REFRESH_SKEW_MILLIS: u64 = 30_000;
pub(crate) const OAUTH_REQUEST_REFRESH_FAILED_ERROR: &str =
"OAuth access token refresh failed before MCP request";

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct StoredOAuthTokens {
Expand Down Expand Up @@ -362,15 +365,20 @@ impl OAuthPersistor {
{
let manager = self.inner.authorization_manager.clone();
let guard = manager.lock().await;
guard.refresh_token().await.with_context(|| {
format!(
"failed to refresh OAuth tokens for server {}",
self.inner.server_name
)
})?;
guard
.get_access_token()
.await
.map_err(|_| anyhow!(OAUTH_REQUEST_REFRESH_FAILED_ERROR))?;
}

self.persist_if_needed().await
if let Err(error) = self.persist_if_needed().await {
warn!(
"failed to persist refreshed OAuth tokens for server {}; will retry: {error}",
self.inner.server_name
);
}

Ok(())
}
}

Expand Down
29 changes: 16 additions & 13 deletions codex-rs/rmcp-client/src/rmcp_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -439,7 +439,7 @@ impl RmcpClient {
params: Option<PaginatedRequestParams>,
timeout: Option<Duration>,
) -> Result<ListToolsResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let result = self
.run_service_operation("tools/list", timeout, move |service| {
let params = params.clone();
Expand All @@ -455,7 +455,7 @@ impl RmcpClient {
params: Option<PaginatedRequestParams>,
timeout: Option<Duration>,
) -> Result<ListToolsWithConnectorIdResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let result = self
.run_service_operation("tools/list", timeout, move |service| {
let params = params.clone();
Expand Down Expand Up @@ -500,7 +500,7 @@ impl RmcpClient {
params: Option<PaginatedRequestParams>,
timeout: Option<Duration>,
) -> Result<ListResourcesResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let result = self
.run_service_operation("resources/list", timeout, move |service| {
let params = params.clone();
Expand All @@ -516,7 +516,7 @@ impl RmcpClient {
params: Option<PaginatedRequestParams>,
timeout: Option<Duration>,
) -> Result<ListResourceTemplatesResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let result = self
.run_service_operation("resources/templates/list", timeout, move |service| {
let params = params.clone();
Expand All @@ -532,7 +532,7 @@ impl RmcpClient {
params: ReadResourceRequestParams,
timeout: Option<Duration>,
) -> Result<ReadResourceResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let result = self
.run_service_operation("resources/read", timeout, move |service| {
let params = params.clone();
Expand All @@ -550,7 +550,7 @@ impl RmcpClient {
meta: Option<serde_json::Value>,
timeout: Option<Duration>,
) -> Result<CallToolResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let arguments = match arguments {
Some(Value::Object(map)) => Some(map),
Some(other) => {
Expand Down Expand Up @@ -606,7 +606,7 @@ impl RmcpClient {
method: &str,
params: Option<serde_json::Value>,
) -> Result<()> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
self.run_service_operation(
"notifications/custom",
/*timeout*/ None,
Expand Down Expand Up @@ -636,7 +636,7 @@ impl RmcpClient {
method: &str,
params: Option<serde_json::Value>,
) -> Result<ServerResult> {
self.refresh_oauth_if_needed().await;
self.refresh_oauth_if_needed().await?;
let response = self
.run_service_operation("requests/custom", /*timeout*/ None, move |service| {
let params = params.clone();
Expand Down Expand Up @@ -700,12 +700,11 @@ impl RmcpClient {
}
}

async fn refresh_oauth_if_needed(&self) {
if let Some(runtime) = self.oauth_persistor().await
&& let Err(error) = runtime.refresh_if_needed().await
{
warn!("failed to refresh OAuth tokens: {error}");
async fn refresh_oauth_if_needed(&self) -> Result<()> {
if let Some(runtime) = self.oauth_persistor().await {
runtime.refresh_if_needed().await?;
}
Ok(())
}

async fn create_pending_transport(
Expand Down Expand Up @@ -1061,6 +1060,10 @@ async fn create_oauth_transport_and_runtime(
Ok((transport, runtime))
}

#[cfg(test)]
#[path = "rmcp_client_oauth_tests.rs"]
mod oauth_tests;

#[cfg(test)]
mod tests {
use std::time::Duration;
Expand Down
Loading
Loading