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
67 changes: 67 additions & 0 deletions crates/agent-lifecycle/tests/adr_review_followups.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
//! Guard for the ADR-1 review follow-ups (openabdev/studio#3).
//!
//! The wording tightened in `docs/adr/agent-lifecycle.md` during the
//! wontfix-or-fold review pass must stay in place: the distinct hard-loss
//! `Stopping → Stopped` edge, the §9 lock-in/reversibility note, the
//! `State.Paused` wire-naming deferral to the RuntimeDriver-contract ADR, and
//! the instance-level `superseded` phrasing. The doc is read relative to this
//! crate's manifest so the test runs in any checkout layout.

use std::path::Path;

fn adr() -> String {
let path = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../docs/adr/agent-lifecycle.md");
std::fs::read_to_string(&path).unwrap_or_else(|e| panic!("read {}: {e}", path.display()))
}

/// `Stopping → Stopped` edges in the state diagram, as trimmed source lines.
fn stopping_to_stopped_edges(doc: &str) -> Vec<&str> {
doc.lines()
.map(str::trim)
.filter(|l| l.starts_with("Stopping") && l.contains("--> Stopped"))
.collect()
}

#[test]
fn stopping_has_graceful_and_hard_loss_edges_to_stopped() {
let doc = adr();
let edges = stopping_to_stopped_edges(&doc);
assert!(
edges.iter().any(|l| l.contains("state saved")),
"graceful Stopping → Stopped edge (state saved) missing; edges: {edges:?}"
);
assert!(
edges.iter().any(|l| l.contains("hard loss")),
"hard-loss Stopping → Stopped edge missing — a mid-drain kill must be \
drawn separately from the graceful 'state saved' edge; edges: {edges:?}"
);
}

#[test]
fn consequences_record_lock_in_and_reversibility() {
let doc = adr();
assert!(
doc.contains("Lock-in / reversibility"),
"§9 must name which parts of the 6-state surface are expensive to reverse"
);
}

#[test]
fn paused_wire_naming_is_deferred_to_the_driver_contract_adr() {
let doc = adr();
assert!(
doc.contains("semantics, not the identifier"),
"the `Paused` enum/wire naming decision must be recorded as deferred to \
the RuntimeDriver-contract ADR"
);
}

#[test]
fn superseded_is_instance_level_and_never_pins_the_instance() {
let doc = adr();
assert!(
doc.contains("never pins an instance"),
"the superseded note must say the instance keeps its normal edges \
(→ Unhealthy, → Stopping) rather than being pinned in Paused"
);
}
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