Skip to content
Open
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
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
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