Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
68 changes: 66 additions & 2 deletions crates/oab-mcp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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`
//!
Expand Down Expand Up @@ -115,6 +115,21 @@ pub fn tools() -> Vec<Tool> {
}
})),
),
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.",
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -648,6 +664,53 @@ impl OabMcp {
}))
}

async fn t_k8s_logs(&self, args: &Map<String, Value>) -> Result<Value> {
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<String, Value>) -> Result<Value> {
let namespace = args
.get("namespace")
Expand Down Expand Up @@ -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",
Expand Down
104 changes: 97 additions & 7 deletions crates/studio-cp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Option<Deployment>> {
) -> anyhow::Result<Option<(String, i32, Vec<k8s_openapi::api::core::v1::Pod>)>> {
use k8s_openapi::api::apps::v1::Deployment as K8sDeployment;
use k8s_openapi::api::core::v1::Pod;
use kube::api::{Api, ListParams};
Expand Down Expand Up @@ -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<Option<Deployment>> {
let Some((name, desired, pods)) = find_k8s_pods(context, namespace, service).await? else {
return Ok(None);
};

let instances: Vec<InstancePhase> = pods
.items
.iter()
.map(|p| InstancePhase {
id: p.metadata.uid.clone().unwrap_or_default(),
Expand All @@ -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<Option<String>> {
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<k8s_openapi::api::core::v1::Pod> = 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
Expand Down
Loading