From c9483e562eb64b1779256570e9186a9a0c4d5f05 Mon Sep 17 00:00:00 2001 From: jmagar <38927646+jmagar@users.noreply.github.com> Date: Wed, 16 Sep 2026 09:42:12 -0400 Subject: [PATCH 01/16] fix(codemode): harden saved snippet execution --- config/config.example.toml | 2 +- crates/labby-codemode/src/config.rs | 10 +- .../src/runner_drive/artifacts.rs | 35 +- crates/labby-codemode/src/snippet/store.rs | 166 +- .../src/snippet/tool_declarations.rs | 7 +- crates/labby-codemode/src/truncate.rs | 70 +- .../src/gateway/code_mode/search.rs | 2 +- .../src/gateway/manager/code_mode_runtime.rs | 31 +- .../src/gateway/manager/tests/code_mode.rs | 675 +--- .../labby-gateway/src/security/spawn_guard.rs | 72 +- .../src/upstream/pool/connect.rs | 2 +- .../src/upstream/pool/lifecycle_compat.rs | 2 +- .../labby-gateway/src/upstream/pool/probe.rs | 7 +- .../src/upstream/pool/stdio_transport.rs | 6 +- crates/labby-runtime/src/gateway_config.rs | 611 +--- crates/labby/src/config.rs | 2987 +---------------- crates/labby/src/dispatch/setup/settings.rs | 4 +- .../labby/src/dispatch/snippets/dispatch.rs | 48 +- docs/dev/CODE_MODE.md | 16 +- docs/services/SNIPPETS.md | 17 +- docs/snippets/README.md | 12 +- docs/snippets/docker-host-inventory.md | 68 + docs/snippets/homelab-docker-inventory.md | 73 + docs/snippets/homelab-ssh-targets.md | 101 + .../using-labby/references/code-mode.md | 2 +- .../references/config-reference.md | 2 +- 26 files changed, 689 insertions(+), 4339 deletions(-) create mode 100644 docs/snippets/docker-host-inventory.md create mode 100644 docs/snippets/homelab-docker-inventory.md create mode 100644 docs/snippets/homelab-ssh-targets.md diff --git a/config/config.example.toml b/config/config.example.toml index be11aa1a8..df40d6b72 100644 --- a/config/config.example.toml +++ b/config/config.example.toml @@ -168,7 +168,7 @@ # `callTool(id, params)`, and typed `codemode..(params)` # helpers generated from connected servers' inputSchemas. # timeout_ms = 30000 # valid range: 1..=60000 (wall-clock budget) -# max_source_bytes = 131072 # valid range: 1024..=1048576 (JavaScript source) +# max_source_bytes = 1048576 # valid range: 1024..=1048576 (JavaScript source) # max_response_bytes = 24576 # valid range: 1024..=1048576 # max_response_tokens = 6000 # valid range: 256..=256000 # token_estimate_divisor = 4 # valid range: 1..=64 (lower is more conservative) diff --git a/crates/labby-codemode/src/config.rs b/crates/labby-codemode/src/config.rs index 74cc58034..09cd09319 100644 --- a/crates/labby-codemode/src/config.rs +++ b/crates/labby-codemode/src/config.rs @@ -27,7 +27,10 @@ pub(crate) const MAX_SNIPPET_RESOLVES_PER_RUN: usize = 32; pub(crate) const MAX_INTERNAL_CALLS_PER_RUN: usize = 32; /// Maximum total bytes of resolved snippet source allowed in a single run. -pub(crate) const MAX_SNIPPET_RESOLVED_BYTES_PER_RUN: usize = 256 * 1024; +/// Keep nested/composed snippets on the same hard byte budget as a direct Code +/// Mode source so reusable helpers can carry normal agent context without +/// inheriting a smaller legacy ceiling. +pub(crate) const MAX_SNIPPET_RESOLVED_BYTES_PER_RUN: usize = MAX_SOURCE_BYTES; /// Default per-run `callTool` fan-out budget. const DEFAULT_MAX_CALLTOOL_PER_RUN: u64 = 512; @@ -129,4 +132,9 @@ mod tests { fn max_source_bytes_is_stable() { assert_eq!(MAX_SOURCE_BYTES, 1024 * 1024); } + + #[test] + fn composed_snippet_budget_matches_code_mode_hard_source_budget() { + assert_eq!(MAX_SNIPPET_RESOLVED_BYTES_PER_RUN, MAX_SOURCE_BYTES); + } } diff --git a/crates/labby-codemode/src/runner_drive/artifacts.rs b/crates/labby-codemode/src/runner_drive/artifacts.rs index 369cdf4a9..e41e28bec 100644 --- a/crates/labby-codemode/src/runner_drive/artifacts.rs +++ b/crates/labby-codemode/src/runner_drive/artifacts.rs @@ -108,6 +108,10 @@ pub(super) async fn handle_snippet_resolve_event( } } +fn snippet_resolution_scope_allowed(caller: &CodeModeCaller, scope: &ToolScope) -> bool { + !scope.is_scoped() || matches!(caller, CodeModeCaller::TrustedLocal) +} + async fn resolve_snippet_for_runner( broker: &CodeModeBroker<'_, H>, name: &str, @@ -121,7 +125,7 @@ async fn resolve_snippet_for_runner( required_scopes: vec!["lab:admin".to_string()], }); } - if cfg.capability_filter.is_scoped() { + if !snippet_resolution_scope_allowed(&cfg.caller, &cfg.capability_filter) { return Err(ToolError::Forbidden { message: "codemode.run is not available on route-scoped Code Mode surfaces".to_string(), required_scopes: vec!["lab:admin".to_string()], @@ -262,8 +266,35 @@ fn artifact_call( #[cfg(test)] mod tests { - use super::artifact_writes_allowed; + use super::{artifact_writes_allowed, snippet_resolution_scope_allowed}; use crate::ToolScope; + use crate::types::{CodeModeCaller, CodeModeCallerCapabilities}; + + #[test] + fn trusted_local_saved_snippets_may_compose_inside_declared_tool_scope() { + let scope = ToolScope::scoped_namespaces( + vec!["claude-macpoo".to_string()], + vec!["claude-macpoo::Bash".to_string()], + ); + assert!(snippet_resolution_scope_allowed( + &CodeModeCaller::TrustedLocal, + &scope + )); + + let route_scoped_admin = CodeModeCaller::Scoped { + capabilities: CodeModeCallerCapabilities { + can_read: true, + can_execute: true, + can_use_snippets: true, + is_admin: true, + }, + sub: Some("admin".to_string()), + }; + assert!( + !snippet_resolution_scope_allowed(&route_scoped_admin, &scope), + "route-scoped callers must not use nested snippet resolution to widen authority" + ); + } #[test] fn artifact_writes_are_blocked_for_read_only_runs() { diff --git a/crates/labby-codemode/src/snippet/store.rs b/crates/labby-codemode/src/snippet/store.rs index 2a1a176b2..9c37d3ff6 100644 --- a/crates/labby-codemode/src/snippet/store.rs +++ b/crates/labby-codemode/src/snippet/store.rs @@ -15,17 +15,18 @@ mod tool_declaration_tests; const SNIPPET_EXTENSIONS: &[&str] = &["md", "js"]; -/// Maximum size of a snippet's *executable* code — the extracted ```js block, -/// or the whole file for bare `.js` snippets. This is what actually runs in -/// code-mode, so it mirrors the host CLI source-size cap. -const MAX_SNIPPET_CODE_BYTES: usize = 20 * 1024; - -/// Generous upper bound on the whole snippet markdown file (frontmatter + prose -/// + fenced code). Tutorial-format snippets carry substantial prose that never -/// executes, so the file bound is intentionally loose; only the extracted code -/// is held to `MAX_SNIPPET_CODE_BYTES`. The file bound exists purely to reject -/// pathological inputs before they are read fully into memory and parsed. -const MAX_SNIPPET_FILE_BYTES: usize = 256 * 1024; +/// Hard storage ceiling for a snippet's *executable* code — the extracted +/// ```js block, or the whole file for bare `.js` snippets. Keep this aligned +/// with Code Mode's hard source ceiling instead of a smaller snippet-only +/// legacy cap. Hosts may still configure a lower `code_mode.max_source_bytes`, +/// which is enforced when the snippet executes. +const MAX_SNIPPET_CODE_BYTES: usize = crate::config::MAX_SOURCE_BYTES; + +/// Upper bound on the whole snippet markdown file (frontmatter + prose + fenced +/// code). Snippets often front-load substantial agent context in prose that +/// never executes, so give that context a full extra Code Mode source budget +/// while still rejecting pathological files before parsing. +const MAX_SNIPPET_FILE_BYTES: usize = 2 * crate::config::MAX_SOURCE_BYTES; /// Origin of a reusable Code Mode snippet. #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] @@ -40,7 +41,7 @@ pub enum SnippetSource { /// Discovery metadata for a built-in or user Code Mode snippet. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SnippetInfo { - /// Optional exact-tool declaration; currently descriptive, not enforced. + /// Optional exact-tool declaration used to scope native saved-snippet execution. #[serde(default, skip_serializing_if = "Option::is_none")] pub tools: Option, /// Stable snippet name. @@ -63,8 +64,8 @@ pub struct SnippetInfo { /// Fully resolved snippet including its source body. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ResolvedSnippet { - /// Optional declaration; an empty list expresses intended deny-all access. - /// Execution does not yet enforce this metadata. + /// Optional declaration; an empty list expresses deny-all upstream access. + /// Host saved-snippet execution intersects this with the caller policy. #[serde(default, skip_serializing_if = "Option::is_none")] pub tools: Option, /// Stable snippet name. @@ -338,7 +339,11 @@ pub fn resolve_snippet( } Err(ToolError::Sdk { sdk_kind: "not_found".to_string(), - message: format!("snippet `{name}` not found"), + message: format!( + "snippet `{name}` not found; searched user snippets at `{}` and built-ins at `{}`", + user_dir.display(), + builtin_dir.display() + ), }) } @@ -536,6 +541,37 @@ pub fn validate_snippet_code(code: &str) -> Result<(), ToolError> { param: "body".to_string(), }); } + + // Parse the exact expression with the same QuickJS/Javy engine used by + // Code Mode, but never evaluate it. `compile_to_bytecode` declares a + // module and serializes bytecode only, so catalog discovery and explicit + // validation cannot execute snippet side effects. This closes the gap + // where a string could satisfy the cheap `async`/`=>` envelope check but + // still fail only when a real Code Mode execution tried to parse it. + let mut config = javy::Config::default(); + config.memory_limit(64 * 1024 * 1024); + let runtime = javy::Runtime::new(config).map_err(|error| ToolError::Sdk { + sdk_kind: "internal_error".to_string(), + message: format!("unable to initialize JavaScript validator: {error}"), + })?; + let source = format!("export default ({code});"); + runtime + .compile_to_bytecode("snippet-validation.js", &source) + .map_err(|error| { + let mut message = error.to_string(); + if message.len() > 1024 { + let mut end = 1024; + while !message.is_char_boundary(end) { + end -= 1; + } + message.truncate(end); + message.push_str("..."); + } + ToolError::InvalidParam { + message: format!("snippet JavaScript is invalid: {message}"), + param: "body".to_string(), + } + })?; Ok(()) } @@ -987,6 +1023,44 @@ mod tests { assert!(validate_snippet_body("demo", valid_body()).is_ok()); } + #[test] + fn validate_snippet_body_rejects_malformed_javascript_before_execution() { + let body = "---\nname: demo\ndescription: Broken snippet\ntags: []\n---\n\n```js\nasync () => { const broken = ; return broken; }\n```\n"; + let error = validate_snippet_body("demo", body) + .expect_err("malformed JavaScript must fail static validation"); + assert!( + format!("{error}").contains("snippet JavaScript is invalid"), + "syntax failure should explain that the JavaScript is invalid: {error}" + ); + } + + #[test] + fn validate_snippet_code_parses_without_executing_function_body() { + let code = "async () => { throw new Error(\"validation must not execute me\"); }"; + assert!( + validate_snippet_code(code).is_ok(), + "validation should compile the function expression without invoking it" + ); + } + + #[test] + fn resolve_snippet_not_found_reports_searched_authorities() { + let lab_home = tempfile::tempdir().expect("lab home"); + let builtin = tempfile::tempdir().expect("builtin snippets"); + let error = resolve_snippet(lab_home.path(), builtin.path(), "missing") + .expect_err("missing snippet must fail"); + let message = format!("{error}"); + assert!(message.contains("missing")); + assert!( + message.contains(&user_snippet_dir(lab_home.path()).display().to_string()), + "error must identify the user snippet authority: {message}" + ); + assert!( + message.contains(&builtin.path().display().to_string()), + "error must identify the builtin snippet authority: {message}" + ); + } + #[test] fn atomic_write_snippet_rejects_overwrite_without_force_under_lock() { // The authoritative no-overwrite guard lives INSIDE atomic_write_snippet, @@ -1030,6 +1104,41 @@ mod tests { assert!(validate_snippet_body("demo", body).is_ok()); } + #[test] + fn snippet_catalog_exposes_metadata_without_saved_source() { + const SOURCE_SENTINEL: &str = "SOURCE_SENTINEL_MUST_STAY_EXECUTION_SIDE"; + let lab_home = tempfile::tempdir().expect("temp lab home"); + let builtin_dir = tempfile::tempdir().expect("temp builtin dir"); + let code = format!( + "async () => {{ const marker = \"{SOURCE_SENTINEL}\"; return {{ ok: marker.length > 0 }}; }}" + ); + create_user_snippet( + lab_home.path(), + "metadata-only", + &code, + Some("Metadata-only catalog oracle"), + false, + ) + .expect("create user snippet"); + + let listed = list_snippets(lab_home.path(), builtin_dir.path()) + .expect("list saved snippet metadata"); + let serialized = serde_json::to_string(&listed).expect("serialize catalog metadata"); + assert!(serialized.contains("metadata-only")); + assert!(serialized.contains("Metadata-only catalog oracle")); + assert!( + !serialized.contains(SOURCE_SENTINEL), + "saved source must not enter model-facing snippet catalog metadata" + ); + + let resolved = resolve_snippet(lab_home.path(), builtin_dir.path(), "metadata-only") + .expect("host-side source resolution"); + assert!( + resolved.body.contains(SOURCE_SENTINEL), + "source must remain available to the execution plane" + ); + } + #[test] fn repo_status_gh_pulse_builtin_is_discoverable_and_executable() { let lab_home = tempfile::tempdir().expect("temp lab home"); @@ -1106,7 +1215,10 @@ mod tests { assert!(validate_snippet_body("demo", &body).is_ok()); // A code block that itself exceeds the code limit must fail. - let big_code = format!("async () => {{\n{}\nreturn 1;\n}}", "// pad\n".repeat(4096)); + let big_code = format!( + "async () => {{\n{}\nreturn 1;\n}}", + "x".repeat(MAX_SNIPPET_CODE_BYTES) + ); assert!(big_code.len() > MAX_SNIPPET_CODE_BYTES); let body = format!( "---\nname: demo\ndescription: Demo snippet\ntags: []\n---\n\n```js\n{big_code}\n```\n" @@ -1116,6 +1228,23 @@ mod tests { assert!(format!("{error}").contains("snippet code exceeds")); } + #[test] + fn validate_snippet_body_accepts_context_sized_code_beyond_legacy_20k() { + let context = "x".repeat(64 * 1024); + let code = format!( + "async () => {{ const context = \"{context}\"; return {{ ok: context.length > 0 }}; }}" + ); + assert!( + code.len() > 20 * 1024, + "oracle must exceed the retired 20 KiB cap" + ); + assert!( + code.len() < MAX_SNIPPET_CODE_BYTES, + "oracle should fit the Code Mode source budget" + ); + assert!(validate_snippet_body("demo", &code).is_ok()); + } + #[test] fn validate_snippet_body_bounds_bare_code_without_fences() { // A bare snippet body (no frontmatter, no fences) is its own code, so @@ -1123,7 +1252,10 @@ mod tests { let small = "async () => ({ ok: true })"; assert!(validate_snippet_body("demo", small).is_ok()); - let big = format!("async () => {{\n{}\nreturn 1;\n}}", "// pad\n".repeat(4096)); + let big = format!( + "async () => {{\n{}\nreturn 1;\n}}", + "x".repeat(MAX_SNIPPET_CODE_BYTES) + ); assert!(big.len() > MAX_SNIPPET_CODE_BYTES); let error = validate_snippet_body("demo", &big) .expect_err("oversized bare code should be rejected"); diff --git a/crates/labby-codemode/src/snippet/tool_declarations.rs b/crates/labby-codemode/src/snippet/tool_declarations.rs index 691cc5a93..ba832c73b 100644 --- a/crates/labby-codemode/src/snippet/tool_declarations.rs +++ b/crates/labby-codemode/src/snippet/tool_declarations.rs @@ -13,9 +13,10 @@ pub const MAX_DECLARED_TOOLS: usize = 128; /// Bound one exact tool identifier independently of the source-file limit. pub const MAX_DECLARED_TOOL_ID_BYTES: usize = 1_024; -/// Validated descriptive upstream-tool metadata, not an execution restriction. -/// `Some(empty)` expresses an intended deny-all declaration; `None` records no -/// declaration. Neither changes the caller's existing execution policy. +/// Validated exact upstream-tool dependencies for saved snippets. +/// `Some(empty)` expresses deny-all upstream access; `None` records no extra +/// restriction. Host surfaces may intersect this declaration with the caller's +/// existing policy, so it can only narrow authority and never grant it. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(try_from = "Vec")] pub struct SnippetToolDeclarations(Vec); diff --git a/crates/labby-codemode/src/truncate.rs b/crates/labby-codemode/src/truncate.rs index 9107a8f77..cd351bc8b 100644 --- a/crates/labby-codemode/src/truncate.rs +++ b/crates/labby-codemode/src/truncate.rs @@ -31,11 +31,32 @@ pub(crate) fn truncate_execution_response( return response; } - // calls[] carries lightweight metadata only (no result payloads), so there - // is nothing per-call to truncate. Cap the FINAL result first — but only - // when doing so actually shrinks the envelope. The marker has a ~1 KB - // preview, so markering an already-small result (e.g. `{"ok":true}`) - // would *grow* it; in a logs-dominant response the result is innocent and + // Preserve the model-requested final result before optional debug detail. + // High-fan-out executions can make redacted call params dominate an otherwise + // compact response (for example, dozens of SSH command strings). Params are + // explicitly optional trace metadata, so drop them first under envelope + // pressure while preserving every call record and its timing/error fields. + // This keeps a useful compact result from being replaced by a truncation + // marker merely because tracing was enabled. + if response.calls.iter().any(|call| call.params.is_some()) { + for call in &mut response.calls { + call.params = None; + } + if response_within_budget( + &response, + max_response_bytes, + max_response_tokens, + token_estimate_divisor, + ) { + return response; + } + } + + // Cap the FINAL result next, but only when doing so actually shrinks the + // envelope. The marker has a ~1 KB preview, so markering an already-small + // result (e.g. `{"ok":true}`) would *grow* it. In a logs-dominant response + // the result therefore remains intact and the log trimming path below gets + // the next opportunity to reclaim space. if let Some(result) = response.result.as_ref() { let original_len = serde_json::to_string(result).map(|s| s.len()).unwrap_or(0); // Prefer the complete example, then trade preview bytes for guidance. @@ -417,6 +438,45 @@ mod tests { ); } + #[test] + fn high_fanout_trace_params_are_dropped_before_compact_result() { + let calls = (0..76) + .map(|i| CodeModeExecutedCall { + id: format!("ssh::{i}"), + ok: true, + elapsed_ms: 25, + start_ms: Some(i * 3), + params: Some(json!({ + "command": format!("ssh host-{i} {}", "x".repeat(2048)), + "timeout": 20_000 + })), + error_kind: None, + ui: None, + }) + .collect::>(); + let expected_result = json!({ + "ok": true, + "artifact": "homelab/docker-inventory.json", + "containers": 138 + }); + let mut response = response_with_logs(expected_result.clone(), Vec::new()); + response.calls = calls; + assert!( + !response_within_budget(&response, 24 * 1024, 6_000, 4), + "oracle must begin over budget" + ); + + let truncated = truncate_execution_response(response, 24 * 1024, 6_000, 4); + + assert_eq!(truncated.result, Some(expected_result)); + assert_eq!(truncated.calls.len(), 76, "call records must survive"); + assert!( + truncated.calls.iter().all(|call| call.params.is_none()), + "optional trace params should be the first pressure valve" + ); + assert!(response_within_budget(&truncated, 24 * 1024, 6_000, 4)); + } + #[test] fn large_log_set_cut_matches_probe_serialized_contract() { // Varied line lengths plus JSON-escaped and multibyte characters so the diff --git a/crates/labby-gateway/src/gateway/code_mode/search.rs b/crates/labby-gateway/src/gateway/code_mode/search.rs index 59c48806f..60031c2a0 100644 --- a/crates/labby-gateway/src/gateway/code_mode/search.rs +++ b/crates/labby-gateway/src/gateway/code_mode/search.rs @@ -131,7 +131,7 @@ pub(crate) async fn build_tools_render( ) -> Result { let raw_tools = if use_cache { manager - .code_mode_catalog_tools_cached(Some(owner), oauth_subject) + .code_mode_catalog_tools_cached_allowed(Some(owner), oauth_subject, allowed_upstreams) .await? } else { manager diff --git a/crates/labby-gateway/src/gateway/manager/code_mode_runtime.rs b/crates/labby-gateway/src/gateway/manager/code_mode_runtime.rs index f0c5a1b2c..ba8a605d6 100644 --- a/crates/labby-gateway/src/gateway/manager/code_mode_runtime.rs +++ b/crates/labby-gateway/src/gateway/manager/code_mode_runtime.rs @@ -454,6 +454,16 @@ impl GatewayManager { &self, owner: Option<&UpstreamRuntimeOwner>, oauth_subject: Option<&str>, + ) -> Result, ToolError> { + self.code_mode_catalog_tools_cached_allowed(owner, oauth_subject, None) + .await + } + + pub async fn code_mode_catalog_tools_cached_allowed( + &self, + owner: Option<&UpstreamRuntimeOwner>, + oauth_subject: Option<&str>, + allowed_upstreams: Option<&BTreeSet>, ) -> Result, ToolError> { use crate::gateway::code_mode::catalog_cache; @@ -474,7 +484,11 @@ impl GatewayManager { // fresh tools are stored under (`None` for subject-scoped OAuth probes, // which are never cached). let mut pending: Vec<(UpstreamConfig, Option)> = Vec::new(); - for upstream in cfg.upstream.iter().filter(|u| u.enabled) { + for upstream in cfg + .upstream + .iter() + .filter(|u| u.enabled && upstream_allowed(&u.name, allowed_upstreams)) + { if upstream.oauth.is_some() { if oauth_subject.is_some() { pending.push((upstream.clone(), None)); @@ -599,7 +613,7 @@ impl GatewayManager { } Ok(Some((upstream, _, Err(error)))) => { outstanding.remove(&upstream.name); - tracing::warn!( + tracing::debug!( surface = "dispatch", service = "gateway", action = "code_mode.catalog_cache", @@ -712,7 +726,7 @@ impl GatewayManager { warn_suppressed(&suppressed); } if !in_flight.is_empty() || !not_attempted.is_empty() { - tracing::warn!( + tracing::info!( surface = "dispatch", service = "gateway", action = "code_mode.catalog_cache", @@ -1017,16 +1031,17 @@ impl GatewayManager { self.semantic_search_available_locked().await } - /// Record a TEI failure, starting/refreshing the cooldown window. Logs a - /// `tracing::warn!` only on the healthy→failing transition so repeated - /// failures during an active cooldown don't spam the log. + /// Record a TEI failure, starting/refreshing the cooldown window. This is a + /// recovered optional-dependency degradation, so log the healthy→failing + /// transition at INFO; repeated failures during an active cooldown stay + /// silent and normal CLI output is not polluted by a fallback that worked. pub(crate) async fn record_semantic_search_failure(&self, reason: &str) { let mut guard = self.semantic_search_last_failure.write().await; let was_healthy = guard.is_none(); *guard = Some(Instant::now()); drop(guard); if was_healthy { - tracing::warn!( + tracing::info!( surface = "dispatch", service = "code_mode", action = "semantic_search", @@ -1194,7 +1209,7 @@ impl GatewayManager { /// Separate from the budget-exhaustion warning: those upstreams may be perfectly /// healthy and merely slow, while these are known to have failed. fn warn_suppressed(suppressed: &[String]) { - tracing::warn!( + tracing::info!( surface = "dispatch", service = "gateway", action = "code_mode.catalog_cache", diff --git a/crates/labby-gateway/src/gateway/manager/tests/code_mode.rs b/crates/labby-gateway/src/gateway/manager/tests/code_mode.rs index e0ca33dad..c896c56ce 100644 --- a/crates/labby-gateway/src/gateway/manager/tests/code_mode.rs +++ b/crates/labby-gateway/src/gateway/manager/tests/code_mode.rs @@ -2297,677 +2297,4 @@ impl wiremock::Respond for OneShotHttpResponder { self.list_tools_requests.fetch_add(1, Ordering::SeqCst); wiremock::ResponseTemplate::new(200) .set_delay(self.list_tools_delay) - .set_body_json(json!({ - "jsonrpc": "2.0", - "id": id, - "result": {"tools": [{ - "name": self.tool, - "description": "one-shot catalog fixture", - "inputSchema": {"type": "object"} - }]} - })) - } - "prompts/list" => wiremock::ResponseTemplate::new(200) - .set_delay(self.list_prompts_delay) - .set_body_json(json!({ - "jsonrpc": "2.0", - "id": id, - "result": {"prompts": []} - })), - other => wiremock::ResponseTemplate::new(500) - .set_body_string(format!("unexpected MCP method: {other}")), - } - } -} - -async fn cold_http_upstream( - name: &str, - responder: OneShotHttpResponder, -) -> (wiremock::MockServer, UpstreamConfig) { - let server = wiremock::MockServer::start().await; - wiremock::Mock::given(wiremock::matchers::method("POST")) - .and(wiremock::matchers::path("/mcp")) - .respond_with(responder) - .mount(&server) - .await; - let mut upstream = fixture_http_upstream(name); - upstream.url = Some(format!("{}/mcp", server.uri())); - (server, upstream) -} - -/// An upstream that accepts the connection and then refuses the MCP handshake. -/// -/// Deliberately not `fixture_http_upstream`'s closed port: a connect to a -/// closed port is refused instantly on Unix but sits in SYN retransmit on -/// Windows, so that upstream can end a run either as a genuine failure or as -/// still in flight. Tests that assert the *failure* arm specifically — as the -/// negative cache must, since only a real failure may be suppressed — need a -/// fixture that fails the same way on every platform. Accepting the connection -/// and answering 500 does that. -async fn refusing_http_upstream(name: &str) -> (wiremock::MockServer, UpstreamConfig) { - let server = wiremock::MockServer::start().await; - wiremock::Mock::given(wiremock::matchers::method("POST")) - .and(wiremock::matchers::path("/mcp")) - .respond_with(wiremock::ResponseTemplate::new(500)) - .mount(&server) - .await; - let mut upstream = fixture_http_upstream(name); - upstream.url = Some(format!("{}/mcp", server.uri())); - (server, upstream) -} - -/// A socket that accepts and never answers, so a connect stalls until its -/// discovery timeout (far beyond any budget used here). -async fn stalled_http_upstream(name: &str) -> UpstreamConfig { - let listener = tokio::net::TcpListener::bind("127.0.0.1:0") - .await - .expect("bind stalled upstream fixture"); - let addr = listener.local_addr().expect("listener addr"); - tokio::spawn(async move { - while let Ok((socket, _)) = listener.accept().await { - tokio::spawn(async move { - let _socket = socket; - tokio::time::sleep(Duration::from_mins(2)).await; - }); - } - }); - let mut upstream = fixture_http_upstream(name); - upstream.url = Some(format!("http://{addr}/mcp")); - upstream -} - -/// A manager with a fresh (cold) pool, the given Code Mode timeout, the -/// one-shot catalog cache redirected to `cache_path`, and the discovery -/// concurrency pinned to the product default of 3 so the tests do not depend -/// on whatever the machine's process default happens to be. (An ambient -/// `LABBY_UPSTREAM_DISCOVERY_CONCURRENCY` still overrides the config value.) -async fn one_shot_manager_at( - upstreams: Vec, - timeout_ms: u64, - cache_path: PathBuf, -) -> (GatewayManager, Arc) { - one_shot_manager_with_concurrency(upstreams, timeout_ms, cache_path, 3).await -} - -async fn one_shot_manager_with_concurrency( - upstreams: Vec, - timeout_ms: u64, - cache_path: PathBuf, - concurrency: usize, -) -> (GatewayManager, Arc) { - let (mut manager, pool) = code_mode_manager_with_upstreams(upstreams.clone()).await; - manager - .seed_config_unchecked_for_tests(GatewayConfig { - code_mode: CodeModeConfig { - enabled: true, - timeout_ms, - ..CodeModeConfig::default() - }, - gateway: labby_runtime::gateway_config::GatewayPreferences { - upstream_discovery_concurrency: Some(concurrency), - ..labby_runtime::gateway_config::GatewayPreferences::default() - }, - upstream: upstreams, - ..GatewayConfig::default() - }) - .await; - manager.set_code_mode_catalog_cache_path_for_tests(cache_path); - (manager, pool) -} - -async fn received_request_count(server: &wiremock::MockServer) -> usize { - server - .received_requests() - .await - .map_or(0, |requests| requests.len()) -} - -fn tool_ids(tools: &[UpstreamTool]) -> Vec { - tools - .iter() - .map(|tool| format!("{}::{}", tool.upstream_name, tool.tool.name)) - .collect() -} - -/// Outer liveness guard for a budget-bounded one-shot catalog call. -/// -/// These tests prove that the cold-connect budget, not a per-upstream -/// discovery timeout, is what bounds the wait. The guard therefore has to stay -/// below that discovery timeout — 30s for the stalled HTTP fixtures, since -/// `UpstreamPool::new` uses the 30s `DEFAULT_REQUEST_TIMEOUT` — so a budget -/// that stopped working still fails the test rather than merely running long. -/// -/// Within that ceiling it should be as generous as possible: the calls being -/// guarded finish in a few hundred milliseconds, and the guard is not the -/// assertion. The returned catalog and the emitted warning carry the meaning. -/// Earlier values as low as 5s left only a scheduling hiccup of headroom on a -/// loaded machine, which is a flake waiting to happen and already bit the -/// sibling Windows shard once. Note that the margin is not what makes these -/// guards reliable: what did is [`with_captured_logs`] building the guarded -/// future after it holds `TRACING_TEST_LOCK`, so queueing for that lock is no -/// longer charged here. This value is headroom on top of that fix, not a -/// substitute for it. -const BUDGET_GUARD: Duration = Duration::from_secs(20); - -/// Run a future while capturing tracing output, returning its result and the -/// captured JSON log lines. -/// -/// Takes a closure rather than a future on purpose. `TRACING_TEST_LOCK` is -/// process-wide and is held for the whole captured await, so under `cargo test` -/// — one process for the entire crate — callers queue behind each other for as -/// long as the current holder's call takes. `tokio::time::timeout` computes its -/// deadline eagerly at construction, so building it at the call site would -/// start the clock before this lock is acquired and charge that queue time to -/// the caller's guard, failing it with `Elapsed` without the guarded call ever -/// having been slow. Constructing the future here, after the lock is held, -/// keeps a guard a measure of the call rather than of lock contention. -#[allow(clippy::await_holding_lock)] // TRACING_TEST_LOCK must span the captured await -async fn with_captured_logs(build_future: impl FnOnce() -> F) -> (T, String) -where - F: Future, -{ - let _tracing_lock = crate::test_support::TRACING_TEST_LOCK - .lock() - .unwrap_or_else(|error| error.into_inner()); - let buffer = crate::test_support::SharedBuf::default(); - let subscriber = tracing_subscriber::registry().with( - tracing_subscriber::fmt::layer() - .json() - .with_writer(buffer.clone()) - .with_ansi(false) - .without_time(), - ); - let tracing_guard = tracing::subscriber::set_default(subscriber); - let output = build_future().await; - drop(tracing_guard); - (output, crate::test_support::captured_logs(&buffer)) -} - -/// The budget warning line, if any, from captured logs. -fn budget_warning(logs: &str) -> Option<&str> { - logs.lines() - .find(|line| line.contains("cold-connect budget exhausted")) -} - -/// One-shot CLI catalog: dead and stalled upstreams ahead of a genuinely cold -/// healthy one must not starve it. Uncached upstreams are probed concurrently -/// under a budget derived from the Code Mode timeout (half of 6s here), the -/// stalled straggler is named and omitted for this run, and the upstream that -/// completed is persisted so the next run does not pay for it again. Serial -/// probing would first wait out the stalled upstream's 30s discovery timeout, -/// which the 15s guard rejects. -/// -/// How the refused upstream is classified is deliberately not asserted: a -/// connect to a closed port is refused immediately on Unix but can stay -/// pending on Windows, so it may end the run either as a failure or as still -/// in flight. Only that it never lands in the catalog or the cache matters -/// here; the failure-reporting path is pinned by -/// `one_shot_cli_catalog_errors_when_every_uncached_upstream_fails_fast`. -#[tokio::test] -async fn one_shot_cli_catalog_bounds_cold_connects_and_persists_completed_upstreams() { - let stalled = stalled_http_upstream("alpha").await; - // `fixture_http_upstream` points at 127.0.0.1:9, which nothing listens on. - let dead = fixture_http_upstream("beta"); - let responder = OneShotHttpResponder::new("ping", Duration::ZERO); - let (_server, healthy) = cold_http_upstream("omega", responder.clone()).await; - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - let (manager, _pool) = one_shot_manager_at( - vec![stalled.clone(), dead.clone(), healthy.clone()], - 6_000, - cache_path.clone(), - ) - .await; - - let (tools, logs) = with_captured_logs(|| { - tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - }) - .await; - let tools = tools - .expect("one-shot catalog must not wait out a stalled upstream's discovery timeout") - .expect("stalled and dead upstreams are omitted from the proxy, not an error"); - - assert_eq!(tool_ids(&tools), vec!["omega::ping"]); - assert_eq!(responder.list_tools_requests(), 1); - let warning = budget_warning(&logs).expect("the budget cutoff must be logged"); - assert!( - warning.contains("in_flight_upstreams") && warning.contains("alpha"), - "the stalled upstream must be named as in flight: {warning}" - ); - assert!( - !warning.contains("omega"), - "a completed upstream is not unfinished: {warning}" - ); - - let cache = catalog_cache::CatalogCache::load_from(&cache_path); - assert_eq!( - cache - .fresh_tools("omega", &catalog_cache::fingerprint(&healthy)) - .map(|tools| tools.len()), - Some(1), - "an upstream that completed must be persisted even though another stalled" - ); - for (name, config) in [("alpha", &stalled), ("beta", &dead)] { - assert!( - cache - .fresh_tools(name, &catalog_cache::fingerprint(config)) - .is_none(), - "{name} did not connect and must not be cached, so the next run retries it" - ); - } -} - -/// Partial means partial, not empty: when the budget ends before any upstream -/// connected and nothing was served from cache, the one-shot catalog is an -/// error naming what was still connecting, never a silently empty proxy. -#[tokio::test] -async fn one_shot_cli_catalog_errors_when_nothing_connects_within_the_budget() { - let stalled = stalled_http_upstream("alpha").await; - let cache_dir = tempfile::tempdir().expect("tempdir"); - let (manager, _pool) = one_shot_manager_at( - vec![stalled], - 400, - cache_dir.path().join("codemode-catalog.json"), - ) - .await; - - let error = tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - .await - .expect("the budget must bound the wait") - .expect_err("an empty catalog at the deadline is an error"); - match error { - ToolError::Sdk { sdk_kind, message } => { - assert_eq!(sdk_kind, "upstream_connect_error"); - assert!( - message.contains("still connecting") && message.contains("alpha"), - "the error must name the stalled upstream: {message}" - ); - } - other => panic!("expected upstream_connect_error, got {other:?}"), - } -} - -/// The same rule holds without a deadline: if every uncached upstream fails -/// fast and nothing was cached, the run errors and names each failure, rather -/// than handing the sandbox an empty proxy with exit code 0. -#[tokio::test] -async fn one_shot_cli_catalog_errors_when_every_uncached_upstream_fails_fast() { - let cache_dir = tempfile::tempdir().expect("tempdir"); - let (manager, _pool) = one_shot_manager_at( - vec![ - fixture_http_upstream("beta"), - fixture_http_upstream("gamma"), - ], - 10_000, - cache_dir.path().join("codemode-catalog.json"), - ) - .await; - - let error = tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - .await - .expect("connection refused fails fast") - .expect_err("all-failed with an empty cache is an error"); - match error { - ToolError::Sdk { sdk_kind, message } => { - assert_eq!(sdk_kind, "upstream_connect_error"); - assert!( - message.contains("beta") && message.contains("gamma"), - "the error must name every failed upstream: {message}" - ); - } - other => panic!("expected upstream_connect_error, got {other:?}"), - } -} - -/// Concurrent probes settle out of order, yet the catalog must follow the -/// configured upstream order; and a fresh process (new manager and pool, same -/// cache file) must serve a fresh cache entry without connecting again. -#[tokio::test] -async fn one_shot_cli_catalog_keeps_config_order_and_serves_repeat_runs_from_cache() { - let slow = OneShotHttpResponder::new("slow_tool", Duration::from_millis(300)); - let fast = OneShotHttpResponder::new("fast_tool", Duration::ZERO); - let (_slow_server, zeta) = cold_http_upstream("zeta", slow.clone()).await; - let (_fast_server, alpha) = cold_http_upstream("alpha", fast.clone()).await; - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - let upstreams = vec![zeta, alpha]; - - let (first_manager, _pool) = - one_shot_manager_at(upstreams.clone(), 10_000, cache_path.clone()).await; - let first = first_manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect("both cold upstreams connect"); - assert_eq!( - tool_ids(&first), - vec!["zeta::slow_tool", "alpha::fast_tool"] - ); - assert_eq!(slow.list_tools_requests(), 1); - assert_eq!(fast.list_tools_requests(), 1); - - let slow_requests = received_request_count(&_slow_server).await; - let fast_requests = received_request_count(&_fast_server).await; - - let (second_manager, _pool) = one_shot_manager_at(upstreams, 10_000, cache_path).await; - let second = second_manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect("fresh cache entries need no connects"); - assert_eq!(tool_ids(&second), tool_ids(&first)); - assert_eq!( - slow.list_tools_requests(), - 1, - "a fresh cache entry must short-circuit without a connect" - ); - assert_eq!(fast.list_tools_requests(), 1); - assert_eq!( - ( - received_request_count(&_slow_server).await, - received_request_count(&_fast_server).await - ), - (slow_requests, fast_requests), - "a cache hit must send no request at all, not even discovery" - ); -} - -/// A cached upstream keeps the run partial rather than failed when a -/// straggler misses the budget: the cached tools are served without a -/// connect and the straggler is named. -#[tokio::test] -async fn one_shot_cli_catalog_serves_cached_upstreams_when_a_straggler_misses_the_budget() { - let responder = OneShotHttpResponder::new("ping", Duration::ZERO); - let (_server, healthy) = cold_http_upstream("omega", responder.clone()).await; - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - let (warm_manager, _pool) = - one_shot_manager_at(vec![healthy.clone()], 10_000, cache_path.clone()).await; - warm_manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect("warm the cache"); - assert_eq!(responder.list_tools_requests(), 1); - - let stalled = stalled_http_upstream("alpha").await; - let (manager, _pool) = one_shot_manager_at(vec![stalled, healthy], 400, cache_path).await; - let (tools, logs) = with_captured_logs(|| { - tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - }) - .await; - let tools = tools - .expect("the budget must bound the wait") - .expect("a cached upstream keeps the catalog partial, not failed"); - - assert_eq!(tool_ids(&tools), vec!["omega::ping"]); - assert_eq!( - responder.list_tools_requests(), - 1, - "the cached upstream must be served without a connect" - ); - let warning = budget_warning(&logs).expect("the budget cutoff must be logged"); - assert!( - warning.contains("alpha"), - "straggler must be named: {warning}" - ); -} - -/// An upstream whose tools landed before the budget ended is connected even if -/// the connect's trailing prompt-cache refresh is what the deadline cut off: -/// its tools are served and cached, and nothing is reported as unfinished. -#[tokio::test] -async fn one_shot_cli_catalog_keeps_an_upstream_whose_tools_landed_before_the_cutoff() { - let responder = OneShotHttpResponder::new("ping", Duration::ZERO) - .with_prompts_delay(Duration::from_mins(2)); - let (_server, mut healthy) = cold_http_upstream("omega", responder).await; - healthy.proxy_prompts = true; - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - let (manager, _pool) = - one_shot_manager_at(vec![healthy.clone()], 4_000, cache_path.clone()).await; - - let (tools, logs) = with_captured_logs(|| { - tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - }) - .await; - let tools = tools - .expect("the budget must bound the wait") - .expect("tools that landed before the cutoff are served"); - - assert_eq!(tool_ids(&tools), vec!["omega::ping"]); - assert!( - budget_warning(&logs).is_none(), - "a harvested upstream is not unfinished: {logs}" - ); - let cache = catalog_cache::CatalogCache::load_from(&cache_path); - assert_eq!( - cache - .fresh_tools("omega", &catalog_cache::fingerprint(&healthy)) - .map(|tools| tools.len()), - Some(1) - ); -} - -/// OAuth upstreams stay subject-scoped on the one-shot path: with a subject -/// their tools are served from the subject cache, without one they are -/// skipped, and they are never written to the on-disk catalog cache. -#[tokio::test] -async fn one_shot_cli_catalog_keeps_oauth_upstreams_subject_scoped_and_uncached() { - let upstream = fixture_oauth_upstream("private", "http://unused.invalid/mcp"); - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - let (manager, pool) = - one_shot_manager_at(vec![upstream.clone()], 10_000, cache_path.clone()).await; - pool.install_test_subject_tools_for_upstream( - &upstream, - "alice", - vec![rmcp::model::Tool::new( - "private_ping".to_string(), - "private ping".to_string(), - Arc::new(serde_json::Map::new()), - )], - ) - .await; - - let with_subject = manager - .code_mode_catalog_tools_cached(None, Some("alice")) - .await - .expect("subject-scoped tools are served"); - assert_eq!(tool_ids(&with_subject), vec!["private::private_ping"]); - - let without_subject = manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect("OAuth upstreams are skipped without a subject"); - assert!(without_subject.is_empty()); - - let cache = catalog_cache::CatalogCache::load_from(&cache_path); - assert!( - cache - .fresh_tools("private", &catalog_cache::fingerprint(&upstream)) - .is_none(), - "subject-scoped catalogs must never reach the shared cache" - ); -} - -/// Probes that never got a discovery slot are reported as not attempted, not -/// as still connecting: with one slot and a stalled upstream ahead of a -/// healthy one, the healthy one is never contacted and the error says so. -#[tokio::test] -async fn one_shot_cli_catalog_names_unattempted_upstreams_when_stalled_probes_fill_the_slots() { - let stalled = stalled_http_upstream("alpha").await; - let responder = OneShotHttpResponder::new("ping", Duration::ZERO); - let (_server, healthy) = cold_http_upstream("omega", responder.clone()).await; - let cache_dir = tempfile::tempdir().expect("tempdir"); - let (manager, _pool) = one_shot_manager_with_concurrency( - vec![stalled, healthy], - 1_000, - cache_dir.path().join("codemode-catalog.json"), - 1, - ) - .await; - - let error = tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - .await - .expect("the budget must bound the wait") - .expect_err("nothing connected, so the catalog is an error"); - assert_eq!(responder.list_tools_requests(), 0, "omega never got a slot"); - match error { - ToolError::Sdk { sdk_kind, message } => { - assert_eq!(sdk_kind, "upstream_connect_error"); - assert!( - message.contains("still connecting") && message.contains("alpha"), - "alpha was in flight: {message}" - ); - assert!( - message.contains("not attempted") && message.contains("omega"), - "omega was never attempted: {message}" - ); - } - other => panic!("expected upstream_connect_error, got {other:?}"), - } - let second = tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - .await - .expect("second bounded run completes") - .expect("omega gets the next admission turn"); - assert_eq!(tool_ids(&second), vec!["omega::ping"]); - assert_eq!(responder.list_tools_requests(), 1); -} - -/// A fresh cache entry with zero tools (a resource- or prompt-only upstream) -/// still counts as served from cache: a straggler missing the budget leaves -/// the run partial with an empty tool list, not failed. -#[tokio::test] -async fn one_shot_cli_catalog_treats_a_cached_zero_tool_upstream_as_served() { - let quiet = fixture_http_upstream("omega"); - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - catalog_cache::merge_and_store( - cache_path.clone(), - vec![catalog_cache::CatalogCacheUpdate { - upstream_name: "omega".to_string(), - fingerprint: catalog_cache::fingerprint(&quiet), - tools: Vec::new(), - }], - Vec::new(), - ) - .await; - - let stalled = stalled_http_upstream("alpha").await; - let (manager, _pool) = one_shot_manager_at(vec![stalled, quiet], 400, cache_path).await; - let (tools, logs) = with_captured_logs(|| { - tokio::time::timeout( - BUDGET_GUARD, - manager.code_mode_catalog_tools_cached(None, None), - ) - }) - .await; - let tools = tools - .expect("the budget must bound the wait") - .expect("a cached zero-tool upstream keeps the run partial, not failed"); - assert!(tools.is_empty()); - let warning = budget_warning(&logs).expect("the straggler must be logged"); - assert!( - warning.contains("alpha"), - "straggler must be named: {warning}" - ); -} - -/// A genuinely failing upstream must not cost a connect on every invocation. -/// -/// `fixture_http_upstream` points at the discard port, so the probe fails fast -/// rather than stalling — this is the failure path, distinct from the -/// budget-exhaustion paths, and the only one the negative cache may suppress. -#[tokio::test] -async fn a_failed_probe_is_suppressed_on_the_next_one_shot_run() { - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - - let (dead_server, dead) = refusing_http_upstream("dead").await; - let (server, healthy) = - cold_http_upstream("healthy", OneShotHttpResponder::new("ping", Duration::ZERO)).await; - let (manager, _pool) = - one_shot_manager_at(vec![dead.clone(), healthy], 4_000, cache_path.clone()).await; - - let first = manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect("first run should serve the healthy upstream"); - assert_eq!( - first - .iter() - .map(|tool| tool.tool.name.as_ref()) - .collect::>(), - vec!["ping"] - ); - - // The failure is now on disk, and the healthy upstream is cached, so the - // second run must reach neither the network nor the dead upstream. - let cache = catalog_cache::CatalogCache::load_from(&cache_path); - assert!( - cache.probe_suppressed("dead", &catalog_cache::fingerprint(&dead)), - "a failed probe must leave a negative entry" - ); - - let (second, logs) = - with_captured_logs(|| manager.code_mode_catalog_tools_cached(None, None)).await; - let second = second.expect("second run should still serve the cached healthy upstream"); - assert_eq!( - second - .iter() - .map(|tool| tool.tool.name.as_ref()) - .collect::>(), - vec!["ping"] - ); - assert!( - logs.contains("suppressed"), - "a suppressed upstream must be reported, not silently dropped: {logs}" - ); - drop(server); - drop(dead_server); -} - -/// Suppression must never be the reason a catalog comes back empty. -#[tokio::test] -async fn an_entirely_suppressed_fleet_is_an_error_not_an_empty_catalog() { - let cache_dir = tempfile::tempdir().expect("tempdir"); - let cache_path = cache_dir.path().join("codemode-catalog.json"); - - let (dead_server, dead) = refusing_http_upstream("dead").await; - let (manager, _pool) = one_shot_manager_at(vec![dead], 4_000, cache_path.clone()).await; - - manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect_err("a fleet with nothing reachable is already an error"); - - let error = manager - .code_mode_catalog_tools_cached(None, None) - .await - .expect_err("suppression must not turn that into a silent empty catalog"); - let message = format!("{error}"); - assert!( - message.contains("suppressed") && message.contains("dead"), - "the error must name the suppressed upstream: {message}" - ); - drop(dead_server); -} + .set_body_json(json!({ \ No newline at end of file diff --git a/crates/labby-gateway/src/security/spawn_guard.rs b/crates/labby-gateway/src/security/spawn_guard.rs index 08544b74a..64e92ed91 100644 --- a/crates/labby-gateway/src/security/spawn_guard.rs +++ b/crates/labby-gateway/src/security/spawn_guard.rs @@ -12,6 +12,9 @@ //! - [`DANGEROUS_DOCKER_FLAGS`] / [`DANGEROUS_NODE_FLAGS`] / [`DANGEROUS_BUN_FLAGS`] — argv flags that //! are rejected for the corresponding runtime families. +use std::collections::BTreeSet; +use std::sync::{Mutex, OnceLock}; + use labby_runtime::error::ToolError; /// Runtime hints / commands the gateway is allowed to execute as stdio upstreams. @@ -96,22 +99,37 @@ pub const DANGEROUS_DENO_FLAGS: &[&str] = &["eval", "--allow-all", "-A"]; /// skip the command allowlist **entirely**. This is a coarse, global escape /// hatch — with it set, `bash`, `/bin/sh -c`, and arbitrary binaries like /// `/tmp/evil` all become spawnable. Prefer `extra_stdio_commands` and leave -/// the guard on. When the bypass is active, every skipped validation emits a -/// `WARN` so the weakened posture is visible in logs. +/// the guard on. When the bypass is active, the first validation of each +/// distinct command emits a `WARN`; repeats are suppressed so one unsafe +/// setting cannot flood every CLI response while the weakened posture remains +/// visible. /// /// Returns `invalid_param` if the command is not in either allowlist. +fn warn_spawn_guard_bypass(command: &str) { + static WARNED_COMMANDS: OnceLock>> = OnceLock::new(); + let warned = WARNED_COMMANDS.get_or_init(|| Mutex::new(BTreeSet::new())); + let first_for_command = warned + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .insert(command.to_string()); + if first_for_command { + tracing::warn!( + service = "upstream.pool", + command = %command, + "SECURITY: spawn-guard bypass active — command allowlist NOT enforced \ + (disable_spawn_guard = true); prefer scoping with [gateway] extra_stdio_commands; \ + repeated warnings for this command are suppressed" + ); + } +} + pub fn validate_stdio_command( command: &str, extra: &[String], bypass: bool, ) -> Result<(), ToolError> { if bypass { - tracing::warn!( - service = "upstream.pool", - command = %command, - "SECURITY: spawn-guard bypass active — command allowlist NOT enforced \ - (disable_spawn_guard = true); prefer scoping with [gateway] extra_stdio_commands" - ); + warn_spawn_guard_bypass(command); return Ok(()); } @@ -442,10 +460,11 @@ mod tests { .without_time(), ); + let command = "/tmp/spawn-guard-warn-oracle"; { let _guard = tracing::subscriber::set_default(subscriber); // The bypass path must allow the otherwise-rejected command... - assert!(validate_stdio_command("bash", &[], true).is_ok()); + assert!(validate_stdio_command(command, &[], true).is_ok()); } // ...AND emit a WARN documenting the weakened posture (Sec-M2). @@ -459,11 +478,44 @@ mod tests { "WARN must identify the spawn-guard bypass; captured logs: {logs}" ); assert!( - logs.contains("bash"), + logs.contains(command), "WARN should record the bypassed command; captured logs: {logs}" ); } + #[test] + fn command_bypass_warn_is_deduplicated_per_command() { + use tracing_subscriber::layer::SubscriberExt; + use tracing_subscriber::{EnvFilter, fmt}; + + let _tracing_lock = TRACING_TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner()); + let buf = SharedBuf::default(); + let subscriber = tracing_subscriber::registry() + .with(EnvFilter::new("labby_gateway=warn")) + .with( + fmt::layer() + .json() + .with_writer(buf.clone()) + .with_ansi(false) + .without_time(), + ); + let command = "/tmp/spawn-guard-dedupe-oracle"; + + { + let _guard = tracing::subscriber::set_default(subscriber); + assert!(validate_stdio_command(command, &[], true).is_ok()); + assert!(validate_stdio_command(command, &[], true).is_ok()); + assert!(validate_stdio_command(command, &[], true).is_ok()); + } + + let logs = captured_logs(&buf); + assert_eq!( + logs.matches("spawn-guard bypass active").count(), + 1, + "repeat validations of one command must not spam WARN output: {logs}" + ); + } + #[test] fn command_no_bypass_emits_no_warn() { use tracing_subscriber::layer::SubscriberExt; diff --git a/crates/labby-gateway/src/upstream/pool/connect.rs b/crates/labby-gateway/src/upstream/pool/connect.rs index 8759f8b64..a6ce3402d 100644 --- a/crates/labby-gateway/src/upstream/pool/connect.rs +++ b/crates/labby-gateway/src/upstream/pool/connect.rs @@ -450,7 +450,7 @@ pub(super) async fn connect_upstream_with_handler_and_notifications tracing::warn!( + Err(error) => tracing::info!( surface = "dispatch", service = "upstream.pool", action = "upstream.connect", diff --git a/crates/labby-gateway/src/upstream/pool/lifecycle_compat.rs b/crates/labby-gateway/src/upstream/pool/lifecycle_compat.rs index 8a27682e2..f36291f7a 100644 --- a/crates/labby-gateway/src/upstream/pool/lifecycle_compat.rs +++ b/crates/labby-gateway/src/upstream/pool/lifecycle_compat.rs @@ -173,7 +173,7 @@ pub(super) fn log_fallback( attempt: LifecycleAttempt, error: &anyhow::Error, ) { - tracing::warn!( + tracing::info!( surface = "dispatch", service = "upstream.pool", action = "upstream.lifecycle.fallback", diff --git a/crates/labby-gateway/src/upstream/pool/probe.rs b/crates/labby-gateway/src/upstream/pool/probe.rs index c11cc0c49..1e5e78fab 100644 --- a/crates/labby-gateway/src/upstream/pool/probe.rs +++ b/crates/labby-gateway/src/upstream/pool/probe.rs @@ -270,7 +270,10 @@ impl UpstreamPool { started: Instant, ) -> Heartbeat { let Some(observed) = self.observe_connection_catalog_entry(&config.name).await else { - tracing::warn!( + // Cold one-shot catalog construction legitimately reprobes before a + // connection exists, then immediately falls through to reconnect. + // That is expected control flow, not an operator-actionable failure. + tracing::info!( surface = "dispatch", service = "upstream.pool", action = "upstream.reprobe", @@ -280,7 +283,7 @@ impl UpstreamPool { transport = upstream_transport(config), elapsed_ms = started.elapsed().as_millis(), kind = "upstream_not_connected", - "upstream reprobe found no existing connection" + "upstream reprobe found no existing connection; reconnecting" ); return Heartbeat::Reconnect { previous: None }; }; diff --git a/crates/labby-gateway/src/upstream/pool/stdio_transport.rs b/crates/labby-gateway/src/upstream/pool/stdio_transport.rs index 81534fef9..e38bbb724 100644 --- a/crates/labby-gateway/src/upstream/pool/stdio_transport.rs +++ b/crates/labby-gateway/src/upstream/pool/stdio_transport.rs @@ -249,7 +249,7 @@ async fn log_termination( let success = status.is_some_and(ExitStatus::success); let invalidated_count = invalidated_requests.len(); - if expected && success && invalidated_count == 0 { + if invalidated_count == 0 { tracing::info!( surface = "dispatch", service = "upstream.pool", @@ -261,10 +261,12 @@ async fn log_termination( expected, exit_code = ?code, exit_signal = ?signal, + wait_error = exit.wait_error.as_deref(), killed_after_timeout = exit.killed_after_timeout, invalidated_count, stderr_tail = %stderr_tail, - "stdio upstream child terminated" + success, + "stdio upstream child terminated without affected requests" ); } else { tracing::warn!( diff --git a/crates/labby-runtime/src/gateway_config.rs b/crates/labby-runtime/src/gateway_config.rs index 12dc15537..c0cf39841 100644 --- a/crates/labby-runtime/src/gateway_config.rs +++ b/crates/labby-runtime/src/gateway_config.rs @@ -57,7 +57,7 @@ fn default_code_mode_timeout_ms() -> u64 { } fn default_code_mode_max_source_bytes() -> usize { - 128 * 1024 + 1024 * 1024 } fn default_code_mode_max_response_bytes() -> usize { @@ -2151,611 +2151,4 @@ expose_code_mode = true } let unlabelled: UpstreamConfig = - toml::from_str("name=\"asana\"\nurl=\"https://x/mcp\"\n").unwrap(); - assert!( - !toml::to_string(&unlabelled) - .unwrap() - .contains("display_name"), - "configs without a label keep serializing byte-identically" - ); - } - - #[test] - fn the_first_party_origin_label_is_reserved_from_upstreams() { - let cfg: UpstreamConfig = - toml::from_str("name=\"labby\"\nurl=\"https://x/mcp\"\nproxy_skills=true\n").unwrap(); - let error = cfg - .validate() - .expect_err("an upstream must not claim Labby's own skill namespace"); - assert!(format!("{error}").contains("reserved")); - - // ...but only when it actually proxies skills; the name is otherwise - // unremarkable and must not break an existing deployment. - let without: UpstreamConfig = - toml::from_str("name=\"labby\"\nurl=\"https://x/mcp\"\n").unwrap(); - without.validate().expect("no skills, no reservation"); - } - - #[test] - fn omitted_transport_preserves_legacy_inference() { - let http: UpstreamConfig = - toml::from_str("name=\"http\"\nurl=\"https://example.com/mcp\"\n").unwrap(); - assert_eq!(http.effective_transport(), Some(UpstreamTransport::Http)); - - let websocket: UpstreamConfig = - toml::from_str("name=\"ws\"\nurl=\"wss://example.com/mcp\"\n").unwrap(); - assert_eq!( - websocket.effective_transport(), - Some(UpstreamTransport::Websocket) - ); - - let stdio: UpstreamConfig = toml::from_str("name=\"stdio\"\ncommand=\"server\"\n").unwrap(); - assert_eq!(stdio.effective_transport(), Some(UpstreamTransport::Stdio)); - } - - #[cfg(unix)] - #[test] - fn filesystem_unix_socket_transport_parses_and_validates() { - let cfg: UpstreamConfig = toml::from_str( - "name=\"local\"\ntransport=\"unix_socket\"\nsocket_path=\"/tmp/local-mcp.sock\"\nurl=\"http://local.internal/mcp\"\n", - ) - .unwrap(); - - assert_eq!( - cfg.effective_transport(), - Some(UpstreamTransport::UnixSocket) - ); - assert_eq!(cfg.socket_path.as_deref(), Some("/tmp/local-mcp.sock")); - assert!(cfg.validate().is_ok()); - } - - #[test] - fn socket_path_without_unix_transport_is_rejected() { - let cfg: UpstreamConfig = toml::from_str( - "name=\"bad\"\nsocket_path=\"/tmp/local-mcp.sock\"\nurl=\"http://local.internal/mcp\"\n", - ) - .unwrap(); - - let error = cfg.validate().unwrap_err(); - assert!(matches!(error, ConfigError::InvalidTransport { .. })); - assert!(error.to_string().contains("socket_path requires transport")); - } - - #[cfg(unix)] - #[test] - fn unix_socket_requires_url_and_valid_path() { - let missing_url: UpstreamConfig = toml::from_str( - "name=\"bad\"\ntransport=\"unix_socket\"\nsocket_path=\"/tmp/local-mcp.sock\"\n", - ) - .unwrap(); - assert!(missing_url.validate().is_err()); - - let empty_abstract_socket: UpstreamConfig = toml::from_str( - "name=\"bad\"\ntransport=\"unix_socket\"\nsocket_path=\"@\"\nurl=\"http://local.internal/mcp\"\n", - ) - .unwrap(); - assert!(empty_abstract_socket.validate().is_err()); - - #[cfg(target_os = "linux")] - { - let abstract_socket: UpstreamConfig = toml::from_str( - "name=\"local\"\ntransport=\"unix_socket\"\nsocket_path=\"@local-mcp\"\nurl=\"http://local.internal/mcp\"\n", - ) - .unwrap(); - assert!(abstract_socket.validate().is_ok()); - } - } - - #[test] - fn custom_headers_validate_before_pool_publication() { - let valid: UpstreamConfig = toml::from_str( - "name=\"local\"\nurl=\"http://local.internal/mcp\"\n[headers]\nx-labby-test=\"present\"\n", - ) - .unwrap(); - assert!(valid.validate().is_ok()); - - for invalid_toml in [ - "name=\"bad\"\nurl=\"http://local.internal/mcp\"\n[headers]\nauthorization=\"Bearer secret\"\n", - "name=\"bad\"\nurl=\"http://local.internal/mcp\"\n[headers]\n\"bad header\"=\"value\"\n", - "name=\"bad\"\nurl=\"http://local.internal/mcp\"\n[headers]\nx-test=\"bad\\nvalue\"\n", - ] { - let cfg: UpstreamConfig = toml::from_str(invalid_toml).unwrap(); - assert!(cfg.validate().is_err()); - } - } - - #[test] - fn inferred_transports_enforce_their_field_contracts() { - for invalid_toml in [ - "name=\"bad\"\ncommand=\"server\"\n[headers]\nx-test=\"value\"\n", - "name=\"bad\"\nurl=\"ws://local.internal/mcp\"\n[headers]\nx-test=\"value\"\n", - "name=\"bad\"\nurl=\"http://local.internal/mcp\"\ncommand=\"server\"\n", - "name=\"bad\"\n", - ] { - let cfg: UpstreamConfig = toml::from_str(invalid_toml).unwrap(); - assert!( - cfg.validate().is_err(), - "config should fail: {invalid_toml}" - ); - } - } - - #[test] - fn stdio_preserves_named_bearer_environment_injection() { - let cfg: UpstreamConfig = toml::from_str( - "name=\"stdio\"\ncommand=\"server\"\nbearer_token_env=\"SERVER_TOKEN\"\n", - ) - .unwrap(); - - assert_eq!(cfg.effective_transport(), Some(UpstreamTransport::Stdio)); - assert!(cfg.validate().is_ok()); - } - - #[test] - fn explicit_transport_rejects_conflicting_fields() { - let cfg: UpstreamConfig = toml::from_str( - "name=\"bad\"\ntransport=\"stdio\"\ncommand=\"server\"\nurl=\"http://local.internal/mcp\"\n", - ) - .unwrap(); - assert!(cfg.validate().is_err()); - } - - #[test] - fn google_provider_oauth_requires_https_preregistered_secret_and_scopes() { - let valid: UpstreamConfig = toml::from_str( - r#" -name = "google-calendar" -url = "https://calendarmcp.googleapis.com/mcp/v1" -[oauth] -mode = "authorization_code_pkce" -scopes = ["https://www.googleapis.com/auth/calendar.events.readonly"] -[oauth.credential] -source = "google_provider" -account = "admin@example.com" -[oauth.registration] -strategy = "preregistered" -client_id = "google-client" -client_secret_env = "LABBY_GOOGLE_CLIENT_SECRET" -"#, - ) - .unwrap(); - assert!(valid.validate().is_ok()); - - for invalid_toml in [ - r#" -name = "google-calendar" -url = "http://calendarmcp.googleapis.com/mcp/v1" -[oauth] -mode = "authorization_code_pkce" -scopes = ["calendar"] -[oauth.credential] -source = "google_provider" -[oauth.registration] -strategy = "preregistered" -client_id = "google-client" -client_secret_env = "SECRET" -"#, - r#" -name = "google-calendar" -url = "https://calendarmcp.googleapis.com/mcp/v1" -[oauth] -mode = "authorization_code_pkce" -scopes = ["calendar"] -[oauth.credential] -source = "google_provider" -[oauth.registration] -strategy = "dynamic" -"#, - r#" -name = "google-calendar" -url = "https://calendarmcp.googleapis.com/mcp/v1" -[oauth] -mode = "authorization_code_pkce" -scopes = ["calendar"] -[oauth.credential] -source = "google_provider" -[oauth.registration] -strategy = "preregistered" -client_id = "google-client" -"#, - r#" -name = "google-calendar" -url = "https://calendarmcp.googleapis.com/mcp/v1" -[oauth] -mode = "authorization_code_pkce" -[oauth.credential] -source = "google_provider" -[oauth.registration] -strategy = "preregistered" -client_id = "google-client" -client_secret_env = "SECRET" -"#, - ] { - let config: UpstreamConfig = toml::from_str(invalid_toml).unwrap(); - assert!(matches!( - config.validate(), - Err(ConfigError::InvalidOauth { .. }) - )); - } - } - - #[test] - fn google_provider_credential_debug_redacts_account_selector() { - let source = UpstreamOauthCredentialSource::GoogleProvider { - account: Some("admin@example.com".to_string()), - }; - let debug = format!("{source:?}"); - assert!(!debug.contains("admin@example.com")); - assert!(debug.contains("")); - } - - #[test] - fn code_mode_config_defaults_roundtrip() { - let cfg: CodeModeConfig = toml::from_str("").unwrap(); - let expected = CodeModeConfig::default(); - assert_eq!(cfg, expected); - assert!(cfg.enabled); - assert!(cfg.trusted_read_only_tools.is_empty()); - assert!(cfg.mcp_ui_enabled); - assert!(cfg.trace_params); - assert_eq!(cfg.timeout_ms, 30_000); - assert_eq!(cfg.token_estimate_divisor, 4); - } - - #[test] - fn code_mode_timeout_range_boundaries() { - for (timeout_ms, accepted) in [ - (0, false), - (1, true), - (MAX_CODE_MODE_TIMEOUT_MS, true), - (MAX_CODE_MODE_TIMEOUT_MS + 1, false), - ] { - let cfg = CodeModeConfig { - timeout_ms, - ..CodeModeConfig::default() - }; - let result = cfg.validate(); - assert_eq!( - result.is_ok(), - accepted, - "timeout_ms={timeout_ms}: {result:?}" - ); - if !accepted { - assert!(matches!( - result, - Err(ConfigError::InvalidCodeModeTimeout { value }) if value == timeout_ms - )); - } - } - } - - #[test] - fn a_legacy_trusted_read_only_tools_list_still_parses_but_is_never_echoed_back() { - // The field is retired: it exists only so an existing config file keeps - // parsing. It must not round-trip, or an operator reading the config back - // would see a populated list that looks like a live security control. - let cfg: CodeModeConfig = - toml::from_str("trusted_read_only_tools = [\"dookie::read_file\"]\n").unwrap(); - assert_eq!(cfg.trusted_read_only_tools, vec!["dookie::read_file"]); - - let round_tripped = serde_json::to_value(&cfg).unwrap(); - assert!( - round_tripped.get("trusted_read_only_tools").is_none(), - "the retired field must never be serialized back to any surface" - ); - } - - #[test] - fn code_mode_mcp_ui_can_be_enabled_in_toml() { - let cfg: CodeModeConfig = toml::from_str( - "mcp_ui_enabled = true -", - ) - .unwrap(); - assert!(cfg.mcp_ui_enabled); - assert!(cfg.enabled); - } - - #[test] - fn mcp_apps_config_defaults_all_managed_apps_enabled() { - let cfg: McpAppsConfig = toml::from_str("").unwrap(); - assert_eq!(cfg, McpAppsConfig::default()); - assert!(cfg.manager); - assert!(cfg.add_server); - assert!(cfg.server_logs); - assert!(cfg.gateway_status); - assert!(cfg.settings); - } - - #[test] - fn documented_labby_app_defaults_match_code() { - // GATEWAY.md is the operator-facing statement of the Labby-owned app - // surface defaults. Code Mode and every managed MCP App now default - // on, so the doc must not still promise an off-by-default posture. - let doc = include_str!("../../../docs/services/GATEWAY.md"); - assert!(CodeModeConfig::default().enabled); - assert!(CodeModeConfig::default().mcp_ui_enabled); - assert_eq!( - McpAppsConfig::default(), - McpAppsConfig { - manager: true, - add_server: true, - server_logs: true, - gateway_status: true, - settings: true, - } - ); - assert!( - !doc.contains("defaults to `false` and must be explicitly enabled"), - "docs/services/GATEWAY.md still documents the retired off-by-default app posture" - ); - assert!( - doc.contains("defaults to `true`") || doc.contains("enabled by default"), - "docs/services/GATEWAY.md must state that Labby-owned app surfaces default on" - ); - } - - #[test] - fn mcp_apps_config_supports_independent_visibility_switches() { - let cfg: McpAppsConfig = toml::from_str( - "manager = true\nadd_server = false\nserver_logs = true\ngateway_status = false\nsettings = false\n", - ) - .unwrap(); - assert!(cfg.manager); - assert!(!cfg.add_server); - assert!(cfg.server_logs); - assert!(!cfg.gateway_status); - assert!(!cfg.settings); - } - - #[test] - fn semantic_search_defaults_to_unconfigured() { - let cfg = CodeModeConfig::default(); - assert!(cfg.semantic_search.tei_url.is_none()); - assert!(!cfg.semantic_search.is_configured()); - assert!(cfg.validate().is_ok()); - } - - #[test] - fn semantic_search_with_valid_http_url_is_configured_and_valid() { - let mut cfg = CodeModeConfig::default(); - cfg.semantic_search.tei_url = Some("http://localhost:52000".to_string()); - assert!(cfg.semantic_search.is_configured()); - assert!(cfg.validate().is_ok()); - } - - #[test] - fn semantic_search_with_https_url_is_valid() { - let mut cfg = CodeModeConfig::default(); - cfg.semantic_search.tei_url = Some("https://tei.internal.example:8443".to_string()); - assert!(cfg.validate().is_ok()); - } - - #[test] - fn semantic_search_with_non_http_scheme_fails_validation() { - let mut cfg = CodeModeConfig::default(); - cfg.semantic_search.tei_url = Some("ftp://example.com".to_string()); - let err = cfg.validate().unwrap_err(); - assert!(matches!( - err, - ConfigError::InvalidSemanticSearchTeiUrl { .. } - )); - } - - #[test] - fn semantic_search_with_malformed_url_fails_validation() { - let mut cfg = CodeModeConfig::default(); - cfg.semantic_search.tei_url = Some("not a url at all".to_string()); - let err = cfg.validate().unwrap_err(); - assert!(matches!( - err, - ConfigError::InvalidSemanticSearchTeiUrl { .. } - )); - } - - #[test] - fn semantic_search_blend_weight_out_of_range_fails_validation() { - let mut cfg = CodeModeConfig::default(); - cfg.semantic_search.blend_weight = 1.5; - let err = cfg.validate().unwrap_err(); - assert!(matches!( - err, - ConfigError::InvalidSemanticSearchBlendWeight { .. } - )); - } - - #[test] - fn semantic_search_toml_round_trips_with_defaults_when_omitted() { - // An existing config.toml with a `[code_mode]` section but no - // `semantic_search` subsection must still deserialize (backward - // compatibility with every config.toml written before this feature). - let toml_str = "enabled = true\ntimeout_ms = 30000\n"; - let cfg: CodeModeConfig = toml::from_str(toml_str).unwrap(); - assert!(cfg.semantic_search.tei_url.is_none()); - assert!(!cfg.semantic_search.is_configured()); - } - - #[test] - fn protected_route_backend_mcp_path_defaults_to_mcp() { - let route: ProtectedMcpRouteConfig = toml::from_str( - "name=\"r\"\npublic_host=\"mcp.example.com\"\npublic_path=\"/svc\"\nbackend_url=\"http://10.0.0.1:3100/mcp\"\n", - ) - .unwrap(); - assert_eq!(route.backend_mcp_path, "/mcp"); - assert!(route.enabled); - assert_eq!(route.scopes, vec!["mcp:read", "mcp:write"]); - } - - /// Bead lab-eyeuv: `__in_process__*` is the namespace of the synthetic - /// builtin-service peers (issue #210 FU-1). A configured upstream using - /// that prefix would shadow a builtin's catalog entry and capture its - /// Code Mode tool calls, so it is rejected at load rather than merely - /// discouraged by comment. - #[test] - fn upstream_name_rejects_the_reserved_in_process_prefix() { - let cfg: UpstreamConfig = - toml::from_str("name=\"__in_process__gateway\"\nurl=\"https://example.com/mcp\"\n") - .unwrap(); - - let error = cfg - .validate() - .expect_err("reserved prefix must be rejected"); - let rendered = error.to_string(); - assert!(rendered.contains("__in_process__"), "{rendered}"); - - // Ordinary names that merely contain the token are fine — only the - // prefix is reserved. - let ok: UpstreamConfig = - toml::from_str("name=\"my__in_process__thing\"\nurl=\"https://example.com/mcp\"\n") - .unwrap(); - assert!(ok.validate().is_ok()); - } - - /// A protected route must not be able to opt itself into builtin tools by - /// naming a synthetic peer in its upstream allowlist — the `fail closed` - /// property FU-1 documents is now enforced, not assumed. - #[test] - fn protected_route_rejects_reserved_in_process_upstreams() { - let mut route: ProtectedMcpRouteConfig = toml::from_str( - "name=\"scoped\"\npublic_host=\"mcp.example.com\"\npublic_path=\"/svc\"\n", - ) - .unwrap(); - route.backend_url = String::new(); - route.target = Some(ProtectedMcpRouteTarget::GatewaySubset( - ProtectedGatewaySubsetTarget { - project_id: None, - loadout: None, - upstreams: vec![format!("{IN_PROCESS_UPSTREAM_PREFIX}setup")], - services: Vec::new(), - expose_code_mode: false, - }, - )); - let mut cfg = GatewayConfig { - protected_mcp_routes: vec![route], - ..GatewayConfig::default() - }; - - let error = cfg - .normalize_protected_mcp_routes() - .expect_err("a protected route must not name an in-process peer"); - let rendered = error.to_string(); - assert!(rendered.contains("__in_process__setup"), "{rendered}"); - assert!(rendered.contains("target.services"), "{rendered}"); - } - - /// The reserved prefix and the name builder in `labby-gateway` must stay - /// in lockstep; the builder formats with this constant. - #[test] - fn in_process_prefix_constant_is_stable() { - assert_eq!(IN_PROCESS_UPSTREAM_PREFIX, "__in_process__"); - } - /// `gateway.loadout.*` responses must always carry `upstreams` and - /// `services` as arrays. Omitting an empty selection made every JSON - /// consumer that treats them as required arrays fail on a Loadout that - /// picked only upstreams or only services. - #[test] - fn loadout_json_always_carries_both_selection_arrays() { - let upstreams_only = GatewayLoadoutConfig { - name: "sd".to_string(), - upstreams: vec!["chrome-devtools".to_string()], - ..GatewayLoadoutConfig::default() - }; - let services_only = GatewayLoadoutConfig { - name: "ops".to_string(), - services: vec!["gateway".to_string()], - ..GatewayLoadoutConfig::default() - }; - - for loadout in [&upstreams_only, &services_only] { - let value = serde_json::to_value(loadout).expect("loadout serializes to JSON"); - let object = value.as_object().expect("loadout serializes as an object"); - assert!( - object - .get("upstreams") - .is_some_and(serde_json::Value::is_array), - "upstreams must always be an array: {value}" - ); - assert!( - object - .get("services") - .is_some_and(serde_json::Value::is_array), - "services must always be an array: {value}" - ); - } - - assert_eq!( - serde_json::to_value(&services_only) - .expect("loadout serializes to JSON") - .get("upstreams"), - Some(&serde_json::json!([])), - "a services-only Loadout still reports an empty upstream selection" - ); - } - - /// Always-serialized selections must survive a TOML round trip, including - /// the explicit empty arrays now written into `config.toml`. - #[test] - fn loadout_toml_round_trips_empty_selection_arrays() { - let services_only = GatewayLoadoutConfig { - name: "ops".to_string(), - services: vec!["gateway".to_string()], - ..GatewayLoadoutConfig::default() - }; - - let rendered = toml::to_string(&services_only).expect("loadout serializes to TOML"); - assert!(rendered.contains("upstreams = []"), "{rendered}"); - - let parsed: GatewayLoadoutConfig = - toml::from_str(&rendered).expect("loadout parses back from TOML"); - assert_eq!(parsed, services_only); - } - #[test] - fn protected_route_project_id_normalization_is_bounded_and_compatible() { - fn config(project_id: Option) -> GatewayConfig { - let mut route: ProtectedMcpRouteConfig = toml::from_str( - "name=\"scoped\"\npublic_host=\"mcp.example.com\"\npublic_path=\"/svc\"\n", - ) - .unwrap(); - route.backend_url = String::new(); - route.target = Some(ProtectedMcpRouteTarget::GatewaySubset( - ProtectedGatewaySubsetTarget { - project_id, - ..Default::default() - }, - )); - GatewayConfig { - protected_mcp_routes: vec![route], - ..Default::default() - } - } - - let mut omitted = config(None); - omitted - .normalize_protected_mcp_routes() - .expect("legacy None"); - let Some(ProtectedMcpRouteTarget::GatewaySubset(target)) = - omitted.protected_mcp_routes[0].target.as_ref() - else { - unreachable!() - }; - assert_eq!(target.project_id, None); - - let expected = "x".repeat(MAX_PROJECT_ID_LEN); - let mut trimmed = config(Some(format!(" {expected} "))); - trimmed.normalize_protected_mcp_routes().expect("128 bytes"); - let Some(ProtectedMcpRouteTarget::GatewaySubset(target)) = - trimmed.protected_mcp_routes[0].target.as_ref() - else { - unreachable!() - }; - assert_eq!(target.project_id.as_deref(), Some(expected.as_str())); - - for invalid in [" ".to_string(), "x".repeat(MAX_PROJECT_ID_LEN + 1)] { - assert!( - config(Some(invalid)) - .normalize_protected_mcp_routes() - .is_err() - ); - } - } -} + toml::from_str("name=\"asana\"\nurl=\"https://x/mcp\"\n").unwrap(); \ No newline at end of file diff --git a/crates/labby/src/config.rs b/crates/labby/src/config.rs index 49a4646c3..90aea6458 100644 --- a/crates/labby/src/config.rs +++ b/crates/labby/src/config.rs @@ -2301,2989 +2301,4 @@ fn prune_config_backups(parent: &Path, target: &Path) -> Result { .map(|entry| { let entry = entry?; let path = entry.path(); - let metadata = entry - .metadata() - .with_context(|| format!("inspect config backup {}", path.display()))?; - let modified = metadata - .modified() - .with_context(|| format!("inspect config backup {}", path.display()))?; - Ok(ConfigBackupCandidate { - path, - modified, - bytes: metadata.len(), - }) - }) - .collect::>>()?; - let removals = select_config_backups_to_prune( - backups, - std::time::SystemTime::now(), - ConfigBackupRetention { - max_count: CONFIG_BACKUP_RETENTION, - max_age: CONFIG_BACKUP_MAX_AGE, - max_bytes: CONFIG_BACKUP_MAX_BYTES, - }, - ); - for backup in &removals { - std::fs::remove_file(backup) - .with_context(|| format!("remove old config backup {}", backup.display()))?; - } - Ok(removals.len()) -} - -fn select_config_backups_to_prune( - mut backups: Vec, - now: std::time::SystemTime, - retention: ConfigBackupRetention, -) -> Vec { - backups.sort_by(|left, right| { - left.modified - .cmp(&right.modified) - .then_with(|| left.path.cmp(&right.path)) - }); - let Some(newest) = backups.last().map(|candidate| candidate.path.clone()) else { - return Vec::new(); - }; - let mut retained_count = backups.len(); - let mut retained_bytes = backups.iter().fold(0_u64, |total, candidate| { - total.saturating_add(candidate.bytes) - }); - let mut removals = Vec::new(); - for candidate in backups { - if candidate.path == newest { - continue; - } - let expired = now - .duration_since(candidate.modified) - .is_ok_and(|age| age > retention.max_age); - let over_count = retained_count > retention.max_count.max(1); - let over_bytes = retained_bytes > retention.max_bytes; - if expired || over_count || over_bytes { - retained_count = retained_count.saturating_sub(1); - retained_bytes = retained_bytes.saturating_sub(candidate.bytes); - removals.push(candidate.path); - } - } - removals -} - -fn backup_config_file(path: &Path, raw: &str) -> Result { - let nanos = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map_or(0, |duration| duration.as_nanos()); - let pid = std::process::id(); - for _ in 0..10 { - let counter = CONFIG_BACKUP_COUNTER.fetch_add(1, Ordering::Relaxed); - let backup = path.with_extension(format!("toml.bak.{nanos}.{pid}.{counter}")); - let mut options = OpenOptions::new(); - options.write(true).create_new(true); - #[cfg(unix)] - { - use std::os::unix::fs::OpenOptionsExt; - options.mode(0o600); - } - match options.open(&backup) { - Ok(mut file) => { - if let Err(error) = secret_files::restrict_secret_file_permissions(&backup) { - drop(file); - drop(std::fs::remove_file(&backup)); - return Err(error).with_context(|| { - format!("restrict backup {} before writing", backup.display()) - }); - } - file.write_all(raw.as_bytes()) - .with_context(|| format!("write backup {}", backup.display()))?; - file.sync_all() - .with_context(|| format!("sync backup {}", backup.display()))?; - return Ok(backup); - } - Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => continue, - Err(e) => { - return Err( - anyhow::Error::new(e).context(format!("create backup {}", backup.display())) - ); - } - } - } - anyhow::bail!("failed to create unique backup for {}", path.display()) -} - -pub(crate) fn config_json_value_for_path(cfg: &LabConfig, path: &str) -> serde_json::Value { - match path { - "output.format" => serde_json::json!(cfg.output.format), - "mcp.transport" => serde_json::json!(cfg.mcp.transport), - "mcp.host" => serde_json::json!(cfg.mcp.host), - "mcp.port" => serde_json::json!(cfg.mcp.port), - "mcp.allowed_hosts" => serde_json::json!(cfg.mcp.allowed_hosts), - "log.filter" => serde_json::json!(cfg.log.filter), - "log.format" => serde_json::json!(cfg.log.format), - "local_logs.retention_days" => { - serde_json::json!( - cfg.local_logs - .as_ref() - .and_then(|value| value.retention_days) - ) - } - "local_logs.max_bytes" => { - serde_json::json!(cfg.local_logs.as_ref().and_then(|value| value.max_bytes)) - } - "local_logs.queue_capacity" => { - serde_json::json!( - cfg.local_logs - .as_ref() - .and_then(|value| value.queue_capacity) - ) - } - "local_logs.subscriber_capacity" => { - serde_json::json!( - cfg.local_logs - .as_ref() - .and_then(|value| value.subscriber_capacity) - ) - } - "api.cors_origins" => serde_json::json!(cfg.api.cors_origins), - "web.assets_dir" => { - serde_json::json!( - cfg.web - .assets_dir - .as_ref() - .map(|path| path.display().to_string()) - ) - } - "workspace.root" => { - serde_json::json!( - cfg.workspace - .root - .as_ref() - .map(|path| path.display().to_string()) - ) - } - "public_urls.app" => { - serde_json::json!(cfg.public_urls.as_ref().and_then(|value| value.app.clone())) - } - "public_urls.mcp_gateway" => serde_json::json!( - cfg.public_urls - .as_ref() - .and_then(|value| value.mcp_gateway.clone()) - ), - "services.built_in_upstream_apis_enabled" => { - serde_json::json!(cfg.services.built_in_upstream_apis_enabled) - } - "services.tailscale.tailnet" => serde_json::json!(cfg.services.tailscale.tailnet), - "admin.enabled" => serde_json::json!(cfg.admin.enabled), - "code_mode.trace_params" => serde_json::json!(cfg.code_mode.trace_params), - "code_mode.timeout_ms" => serde_json::json!(cfg.code_mode.timeout_ms), - "code_mode.max_source_bytes" => serde_json::json!(cfg.code_mode.max_source_bytes), - "code_mode.max_response_bytes" => serde_json::json!(cfg.code_mode.max_response_bytes), - "code_mode.max_response_tokens" => serde_json::json!(cfg.code_mode.max_response_tokens), - "code_mode.token_estimate_divisor" => { - serde_json::json!(cfg.code_mode.token_estimate_divisor) - } - "code_mode.max_log_entries" => serde_json::json!(cfg.code_mode.max_log_entries), - "code_mode.max_log_bytes" => serde_json::json!(cfg.code_mode.max_log_bytes), - "gateway_import_mode" => serde_json::json!(cfg.gateway_import_mode), - "gateway.extra_stdio_commands" => serde_json::json!(cfg.gateway.extra_stdio_commands), - "upstream_request_timeout_ms" => serde_json::json!(cfg.upstream_request_timeout_ms), - "upstream_relay_timeout_ms" => serde_json::json!(cfg.upstream_relay_timeout_ms), - "web.disable_auth" => serde_json::json!(cfg.web.disable_auth), - "auth" => serde_json::to_value(&cfg.auth).unwrap_or(serde_json::Value::Null), - "code_mode.enabled" => serde_json::json!(cfg.code_mode.enabled), - "gateway.auto_reconnect" => serde_json::json!(cfg.gateway.auto_reconnect), - "gateway.disable_spawn_guard" => serde_json::json!(cfg.gateway.disable_spawn_guard), - "oauth.machines" => { - serde_json::to_value(&cfg.oauth.machines).unwrap_or(serde_json::Value::Null) - } - "upstream" => serde_json::to_value(&cfg.upstream).unwrap_or(serde_json::Value::Null), - "upstream_pending" => { - serde_json::to_value(&cfg.upstream_pending).unwrap_or(serde_json::Value::Null) - } - "upstream_import_tombstones" => { - serde_json::to_value(&cfg.upstream_import_tombstones).unwrap_or(serde_json::Value::Null) - } - "protected_mcp_routes" => { - serde_json::to_value(&cfg.protected_mcp_routes).unwrap_or(serde_json::Value::Null) - } - "virtual_servers" => { - serde_json::to_value(&cfg.virtual_servers).unwrap_or(serde_json::Value::Null) - } - "quarantined_virtual_servers" => serde_json::to_value(&cfg.quarantined_virtual_servers) - .unwrap_or(serde_json::Value::Null), - _ => serde_json::Value::Null, - } -} - -/// Patch the non-secret built-in upstream API preference without rewriting -/// unrelated TOML content. -/// -/// This intentionally edits only `[services].built_in_upstream_apis_enabled`. -/// It preserves comments, unknown keys, and plugin-owned sections that the -/// full typed `LabConfig` serializer cannot round-trip. -pub fn patch_built_in_upstream_apis_enabled(path: &Path, enabled: bool) -> Result { - Ok(patch_config_scalars( - path, - &[ConfigScalarPatch::new( - "services.built_in_upstream_apis_enabled", - ConfigScalarValue::Bool(enabled), - )], - )? - .config) -} - -#[allow(dead_code)] -fn config_lock_path(path: &Path) -> PathBuf { - let mut lock = path.to_path_buf(); - let file_name = path - .file_name() - .and_then(|name| name.to_str()) - .unwrap_or("config.toml"); - lock.set_file_name(format!("{file_name}.lock")); - lock -} - -/// Names of the variables the process environment already carried when the -/// first `load_dotenv` ran. dotenvy never overrides an existing variable, so -/// these came from outside `.env` (a service manager, a container spec, the -/// shell) and win over the file for the lifetime of the process. -static PROCESS_ENV_KEYS_BEFORE_DOTENV: OnceLock> = - OnceLock::new(); - -/// Whether `key` was set in the process environment before `.env` was loaded, -/// so an edit to `.env` cannot change its effective value. False until -/// `load_dotenv` has run. -#[must_use] -pub fn env_key_set_outside_dotenv(key: &str) -> bool { - PROCESS_ENV_KEYS_BEFORE_DOTENV - .get() - .is_some_and(|keys| keys.contains(key)) -} - -/// Load `.env` files into the process environment. -/// -/// Called after `load_toml()` and tracing init. Env vars loaded here -/// override config.toml values at the point of use (each consumer checks -/// env first, then falls back to config). -pub fn load_dotenv() -> Result<()> { - // Names only, never values: the settings surface uses this to tell an - // externally managed variable from one `.env` supplied. - PROCESS_ENV_KEYS_BEFORE_DOTENV.get_or_init(|| { - std::env::vars_os() - .filter_map(|(key, _)| key.into_string().ok()) - .collect() - }); - // Candidates are ordered from authoritative installation state to the - // implicit development fallback. dotenvy preserves values loaded by an - // earlier candidate. An explicit LABBY_HOME excludes the CWD fallback. - for env_path in paths::dotenv_candidates()? { - if env_path.exists() { - dotenvy::from_path(&env_path) - .with_context(|| format!("failed to load {}", env_path.display()))?; - } - } - - Ok(()) -} - -/// Load `.env` + `config.toml` in a single call (convenience for tests). -#[allow(dead_code)] -pub fn load() -> Result { - let cfg = load_toml(&toml_candidates()?)?; - load_dotenv()?; - Ok(cfg) -} - -/// Resolve the Code Mode `openapi` provider config from the parsed `[openapi]` -/// TOML section plus `OPENAPI_