From be1a36a4e7c4e722d2fa5b9061332e70cb4203fc Mon Sep 17 00:00:00 2001 From: Brett Chien Date: Sun, 13 Sep 2026 21:28:33 +0800 Subject: [PATCH] =?UTF-8?q?feat(mcp):=20add=20k8s=5Flogs=20=E2=80=94=20pod?= =?UTF-8?q?=20container=20logs=20for=20k8s-runtime=20fleets?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit deploy_events (ECS control-plane events, archived via EventBridge) has no k8s equivalent and explicitly refuses k8s-runtime fleets — there's no archival system to read because the k8s API serves pod logs directly. Adds a k8s_logs tool built on the kube client k8s_client_for() already wires (list_k8s_contexts/list_namespaces/fleet_config's k8s dispatch): - studio-cp: extract find_k8s_pods() (deployment lookup + live pod list) out of observe_k8s_deployment so both it and the new fetch_k8s_pod_logs() share the same selector logic. fetch_k8s_pod_logs() disambiguates by instance_id (the same pod uid deploy_get/get_agent_states already surface) when more than one pod matches — the shape hit debugging seaturtle's image swap, where a crashed pod sat next to its replacement mid-rollout — and supports `previous` (kubectl logs -p) to read a CrashLoopBackOff pod's last terminated container, since its current log is empty post-restart. - oab-mcp: new k8s_logs tool, dispatched the same way deploy_get/deploy_list require `fleet` for k8s (no bare-cluster k8s path exists for those either). Test plan: - cargo check -p studio-cp, -p oab-mcp: clean - cargo test -p studio-cp/-p oab-mcp --lib: hits the same aws-sdk-ec2 test-cfg OOM on this box PR #153 already documented and deferred to CI (unrelated to this change — cargo check compiles the same code cleanly) 🤖 Generated with Claude Code --- crates/oab-mcp/src/lib.rs | 68 ++++++++++++++++++++++- crates/studio-cp/src/lib.rs | 104 +++++++++++++++++++++++++++++++++--- 2 files changed, 163 insertions(+), 9 deletions(-) diff --git a/crates/oab-mcp/src/lib.rs b/crates/oab-mcp/src/lib.rs index bf9b5fa..1ccc5d0 100644 --- a/crates/oab-mcp/src/lib.rs +++ b/crates/oab-mcp/src/lib.rs @@ -5,7 +5,7 @@ //! tools: //! //! - read: `deploy_list`, `deploy_get`, `get_agent_states`, `deploy_events`, -//! `runtime_context`, `fleet_config` +//! `k8s_logs`, `runtime_context`, `fleet_config` //! - write: `deploy_apply`, `deploy_provision`, `deploy_scale`, `deploy_delete`, //! `fleet_config_write` //! @@ -115,6 +115,21 @@ pub fn tools() -> Vec { } })), ), + Tool::new( + "k8s_logs", + "Fetch one Pod's container log for a k8s-runtime fleet — the k8s counterpart to `deploy_events`, which only supports ECS (control-plane events aren't container stdout/stderr, and k8s has no archival system to read; the k8s API serves pod logs directly). Requires `fleet` naming a k8s-runtime fleet (see fleet_config), same as deploy_get/deploy_list's k8s dispatch.", + as_map(json!({ + "type": "object", + "properties": { + "service": { "type": "string", "description": "k8s reconstructed oab-{namespace}-{name}, or bare agent name." }, + "fleet": { "type": "string", "description": "Fleet name (see fleet_config) naming a k8s-runtime fleet; targets its context/namespace. Required — no bare-cluster k8s dispatch exists, same as deploy_get/deploy_list." }, + "instance_id": { "type": "string", "description": "Pod uid to select, from deploy_get/get_agent_states's InstancePhase.id. Required when more than one Pod matches (a rollout in flight, or a crashed Pod next to its replacement); optional with exactly one live Pod." }, + "tail_lines": { "type": "integer", "description": "Max lines from the end of the log (default 200, max 10000)." }, + "previous": { "type": "boolean", "description": "Read the last terminated container's log instead of the current one (like `kubectl logs -p`) — the only way to see why a CrashLoopBackOff pod died, since its current log is empty post-restart. Default false." } + }, + "required": ["service"] + })), + ), Tool::new( "deploy_apply", "Apply an OABService/OABFleet manifest (create or update). Returns the number of services reconciled.", @@ -416,6 +431,7 @@ impl OabMcp { "deploy_get" => self.t_get(args).await, "get_agent_states" => self.t_states(args).await, "deploy_events" => self.t_events(args).await, + "k8s_logs" => self.t_k8s_logs(args).await, "deploy_apply" => self.t_apply(args).await, "deploy_provision" => self.t_provision(args).await, "deploy_provision_agent" => self.t_provision_agent(args).await, @@ -648,6 +664,53 @@ impl OabMcp { })) } + async fn t_k8s_logs(&self, args: &Map) -> Result { + let service = args + .get("service") + .and_then(Value::as_str) + .ok_or_else(|| anyhow::anyhow!("missing required arg: service"))?; + let Some(b) = self.named_fleet(args)? else { + anyhow::bail!( + "k8s_logs requires `fleet` naming a k8s-runtime fleet (see fleet_config) — no bare-cluster k8s dispatch exists, same as deploy_get/deploy_list" + ); + }; + if b.runtime != scp::FleetRuntime::K8s { + anyhow::bail!( + "fleet {:?} is not a k8s-runtime fleet; use deploy_events for an ecs fleet's logs", + b.name + ); + } + let namespace = b.namespace.clone().unwrap_or_else(|| "default".to_string()); + let instance_id = args.get("instance_id").and_then(Value::as_str); + let tail_lines = args + .get("tail_lines") + .and_then(Value::as_i64) + .unwrap_or(200) + .clamp(1, 10_000); + let previous = args + .get("previous") + .and_then(Value::as_bool) + .unwrap_or(false); + match scp::fetch_k8s_pod_logs( + b.context.as_deref(), + &namespace, + service, + instance_id, + tail_lines, + previous, + ) + .await? + { + Some(logs) => Ok(json!({ + "service": service, + "namespace": namespace, + "previous": previous, + "logs": logs, + })), + None => Ok(json!({ "found": false, "service": service })), + } + } + async fn t_provision(&self, args: &Map) -> Result { let namespace = args .get("namespace") @@ -1165,12 +1228,13 @@ mod tests { .iter() .map(|t| t["name"].as_str().expect("tool has a name").to_string()) .collect(); - assert_eq!(names.len(), 17); + assert_eq!(names.len(), 18); for expected in [ "deploy_list", "deploy_get", "get_agent_states", "deploy_events", + "k8s_logs", "deploy_apply", "deploy_provision", "deploy_provision_agent", diff --git a/crates/studio-cp/src/lib.rs b/crates/studio-cp/src/lib.rs index 27be58a..b733852 100644 --- a/crates/studio-cp/src/lib.rs +++ b/crates/studio-cp/src/lib.rs @@ -312,15 +312,18 @@ pub fn k8s_instance_phase(pod: &k8s_openapi::api::core::v1::Pod, verified_before K8sDriver.project(&native, verified_before).classify() } -/// Observe one k8s Deployment end-to-end: replica counters + per-Pod phase — -/// the k8s counterpart to [`observe_deployment`]. `service` matches either -/// the reconstructed `oab-{namespace}-{name}` form or the bare agent name -/// (same dual-match spirit as ECS's `resolve_service`). -pub async fn observe_k8s_deployment( +/// Resolve `service` (either the reconstructed `oab-{namespace}-{name}` form +/// or the bare agent name) to its k8s Deployment's oab short name, desired +/// replica count, and the live Pods its selector currently matches. Shared by +/// [`observe_k8s_deployment`] (which only needs `phase`) and +/// [`fetch_k8s_pod_logs`] (which needs the raw `Pod` — specifically +/// `metadata.name`, which `InstancePhase` deliberately doesn't carry since +/// ECS's `InstancePhase.id` has no equivalent "log by this name" use). +async fn find_k8s_pods( context: Option<&str>, namespace: &str, service: &str, -) -> anyhow::Result> { +) -> anyhow::Result)>> { use k8s_openapi::api::apps::v1::Deployment as K8sDeployment; use k8s_openapi::api::core::v1::Pod; use kube::api::{Api, ListParams}; @@ -357,8 +360,23 @@ pub async fn observe_k8s_deployment( .await .map_err(|e| anyhow::anyhow!("failed to list pods for k8s deployment '{name}': {e}"))?; + Ok(Some((name, desired, pods.items))) +} + +/// Observe one k8s Deployment end-to-end: replica counters + per-Pod phase — +/// the k8s counterpart to [`observe_deployment`]. `service` matches either +/// the reconstructed `oab-{namespace}-{name}` form or the bare agent name +/// (same dual-match spirit as ECS's `resolve_service`). +pub async fn observe_k8s_deployment( + context: Option<&str>, + namespace: &str, + service: &str, +) -> anyhow::Result> { + let Some((name, desired, pods)) = find_k8s_pods(context, namespace, service).await? else { + return Ok(None); + }; + let instances: Vec = pods - .items .iter() .map(|p| InstancePhase { id: p.metadata.uid.clone().unwrap_or_default(), @@ -380,6 +398,78 @@ pub async fn observe_k8s_deployment( })) } +/// Fetch one Pod's container log for a k8s-runtime deployment — the k8s +/// counterpart ECS gets from `deploy_events`'s CloudWatch-archived stream, +/// which doesn't exist for k8s (control-plane events aren't the same thing +/// as container stdout/stderr, and there's no archival system to read here — +/// the k8s API serves pod logs directly). +/// +/// `instance_id` selects a specific Pod by the same `uid` `deploy_get`/ +/// `get_agent_states` already surface as `InstancePhase.id` — required +/// whenever more than one Pod matches (e.g. a rollout in flight, or a crashed +/// Pod sitting next to its replacement, exactly the shape hit debugging +/// "seaturtle"'s `0.9.0` → `0.10.0-beta.4` image swap). With exactly one Pod, +/// `instance_id` may be omitted. `previous` reads the last terminated +/// container's log (`kubectl logs -p`) — the only way to see why a +/// CrashLoopBackOff pod died, since its current log is empty post-restart. +pub async fn fetch_k8s_pod_logs( + context: Option<&str>, + namespace: &str, + service: &str, + instance_id: Option<&str>, + tail_lines: i64, + previous: bool, +) -> anyhow::Result> { + use kube::api::{Api, LogParams}; + + let Some((name, _desired, pods)) = find_k8s_pods(context, namespace, service).await? else { + return Ok(None); + }; + let pod = match instance_id { + Some(id) => pods + .iter() + .find(|p| p.metadata.uid.as_deref() == Some(id)) + .ok_or_else(|| { + let known: Vec<&str> = pods + .iter() + .filter_map(|p| p.metadata.uid.as_deref()) + .collect(); + anyhow::anyhow!( + "no pod with instance_id {id:?} for '{name}'; live instance ids: [{}]", + known.join(", ") + ) + })?, + None => match pods.len() { + 1 => &pods[0], + 0 => return Ok(None), + n => anyhow::bail!( + "'{name}' has {n} live pods — pass instance_id to pick one (see deploy_get/get_agent_states)" + ), + }, + }; + let pod_name = pod + .metadata + .name + .clone() + .ok_or_else(|| anyhow::anyhow!("pod for '{name}' has no metadata.name"))?; + + let client = k8s_client_for(context).await?; + let pod_api: Api = Api::namespaced(client, namespace); + let logs = pod_api + .logs( + &pod_name, + &LogParams { + tail_lines: Some(tail_lines), + previous, + timestamps: true, + ..Default::default() + }, + ) + .await + .map_err(|e| anyhow::anyhow!("failed to fetch logs for pod '{pod_name}': {e}"))?; + Ok(Some(logs)) +} + // ---- Effective runtime identity/context (ADR: Per-Fleet managing identity) -- // // Read-only observation of *who this control plane is actually acting as*. The