Skip to content
62 changes: 46 additions & 16 deletions crates/acp-tunnel/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -227,7 +227,9 @@ impl AgentRegistry {
return Err("every [[agent]] needs a name".to_string());
}
if !seen.insert(name) {
return Err(format!("duplicate agent name {name:?} — names must be unique"));
return Err(format!(
"duplicate agent name {name:?} — names must be unique"
));
}
if a.management {
management += 1;
Expand Down Expand Up @@ -272,17 +274,23 @@ cwd = "/work"

#[test]
fn cwd_defaults_when_omitted() {
let cfg = RemoteConfig::parse(r#"url = "wss://x/acp"
token = "t""#)
.expect("parse");
let cfg = RemoteConfig::parse(
r#"url = "wss://x/acp"
token = "t""#,
)
.expect("parse");
assert_eq!(cfg.cwd, "/");
}

#[test]
fn validate_rejects_missing_url_bad_scheme_and_missing_token() {
assert!(RemoteConfig { url: "".into(), token: "t".into(), cwd: "/".into() }
.validate()
.is_err());
assert!(RemoteConfig {
url: "".into(),
token: "t".into(),
cwd: "/".into()
}
.validate()
.is_err());
assert!(RemoteConfig {
url: "http://x/acp".into(),
token: "t".into(),
Expand Down Expand Up @@ -339,7 +347,7 @@ token = "s2"
let mira = reg.get("mira").expect("mira present");
assert_eq!(mira.cwd, "/");
assert!(!mira.management); // management defaults off (least privilege)
// management() returns the one flagged entry.
// management() returns the one flagged entry.
assert_eq!(reg.management().map(|a| a.name.as_str()), Some("orca"));
}

Expand Down Expand Up @@ -372,23 +380,42 @@ token = "s2"
#[test]
fn validate_rejects_dup_names_missing_names_and_two_managements() {
assert!(AgentRegistry {
agents: vec![AgentEndpoint { name: "a".into(), ..Default::default() },
AgentEndpoint { name: "a".into(), ..Default::default() }],
agents: vec![
AgentEndpoint {
name: "a".into(),
..Default::default()
},
AgentEndpoint {
name: "a".into(),
..Default::default()
}
],
}
.validate()
.unwrap_err()
.contains("unique"));

assert!(AgentRegistry {
agents: vec![AgentEndpoint { name: " ".into(), ..Default::default() }],
agents: vec![AgentEndpoint {
name: " ".into(),
..Default::default()
}],
}
.validate()
.is_err());

assert!(AgentRegistry {
agents: vec![
AgentEndpoint { name: "a".into(), management: true, ..Default::default() },
AgentEndpoint { name: "b".into(), management: true, ..Default::default() },
AgentEndpoint {
name: "a".into(),
management: true,
..Default::default()
},
AgentEndpoint {
name: "b".into(),
management: true,
..Default::default()
},
],
}
.validate()
Expand All @@ -409,9 +436,12 @@ token = "s2"
.unwrap_err()
.contains("name"));
// Named but unconfigured conn → the RemoteConfig validation fires.
assert!(AgentEndpoint { name: "x".into(), ..Default::default() }
.validate()
.is_err());
assert!(AgentEndpoint {
name: "x".into(),
..Default::default()
}
.validate()
.is_err());
}

#[test]
Expand Down
11 changes: 2 additions & 9 deletions crates/agent-lifecycle/src/k8s.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,8 +67,7 @@ impl RuntimeDriver for K8sDriver {

// `identity_verified` latches once the pod has ever reported phase
// Running with Ready true.
let identity_verified =
verified_before || (pod.phase == PodPhase::Running && pod.ready);
let identity_verified = verified_before || (pod.phase == PodPhase::Running && pod.ready);

// Node-unreachable (Unknown phase), a crash-loop, or a lost lease all
// fault outright. A *declared* readiness probe failing faults too —
Expand Down Expand Up @@ -100,13 +99,7 @@ mod tests {
use super::*;
use crate::AgentState;

fn pod(
phase: PodPhase,
deleting: bool,
ready: bool,
lease: bool,
accepting: bool,
) -> K8sPod {
fn pod(phase: PodPhase, deleting: bool, ready: bool, lease: bool, accepting: bool) -> K8sPod {
// Default: no readiness probe declared (the common case for OAB agents).
pod_probe(phase, deleting, ready, lease, accepting, false, false)
}
Expand Down
20 changes: 20 additions & 0 deletions crates/agent-lifecycle/tests/rustfmt_clean.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
//! Repo-hygiene regression: the whole workspace stays `cargo fmt` clean.
//! `cargo fmt --all -- --check` is part of the repo verification suite; this
//! test fails fast when formatting drifts in any crate instead of letting it
//! accumulate silently.

use std::process::Command;

#[test]
fn workspace_is_rustfmt_clean() {
let output = Command::new(env!("CARGO"))
.args(["fmt", "--all", "--", "--check"])
.output()
.expect("spawn cargo fmt");
assert!(
output.status.success(),
"cargo fmt --all -- --check failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
44 changes: 34 additions & 10 deletions crates/oab-mcp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -467,7 +467,10 @@ impl OabMcp {
let guard = self.bindings.read().unwrap();
let binding = guard.get(name).cloned().ok_or_else(|| {
let known: Vec<&str> = guard.fleets.iter().map(|f| f.name.as_str()).collect();
anyhow::anyhow!("unknown fleet {name:?}; configured fleets: [{}]", known.join(", "))
anyhow::anyhow!(
"unknown fleet {name:?}; configured fleets: [{}]",
known.join(", ")
)
})?;
let cluster = binding.cluster.clone().ok_or_else(|| {
anyhow::anyhow!(
Expand Down Expand Up @@ -501,7 +504,10 @@ impl OabMcp {
let guard = self.bindings.read().unwrap();
let binding = guard.get(name).cloned().ok_or_else(|| {
let known: Vec<&str> = guard.fleets.iter().map(|f| f.name.as_str()).collect();
anyhow::anyhow!("unknown fleet {name:?}; configured fleets: [{}]", known.join(", "))
anyhow::anyhow!(
"unknown fleet {name:?}; configured fleets: [{}]",
known.join(", ")
)
})?;
Ok(Some(binding))
}
Expand Down Expand Up @@ -542,7 +548,9 @@ impl OabMcp {
.filter(|s| b.includes(&s.service_name, &s.name))
.map(|s| json!({ "name": s.name, "namespace": s.namespace, "desired": s.desired }))
.collect();
return Ok(json!({ "context": b.context, "namespace": namespace, "deployments": deployments }));
return Ok(
json!({ "context": b.context, "namespace": namespace, "deployments": deployments }),
);
}
}
let t = self.target(args)?;
Expand Down Expand Up @@ -575,7 +583,9 @@ impl OabMcp {
if let Some(b) = self.named_fleet(args)? {
if b.runtime == scp::FleetRuntime::K8s {
let namespace = b.namespace.clone().unwrap_or_else(|| "default".to_string());
return match scp::observe_k8s_deployment(b.context.as_deref(), &namespace, service).await? {
return match scp::observe_k8s_deployment(b.context.as_deref(), &namespace, service)
.await?
{
Some(d) => Ok(deployment_json(&d)),
None => Ok(json!({ "found": false, "service": service })),
};
Expand Down Expand Up @@ -817,18 +827,33 @@ impl OabMcp {
.and_then(Value::as_str)
.ok_or_else(|| anyhow::anyhow!("missing required arg: image"))?;
let input = scp::AgentWizardInput {
api_key: args.get("api_key").and_then(Value::as_str).map(str::to_string),
chat_platform: args.get("chat_platform").and_then(Value::as_str).map(str::to_string),
chat_bot_token: args.get("chat_bot_token").and_then(Value::as_str).map(str::to_string),
api_key: args
.get("api_key")
.and_then(Value::as_str)
.map(str::to_string),
chat_platform: args
.get("chat_platform")
.and_then(Value::as_str)
.map(str::to_string),
chat_bot_token: args
.get("chat_bot_token")
.and_then(Value::as_str)
.map(str::to_string),
chat_channel_secret: args
.get("chat_channel_secret")
.and_then(Value::as_str)
.map(str::to_string),
// studio#119: default on when the caller doesn't say — matches
// build_default_manifest's original hardcoded behavior before
// studio#128 made it caller-controlled.
acp_enabled: args.get("acp_enabled").and_then(Value::as_bool).unwrap_or(true),
acp_token: args.get("acp_token").and_then(Value::as_str).map(str::to_string),
acp_enabled: args
.get("acp_enabled")
.and_then(Value::as_bool)
.unwrap_or(true),
acp_token: args
.get("acp_token")
.and_then(Value::as_str)
.map(str::to_string),
};
// studio#135: Brett's explicit ordering — write local first, S3
// (via provision_agent[_k8s]'s existing upload) after. Optional:
Expand Down Expand Up @@ -1160,7 +1185,6 @@ impl OabMcp {
Ok(json!({ "service_accounts": service_accounts }))
}


async fn t_delete(&self, args: &Map<String, Value>) -> Result<Value> {
let t = self.target(args)?;
let cluster = t.cluster.clone();
Expand Down
Loading