feat(llm): send x-session-id trace header with the run ID

Tag every LLM request in a run with an x-session-id header carrying the
run ID, so gateways that understand session tracing (e.g. OpenRouter
broadcast) can group a run's requests into one session.

Adds ExtraHeadersCredentialSource to fabro-auth: a CredentialSource
decorator that appends fixed headers to every resolved credential,
leaving operator-configured extra_headers untouched. The run pipeline
wraps its vault/env source with it, so agent stages, prompt stages,
hooks, and PR-content generation all pick up the header through the
existing extra_headers plumbing with no fabro-llm changes.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Release Repro 2026-07-23 13:26:20 -04:00
parent 30d770046a
commit 1448d996e2
No known key found for this signature in database
3 changed files with 210 additions and 5 deletions

View file

@ -0,0 +1,164 @@
use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use fabro_model::{Catalog, ProviderId};
use crate::credential_source::{CredentialSource, ResolvedCredentials};
/// Decorates another [`CredentialSource`] by appending fixed extra headers to
/// every credential it resolves.
///
/// Headers already present on a credential (for example from explicit
/// provider configuration) are left untouched.
pub struct ExtraHeadersCredentialSource {
inner: Arc<dyn CredentialSource>,
headers: HashMap<String, String>,
}
impl ExtraHeadersCredentialSource {
#[must_use]
pub fn new(inner: Arc<dyn CredentialSource>, headers: HashMap<String, String>) -> Self {
Self { inner, headers }
}
}
#[async_trait]
impl CredentialSource for ExtraHeadersCredentialSource {
async fn resolve(&self, catalog: &Catalog) -> anyhow::Result<ResolvedCredentials> {
let mut resolved = self.inner.resolve(catalog).await?;
for credential in &mut resolved.credentials {
for (name, value) in &self.headers {
credential
.extra_headers
.entry(name.clone())
.or_insert_with(|| value.clone());
}
}
Ok(resolved)
}
async fn configured_providers(&self, catalog: &Catalog) -> Vec<ProviderId> {
self.inner.configured_providers(catalog).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{ApiCredential, ResolveError};
struct StubSource {
credentials: Vec<ApiCredential>,
auth_issues: Vec<(ProviderId, ResolveError)>,
}
#[async_trait]
impl CredentialSource for StubSource {
async fn resolve(&self, _catalog: &Catalog) -> anyhow::Result<ResolvedCredentials> {
Ok(ResolvedCredentials {
credentials: self.credentials.clone(),
auth_issues: self
.auth_issues
.iter()
.map(|(provider, _)| {
(
provider.clone(),
ResolveError::NotConfigured(provider.clone()),
)
})
.collect(),
})
}
async fn configured_providers(&self, _catalog: &Catalog) -> Vec<ProviderId> {
self.credentials
.iter()
.map(|c| c.provider.clone())
.collect()
}
}
fn credential(provider: &str, extra_headers: HashMap<String, String>) -> ApiCredential {
ApiCredential {
provider: ProviderId::new(provider),
auth_header: None,
extra_headers,
base_url: None,
codex_mode: false,
org_id: None,
project_id: None,
}
}
fn catalog() -> Catalog {
Catalog::from_builtin().unwrap()
}
#[tokio::test]
async fn appends_headers_to_every_resolved_credential() {
let source = ExtraHeadersCredentialSource::new(
Arc::new(StubSource {
credentials: vec![
credential("anthropic", HashMap::new()),
credential("openai", HashMap::new()),
],
auth_issues: Vec::new(),
}),
HashMap::from([("x-session-id".to_string(), "run-123".to_string())]),
);
let resolved = source.resolve(&catalog()).await.unwrap();
assert_eq!(resolved.credentials.len(), 2);
for credential in &resolved.credentials {
assert_eq!(
credential.extra_headers.get("x-session-id"),
Some(&"run-123".to_string())
);
}
}
#[tokio::test]
async fn preserves_headers_already_set_on_a_credential() {
let source = ExtraHeadersCredentialSource::new(
Arc::new(StubSource {
credentials: vec![credential(
"openrouter",
HashMap::from([("x-session-id".to_string(), "configured".to_string())]),
)],
auth_issues: Vec::new(),
}),
HashMap::from([("x-session-id".to_string(), "run-123".to_string())]),
);
let resolved = source.resolve(&catalog()).await.unwrap();
assert_eq!(
resolved.credentials[0].extra_headers.get("x-session-id"),
Some(&"configured".to_string())
);
}
#[tokio::test]
async fn passes_through_auth_issues_and_configured_providers() {
let provider = ProviderId::new("anthropic");
let source = ExtraHeadersCredentialSource::new(
Arc::new(StubSource {
credentials: vec![credential("openai", HashMap::new())],
auth_issues: vec![(
provider.clone(),
ResolveError::NotConfigured(provider.clone()),
)],
}),
HashMap::from([("x-session-id".to_string(), "run-123".to_string())]),
);
let resolved = source.resolve(&catalog()).await.unwrap();
assert_eq!(resolved.auth_issues.len(), 1);
assert_eq!(resolved.auth_issues[0].0, provider);
let providers = source.configured_providers(&catalog()).await;
assert_eq!(providers, vec![ProviderId::new("openai")]);
}
}

View file

@ -2,6 +2,7 @@ mod context;
mod credential;
mod credential_source;
mod env_source;
mod extra_headers_source;
mod refresh;
mod resolve;
mod sql_vault_source;
@ -15,6 +16,7 @@ pub use context::{AuthContextRequest, AuthContextResponse};
pub use credential::{ApiKeyHeader, OAuthConfig, OAuthCredential, OAuthTokens};
pub use credential_source::{CredentialSource, ResolvedCredentials};
pub use env_source::EnvCredentialSource;
pub use extra_headers_source::ExtraHeadersCredentialSource;
pub use refresh::refresh_oauth_credential;
pub use resolve::{
ApiCredential, CredentialResolver, CredentialUsage, EnvLookup, ResolveError,

View file

@ -5,7 +5,8 @@ use std::time::Instant;
use fabro_agent::{Sandbox, ToolSecrets};
use fabro_auth::{
CredentialSource, EnvCredentialSource, VaultCredentialSource, auth_issue_message,
CredentialSource, EnvCredentialSource, ExtraHeadersCredentialSource, VaultCredentialSource,
auth_issue_message,
};
use fabro_graphviz::graph;
use fabro_hooks::{HookContext, HookDecision, HookEvent, HookExecutionContext, HookRunner};
@ -260,11 +261,23 @@ fn graph_needs_api_backend(graph: &graph::Graph) -> bool {
graph.nodes.values().any(routing::node_needs_api_backend)
}
fn build_llm_source(vault: Option<Arc<AsyncRwLock<Vault>>>) -> Arc<dyn CredentialSource> {
match vault {
/// Trace header attached to every LLM request in a run so gateways that
/// understand it (e.g. OpenRouter broadcast) can group the run's requests
/// into one session. Explicit `extra_headers` provider configuration wins.
const SESSION_ID_HEADER: &str = "x-session-id";
fn build_llm_source(
vault: Option<Arc<AsyncRwLock<Vault>>>,
run_id: fabro_types::RunId,
) -> Arc<dyn CredentialSource> {
let inner: Arc<dyn CredentialSource> = match vault {
Some(vault) => Arc::new(VaultCredentialSource::new(vault)),
None => Arc::new(EnvCredentialSource::new()),
}
};
Arc::new(ExtraHeadersCredentialSource::new(
inner,
HashMap::from([(SESSION_ID_HEADER.to_string(), run_id.to_string())]),
))
}
/// INITIALIZE phase: prepare the sandbox, env, and handlers for execution.
@ -277,7 +290,7 @@ pub async fn initialize(
options.run_options.run_dir = run_dir.clone();
options.run_options.git = options.git.clone();
let llm_source = build_llm_source(options.vault.clone());
let llm_source = build_llm_source(options.vault.clone(), options.run_options.run_id);
let tool_secrets = tool_secrets_from_configured_sources(options.vault.as_ref()).await;
let catalog = Arc::clone(&options.catalog);
let sandbox_git = Arc::new(SandboxGitRuntime::new());
@ -1028,6 +1041,32 @@ mod tests {
assert!(!effective_dry_run);
}
#[tokio::test]
async fn build_llm_source_appends_run_session_trace_header() {
let dir = tempfile::tempdir().unwrap();
let mut vault = Vault::load(dir.path().join("secrets.json")).unwrap();
vault
.set(
"ANTHROPIC_API_KEY",
"anthropic-key",
SecretType::Token,
None,
)
.unwrap();
let vault = Arc::new(AsyncRwLock::new(vault));
let source = build_llm_source(Some(vault), test_run_id());
let resolved = source.resolve(test_catalog().as_ref()).await.unwrap();
assert!(!resolved.credentials.is_empty());
for credential in &resolved.credentials {
assert_eq!(
credential.extra_headers.get(SESSION_ID_HEADER),
Some(&test_run_id().to_string())
);
}
}
#[tokio::test]
async fn initialize_executes_acp_backend_node_from_registry() {
let temp = tempfile::tempdir().unwrap();