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
8 changes: 8 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,14 @@ Expected config growth:

Keep persistent config access behind `UserConfigStore`. Do not push Redis details into routing code.

## Logging

- Use consistent formatted tracing messages for related events so logs are easy to grep across request paths.
- Prefer captured variables inside the message, for example `level!("method_name - event text field = {field} other_field = {other_field}")`, over structured field syntax for dataplane logs.
- Keep the method/event prefix stable and reuse the same field names/order for related events.
- Keep warning logs for unexpected conditions that likely need operator attention. Expected user/config misses should be debug or info unless they indicate a platform problem.
- Do not log tokens, authorization headers, secrets, Redis key/value bytes, full `UserConfig`, or backend credentials.

## Backend Sessions

Initialization fans out:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ use rmcp::{
ErrorData, RoleServer, model::ErrorCode, service::RequestContext,
transport::streamable_http_server::tower::DownstreamSessionId,
};
use tracing::info;
use tracing::debug;

use crate::{
common::ContextForgeClaims,
Expand All @@ -28,9 +28,14 @@ impl<'a> AuthorizedCallValidator<'a> {
let maybe_claims = maybe_parts.and_then(|parts| parts.extensions.get::<ContextForgeClaims>());

let maybe_virtual_host_id = maybe_parts.and_then(|parts| parts.extensions.get::<VirtualHostId>());
info!(
"{} user_config = {maybe_user_config:#?} session_id = {maybe_session_id:#?} virtual_host_id = {maybe_virtual_host_id:#?}",
self.call_name
let call_name = self.call_name;
let has_user_config = maybe_user_config.is_some();
let virtual_hosts = maybe_user_config.map_or(0, |user_config| user_config.virtual_hosts.len());
let has_session_id = maybe_session_id.is_some();
let has_claims = maybe_claims.is_some();
let virtual_host_id = maybe_virtual_host_id.map_or("<missing>", |id| id.value().as_str());
debug!(
"AuthorizedCallValidator::validate - mcp call validation call_name = {call_name} has_user_config = {has_user_config} virtual_hosts = {virtual_hosts} has_session_id = {has_session_id} has_claims = {has_claims} virtual_host_id = {virtual_host_id}"
);

let Some(session_id) = maybe_session_id else {
Expand Down Expand Up @@ -58,6 +63,12 @@ impl<'a> AuthorizedCallValidator<'a> {
};

let Some(virtual_host) = user_config.virtual_hosts.get(virtual_host_id.value()) else {
let call_name = self.call_name;
let virtual_host_id = virtual_host_id.value();
let virtual_hosts = user_config.virtual_hosts.len();
debug!(
"AuthorizedCallValidator::validate - mcp virtual host config missing call_name = {call_name} virtual_host_id = {virtual_host_id} virtual_hosts = {virtual_hosts}"
);
return Err(ErrorData {
code: ErrorCode::RESOURCE_NOT_FOUND,
message: "No configuration".into(),
Expand Down Expand Up @@ -92,8 +103,14 @@ impl<'a> InitializeCallValidator<'a> {
let maybe_user_config = maybe_parts.and_then(|parts| parts.extensions.get::<UserConfig>());
let maybe_virtual_host_id = maybe_parts.and_then(|parts| parts.extensions.get::<VirtualHostId>());
let maybe_claims = maybe_parts.and_then(|parts| parts.extensions.get::<ContextForgeClaims>());
info!(
"intialize user_config = {maybe_user_config:#?} downstream_session_id = {maybe_downstream_session:#?} virtual_host_id = {maybe_virtual_host_id:#?}"
let call_name = "initialize";
let has_user_config = maybe_user_config.is_some();
let virtual_hosts = maybe_user_config.map_or(0, |user_config| user_config.virtual_hosts.len());
let has_session_id = maybe_downstream_session.is_some();
let has_claims = maybe_claims.is_some();
let virtual_host_id = maybe_virtual_host_id.map_or("<missing>", |id| id.value().as_str());
debug!(
"InitializeCallValidator::validate - mcp call validation call_name = {call_name} has_user_config = {has_user_config} virtual_hosts = {virtual_hosts} has_session_id = {has_session_id} has_claims = {has_claims} virtual_host_id = {virtual_host_id}"
);

let Some(downstream_session_id) = maybe_downstream_session else {
Expand Down Expand Up @@ -121,6 +138,12 @@ impl<'a> InitializeCallValidator<'a> {
};

let Some(virtual_host) = user_config.virtual_hosts.get(virtual_host_id.value()) else {
let call_name = "initialize";
let virtual_host_id = virtual_host_id.value();
let virtual_hosts = user_config.virtual_hosts.len();
debug!(
"InitializeCallValidator::validate - mcp virtual host config missing call_name = {call_name} virtual_host_id = {virtual_host_id} virtual_hosts = {virtual_hosts}"
);
return Err(ErrorData {
code: ErrorCode::RESOURCE_NOT_FOUND,
message: "No configuration".into(),
Expand Down
43 changes: 30 additions & 13 deletions crates/contextforge-gateway-rs-lib/src/layers/user_config_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,31 +14,48 @@ pub async fn user_config_store_layer(
mut request: http::Request<axum::body::Body>,
next: Next,
) -> Response {
let method = request.method().clone();
let path = request.uri().path().to_owned();
let maybe_claims = request.extensions().get::<ContextForgeClaims>();
if let Some(claims) = maybe_claims {
let subject = claims.sub.clone();
debug!("Getting user config for {subject:?}");
debug!(
"user_config_store_layer - getting user config for request subject = {subject} method = {method} path = {path}"
);
match state.config_store.get_config(&User::new(&subject)).await {
Ok(user_config) => {
info!(subject, virtual_hosts = user_config.virtual_hosts.len(), "loaded user config");
let virtual_hosts = user_config.virtual_hosts.len();
info!(
"user_config_store_layer - loaded user config subject = {subject} virtual_hosts = {virtual_hosts}"
);
request.extensions_mut().insert(user_config);
next.run(request).await
},

Err(ConfigStoreError::NoDataForKey) => Response::builder()
.status(StatusCode::BAD_REQUEST)
.header(header::CONTENT_TYPE, "text/plain")
.body(Body::from("Problem occurred retrieving the configuration"))
.expect("Expecting this to work"),
Err(ConfigStoreError::NoDataForKey) => {
debug!(
"user_config_store_layer - user config lookup returned no data subject = {subject} method = {method} path = {path}"
);
Response::builder()
.status(StatusCode::BAD_REQUEST)
.header(header::CONTENT_TYPE, "text/plain")
.body(Body::from("Problem occurred retrieving the configuration"))
.expect("Expecting this to work")
},

Err(_) => Response::builder()
.status(StatusCode::INTERNAL_SERVER_ERROR)
.header(header::CONTENT_TYPE, "text/plain")
.body(Body::from("Problem occurred retrieving the configuration"))
.expect("Expecting this to work"),
Err(error) => {
debug!(
"user_config_store_layer - user config lookup failed subject = {subject} method = {method} path = {path} error = {error}"
);
Response::builder()
.status(StatusCode::INTERNAL_SERVER_ERROR)
.header(header::CONTENT_TYPE, "text/plain")
.body(Body::from("Problem occurred retrieving the configuration"))
.expect("Expecting this to work")
},
}
} else {
warn!("No claims");
warn!("user_config_store_layer - no claims found in request extensions method = {method} path = {path}");
Response::builder()
.status(StatusCode::BAD_REQUEST)
.header(header::CONTENT_TYPE, "text/plain")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,13 +14,14 @@ impl VirtualHostId {
}

pub async fn virtual_host_id_layer(mut request: http::Request<axum::body::Body>, next: Next) -> Response {
let uri = request.uri();
let path = request.uri().path().to_owned();

debug!("Extracting virtual host from path {:?}", uri.path());
if let Some(virtual_host_id) = extract_virtual_host_id(uri.path()) {
debug!("virtual_host_id_layer - extracting virtual host from path path = {path}");
if let Some(virtual_host_id) = extract_virtual_host_id(&path) {
request.extensions_mut().insert(virtual_host_id);
next.run(request).await
} else {
debug!("virtual_host_id_layer - failed to extract virtual host id from request path path = {path}");
Response::builder()
.status(StatusCode::BAD_REQUEST)
.header(header::CONTENT_TYPE, "text/plain")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use redis::{
cmd,
};
use tokio::sync::Mutex;
use tracing::{debug, warn};

use super::{ConfigStoreError, UserConfigStore};
use crate::{
Expand All @@ -30,7 +31,10 @@ impl RedisUserConfigStore {
ConnectionManagerConfig::default().set_number_of_retries(REDIS_RETRIES),
)
.await
.map_err(|_| ConfigStoreError::InvalidConnection)?,
.map_err(|error| {
warn!("RedisUserConfigStore::new - failed to create Redis user config connection error = {error}");
ConfigStoreError::InvalidConnection
})?,
cache: Arc::new(Mutex::new(LruCache::with_expiry_duration_and_capacity(
LRU_CACHE_EXPIRY_DURATION,
LRU_CACHE_ENTRIES,
Expand All @@ -42,51 +46,103 @@ impl RedisUserConfigStore {
#[async_trait]
impl UserConfigStore for RedisUserConfigStore {
async fn get_config<'a>(&self, user_key: &'a User) -> Result<UserConfig, ConfigStoreError> {
let has_key = { self.cache.lock().await.contains_key(user_key.key()) };
if has_key {
if let Some(user_config) = self.cache.lock().await.get_mut(user_key.key()) {
Ok(user_config.clone())
} else {
return Err(ConfigStoreError::NoDataForKey);
let subject = user_key.key();

{
let mut cache = self.cache.lock().await;
if let Some(user_config) = cache.get_mut(subject) {
let virtual_hosts = user_config.virtual_hosts.len();
debug!(
"RedisUserConfigStore::get_config - user config cache hit subject = {subject} virtual_hosts = {virtual_hosts}"
);
return Ok(user_config.clone());
}
} else {
let Ok(key) = rmp_serde::encode::to_vec::<User>(user_key) else {
return Err(ConfigStoreError::DataEncoding);
};
}

debug!("RedisUserConfigStore::get_config - user config cache miss subject = {subject}");

let mut connection = self.connection.clone();
let maybe_user_config: Result<Option<Vec<u8>>, RedisError> =
cmd("GET").arg(key).take().query_async(&mut connection).await;
let Ok(key) = rmp_serde::encode::to_vec::<User>(user_key) else {
warn!("RedisUserConfigStore::get_config - failed to encode Redis user config key subject = {subject}");
return Err(ConfigStoreError::DataEncoding);
};

let Ok(Some(user_config)) = maybe_user_config else {
let mut connection = self.connection.clone();
let maybe_user_config: Result<Option<Vec<u8>>, RedisError> =
cmd("GET").arg(key).take().query_async(&mut connection).await;

let user_config = match maybe_user_config {
Ok(Some(user_config)) => {
let bytes = user_config.len();
debug!(
"RedisUserConfigStore::get_config - loaded user config blob from Redis subject = {subject} bytes = {bytes}"
);
user_config
},
Ok(None) => {
debug!("RedisUserConfigStore::get_config - no user config found in Redis subject = {subject}");
return Err(ConfigStoreError::NoDataForKey);
},
Err(error) => {
warn!(
"RedisUserConfigStore::get_config - failed to load user config from Redis subject = {subject} error = {error}"
);
return Err(ConfigStoreError::NoDataForKey);
};
},
};

let Ok(user_config) = rmp_serde::decode::from_slice::<UserConfig>(&user_config) else {
let user_config = match rmp_serde::decode::from_slice::<UserConfig>(&user_config) {
Ok(user_config) => user_config,
Err(error) => {
warn!(
"RedisUserConfigStore::get_config - failed to decode Redis user config blob subject = {subject} error = {error}"
);
return Err(ConfigStoreError::DataWrongFormat);
};
},
};

self.cache.lock().await.insert(user_key.key().to_owned(), user_config.clone());
Ok(user_config)
}
let virtual_hosts = user_config.virtual_hosts.len();
debug!(
"RedisUserConfigStore::get_config - decoded user config subject = {subject} virtual_hosts = {virtual_hosts}"
);

self.cache.lock().await.insert(subject.to_owned(), user_config.clone());
Ok(user_config)
}

async fn set_config<'a>(&self, user_key: &'a User, config: &'a UserConfig) -> Result<(), ConfigStoreError> {
let subject = user_key.key();

let Ok(key) = rmp_serde::encode::to_vec::<User>(user_key) else {
warn!("RedisUserConfigStore::set_config - failed to encode Redis user config key subject = {subject}");
return Err(ConfigStoreError::DataEncoding);
};

let Ok(encoded) = rmp_serde::encode::to_vec::<UserConfig>(config) else {
let virtual_hosts = config.virtual_hosts.len();
warn!(
"RedisUserConfigStore::set_config - failed to encode user config subject = {subject} virtual_hosts = {virtual_hosts}"
);
return Err(ConfigStoreError::DataEncoding);
};

let mut connection = self.connection.clone();

if connection.set::<&[u8], &[u8], String>(&key, &encoded).await.is_ok() {
self.cache.lock().await.insert(user_key.key().to_owned(), config.clone());
Ok(())
} else {
return Err(ConfigStoreError::CantWriteData);
match connection.set::<&[u8], &[u8], String>(&key, &encoded).await {
Ok(_) => {
let bytes = encoded.len();
let virtual_hosts = config.virtual_hosts.len();
debug!(
"RedisUserConfigStore::set_config - wrote user config to Redis subject = {subject} bytes = {bytes} virtual_hosts = {virtual_hosts}"
);
self.cache.lock().await.insert(subject.to_owned(), config.clone());
Ok(())
},
Err(error) => {
warn!(
"RedisUserConfigStore::set_config - failed to write user config to Redis subject = {subject} error = {error}"
);
Err(ConfigStoreError::CantWriteData)
},
}
}
}