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
46 changes: 35 additions & 11 deletions crates/oab-mcp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -195,7 +195,7 @@ pub fn tools() -> Vec<Tool> {
"chat_channel_secret": { "type": "string", "description": "LINE only." },
"acp_enabled": { "type": "boolean", "description": "Enable the reverse-MCP-over-ACP tunnel on this agent. Defaults to true when omitted (studio#119: Studio-deployed agents default to ACP on). Not honorable for every vendor — the caller is responsible for not setting this true for a vendor that can't support it (e.g. agy, whose bridge bypasses /acp entirely)." },
"acp_token": { "type": "string", "description": "Optional (studio#136), only meaningful when acp_enabled is true. The OPENAB_ACP_AUTH_KEY to use, instead of generating a random one." },
"local_config_folder": { "type": "string", "description": "Optional local directory (studio#135) — when set, config.toml is written to <local_config_folder>/<name>/config.toml *before* anything touches S3. Omit to skip the local mirror entirely." },
"local_config_folder": { "type": "string", "description": "Optional local directory (studio#135) — when set, config.toml is written to <local_config_folder>/<name>/config.toml *before* anything touches S3. `name` must be a single path component; path-like names (`..`, separators, absolute) are refused rather than allowed to escape the folder. Omit to skip the local mirror entirely." },
"provider": { "type": "string", "description": "\"aws\" (default) or \"k8s\" — which driver applies the result." },
"fleet": { "type": "string", "description": "AWS only. Fleet name (see fleet_config): targets the fleet's cluster and managing credential; a write to a service outside the fleet's members is refused. Overrides the cluster arg." },
"cluster": { "type": "string", "description": "AWS only. ECS cluster (defaults to the server's configured cluster)." },
Expand Down 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