From 1448d996e2bf74268cc969d88580768a149ca036 Mon Sep 17 00:00:00 2001 From: Release Repro Date: Thu, 23 Jul 2026 13:26:20 -0400 Subject: [PATCH] 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 --- .../fabro-auth/src/extra_headers_source.rs | 164 ++++++++++++++++++ lib/crates/fabro-auth/src/lib.rs | 2 + .../fabro-workflow/src/pipeline/initialize.rs | 49 +++++- 3 files changed, 210 insertions(+), 5 deletions(-) create mode 100644 lib/crates/fabro-auth/src/extra_headers_source.rs diff --git a/lib/crates/fabro-auth/src/extra_headers_source.rs b/lib/crates/fabro-auth/src/extra_headers_source.rs new file mode 100644 index 000000000..2e57ee6ea --- /dev/null +++ b/lib/crates/fabro-auth/src/extra_headers_source.rs @@ -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, + headers: HashMap, +} + +impl ExtraHeadersCredentialSource { + #[must_use] + pub fn new(inner: Arc, headers: HashMap) -> Self { + Self { inner, headers } + } +} + +#[async_trait] +impl CredentialSource for ExtraHeadersCredentialSource { + async fn resolve(&self, catalog: &Catalog) -> anyhow::Result { + 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 { + self.inner.configured_providers(catalog).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::{ApiCredential, ResolveError}; + + struct StubSource { + credentials: Vec, + auth_issues: Vec<(ProviderId, ResolveError)>, + } + + #[async_trait] + impl CredentialSource for StubSource { + async fn resolve(&self, _catalog: &Catalog) -> anyhow::Result { + 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 { + self.credentials + .iter() + .map(|c| c.provider.clone()) + .collect() + } + } + + fn credential(provider: &str, extra_headers: HashMap) -> 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")]); + } +} diff --git a/lib/crates/fabro-auth/src/lib.rs b/lib/crates/fabro-auth/src/lib.rs index 1bc135149..77c217317 100644 --- a/lib/crates/fabro-auth/src/lib.rs +++ b/lib/crates/fabro-auth/src/lib.rs @@ -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, diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs index 54a6f1f45..ffe53e180 100644 --- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs @@ -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 { - 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>>, + run_id: fabro_types::RunId, +) -> Arc { + let inner: Arc = 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();