From 091bfeb8109e48763a3121545a8e1837df5a6387 Mon Sep 17 00:00:00 2001 From: Devin Date: Wed, 30 Sep 2026 13:03:06 +0800 Subject: [PATCH 1/7] test(agent-lifecycle): guard repo-wide rustfmt cleanliness cargo fmt --all -- --check is part of the workspace verification suite, but the repo had drifted out of rustfmt-clean (the authoring runtime runs no rustfmt and CI only checks -p agent-lifecycle). Add a regression test that runs the same check so drift fails fast in cargo test instead of accumulating silently across crates. --- crates/agent-lifecycle/tests/rustfmt_clean.rs | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) create mode 100644 crates/agent-lifecycle/tests/rustfmt_clean.rs diff --git a/crates/agent-lifecycle/tests/rustfmt_clean.rs b/crates/agent-lifecycle/tests/rustfmt_clean.rs new file mode 100644 index 0000000..cb298d1 --- /dev/null +++ b/crates/agent-lifecycle/tests/rustfmt_clean.rs @@ -0,0 +1,19 @@ +//! 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:\n{}", + String::from_utf8_lossy(&output.stdout) + ); +} From 3678e6ed5c2f98c3facf01aae41e287bfcf2ccb9 Mon Sep 17 00:00:00 2001 From: Devin Date: Wed, 30 Sep 2026 13:03:09 +0800 Subject: [PATCH 2/7] style: apply cargo fmt --all workspace-wide MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Mechanical rustfmt pass — the repo had drifted from rustfmt-clean across acp-tunnel, agent-lifecycle, oab-mcp and oabctl. No semantic changes. --- crates/acp-tunnel/src/config.rs | 62 ++++-- crates/agent-lifecycle/src/k8s.rs | 11 +- crates/oab-mcp/src/lib.rs | 44 ++++- crates/oabctl/src/bootstrap.rs | 295 +++++++++++++++++++++++------ crates/oabctl/src/config.rs | 12 +- crates/oabctl/src/create.rs | 269 +++++++++++++++++++------- crates/oabctl/src/delete.rs | 11 +- crates/oabctl/src/driver.rs | 35 +++- crates/oabctl/src/events.rs | 9 +- crates/oabctl/src/k8s_driver.rs | 92 +++++++-- crates/oabctl/src/manifest.rs | 135 ++++++++----- crates/oabctl/src/scale.rs | 4 +- crates/oabctl/src/secrets.rs | 53 ++++-- crates/oabctl/src/status.rs | 7 +- crates/oabctl/src/studio_api.rs | 67 +++++-- crates/oabctl/src/vendor_images.rs | 31 ++- 16 files changed, 852 insertions(+), 285 deletions(-) diff --git a/crates/acp-tunnel/src/config.rs b/crates/acp-tunnel/src/config.rs index 22c5965..67951e6 100644 --- a/crates/acp-tunnel/src/config.rs +++ b/crates/acp-tunnel/src/config.rs @@ -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; @@ -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(), @@ -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")); } @@ -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() @@ -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] diff --git a/crates/agent-lifecycle/src/k8s.rs b/crates/agent-lifecycle/src/k8s.rs index 4514b0a..c85b0a7 100644 --- a/crates/agent-lifecycle/src/k8s.rs +++ b/crates/agent-lifecycle/src/k8s.rs @@ -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 — @@ -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) } diff --git a/crates/oab-mcp/src/lib.rs b/crates/oab-mcp/src/lib.rs index 6469f57..d0ac5e9 100644 --- a/crates/oab-mcp/src/lib.rs +++ b/crates/oab-mcp/src/lib.rs @@ -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!( @@ -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)) } @@ -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)?; @@ -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 })), }; @@ -817,9 +827,18 @@ 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) @@ -827,8 +846,14 @@ impl OabMcp { // 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: @@ -1160,7 +1185,6 @@ impl OabMcp { Ok(json!({ "service_accounts": service_accounts })) } - async fn t_delete(&self, args: &Map) -> Result { let t = self.target(args)?; let cluster = t.cluster.clone(); diff --git a/crates/oabctl/src/bootstrap.rs b/crates/oabctl/src/bootstrap.rs index ae98c10..1982328 100644 --- a/crates/oabctl/src/bootstrap.rs +++ b/crates/oabctl/src/bootstrap.rs @@ -70,7 +70,12 @@ pub struct ImportOptions { pub task_role: Option, } -pub async fn run(config: &aws_config::SdkConfig, delete: bool, status: bool, imports: ImportOptions) -> Result<()> { +pub async fn run( + config: &aws_config::SdkConfig, + delete: bool, + status: bool, + imports: ImportOptions, +) -> Result<()> { if status { return show_status(config).await; } @@ -84,7 +89,10 @@ async fn get_account_and_region(config: &aws_config::SdkConfig) -> Result<(Strin let sts = StsClient::new(config); let identity = sts.get_caller_identity().send().await?; let account = identity.account().context("no account ID")?.to_string(); - let region = config.region().map(|r| r.to_string()).unwrap_or_else(|| "us-east-1".to_string()); + let region = config + .region() + .map(|r| r.to_string()) + .unwrap_or_else(|| "us-east-1".to_string()); Ok((account, region)) } @@ -139,44 +147,107 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul let bucket_exists = s3.head_bucket().bucket(&bucket).send().await.is_ok(); let cluster_name = imports.cluster.as_deref().unwrap_or(CLUSTER_NAME); - let cluster_exists = ecs.describe_clusters().clusters(cluster_name).send().await - .map(|r| r.clusters().first().is_some_and(|c| c.status() == Some("ACTIVE"))) + let cluster_exists = ecs + .describe_clusters() + .clusters(cluster_name) + .send() + .await + .map(|r| { + r.clusters() + .first() + .is_some_and(|c| c.status() == Some("ACTIVE")) + }) .unwrap_or(false); let exec_role_exists = imports.execution_role.is_some() - || iam.get_role().role_name(EXECUTION_ROLE).send().await.is_ok(); - let task_role_exists = imports.task_role.is_some() - || iam.get_role().role_name(TASK_ROLE).send().await.is_ok(); + || iam + .get_role() + .role_name(EXECUTION_ROLE) + .send() + .await + .is_ok(); + let task_role_exists = + imports.task_role.is_some() || iam.get_role().role_name(TASK_ROLE).send().await.is_ok(); let vpc_id_for_check = if let Some(ref v) = imports.vpc { v.clone() } else { ec2.describe_vpcs() - .filters(aws_sdk_ec2::types::Filter::builder().name("isDefault").values("true").build()) - .send().await.ok() - .and_then(|r| r.vpcs().first().and_then(|v| v.vpc_id()).map(|s| s.to_string())) + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("isDefault") + .values("true") + .build(), + ) + .send() + .await + .ok() + .and_then(|r| { + r.vpcs() + .first() + .and_then(|v| v.vpc_id()) + .map(|s| s.to_string()) + }) .unwrap_or_default() }; let sg_exists = imports.security_group.is_some() - || ec2.describe_security_groups() - .filters(aws_sdk_ec2::types::Filter::builder().name("group-name").values(SG_NAME).build()) - .filters(aws_sdk_ec2::types::Filter::builder().name("vpc-id").values(&vpc_id_for_check).build()) - .send().await + || ec2 + .describe_security_groups() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("group-name") + .values(SG_NAME) + .build(), + ) + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("vpc-id") + .values(&vpc_id_for_check) + .build(), + ) + .send() + .await .map(|r| !r.security_groups().is_empty()) .unwrap_or(false); - let log_group_exists = logs.describe_log_groups() + let log_group_exists = logs + .describe_log_groups() .log_group_name_prefix(LOG_GROUP) - .send().await - .map(|r| r.log_groups().iter().any(|g| g.log_group_name() == Some(LOG_GROUP))) + .send() + .await + .map(|r| { + r.log_groups() + .iter() + .any(|g| g.log_group_name() == Some(LOG_GROUP)) + }) .unwrap_or(false); // ─── DISPLAY PLAN ───────────────────────────────────────────────────── eprintln!(" Resource Action"); eprintln!(" ─────────────────────────────────────────"); plan_line("S3 Bucket", &bucket, bucket_exists, true); - plan_line("ECS Cluster", cluster_name, cluster_exists, imports.cluster.is_none()); - plan_line("IAM Execution Role", imports.execution_role.as_deref().unwrap_or(EXECUTION_ROLE), exec_role_exists, imports.execution_role.is_none()); - plan_line("IAM Task Role", imports.task_role.as_deref().unwrap_or(TASK_ROLE), task_role_exists, imports.task_role.is_none()); - plan_line("Security Group", imports.security_group.as_deref().unwrap_or(SG_NAME), sg_exists, imports.security_group.is_none()); + plan_line( + "ECS Cluster", + cluster_name, + cluster_exists, + imports.cluster.is_none(), + ); + plan_line( + "IAM Execution Role", + imports.execution_role.as_deref().unwrap_or(EXECUTION_ROLE), + exec_role_exists, + imports.execution_role.is_none(), + ); + plan_line( + "IAM Task Role", + imports.task_role.as_deref().unwrap_or(TASK_ROLE), + task_role_exists, + imports.task_role.is_none(), + ); + plan_line( + "Security Group", + imports.security_group.as_deref().unwrap_or(SG_NAME), + sg_exists, + imports.security_group.is_none(), + ); plan_line("CloudWatch Log Group", LOG_GROUP, log_group_exists, true); eprintln!(); @@ -216,7 +287,9 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul .restrict_public_buckets(true) .build(), ) - .send().await.ok(); + .send() + .await + .ok(); eprintln!(" ✓ Created S3 bucket: {bucket} (public access blocked)"); managed.bucket = true; } @@ -224,7 +297,9 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul // 2. ECS Cluster — save state incrementally after this point let (cluster_arn, cluster_managed) = if let Some(ref name) = imports.cluster { let resp = ecs.describe_clusters().clusters(name).send().await?; - let arn = resp.clusters().first() + let arn = resp + .clusters() + .first() .and_then(|c| c.cluster_arn()) .context(format!("cluster '{}' not found", name))? .to_string(); @@ -232,13 +307,22 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul (arn, false) } else { match ecs.describe_clusters().clusters(CLUSTER_NAME).send().await { - Ok(resp) if resp.clusters().first().is_some_and(|c| c.status() == Some("ACTIVE")) => { - let arn = resp.clusters()[0].cluster_arn().unwrap_or_default().to_string(); + Ok(resp) + if resp + .clusters() + .first() + .is_some_and(|c| c.status() == Some("ACTIVE")) => + { + let arn = resp.clusters()[0] + .cluster_arn() + .unwrap_or_default() + .to_string(); eprintln!(" ✓ ECS cluster already exists: {CLUSTER_NAME}"); (arn, true) } _ => { - let resp = ecs.create_cluster() + let resp = ecs + .create_cluster() .cluster_name(CLUSTER_NAME) .capacity_providers("FARGATE") .capacity_providers("FARGATE_SPOT") @@ -251,7 +335,11 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul .send() .await .context("failed to create ECS cluster")?; - let arn = resp.cluster().and_then(|c| c.cluster_arn()).unwrap_or_default().to_string(); + let arn = resp + .cluster() + .and_then(|c| c.cluster_arn()) + .unwrap_or_default() + .to_string(); eprintln!(" ✓ Created ECS cluster: {CLUSTER_NAME}"); (arn, true) } @@ -268,7 +356,9 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul iam.attach_role_policy() .role_name(EXECUTION_ROLE) .policy_arn("arn:aws:iam::aws:policy/service-role/AmazonECSTaskExecutionRolePolicy") - .send().await.ok(); + .send() + .await + .ok(); eprintln!(" ✓ IAM execution role: {EXECUTION_ROLE}"); managed.execution_role = true; arn @@ -330,7 +420,9 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul .role_name(TASK_ROLE) .policy_name("oab-s3-artifacts") .policy_document(&artifacts_policy) - .send().await.ok(); + .send() + .await + .ok(); // Secrets Manager access (agent reads its own secrets at runtime) iam.put_role_policy() .role_name(TASK_ROLE) @@ -347,32 +439,58 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul eprintln!(" ✓ Using existing security group: {sg}"); (vpc, sg.clone()) } else { - let default_vpc = ec2.describe_vpcs() - .filters(aws_sdk_ec2::types::Filter::builder().name("isDefault").values("true").build()) - .send().await?; + let default_vpc = ec2 + .describe_vpcs() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("isDefault") + .values("true") + .build(), + ) + .send() + .await?; let vid = imports.vpc.clone().unwrap_or_else(|| { - default_vpc.vpcs().first() + default_vpc + .vpcs() + .first() .and_then(|v| v.vpc_id()) .unwrap_or_default() .to_string() }); - let sid = match ec2.describe_security_groups() - .filters(aws_sdk_ec2::types::Filter::builder().name("group-name").values(SG_NAME).build()) - .filters(aws_sdk_ec2::types::Filter::builder().name("vpc-id").values(&vid).build()) - .send().await + let sid = match ec2 + .describe_security_groups() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("group-name") + .values(SG_NAME) + .build(), + ) + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("vpc-id") + .values(&vid) + .build(), + ) + .send() + .await { Ok(resp) if !resp.security_groups().is_empty() => { - let id = resp.security_groups()[0].group_id().unwrap_or_default().to_string(); + let id = resp.security_groups()[0] + .group_id() + .unwrap_or_default() + .to_string(); eprintln!(" ✓ Security group already exists: {id}"); id } _ => { - let resp = ec2.create_security_group() + let resp = ec2 + .create_security_group() .group_name(SG_NAME) .description("OAB agent containers — managed by oabctl bootstrap") .vpc_id(&vid) - .send().await + .send() + .await .context("failed to create security group")?; let id = resp.group_id().unwrap_or_default().to_string(); managed.security_group = true; @@ -388,17 +506,34 @@ async fn create(config: &aws_config::SdkConfig, imports: ImportOptions) -> Resul eprintln!(" ✓ Using provided subnets: {}", s.join(", ")); s.clone() } else { - let subnets_resp = ec2.describe_subnets() - .filters(aws_sdk_ec2::types::Filter::builder().name("vpc-id").values(&vpc_id).build()) - .send().await?; - subnets_resp.subnets().iter() + let subnets_resp = ec2 + .describe_subnets() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("vpc-id") + .values(&vpc_id) + .build(), + ) + .send() + .await?; + subnets_resp + .subnets() + .iter() .filter_map(|s| s.subnet_id().map(|id| id.to_string())) .collect() }; // 7. CloudWatch Log Group - match logs.create_log_group().log_group_name(LOG_GROUP).send().await { - Ok(_) => { managed.log_group = true; eprintln!(" ✓ Created log group: {LOG_GROUP}"); } + match logs + .create_log_group() + .log_group_name(LOG_GROUP) + .send() + .await + { + Ok(_) => { + managed.log_group = true; + eprintln!(" ✓ Created log group: {LOG_GROUP}"); + } Err(_) => eprintln!(" ✓ Log group already exists: {LOG_GROUP}"), } @@ -431,7 +566,8 @@ async fn teardown(config: &aws_config::SdkConfig) -> Result<()> { let bucket = bucket_name(&account); let s3 = S3Client::new(config); - let state = load_state(&s3, &bucket).await? + let state = load_state(&s3, &bucket) + .await? .context("no bootstrap state found — nothing to delete")?; eprintln!("🗑️ Tearing down OAB bootstrap resources...\n"); @@ -456,7 +592,12 @@ async fn teardown(config: &aws_config::SdkConfig) -> Result<()> { // Reverse order — only delete resources we created (managed) // 1. Log group if state.managed.log_group { - match logs.delete_log_group().log_group_name(&state.resources.log_group).send().await { + match logs + .delete_log_group() + .log_group_name(&state.resources.log_group) + .send() + .await + { Ok(_) => eprintln!(" ✓ Deleted log group: {}", state.resources.log_group), Err(e) => eprintln!(" ⚠ Failed to delete log group: {e}"), } @@ -466,8 +607,16 @@ async fn teardown(config: &aws_config::SdkConfig) -> Result<()> { // 2. Security group if state.managed.security_group { - match ec2.delete_security_group().group_id(&state.resources.security_group_id).send().await { - Ok(_) => eprintln!(" ✓ Deleted security group: {}", state.resources.security_group_id), + match ec2 + .delete_security_group() + .group_id(&state.resources.security_group_id) + .send() + .await + { + Ok(_) => eprintln!( + " ✓ Deleted security group: {}", + state.resources.security_group_id + ), Err(e) => eprintln!(" ⚠ Failed to delete security group: {e}"), } } else { @@ -499,7 +648,12 @@ async fn teardown(config: &aws_config::SdkConfig) -> Result<()> { } // 5. Delete state file (keep bucket for user data) - s3.delete_object().bucket(&bucket).key(STATE_KEY).send().await.ok(); + s3.delete_object() + .bucket(&bucket) + .key(STATE_KEY) + .send() + .await + .ok(); eprintln!(" ✓ Deleted bootstrap state"); eprintln!("\n ℹ️ S3 bucket '{bucket}' preserved (may contain manifests/config)."); eprintln!(" To fully remove: aws s3 rb s3://{bucket} --force"); @@ -542,31 +696,56 @@ async fn show_status(config: &aws_config::SdkConfig) -> Result<()> { async fn ensure_role(iam: &IamClient, name: &str, _account: &str) -> Result { match iam.get_role().role_name(name).send().await { - Ok(resp) => Ok(resp.role().context("no role in response")?.arn().to_string()), + Ok(resp) => Ok(resp + .role() + .context("no role in response")? + .arn() + .to_string()), Err(_) => { - let resp = iam.create_role() + let resp = iam + .create_role() .role_name(name) .assume_role_policy_document(ASSUME_ROLE_POLICY) - .send().await + .send() + .await .with_context(|| format!("failed to create role {name}"))?; - Ok(resp.role().context("no role in response")?.arn().to_string()) + Ok(resp + .role() + .context("no role in response")? + .arn() + .to_string()) } } } async fn delete_role(iam: &IamClient, name: &str) { // Detach managed policies - if let Ok(resp) = iam.list_attached_role_policies().role_name(name).send().await { + if let Ok(resp) = iam + .list_attached_role_policies() + .role_name(name) + .send() + .await + { for p in resp.attached_policies() { if let Some(arn) = p.policy_arn() { - iam.detach_role_policy().role_name(name).policy_arn(arn).send().await.ok(); + iam.detach_role_policy() + .role_name(name) + .policy_arn(arn) + .send() + .await + .ok(); } } } // Delete inline policies if let Ok(resp) = iam.list_role_policies().role_name(name).send().await { for p in resp.policy_names() { - iam.delete_role_policy().role_name(name).policy_name(p).send().await.ok(); + iam.delete_role_policy() + .role_name(name) + .policy_name(p) + .send() + .await + .ok(); } } iam.delete_role().role_name(name).send().await.ok(); diff --git a/crates/oabctl/src/config.rs b/crates/oabctl/src/config.rs index cdd5f1b..88991e3 100644 --- a/crates/oabctl/src/config.rs +++ b/crates/oabctl/src/config.rs @@ -39,8 +39,12 @@ pub struct BootstrapConfig { pub bucket: Option, } -fn default_namespace() -> String { "prod".to_string() } -fn default_cluster() -> String { "oab".to_string() } +fn default_namespace() -> String { + "prod".to_string() +} +fn default_cluster() -> String { + "oab".to_string() +} impl OabConfig { pub fn load() -> Result { @@ -65,7 +69,9 @@ impl OabConfig { /// Get the control plane bucket name (config > env var > account-based default) pub fn bucket(&self) -> Option { - self.bootstrap.bucket.clone() + self.bootstrap + .bucket + .clone() .or_else(|| std::env::var("OAB_CONTROL_PLANE_BUCKET").ok()) } } diff --git a/crates/oabctl/src/create.rs b/crates/oabctl/src/create.rs index 0eb7122..a0e690e 100644 --- a/crates/oabctl/src/create.rs +++ b/crates/oabctl/src/create.rs @@ -15,11 +15,19 @@ const BACKENDS: &[(&str, &str)] = &[ const CHANNELS: &[&str] = &["stable", "beta"]; -pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, auto_apply: bool) -> Result<()> { +pub async fn run( + config: &aws_config::SdkConfig, + name: &str, + namespace: &str, + auto_apply: bool, +) -> Result<()> { eprintln!("🤖 Creating agent: {name}\n"); // 1. Backend - let backend = prompt_select("Backend platform", &BACKENDS.iter().map(|(n, _)| *n).collect::>())?; + let backend = prompt_select( + "Backend platform", + &BACKENDS.iter().map(|(n, _)| *n).collect::>(), + )?; let image_base = BACKENDS.iter().find(|(n, _)| *n == backend).unwrap().1; // 2. Release channel @@ -35,8 +43,8 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au let secret_name = format!("oab/{namespace}/{name}"); // 3b. STT API key (optional) - let stt_key = rpassword::prompt_password(" STT API key (Groq, enter to skip): ") - .unwrap_or_default(); + let stt_key = + rpassword::prompt_password(" STT API key (Groq, enter to skip): ").unwrap_or_default(); let stt_enabled = !stt_key.is_empty(); let mut secret_obj = serde_json::json!({ "DISCORD_BOT_TOKEN": token }); @@ -58,8 +66,15 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au } // 5. Capacity provider - let cap = prompt_select("Capacity provider", &["FARGATE_SPOT (cost-optimized)", "FARGATE (on-demand)"])?; - let capacity_provider = if cap.starts_with("FARGATE_SPOT") { "FARGATE_SPOT" } else { "FARGATE" }; + let cap = prompt_select( + "Capacity provider", + &["FARGATE_SPOT (cost-optimized)", "FARGATE (on-demand)"], + )?; + let capacity_provider = if cap.starts_with("FARGATE_SPOT") { + "FARGATE_SPOT" + } else { + "FARGATE" + }; // 6. VPC let ec2 = Ec2Client::new(config); @@ -75,7 +90,13 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au let subnets = select_subnets(&ec2, &vpc.id).await?; eprintln!(" Subnets (auto-selected):"); for s in &subnets { - eprintln!(" ✓ {} ({}, {}, {})", s.id, s.az, s.kind, if s.has_nat { "NAT ✓" } else { "no NAT" }); + eprintln!( + " ✓ {} ({}, {}, {})", + s.id, + s.az, + s.kind, + if s.has_nat { "NAT ✓" } else { "no NAT" } + ); } eprintln!(); @@ -88,17 +109,23 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au let sg_id = if sg_choice.starts_with("Create new") { let sg_name = format!("oab-{name}"); - let resp = ec2.create_security_group() + let resp = ec2 + .create_security_group() .group_name(&sg_name) .description(format!("OAB agent {name}")) .vpc_id(&vpc.id) - .send().await + .send() + .await .context("failed to create security group")?; let id = resp.group_id().unwrap_or_default().to_string(); eprintln!(" → Created security group: {id}\n"); id } else { - sgs.iter().find(|s| sg_choice.contains(&s.id)).unwrap().id.clone() + sgs.iter() + .find(|s| sg_choice.contains(&s.id)) + .unwrap() + .id + .clone() }; // ─── Generate config.toml ────────────────────────────────────────────── @@ -106,7 +133,8 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au // ─── Resolve bucket for configFrom path ──────────────────────────────── let s3 = S3Client::new(config); - let bucket = resolve_bucket(&s3, config).await + let bucket = resolve_bucket(&s3, config) + .await .unwrap_or_else(|| "oab-control-plane-unknown".to_string()); let config_s3_key = format!("artifacts/{namespace}/{name}/config.toml"); @@ -118,7 +146,15 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au std::fs::write(format!("{dir}/config.toml"), &config_toml)?; let subnet_ids: Vec = subnets.iter().map(|s| s.id.clone()).collect(); - let manifest_yaml = generate_manifest(name, namespace, &image, &config_from, capacity_provider, &subnet_ids, &sg_id); + let manifest_yaml = generate_manifest( + name, + namespace, + &image, + &config_from, + capacity_provider, + &subnet_ids, + &sg_id, + ); std::fs::write(format!("{dir}/manifest.yaml"), &manifest_yaml)?; // ─── Summary ─────────────────────────────────────────────────────────── @@ -138,7 +174,10 @@ pub async fn run(config: &aws_config::SdkConfig, name: &str, namespace: &str, au io::stdout().flush()?; let mut input = String::new(); io::stdin().read_line(&mut input)?; - if !input.trim().is_empty() && !input.trim().eq_ignore_ascii_case("y") && !input.trim().eq_ignore_ascii_case("yes") { + if !input.trim().is_empty() + && !input.trim().eq_ignore_ascii_case("y") + && !input.trim().eq_ignore_ascii_case("yes") + { eprintln!("Aborted."); return Ok(()); } @@ -194,11 +233,21 @@ fn prompt_secret(label: &str) -> Result { /// `OPENAB_ACP_AUTH_KEY`, same create-or-update pattern this CLI wizard /// already uses for the Discord bot token. pub async fn store_secret(sm: &SmClient, name: &str, value: &str) -> Result<()> { - match sm.create_secret().name(name).secret_string(value).send().await { + match sm + .create_secret() + .name(name) + .secret_string(value) + .send() + .await + { Ok(_) => Ok(()), Err(_) => { // Already exists — update - sm.put_secret_value().secret_id(name).secret_string(value).send().await + sm.put_secret_value() + .secret_id(name) + .secret_string(value) + .send() + .await .context("failed to store secret")?; Ok(()) } @@ -210,21 +259,38 @@ pub async fn store_secret(sm: &SmClient, name: &str, value: &str) -> Result<()> /// `pub` (studio#111): reused by the provision-from-scratch path to pick a /// sensible default VPC without prompting, same discovery `oabctl create`'s /// wizard uses interactively. -pub struct VpcInfo { pub id: String, pub label: String, pub is_default: bool } +pub struct VpcInfo { + pub id: String, + pub label: String, + pub is_default: bool, +} pub async fn list_vpcs(ec2: &Ec2Client) -> Result> { let resp = ec2.describe_vpcs().send().await?; - Ok(resp.vpcs().iter().map(|v| { - let id = v.vpc_id().unwrap_or_default().to_string(); - let cidr = v.cidr_block().unwrap_or_default(); - let is_default = v.is_default().unwrap_or(false); - let name = v.tags().iter() - .find(|t| t.key() == Some("Name")) - .and_then(|t| t.value()) - .unwrap_or("unnamed"); - let label = format!("{id} ({name}, {cidr}{})", if is_default { ", default" } else { "" }); - VpcInfo { id, label, is_default } - }).collect()) + Ok(resp + .vpcs() + .iter() + .map(|v| { + let id = v.vpc_id().unwrap_or_default().to_string(); + let cidr = v.cidr_block().unwrap_or_default(); + let is_default = v.is_default().unwrap_or(false); + let name = v + .tags() + .iter() + .find(|t| t.key() == Some("Name")) + .and_then(|t| t.value()) + .unwrap_or("unnamed"); + let label = format!( + "{id} ({name}, {cidr}{})", + if is_default { ", default" } else { "" } + ); + VpcInfo { + id, + label, + is_default, + } + }) + .collect()) } /// Pick a VPC with **no interactive prompt**, for the provision-from-scratch @@ -241,37 +307,59 @@ pub async fn default_vpc(ec2: &Ec2Client) -> Result { match defaults.len() { 1 => Ok(defaults.into_iter().next().unwrap()), 0 => anyhow::bail!("no default VPC in this account/region — configure a VPC explicitly"), - n => anyhow::bail!("{n} VPCs marked default — configure a VPC explicitly, can't pick automatically"), + n => anyhow::bail!( + "{n} VPCs marked default — configure a VPC explicitly, can't pick automatically" + ), } } /// A subnet candidate, already classified private/public + NAT reachability. /// `pub` (studio#111): see `VpcInfo`. -pub struct SubnetInfo { pub id: String, pub az: String, pub kind: String, pub has_nat: bool } +pub struct SubnetInfo { + pub id: String, + pub az: String, + pub kind: String, + pub has_nat: bool, +} /// Auto-select up to 3 subnets (one per AZ), preferring private+NAT > /// private > public — this already has zero interactive prompting, it's the /// exact "sensible default" the provision-from-scratch path (studio#111) /// needs, just needed to be reachable outside `create.rs`. pub async fn select_subnets(ec2: &Ec2Client, vpc_id: &str) -> Result> { - let subnets_resp = ec2.describe_subnets() - .filters(aws_sdk_ec2::types::Filter::builder().name("vpc-id").values(vpc_id).build()) - .send().await?; + let subnets_resp = ec2 + .describe_subnets() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("vpc-id") + .values(vpc_id) + .build(), + ) + .send() + .await?; // Get route tables to determine private vs public + NAT - let rt_resp = ec2.describe_route_tables() - .filters(aws_sdk_ec2::types::Filter::builder().name("vpc-id").values(vpc_id).build()) - .send().await?; + let rt_resp = ec2 + .describe_route_tables() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("vpc-id") + .values(vpc_id) + .build(), + ) + .send() + .await?; // Build subnet → route table mapping - let mut subnet_routes: std::collections::HashMap = std::collections::HashMap::new(); + let mut subnet_routes: std::collections::HashMap = + std::collections::HashMap::new(); for rt in rt_resp.route_tables() { let has_igw = rt.routes().iter().any(|r| { - r.gateway_id().map(|g| g.starts_with("igw-")).unwrap_or(false) - }); - let has_nat = rt.routes().iter().any(|r| { - r.nat_gateway_id().is_some() + r.gateway_id() + .map(|g| g.starts_with("igw-")) + .unwrap_or(false) }); + let has_nat = rt.routes().iter().any(|r| r.nat_gateway_id().is_some()); for assoc in rt.associations() { if let Some(sid) = assoc.subnet_id() { subnet_routes.insert(sid.to_string(), (has_igw, has_nat)); @@ -279,13 +367,26 @@ pub async fn select_subnets(ec2: &Ec2Client, vpc_id: &str) -> Result = subnets_resp.subnets().iter().map(|s| { - let id = s.subnet_id().unwrap_or_default().to_string(); - let az = s.availability_zone().unwrap_or_default().to_string(); - let (has_igw, has_nat) = subnet_routes.get(&id).copied().unwrap_or((false, false)); - let kind = if !has_igw { "private".to_string() } else { "public".to_string() }; - SubnetInfo { id, az, kind, has_nat } - }).collect(); + let mut all: Vec = subnets_resp + .subnets() + .iter() + .map(|s| { + let id = s.subnet_id().unwrap_or_default().to_string(); + let az = s.availability_zone().unwrap_or_default().to_string(); + let (has_igw, has_nat) = subnet_routes.get(&id).copied().unwrap_or((false, false)); + let kind = if !has_igw { + "private".to_string() + } else { + "public".to_string() + }; + SubnetInfo { + id, + az, + kind, + has_nat, + } + }) + .collect(); // Priority: private+NAT > private > public, pick 2-3 unique AZs all.sort_by(|a, b| { @@ -303,8 +404,12 @@ pub async fn select_subnets(ec2: &Ec2Client, vpc_id: &str) -> Result= 3 { break; } - if seen_azs.contains(&s.az) { continue; } + if seen_azs.len() >= 3 { + break; + } + if seen_azs.contains(&s.az) { + continue; + } seen_azs.insert(s.az.clone()); selected.push(s); } @@ -316,18 +421,30 @@ pub async fn select_subnets(ec2: &Ec2Client, vpc_id: &str) -> Result Result> { - let resp = ec2.describe_security_groups() - .filters(aws_sdk_ec2::types::Filter::builder().name("vpc-id").values(vpc_id).build()) - .send().await?; - Ok(resp.security_groups().iter().map(|sg| { - SgInfo { + let resp = ec2 + .describe_security_groups() + .filters( + aws_sdk_ec2::types::Filter::builder() + .name("vpc-id") + .values(vpc_id) + .build(), + ) + .send() + .await?; + Ok(resp + .security_groups() + .iter() + .map(|sg| SgInfo { id: sg.group_id().unwrap_or_default().to_string(), name: sg.group_name().unwrap_or_default().to_string(), - } - }).collect()) + }) + .collect()) } /// Pick a security group for a new agent with **no interactive prompt** — @@ -344,11 +461,13 @@ pub async fn default_security_group(ec2: &Ec2Client, vpc_id: &str, name: &str) - if let Some(sg) = existing.iter().find(|sg| sg.name == sg_name) { return Ok(sg.id.clone()); } - let resp = ec2.create_security_group() + let resp = ec2 + .create_security_group() .group_name(&sg_name) .description(format!("OAB agent {name}")) .vpc_id(vpc_id) - .send().await + .send() + .await .context("failed to create security group")?; Ok(resp.group_id().unwrap_or_default().to_string()) } @@ -366,7 +485,10 @@ pub struct DefaultNetworking { pub security_groups: Vec, } -pub async fn default_networking(config: &aws_config::SdkConfig, name: &str) -> Result { +pub async fn default_networking( + config: &aws_config::SdkConfig, + name: &str, +) -> Result { let ec2 = Ec2Client::new(config); let vpc = default_vpc(&ec2).await?; let subnets = select_subnets(&ec2, &vpc.id).await?; @@ -385,7 +507,8 @@ enabled = true api_key = "${secrets.stt_api_key}" model = "whisper-large-v3-turbo" base_url = "https://api.groq.com/openai/v1" -"#.to_string() +"# + .to_string() } else { "[stt]\nenabled = false\n".to_string() }; @@ -436,8 +559,20 @@ usercron_path = "cronjob.toml" ) } -fn generate_manifest(name: &str, namespace: &str, image: &str, config_from: &str, cap: &str, subnets: &[String], sg: &str) -> String { - let subnets_yaml = subnets.iter().map(|s| format!("\"{}\"", s)).collect::>().join(", "); +fn generate_manifest( + name: &str, + namespace: &str, + image: &str, + config_from: &str, + cap: &str, + subnets: &[String], + sg: &str, +) -> String { + let subnets_yaml = subnets + .iter() + .map(|s| format!("\"{}\"", s)) + .collect::>() + .join(", "); format!( r#"apiVersion: oab.dev/v2 kind: OABService @@ -466,6 +601,12 @@ async fn resolve_bucket(_s3: &S3Client, config: &aws_config::SdkConfig) -> Optio return Some(b); } let sts = aws_sdk_sts::Client::new(config); - let account = sts.get_caller_identity().send().await.ok()?.account()?.to_string(); + let account = sts + .get_caller_identity() + .send() + .await + .ok()? + .account()? + .to_string(); Some(format!("oab-control-plane-{account}")) } diff --git a/crates/oabctl/src/delete.rs b/crates/oabctl/src/delete.rs index b8b273f..0144849 100644 --- a/crates/oabctl/src/delete.rs +++ b/crates/oabctl/src/delete.rs @@ -237,15 +237,8 @@ pub(crate) async fn run_with_bucket( }; for object in response.contents() { if let Some(key) = object.key() { - if let Err(error) = s3 - .delete_object() - .bucket(bucket) - .key(key) - .send() - .await - { - cleanup_failures - .push(format!("failed to delete s3://{bucket}/{key}: {error}")); + if let Err(error) = s3.delete_object().bucket(bucket).key(key).send().await { + cleanup_failures.push(format!("failed to delete s3://{bucket}/{key}: {error}")); } } } diff --git a/crates/oabctl/src/driver.rs b/crates/oabctl/src/driver.rs index 01bed27..593d841 100644 --- a/crates/oabctl/src/driver.rs +++ b/crates/oabctl/src/driver.rs @@ -38,7 +38,11 @@ pub struct ProvisionOptions { #[async_trait] pub trait ProvisionDriver { /// Create-or-update the given manifests. - async fn apply(&self, manifests: &[OABServiceManifest], opts: &ProvisionOptions) -> Result; + async fn apply( + &self, + manifests: &[OABServiceManifest], + opts: &ProvisionOptions, + ) -> Result; /// Scale a single service. OAB services carry a single bot token, so /// `size` must be 0 (off) or 1 (on) — enforced by the implementation. @@ -47,7 +51,13 @@ pub trait ProvisionDriver { /// Delete a control-plane resource (`resource` is currently always /// `"oabservice"`). `control_plane_bucket` is the already-resolved bucket /// (see `control_plane::resolve_bucket`) — the driver does not resolve it. - async fn delete(&self, resource: &str, name: &str, namespace: &str, control_plane_bucket: &str) -> Result<()>; + async fn delete( + &self, + resource: &str, + name: &str, + namespace: &str, + control_plane_bucket: &str, + ) -> Result<()>; } /// The ECS implementation. Every method is a thin wrapper over the existing @@ -59,7 +69,11 @@ pub struct EcsDriver<'a> { #[async_trait] impl<'a> ProvisionDriver for EcsDriver<'a> { - async fn apply(&self, manifests: &[OABServiceManifest], opts: &ProvisionOptions) -> Result { + async fn apply( + &self, + manifests: &[OABServiceManifest], + opts: &ProvisionOptions, + ) -> Result { let mut ecs_opts = crate::apply::ApplyOptions::new(self.cluster).with_wait(opts.wait); if let Some(bucket) = &opts.control_plane_bucket { ecs_opts = ecs_opts.with_control_plane_bucket(bucket.clone()); @@ -81,7 +95,13 @@ impl<'a> ProvisionDriver for EcsDriver<'a> { ecsctl::scale::scale_service(&ecs, self.cluster, &service_name, size, false).await } - async fn delete(&self, resource: &str, name: &str, namespace: &str, control_plane_bucket: &str) -> Result<()> { + async fn delete( + &self, + resource: &str, + name: &str, + namespace: &str, + control_plane_bucket: &str, + ) -> Result<()> { crate::delete::run_with_bucket( self.aws_config, resource, @@ -111,7 +131,9 @@ fn ecs_apply_error_to_anyhow(e: crate::apply::ApplyError) -> anyhow::Error { #[cfg(test)] mod tests { use super::*; - use crate::apply::{ApplyAction, AppliedService, ApplyError, ApplyErrorKind, ApplyReport, ServiceTarget}; + use crate::apply::{ + AppliedService, ApplyAction, ApplyError, ApplyErrorKind, ApplyReport, ServiceTarget, + }; #[test] fn ecs_apply_error_to_anyhow_preserves_downcast_and_partial_progress() { @@ -130,7 +152,8 @@ mod tests { name: "mira".to_string(), ecs_service_name: "oab-prod-mira".to_string(), }; - let source = ApplyError::reconciliation(failed.clone(), completed.clone(), anyhow::anyhow!("boom")); + let source = + ApplyError::reconciliation(failed.clone(), completed.clone(), anyhow::anyhow!("boom")); let err = ecs_apply_error_to_anyhow(source); diff --git a/crates/oabctl/src/events.rs b/crates/oabctl/src/events.rs index 2db23b5..62a0ddd 100644 --- a/crates/oabctl/src/events.rs +++ b/crates/oabctl/src/events.rs @@ -49,7 +49,10 @@ pub struct EcsEvent { fn normalize_service(s: &str) -> String { match s.strip_prefix("oab-") { // "prod-mira" → "mira"; "prod-foo-bar" → "foo-bar" (name keeps dashes) - Some(rest) => rest.split_once('-').map_or(rest, |(_, name)| name).to_string(), + Some(rest) => rest + .split_once('-') + .map_or(rest, |(_, name)| name) + .to_string(), None => s.to_string(), } } @@ -162,7 +165,9 @@ pub async fn fetch_ecs_events( let want_service = service.map(normalize_service); let mut out = Vec::new(); for (_, msg) in raw { - let Some(ev) = parse_event(&msg) else { continue }; + let Some(ev) = parse_event(&msg) else { + continue; + }; if let Some(c) = cluster { match &ev.cluster_arn { Some(arn) if arn.ends_with(&format!("/{c}")) => {} diff --git a/crates/oabctl/src/k8s_driver.rs b/crates/oabctl/src/k8s_driver.rs index f3a34db..9896f19 100644 --- a/crates/oabctl/src/k8s_driver.rs +++ b/crates/oabctl/src/k8s_driver.rs @@ -40,7 +40,7 @@ //! mapping table) to land as its own follow-up rather than growing this PR //! further. -use crate::apply::{ApplyAction, AppliedService, ApplyReport}; +use crate::apply::{AppliedService, ApplyAction, ApplyReport}; use crate::driver::{ProvisionDriver, ProvisionOptions}; use crate::manifest::{OABServiceManifest, Runtime}; use anyhow::{Context, Result}; @@ -68,7 +68,13 @@ use std::collections::BTreeMap; pub fn k8s_safe_name(name: &str) -> String { name.to_lowercase() .chars() - .map(|c| if c.is_ascii_alphanumeric() || c == '-' { c } else { '-' }) + .map(|c| { + if c.is_ascii_alphanumeric() || c == '-' { + c + } else { + '-' + } + }) .collect::() .trim_matches('-') .to_string() @@ -112,7 +118,9 @@ impl K8sDriver { } } -fn require_kubernetes_runtime(m: &OABServiceManifest) -> Result<&crate::manifest::KubernetesRuntime> { +fn require_kubernetes_runtime( + m: &OABServiceManifest, +) -> Result<&crate::manifest::KubernetesRuntime> { match &m.spec.runtime { Runtime::Kubernetes(rt) => Ok(rt), Runtime::Ecs(_) => anyhow::bail!( @@ -230,7 +238,11 @@ fn build_deployment(m: &OABServiceManifest) -> Result { // frame, so a separate wrapper here would silently swallow the // parser's own detail (which scheme, which malformed part). Some(Err(e)) => { - anyhow::bail!("{e} — manifest '{}/{}'", m.metadata.namespace, m.metadata.name) + anyhow::bail!( + "{e} — manifest '{}/{}'", + m.metadata.namespace, + m.metadata.name + ) } }; let (command, volumes, volume_mounts) = if let Some((config_map_name, _key)) = configmap_ref { @@ -321,7 +333,11 @@ fn build_deployment(m: &OABServiceManifest) -> Result { #[async_trait] impl ProvisionDriver for K8sDriver { - async fn apply(&self, manifests: &[OABServiceManifest], _opts: &ProvisionOptions) -> Result { + async fn apply( + &self, + manifests: &[OABServiceManifest], + _opts: &ProvisionOptions, + ) -> Result { let mut services = Vec::with_capacity(manifests.len()); for m in manifests { let deployment = build_deployment(m)?; @@ -351,7 +367,11 @@ impl ProvisionDriver for K8sDriver { namespace: m.metadata.namespace.clone(), name: m.metadata.name.clone(), resource_name: name, - action: if existed { ApplyAction::Updated } else { ApplyAction::Created }, + action: if existed { + ApplyAction::Updated + } else { + ApplyAction::Created + }, webhook_urls: Vec::new(), warnings: Vec::new(), }); @@ -375,7 +395,13 @@ impl ProvisionDriver for K8sDriver { Ok(()) } - async fn delete(&self, resource: &str, name: &str, namespace: &str, _control_plane_bucket: &str) -> Result<()> { + async fn delete( + &self, + resource: &str, + name: &str, + namespace: &str, + _control_plane_bucket: &str, + ) -> Result<()> { if resource != "oabservice" { anyhow::bail!("unknown resource type: {resource}. Use 'oabservice'"); } @@ -385,7 +411,9 @@ impl ProvisionDriver for K8sDriver { Ok(_) => Ok(()), // Delete is idempotent — already gone is success, not an error. Err(kube::Error::Api(e)) if e.code == 404 => Ok(()), - Err(e) => Err(e).with_context(|| format!("failed to delete k8s deployment '{dep_name}'")), + Err(e) => { + Err(e).with_context(|| format!("failed to delete k8s deployment '{dep_name}'")) + } } } } @@ -457,12 +485,21 @@ mod tests { // a manifest carrying it builds identically to one without. let with = k8s_manifest(Some("s3://bucket/artifacts/prod/orca/"), &[]); let without = k8s_manifest(None, &[]); - assert_eq!(build_deployment(&with).unwrap(), build_deployment(&without).unwrap()); + assert_eq!( + build_deployment(&with).unwrap(), + build_deployment(&without).unwrap() + ); } #[test] fn build_deployment_wires_secret_key_ref() { - let m = k8s_manifest(None, &[("DISCORD_BOT_TOKEN", "k8s-secret://oab-orca#DISCORD_BOT_TOKEN")]); + let m = k8s_manifest( + None, + &[( + "DISCORD_BOT_TOKEN", + "k8s-secret://oab-orca#DISCORD_BOT_TOKEN", + )], + ); let dep = build_deployment(&m).unwrap(); let pod = dep.spec.unwrap().template.spec.unwrap(); let env = pod.containers[0].env.as_ref().unwrap(); @@ -480,7 +517,13 @@ mod tests { #[test] fn build_deployment_rejects_non_k8s_secret_scheme() { - let m = k8s_manifest(None, &[("DISCORD_BOT_TOKEN", "aws-sm://oab/prod/orca#DISCORD_BOT_TOKEN")]); + let m = k8s_manifest( + None, + &[( + "DISCORD_BOT_TOKEN", + "aws-sm://oab/prod/orca#DISCORD_BOT_TOKEN", + )], + ); let err = build_deployment(&m).unwrap_err(); assert!(err.to_string().contains("k8s-secret://")); } @@ -492,7 +535,10 @@ mod tests { let m = k8s_manifest(None, &[("DISCORD_BOT_TOKEN", "k8s-secret://oab-orca")]); let err = build_deployment(&m).unwrap_err(); let msg = err.to_string(); - assert!(msg.contains("DISCORD_BOT_TOKEN"), "must name the env var: {msg}"); + assert!( + msg.contains("DISCORD_BOT_TOKEN"), + "must name the env var: {msg}" + ); assert!(msg.contains("prod/orca"), "must name the agent: {msg}"); } @@ -525,7 +571,10 @@ mod tests { m.spec.config_from = "k8s-configmap://orca-config".to_string(); // missing #key let err = build_deployment(&m).unwrap_err(); let msg = err.to_string(); - assert!(msg.contains("k8s-configmap://"), "must name the scheme: {msg}"); + assert!( + msg.contains("k8s-configmap://"), + "must name the scheme: {msg}" + ); assert!(msg.contains("prod/orca"), "must name the agent: {msg}"); } @@ -538,7 +587,10 @@ mod tests { let pod = dep.spec.unwrap().template.spec.unwrap(); let container = &pod.containers[0]; - assert_eq!(container.image.as_deref(), Some("ghcr.io/openabdev/openab:latest")); + assert_eq!( + container.image.as_deref(), + Some("ghcr.io/openabdev/openab:latest") + ); assert_eq!( container.command.as_deref(), Some( @@ -564,15 +616,21 @@ mod tests { #[test] fn build_deployment_wires_service_account_and_node_selector() { let mut m = k8s_manifest(None, &[]); - let Runtime::Kubernetes(rt) = &mut m.spec.runtime else { unreachable!() }; + let Runtime::Kubernetes(rt) = &mut m.spec.runtime else { + unreachable!() + }; rt.service_account = Some("orca-sa".to_string()); - rt.node_selector.insert("kubernetes.io/arch".to_string(), "arm64".to_string()); + rt.node_selector + .insert("kubernetes.io/arch".to_string(), "arm64".to_string()); let dep = build_deployment(&m).unwrap(); let pod = dep.spec.unwrap().template.spec.unwrap(); assert_eq!(pod.service_account_name.as_deref(), Some("orca-sa")); assert_eq!( - pod.node_selector.unwrap().get("kubernetes.io/arch").map(String::as_str), + pod.node_selector + .unwrap() + .get("kubernetes.io/arch") + .map(String::as_str), Some("arm64") ); } diff --git a/crates/oabctl/src/manifest.rs b/crates/oabctl/src/manifest.rs index a49dbc7..1774c29 100644 --- a/crates/oabctl/src/manifest.rs +++ b/crates/oabctl/src/manifest.rs @@ -79,7 +79,10 @@ pub struct AgentOverride { impl OABFleetManifest { pub fn validate(&self) -> anyhow::Result<()> { if self.api_version != "oab.dev/v2" { - anyhow::bail!("unsupported apiVersion: {} (expected oab.dev/v2)", self.api_version); + anyhow::bail!( + "unsupported apiVersion: {} (expected oab.dev/v2)", + self.api_version + ); } if self.kind != "OABFleet" { anyhow::bail!("unsupported kind: {}", self.kind); @@ -103,51 +106,69 @@ impl OABFleetManifest { /// Expand fleet into individual OABService manifests pub fn expand(&self) -> Vec { - self.spec.agents.iter().map(|agent| { - let resources = agent.resources.clone() - .or(self.spec.template.resources.clone()) - .unwrap_or(Resources { cpu: "256".into(), memory: "512".into() }); - let base_secrets = agent.secrets.clone() - .unwrap_or_else(|| self.spec.template.secrets.clone()); - // Interpolate ${name} in secret values - let secrets = base_secrets.into_iter().map(|(k, v)| { - (k, v.replace("${name}", &agent.name)) - }).collect(); - - OABServiceManifest { - api_version: self.api_version.clone(), - kind: "OABService".to_string(), - metadata: Metadata { - name: agent.name.clone(), - namespace: self.metadata.namespace.clone(), - generation: 0, - }, - spec: Spec { - image: agent.image.clone() - .unwrap_or_else(|| self.spec.template.image.clone()), - resources, - config_from: agent.config_from.replace("${name}", &agent.name), - bundle_from: agent.bundle_from.clone() - .or(self.spec.template.bundle_from.clone()) - .map(|s| s.replace("${name}", &agent.name)), - bootstrap_from: agent.bootstrap_from.clone() - .or(self.spec.template.bootstrap_from.clone()) - .map(|s| s.replace("${name}", &agent.name)), - secrets, - runtime: self.spec.template.runtime.clone(), - ingress: agent - .ingress - .clone() - .or_else(|| self.spec.template.ingress.clone()), - // Fleets have no per-agent/template ACP concept yet — - // unaffected by studio#119's default-on behavior, which - // only applies to the "+ New fleet" single-agent wizard - // path (build_default_manifest/build_default_k8s_manifest - // in studio-cp). - acp_enabled: None, - }, - } - }).collect() + self.spec + .agents + .iter() + .map(|agent| { + let resources = agent + .resources + .clone() + .or(self.spec.template.resources.clone()) + .unwrap_or(Resources { + cpu: "256".into(), + memory: "512".into(), + }); + let base_secrets = agent + .secrets + .clone() + .unwrap_or_else(|| self.spec.template.secrets.clone()); + // Interpolate ${name} in secret values + let secrets = base_secrets + .into_iter() + .map(|(k, v)| (k, v.replace("${name}", &agent.name))) + .collect(); + + OABServiceManifest { + api_version: self.api_version.clone(), + kind: "OABService".to_string(), + metadata: Metadata { + name: agent.name.clone(), + namespace: self.metadata.namespace.clone(), + generation: 0, + }, + spec: Spec { + image: agent + .image + .clone() + .unwrap_or_else(|| self.spec.template.image.clone()), + resources, + config_from: agent.config_from.replace("${name}", &agent.name), + bundle_from: agent + .bundle_from + .clone() + .or(self.spec.template.bundle_from.clone()) + .map(|s| s.replace("${name}", &agent.name)), + bootstrap_from: agent + .bootstrap_from + .clone() + .or(self.spec.template.bootstrap_from.clone()) + .map(|s| s.replace("${name}", &agent.name)), + secrets, + runtime: self.spec.template.runtime.clone(), + ingress: agent + .ingress + .clone() + .or_else(|| self.spec.template.ingress.clone()), + // Fleets have no per-agent/template ACP concept yet — + // unaffected by studio#119's default-on behavior, which + // only applies to the "+ New fleet" single-agent wizard + // path (build_default_manifest/build_default_k8s_manifest + // in studio-cp). + acp_enabled: None, + }, + } + }) + .collect() } } @@ -324,7 +345,10 @@ const VALID_ECS_CPU: &[&str] = &["256", "512", "1024", "2048", "4096"]; impl OABServiceManifest { pub fn validate(&self) -> anyhow::Result<()> { if self.api_version != "oab.dev/v2" { - anyhow::bail!("unsupported apiVersion: {} (expected oab.dev/v2)", self.api_version); + anyhow::bail!( + "unsupported apiVersion: {} (expected oab.dev/v2)", + self.api_version + ); } if self.kind != "OABService" { anyhow::bail!("unsupported kind: {}", self.kind); @@ -576,7 +600,8 @@ spec: securityGroups: ["sg-1"] "#; let m = parse(yaml); - m.validate().expect("should be valid with default architecture"); + m.validate() + .expect("should be valid with default architecture"); match &m.spec.runtime { Runtime::Ecs(ecs) => assert_eq!(ecs.architecture, "X86_64"), _ => panic!("expected ECS runtime"), @@ -607,7 +632,9 @@ spec: "#; let m = parse(yaml); let err = m.validate().unwrap_err(); - assert!(err.to_string().contains("runtime.architecture must be one of")); + assert!(err + .to_string() + .contains("runtime.architecture must be one of")); } #[test] @@ -634,7 +661,9 @@ spec: "#; let m = parse(yaml); let err = m.validate().unwrap_err(); - assert!(err.to_string().contains("runtime.architecture must be one of")); + assert!(err + .to_string() + .contains("runtime.architecture must be one of")); } #[test] @@ -671,7 +700,11 @@ spec: assert_eq!(expanded.len(), 2); let from_template = &expanded[0]; - let ing = from_template.spec.ingress.as_ref().expect("template ingress"); + let ing = from_template + .spec + .ingress + .as_ref() + .expect("template ingress"); assert_eq!(ing.paths, vec!["/webhook/telegram"]); assert_eq!(ing.cloud_map_namespace, "oab"); diff --git a/crates/oabctl/src/scale.rs b/crates/oabctl/src/scale.rs index f7cd390..ed61005 100644 --- a/crates/oabctl/src/scale.rs +++ b/crates/oabctl/src/scale.rs @@ -363,7 +363,9 @@ pub async fn list_schedules(aws_config: &aws_config::SdkConfig) -> Result<()> { if all_schedules.is_empty() { println!("No schedules found in group '{group_name}'."); - println!(" Use 'oabctl schedule create --expr ' to create one."); + println!( + " Use 'oabctl schedule create --expr ' to create one." + ); return Ok(()); } diff --git a/crates/oabctl/src/secrets.rs b/crates/oabctl/src/secrets.rs index b2ffce7..f4af579 100644 --- a/crates/oabctl/src/secrets.rs +++ b/crates/oabctl/src/secrets.rs @@ -131,7 +131,9 @@ fn split_ecs_json_key_suffix(value: &str) -> Result<(&str, Option<&str>)> { return Ok((value, None)); }; let base = &value[..after_name_start + name_end]; - let fields: Vec<&str> = value[after_name_start + name_end + 1..].split(':').collect(); + let fields: Vec<&str> = value[after_name_start + name_end + 1..] + .split(':') + .collect(); let (json_key, version_stage, version_id) = match fields.as_slice() { [k] => (*k, "", ""), [k, s] => (*k, *s, ""), @@ -210,7 +212,10 @@ mod tests { "arn:aws:secretsmanager:us-east-1:903779448426:secret:oab/telegram/pahudxbot-AC80TP:TELEGRAM_BOT_TOKEN::", ) .unwrap(); - assert_eq!(base, "arn:aws:secretsmanager:us-east-1:903779448426:secret:oab/telegram/pahudxbot-AC80TP"); + assert_eq!( + base, + "arn:aws:secretsmanager:us-east-1:903779448426:secret:oab/telegram/pahudxbot-AC80TP" + ); assert_eq!(key, Some("TELEGRAM_BOT_TOKEN")); } @@ -220,7 +225,10 @@ mod tests { "arn:aws:secretsmanager:us-east-1:903779448426:secret:oab/telegram/pahudxbot-AC80TP", ) .unwrap(); - assert_eq!(base, "arn:aws:secretsmanager:us-east-1:903779448426:secret:oab/telegram/pahudxbot-AC80TP"); + assert_eq!( + base, + "arn:aws:secretsmanager:us-east-1:903779448426:secret:oab/telegram/pahudxbot-AC80TP" + ); assert_eq!(key, None); } @@ -237,9 +245,14 @@ mod tests { // optional fields empty (`::`) must not be misparsed as // json-key="mysecret" — "mysecret" here is part of the secret // name/base ARN, not a suffix field. - let (base, key) = - split_ecs_json_key_suffix("arn:aws:secretsmanager:us-east-1:903779448426:secret:mysecret::").unwrap(); - assert_eq!(base, "arn:aws:secretsmanager:us-east-1:903779448426:secret:mysecret"); + let (base, key) = split_ecs_json_key_suffix( + "arn:aws:secretsmanager:us-east-1:903779448426:secret:mysecret::", + ) + .unwrap(); + assert_eq!( + base, + "arn:aws:secretsmanager:us-east-1:903779448426:secret:mysecret" + ); assert_eq!(key, None); } @@ -284,7 +297,9 @@ mod tests { #[test] fn parse_aws_sm_uri_rejects_missing_hash() { - assert!(parse_aws_sm_uri("aws-sm://oab/telegram/pahudxbot").unwrap().is_err()); + assert!(parse_aws_sm_uri("aws-sm://oab/telegram/pahudxbot") + .unwrap() + .is_err()); } #[test] @@ -295,7 +310,9 @@ mod tests { #[test] fn parse_aws_sm_uri_returns_none_for_other_schemes() { - assert!(parse_aws_sm_uri("arn:aws:secretsmanager:us-east-1:123:secret:oab/x-AbCdEf").is_none()); + assert!( + parse_aws_sm_uri("arn:aws:secretsmanager:us-east-1:123:secret:oab/x-AbCdEf").is_none() + ); assert!(parse_aws_sm_uri("plain-secret-name").is_none()); } @@ -310,13 +327,17 @@ mod tests { #[test] fn parse_k8s_secret_uri_rejects_missing_hash() { - assert!(parse_k8s_secret_uri("k8s-secret://oab-orca").unwrap().is_err()); + assert!(parse_k8s_secret_uri("k8s-secret://oab-orca") + .unwrap() + .is_err()); } #[test] fn parse_k8s_secret_uri_rejects_empty_parts() { assert!(parse_k8s_secret_uri("k8s-secret://#key").unwrap().is_err()); - assert!(parse_k8s_secret_uri("k8s-secret://oab-orca#").unwrap().is_err()); + assert!(parse_k8s_secret_uri("k8s-secret://oab-orca#") + .unwrap() + .is_err()); } #[test] @@ -336,13 +357,19 @@ mod tests { #[test] fn parse_k8s_configmap_uri_rejects_missing_hash() { - assert!(parse_k8s_configmap_uri("k8s-configmap://oab-orca-config").unwrap().is_err()); + assert!(parse_k8s_configmap_uri("k8s-configmap://oab-orca-config") + .unwrap() + .is_err()); } #[test] fn parse_k8s_configmap_uri_rejects_empty_parts() { - assert!(parse_k8s_configmap_uri("k8s-configmap://#config.toml").unwrap().is_err()); - assert!(parse_k8s_configmap_uri("k8s-configmap://oab-orca-config#").unwrap().is_err()); + assert!(parse_k8s_configmap_uri("k8s-configmap://#config.toml") + .unwrap() + .is_err()); + assert!(parse_k8s_configmap_uri("k8s-configmap://oab-orca-config#") + .unwrap() + .is_err()); } #[test] diff --git a/crates/oabctl/src/status.rs b/crates/oabctl/src/status.rs index 9130b6c..ea72dc6 100644 --- a/crates/oabctl/src/status.rs +++ b/crates/oabctl/src/status.rs @@ -242,7 +242,12 @@ async fn task_def_has_health_check( if let Some(v) = cache.get(arn) { return *v; } - let defined = match ecs.describe_task_definition().task_definition(arn).send().await { + let defined = match ecs + .describe_task_definition() + .task_definition(arn) + .send() + .await + { Ok(resp) => resp .task_definition() .map(|td| { diff --git a/crates/oabctl/src/studio_api.rs b/crates/oabctl/src/studio_api.rs index 5097ae0..8eacfb2 100644 --- a/crates/oabctl/src/studio_api.rs +++ b/crates/oabctl/src/studio_api.rs @@ -223,10 +223,13 @@ pub async fn provision( control_plane_bucket: control_plane_bucket.map(str::to_string), wait: false, }; - crate::driver::EcsDriver { aws_config: config, cluster } - .apply(&manifests, &opts) - .await - .context("failed to apply manifest during provision") + crate::driver::EcsDriver { + aws_config: config, + cluster, + } + .apply(&manifests, &opts) + .await + .context("failed to apply manifest during provision") } /// [`provision`], but takes an already-built [`crate::manifest::OABServiceManifest`] @@ -310,7 +313,14 @@ pub async fn load_manifest( Ok(Some(manifest)) } // A missing object is the "not provisioned yet" signal, not an error. - Err(err) if err.as_service_error().map(|e| e.is_no_such_key()).unwrap_or(false) => Ok(None), + Err(err) + if err + .as_service_error() + .map(|e| e.is_no_such_key()) + .unwrap_or(false) => + { + Ok(None) + } Err(err) => { Err(anyhow::Error::new(err).context(format!("failed to fetch stored manifest '{key}'"))) } @@ -338,7 +348,9 @@ pub async fn redeploy( let mut manifest = load_manifest(config, namespace, name, Some(&bucket)) .await? .with_context(|| { - format!("no stored manifest for {namespace}/{name} — create the agent before redeploying") + format!( + "no stored manifest for {namespace}/{name} — create the agent before redeploying" + ) })?; if let Some(img) = image.filter(|s| !s.is_empty()) { @@ -363,9 +375,12 @@ pub async fn scale( name: &str, size: i32, ) -> Result<()> { - crate::driver::EcsDriver { aws_config: config, cluster } - .scale(namespace, name, size) - .await + crate::driver::EcsDriver { + aws_config: config, + cluster, + } + .scale(namespace, name, size) + .await } /// Delete a control-plane resource (currently `oabservice`). @@ -382,9 +397,12 @@ pub async fn delete( control_plane_bucket: Option<&str>, ) -> Result<()> { let bucket = crate::control_plane::resolve_bucket(config, control_plane_bucket).await?; - crate::driver::EcsDriver { aws_config: config, cluster } - .delete(resource, name, namespace, &bucket) - .await + crate::driver::EcsDriver { + aws_config: config, + cluster, + } + .delete(resource, name, namespace, &bucket) + .await } #[cfg(test)] @@ -417,7 +435,8 @@ mod tests { #[test] fn inject_pre_seed_hook_appends_to_existing_config() { let config = b"[agent]\nname = \"orca\"\n"; - let out = inject_pre_seed_hook(config, "s3://bucket/artifacts/prod/orca/bundle.zip").unwrap(); + let out = + inject_pre_seed_hook(config, "s3://bucket/artifacts/prod/orca/bundle.zip").unwrap(); let text = String::from_utf8(out).unwrap(); assert!(text.starts_with("[agent]\nname = \"orca\"\n")); assert!(text.contains("[hooks.pre_seed]")); @@ -428,8 +447,10 @@ mod tests { #[test] fn inject_pre_seed_hook_is_a_true_noop_only_when_source_already_present() { - let config = b"[hooks.pre_seed]\nsources = [\"s3://bucket/artifacts/prod/orca/bundle.zip\"]\n"; - let out = inject_pre_seed_hook(config, "s3://bucket/artifacts/prod/orca/bundle.zip").unwrap(); + let config = + b"[hooks.pre_seed]\nsources = [\"s3://bucket/artifacts/prod/orca/bundle.zip\"]\n"; + let out = + inject_pre_seed_hook(config, "s3://bucket/artifacts/prod/orca/bundle.zip").unwrap(); // unchanged, byte-for-byte — this exact source is already wired, a // second redeploy of the same agent shouldn't touch the file at all assert_eq!(out, config); @@ -442,10 +463,17 @@ mod tests { // bundle) used to make injection silently no-op entirely, so the new // bundle was never wired in at all. It must be added alongside. let config = b"[hooks.pre_seed]\nsources = [\"s3://other/state.zip\"]\n"; - let out = inject_pre_seed_hook(config, "s3://bucket/artifacts/prod/orca/bundle.zip").unwrap(); + let out = + inject_pre_seed_hook(config, "s3://bucket/artifacts/prod/orca/bundle.zip").unwrap(); let text = String::from_utf8(out).unwrap(); - assert!(text.contains("s3://other/state.zip"), "operator's own source must survive: {text}"); - assert!(text.contains("s3://bucket/artifacts/prod/orca/bundle.zip"), "new source must be added: {text}"); + assert!( + text.contains("s3://other/state.zip"), + "operator's own source must survive: {text}" + ); + assert!( + text.contains("s3://bucket/artifacts/prod/orca/bundle.zip"), + "new source must be added: {text}" + ); let reparsed: toml::Value = text.parse().expect("valid toml"); let sources = reparsed["hooks"]["pre_seed"]["sources"].as_array().unwrap(); assert_eq!(sources.len(), 2); @@ -484,7 +512,8 @@ sources = ["s3://a", "s3://b", "s3://c", "s3://d", "s3://e"] let config = b"# a comment worth keeping\n[agent]\nname = \"orca\" # inline comment\n"; let out = inject_pre_seed_hook(config, "s3://x/bundle.zip").unwrap(); let text = String::from_utf8(out).unwrap(); - assert!(text.starts_with("# a comment worth keeping\n[agent]\nname = \"orca\" # inline comment\n")); + assert!(text + .starts_with("# a comment worth keeping\n[agent]\nname = \"orca\" # inline comment\n")); } #[test] diff --git a/crates/oabctl/src/vendor_images.rs b/crates/oabctl/src/vendor_images.rs index b7a9f67..c73cd5c 100644 --- a/crates/oabctl/src/vendor_images.rs +++ b/crates/oabctl/src/vendor_images.rs @@ -100,7 +100,9 @@ async fn fetch_openab_releases(client: &reqwest::Client) -> Result Vec { releases .iter() - .filter(|r| !r.prerelease && r.tag_name.starts_with("openab-") && !r.tag_name.contains("-beta")) + .filter(|r| { + !r.prerelease && r.tag_name.starts_with("openab-") && !r.tag_name.contains("-beta") + }) .map(|r| r.tag_name.trim_start_matches("openab-").to_string()) .collect() } @@ -151,7 +153,10 @@ pub async fn resolve_vendor_image_tags(vendor: &str) -> VendorImageTags { for version in beta_release_versions(&releases) { let candidate = format!("{version}-{vendor}"); - if ghcr_tag_exists(&client, &token, &candidate).await.unwrap_or(false) { + if ghcr_tag_exists(&client, &token, &candidate) + .await + .unwrap_or(false) + { out.beta = Some(candidate); break; } @@ -159,7 +164,10 @@ pub async fn resolve_vendor_image_tags(vendor: &str) -> VendorImageTags { for version in stable_release_versions(&releases) { let candidate = format!("{version}-{vendor}"); - if ghcr_tag_exists(&client, &token, &candidate).await.unwrap_or(false) { + if ghcr_tag_exists(&client, &token, &candidate) + .await + .unwrap_or(false) + { out.stable = Some(candidate); break; } @@ -173,7 +181,10 @@ mod tests { use super::*; fn release(tag_name: &str, prerelease: bool) -> GhRelease { - GhRelease { tag_name: tag_name.to_string(), prerelease } + GhRelease { + tag_name: tag_name.to_string(), + prerelease, + } } // Real-world snapshot (2026-09-08, `gh api repos/openabdev/openab/releases`): @@ -200,13 +211,21 @@ mod tests { fn beta_versions_newest_first_ignores_prerelease_flag() { assert_eq!( beta_release_versions(&sample_releases()), - vec!["0.10.0-beta.3", "0.10.0-beta.2", "0.10.0-beta.1", "0.9.0-beta.12"] + vec![ + "0.10.0-beta.3", + "0.10.0-beta.2", + "0.10.0-beta.1", + "0.9.0-beta.12" + ] ); } #[test] fn both_channels_ignore_non_openab_prefixed_releases() { - let releases = vec![release("oabctl-pre-beta", true), release("pre-seed-utils-v2.35.13-ghp0.3.2", false)]; + let releases = vec![ + release("oabctl-pre-beta", true), + release("pre-seed-utils-v2.35.13-ghp0.3.2", false), + ]; assert!(stable_release_versions(&releases).is_empty()); assert!(beta_release_versions(&releases).is_empty()); } From b2325478bf9c9aa0589d5782dce854493f7cd2de Mon Sep 17 00:00:00 2001 From: Devin Date: Wed, 30 Sep 2026 13:03:14 +0800 Subject: [PATCH 3/7] refactor: satisfy workspace-wide clippy -D warnings MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Pre-existing violations blocking cargo clippy --workspace --all-targets -- -D warnings: - studio-compose: SkillsLibrary::from_iter tripped should_implement_trait — implement FromIterator<(N, Skill)> for real; the five existing SkillsLibrary::from_iter call sites resolve to the trait method unchanged, and .collect() into a library now works too. - studio-cp: derive Default for FleetRuntime (was a manual impl clippy flags as derivable), rewrite a strip_prefix/else-return-None as ? in role_identity, and add targeted allow(clippy::too_many_arguments) on the four provision_* boundary functions whose signatures are already public API — restructuring them ripples outside the workspace. Both files also carry their cargo fmt normalization. --- crates/studio-compose/src/lib.rs | 31 ++-- crates/studio-cp/src/lib.rs | 264 ++++++++++++++++++++++--------- 2 files changed, 209 insertions(+), 86 deletions(-) diff --git a/crates/studio-compose/src/lib.rs b/crates/studio-compose/src/lib.rs index ef8f7b0..db9b791 100644 --- a/crates/studio-compose/src/lib.rs +++ b/crates/studio-compose/src/lib.rs @@ -90,13 +90,9 @@ pub struct SkillsLibrary { pub skills: BTreeMap, } -impl SkillsLibrary { - /// Build a library from `(name, skill)` pairs — convenience for callers/tests. - pub fn from_iter(it: I) -> Self - where - I: IntoIterator, - N: Into, - { +/// Build a library from `(name, skill)` pairs — convenience for callers/tests. +impl> FromIterator<(N, Skill)> for SkillsLibrary { + fn from_iter>(it: T) -> Self { SkillsLibrary { skills: it.into_iter().map(|(n, s)| (n.into(), s)).collect(), } @@ -190,7 +186,9 @@ impl Bundle { // start_file/write on an in-memory Cursor> cannot fail // (no OS I/O involved) — unwrap keeps this fn infallible, matching // every other pure Bundle method (digest, preview, artifact_objects). - writer.start_file(path, options).expect("zip: in-memory write"); + writer + .start_file(path, options) + .expect("zip: in-memory write"); std::io::Write::write_all(&mut writer, bytes).expect("zip: in-memory write"); } writer.finish().expect("zip: in-memory write"); @@ -520,10 +518,7 @@ mod tests { let mut t = tmpl(); t.skills = vec!["s".into()]; let overlay = Overlay { - files: BTreeMap::from([( - ".claude/skills/s/SKILL.md".into(), - "from overlay\n".into(), - )]), + files: BTreeMap::from([(".claude/skills/s/SKILL.md".into(), "from overlay\n".into())]), ..Default::default() }; let b = compose(&t, &overlay, &lib).unwrap(); @@ -542,7 +537,9 @@ mod tests { referenced_by: SkillRef::Template, } ); - assert!(err.to_string().contains("template references skill \"nope\"")); + assert!(err + .to_string() + .contains("template references skill \"nope\"")); } #[test] @@ -601,7 +598,8 @@ mod tests { fn digest_changes_when_a_byte_changes() { let base = compose(&tmpl(), &Overlay::default(), &SkillsLibrary::default()).unwrap(); let mut t = tmpl(); - t.files.insert("config.toml".into(), "[agent]\nname = \"x\"\n".into()); + t.files + .insert("config.toml".into(), "[agent]\nname = \"x\"\n".into()); let changed = compose(&t, &Overlay::default(), &SkillsLibrary::default()).unwrap(); assert_ne!(base.digest(), changed.digest()); } @@ -702,7 +700,10 @@ mod tests { ] ); // bytes travel with the key, unmodified - let (_, cfg) = objs.iter().find(|(k, _)| k.ends_with("/config.toml")).unwrap(); + let (_, cfg) = objs + .iter() + .find(|(k, _)| k.ends_with("/config.toml")) + .unwrap(); assert_eq!(cfg, b"[agent]\nname = \"base\"\n"); } diff --git a/crates/studio-cp/src/lib.rs b/crates/studio-cp/src/lib.rs index b35d885..49f1105 100644 --- a/crates/studio-cp/src/lib.rs +++ b/crates/studio-cp/src/lib.rs @@ -167,7 +167,15 @@ pub async fn observe_events( since_ms: i64, limit: i32, ) -> anyhow::Result> { - oabctl::fetch_ecs_events(aws_config, log_group, Some(cluster), service, since_ms, limit).await + oabctl::fetch_ecs_events( + aws_config, + log_group, + Some(cluster), + service, + since_ms, + limit, + ) + .await } // ---- k8s observe (studio#146) -------------------------------------------- @@ -296,7 +304,10 @@ fn k8s_latched_verified(pod: &k8s_openapi::api::core::v1::Pod) -> bool { /// `instance_phase`'s ECS derivation. `lease_valid`/`accepting_work` are /// CP-level, not yet k8s-observable, so they default to valid/admitting — /// same stance `instance_phase` takes for ECS today. -pub fn k8s_instance_phase(pod: &k8s_openapi::api::core::v1::Pod, verified_before: bool) -> AgentState { +pub fn k8s_instance_phase( + pod: &k8s_openapi::api::core::v1::Pod, + verified_before: bool, +) -> AgentState { use agent_lifecycle::k8s::{K8sDriver, K8sPod}; use agent_lifecycle::RuntimeDriver; @@ -335,7 +346,9 @@ async fn find_k8s_pods( .await .map_err(|e| anyhow::anyhow!("failed to list k8s deployments in '{namespace}': {e}"))?; let Some(dep) = deployments.items.into_iter().find(|d| { - let Some(name) = k8s_oab_name(d) else { return false }; + let Some(name) = k8s_oab_name(d) else { + return false; + }; service == format!("oab-{namespace}-{name}") || service == name }) else { return Ok(None); @@ -591,13 +604,17 @@ pub async fn observe_k8s_identity(context: Option<&str>) -> anyhow::Result format!("{}/{}", c.cluster, c.namespace.as_deref().unwrap_or("default")), + Some(c) => format!( + "{}/{}", + c.cluster, + c.namespace.as_deref().unwrap_or("default") + ), None => String::new(), }; - let client = k8s_client_for(context) - .await - .map_err(|e| anyhow::anyhow!("failed to resolve kubeconfig context '{context_name}': {e}"))?; + let client = k8s_client_for(context).await.map_err(|e| { + anyhow::anyhow!("failed to resolve kubeconfig context '{context_name}': {e}") + })?; let api: Api = Api::all(client); let review = api @@ -868,7 +885,11 @@ pub async fn list_namespaces(context: Option<&str>) -> anyhow::Result) -> anyhow::Result, namespace: &str) -> anyhow::Result> { +pub async fn list_service_accounts( + context: Option<&str>, + namespace: &str, +) -> anyhow::Result> { use k8s_openapi::api::core::v1::ServiceAccount; use kube::api::{Api, ListParams}; @@ -887,7 +911,11 @@ pub async fn list_service_accounts(context: Option<&str>, namespace: &str) -> an .list(&ListParams::default()) .await .map_err(|e| anyhow::anyhow!("failed to list service accounts: {e}"))?; - Ok(list.items.into_iter().filter_map(|sa| sa.metadata.name).collect()) + Ok(list + .items + .into_iter() + .filter_map(|sa| sa.metadata.name) + .collect()) } // ---- Fleet → managing-credential binding (ADR: Per-Fleet managing identity) -- @@ -908,19 +936,14 @@ pub async fn list_service_accounts(context: Option<&str>, namespace: &str) -> an // the existing schema, not a parallel one). `fleets-k8s.toml` is migrated into // `fleets.toml` once at startup — see `migrate_legacy_k8s_bindings`. -#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Deserialize, serde::Serialize)] +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Deserialize, serde::Serialize)] #[serde(rename_all = "lowercase")] pub enum FleetRuntime { + #[default] Ecs, K8s, } -impl Default for FleetRuntime { - fn default() -> Self { - FleetRuntime::Ecs - } -} - /// A declarative binding of a managed fleet to the credential/context that /// should manage it, plus the fleet's members. Profile-first for ECS /// (assume-role is later work); `context`+`namespace` stand in for k8s (no AWS @@ -1152,10 +1175,8 @@ fn role_identity(arn: &str) -> Option<(String, String)> { let resource = parts[5..].join(":"); let name = if let Some(r) = resource.strip_prefix("assumed-role/") { r.split('/').next()?.to_string() - } else if let Some(r) = resource.strip_prefix("role/") { - r.to_string() } else { - return None; + resource.strip_prefix("role/")?.to_string() }; Some((account, name)) } @@ -1422,10 +1443,7 @@ pub fn write_bindings_atomic(path: &std::path::Path, text: &str) -> anyhow::Resu /// Validate `text` parses as a bindings file and, if so, persist it verbatim, /// returning the parsed set for the caller to hot-reload. A parse error is /// returned **before** anything is written, so a bad edit never lands on disk. -pub fn save_bindings_text( - path: &std::path::Path, - text: &str, -) -> anyhow::Result { +pub fn save_bindings_text(path: &std::path::Path, text: &str) -> anyhow::Result { let parsed: FleetBindings = toml::from_str(text)?; write_bindings_atomic(path, text)?; Ok(parsed) @@ -1480,7 +1498,10 @@ pub struct ProvisionOutcome { /// already uploads a copy of the composed config.toml to, and the same /// convention `oabctl create`'s wizard uses for its own generated manifest. fn default_config_from_uri(bucket: &str, namespace: &str, name: &str) -> String { - format!("s3://{bucket}/{}/config.toml", studio_compose::artifacts_prefix(namespace, name)) + format!( + "s3://{bucket}/{}/config.toml", + studio_compose::artifacts_prefix(namespace, name) + ) } /// Stores an `OPENAB_ACP_AUTH_KEY` in Secrets Manager under the same @@ -1501,7 +1522,9 @@ async fn provision_acp_auth_secret( ) -> anyhow::Result { let sm = aws_sdk_secretsmanager::Client::new(aws_config); let secret_name = format!("oab/{namespace}/{name}"); - let key = token.map(str::to_string).unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); + let key = token + .map(str::to_string) + .unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); let secret_obj = serde_json::json!({ "OPENAB_ACP_AUTH_KEY": key }); oabctl::create::store_secret(&sm, &secret_name, &secret_obj.to_string()).await?; Ok(format!("aws-sm://{secret_name}#OPENAB_ACP_AUTH_KEY")) @@ -1532,7 +1555,8 @@ async fn build_default_manifest( let config_from = default_config_from_uri(bucket, namespace, name); let mut secrets = std::collections::HashMap::new(); if acp_enabled { - let acp_auth_ref = provision_acp_auth_secret(aws_config, namespace, name, acp_token).await?; + let acp_auth_ref = + provision_acp_auth_secret(aws_config, namespace, name, acp_token).await?; secrets.insert("OPENAB_ACP_AUTH_KEY".to_string(), acp_auth_ref); } Ok(oabctl::manifest::OABServiceManifest { @@ -1578,6 +1602,7 @@ async fn build_default_manifest( /// built instead ([`build_default_manifest`], studio#111). /// /// `image_override` (when non-empty) wins over the bundle's own default image tag. +#[allow(clippy::too_many_arguments)] pub async fn provision_from_library( aws_config: &aws_config::SdkConfig, cluster: &str, @@ -1648,10 +1673,19 @@ pub async fn provision_from_library( // (compose-library) path; the caller-controlled toggle is // studio#128's wizard-only `provision_agent`. let mut manifest = - build_default_manifest(aws_config, namespace, name, &image, &bucket, true, None).await?; - manifest.spec.bundle_from = Some(oabctl::studio_api::bundle_from_uri(&bucket, namespace, name)); - oabctl::studio_api::provision_manifest(aws_config, cluster, &manifest, &objects, Some(&bucket)) - .await? + build_default_manifest(aws_config, namespace, name, &image, &bucket, true, None) + .await?; + manifest.spec.bundle_from = Some(oabctl::studio_api::bundle_from_uri( + &bucket, namespace, name, + )); + oabctl::studio_api::provision_manifest( + aws_config, + cluster, + &manifest, + &objects, + Some(&bucket), + ) + .await? } }; @@ -1725,17 +1759,26 @@ async fn provision_agent_secrets( ) -> anyhow::Result<()> { let mut obj = serde_json::Map::new(); if let Some(key) = &input.api_key { - obj.insert("VENDOR_API_KEY".to_string(), serde_json::Value::String(key.clone())); + obj.insert( + "VENDOR_API_KEY".to_string(), + serde_json::Value::String(key.clone()), + ); } match input.chat_platform.as_deref() { Some("discord") => { if let Some(t) = &input.chat_bot_token { - obj.insert("DISCORD_BOT_TOKEN".to_string(), serde_json::Value::String(t.clone())); + obj.insert( + "DISCORD_BOT_TOKEN".to_string(), + serde_json::Value::String(t.clone()), + ); } } Some("telegram") => { if let Some(t) = &input.chat_bot_token { - obj.insert("TELEGRAM_BOT_TOKEN".to_string(), serde_json::Value::String(t.clone())); + obj.insert( + "TELEGRAM_BOT_TOKEN".to_string(), + serde_json::Value::String(t.clone()), + ); } } Some("line") => { @@ -1746,7 +1789,10 @@ async fn provision_agent_secrets( ); } if let Some(s) = &input.chat_channel_secret { - obj.insert("LINE_CHANNEL_SECRET".to_string(), serde_json::Value::String(s.clone())); + obj.insert( + "LINE_CHANNEL_SECRET".to_string(), + serde_json::Value::String(s.clone()), + ); } } _ => {} @@ -1756,7 +1802,12 @@ async fn provision_agent_secrets( } let sm = aws_sdk_secretsmanager::Client::new(aws_config); let secret_name = format!("oab/{namespace}/{name}"); - oabctl::create::store_secret(&sm, &secret_name, &serde_json::Value::Object(obj).to_string()).await?; + oabctl::create::store_secret( + &sm, + &secret_name, + &serde_json::Value::Object(obj).to_string(), + ) + .await?; Ok(()) } @@ -2012,9 +2063,17 @@ pub async fn provision_agent( input.acp_token.as_deref(), ) .await?; - manifest.spec.bundle_from = Some(oabctl::studio_api::bundle_from_uri(&bucket, namespace, name)); - oabctl::studio_api::provision_manifest(aws_config, cluster, &manifest, &objects, Some(&bucket)) - .await? + manifest.spec.bundle_from = Some(oabctl::studio_api::bundle_from_uri( + &bucket, namespace, name, + )); + oabctl::studio_api::provision_manifest( + aws_config, + cluster, + &manifest, + &objects, + Some(&bucket), + ) + .await? } }; @@ -2048,6 +2107,7 @@ pub async fn provision_agent( /// unaffected (that key is delivered as a container-level env var via /// `spec.secrets`/`k8s-secret://`, a completely different, already-working /// mechanism — see `provision_acp_auth_k8s_secret`). +#[allow(clippy::too_many_arguments)] pub async fn provision_agent_k8s( aws_config: &aws_config::SdkConfig, context: Option<&str>, @@ -2093,7 +2153,8 @@ pub async fn provision_agent_k8s( // bundle.zip + `hooks.pre_seed` carrier the AWS path uses. A k8s deploy // now never touches S3, and the pod itself never needs AWS credentials // just to boot. - let config_from = provision_config_k8s_configmap(context, namespace, name, &config_toml).await?; + let config_from = + provision_config_k8s_configmap(context, namespace, name, &config_toml).await?; // Content-address of what's actually being applied — reuses // `studio_compose::Bundle::digest()`'s tested hashing rather than @@ -2125,7 +2186,10 @@ pub async fn provision_agent_k8s( .await?; let driver = oabctl::K8sDriver::from_context(context).await?; - let opts = oabctl::ProvisionOptions { control_plane_bucket: None, wait: false }; + let opts = oabctl::ProvisionOptions { + control_plane_bucket: None, + wait: false, + }; let report = { use oabctl::ProvisionDriver; driver.apply(std::slice::from_ref(&manifest), &opts).await? @@ -2181,7 +2245,9 @@ async fn provision_acp_auth_k8s_secret( let client = k8s_client_for(context).await?; let secret_name = format!("{}-acp", oabctl::k8s_safe_name(name)); - let key = token.map(str::to_string).unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); + let key = token + .map(str::to_string) + .unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); let mut string_data = std::collections::BTreeMap::new(); string_data.insert("OPENAB_ACP_AUTH_KEY".to_string(), key); let secret = Secret { @@ -2194,9 +2260,13 @@ async fn provision_acp_auth_k8s_secret( ..Default::default() }; let api: Api = Api::namespaced(client, namespace); - api.patch(&secret_name, &PatchParams::apply("studio-cp"), &Patch::Apply(&secret)) - .await - .map_err(|e| anyhow::anyhow!("failed to create/apply k8s Secret '{secret_name}': {e}"))?; + api.patch( + &secret_name, + &PatchParams::apply("studio-cp"), + &Patch::Apply(&secret), + ) + .await + .map_err(|e| anyhow::anyhow!("failed to create/apply k8s Secret '{secret_name}': {e}"))?; Ok(format!("k8s-secret://{secret_name}#OPENAB_ACP_AUTH_KEY")) } @@ -2222,9 +2292,13 @@ async fn ensure_namespace_k8s(context: Option<&str>, namespace: &str) -> anyhow: ..Default::default() }; let api: Api = Api::all(client); - api.patch(namespace, &PatchParams::apply("studio-cp"), &Patch::Apply(&ns)) - .await - .map_err(|e| anyhow::anyhow!("failed to create/apply k8s namespace '{namespace}': {e}"))?; + api.patch( + namespace, + &PatchParams::apply("studio-cp"), + &Patch::Apply(&ns), + ) + .await + .map_err(|e| anyhow::anyhow!("failed to create/apply k8s namespace '{namespace}': {e}"))?; Ok(()) } @@ -2266,9 +2340,15 @@ async fn provision_config_k8s_configmap( ..Default::default() }; let api: Api = Api::namespaced(client, namespace); - api.patch(&config_map_name, &PatchParams::apply("studio-cp"), &Patch::Apply(&config_map)) - .await - .map_err(|e| anyhow::anyhow!("failed to create/apply k8s ConfigMap '{config_map_name}': {e}"))?; + api.patch( + &config_map_name, + &PatchParams::apply("studio-cp"), + &Patch::Apply(&config_map), + ) + .await + .map_err(|e| { + anyhow::anyhow!("failed to create/apply k8s ConfigMap '{config_map_name}': {e}") + })?; Ok(format!("k8s-configmap://{config_map_name}#config.toml")) } @@ -2285,6 +2365,7 @@ async fn provision_config_k8s_configmap( /// applies (same "unset = use default" contract `default_security_group`- /// style AWS defaults don't have an equivalent of, since k8s already has one /// built in). +#[allow(clippy::too_many_arguments)] async fn build_default_k8s_manifest( context: Option<&str>, namespace: &str, @@ -2297,7 +2378,8 @@ async fn build_default_k8s_manifest( ) -> anyhow::Result { let mut secrets = std::collections::HashMap::new(); if acp_enabled { - let acp_auth_ref = provision_acp_auth_k8s_secret(context, namespace, name, acp_token).await?; + let acp_auth_ref = + provision_acp_auth_k8s_secret(context, namespace, name, acp_token).await?; secrets.insert("OPENAB_ACP_AUTH_KEY".to_string(), acp_auth_ref); } Ok(oabctl::manifest::OABServiceManifest { @@ -2350,6 +2432,7 @@ async fn build_default_k8s_manifest( /// /// `expected_principal` is `FleetBinding.expected_principal` — see /// [`k8s_service_account_from_principal`] for how it maps onto the manifest. +#[allow(clippy::too_many_arguments)] pub async fn provision_from_library_k8s( aws_config: &aws_config::SdkConfig, context: Option<&str>, @@ -2417,10 +2500,13 @@ pub async fn provision_from_library_k8s( .await? } }; - manifest.spec.bundle_from = Some(oabctl::studio_api::bundle_from_uri(&bucket, namespace, name)); + manifest.spec.bundle_from = Some(oabctl::studio_api::bundle_from_uri( + &bucket, namespace, name, + )); let report = - oabctl::studio_api::provision_k8s(aws_config, context, &manifest, &objects, Some(&bucket)).await?; + oabctl::studio_api::provision_k8s(aws_config, context, &manifest, &objects, Some(&bucket)) + .await?; Ok(ProvisionOutcome { image, @@ -2461,7 +2547,10 @@ pub async fn scale_k8s_deployment( size: i32, ) -> anyhow::Result<()> { use oabctl::ProvisionDriver; - oabctl::K8sDriver::from_context(context).await?.scale(namespace, name, size).await + oabctl::K8sDriver::from_context(context) + .await? + .scale(namespace, name, size) + .await } /// Delete a control-plane resource (e.g. an `OABService`). @@ -2697,7 +2786,8 @@ members = ["oab-prod-mira"] assert_eq!(b.get("mira").unwrap().cluster.as_deref(), Some("oab")); // membership routing assert_eq!( - b.fleet_for_service("oab-prod-mira").map(|f| f.name.as_str()), + b.fleet_for_service("oab-prod-mira") + .map(|f| f.name.as_str()), Some("mira") ); assert!(b.fleet_for_service("oab-prod-nope").is_none()); @@ -2732,7 +2822,10 @@ members = ["oab-prod-mira"] #[test] fn empty_config_and_no_fleet_key_parse_to_empty() { - assert!(toml::from_str::("").unwrap().fleets.is_empty()); + assert!(toml::from_str::("") + .unwrap() + .fleets + .is_empty()); assert!(toml::from_str::("# just a comment\n") .unwrap() .fleets @@ -2781,10 +2874,16 @@ namespace = "prod" "#; let b: FleetBindings = toml::from_str(doc).expect("parse"); assert_eq!( - b.get("dev").expect("dev fleet").expected_principal.as_deref(), + b.get("dev") + .expect("dev fleet") + .expected_principal + .as_deref(), Some("system:serviceaccount:dev:oab-agent") ); - assert_eq!(b.get("unset").expect("unset fleet").expected_principal, None); + assert_eq!( + b.get("unset").expect("unset fleet").expected_principal, + None + ); } #[test] @@ -2799,7 +2898,10 @@ namespace = "prod" #[test] fn normalize_bindings_runtime_field_inserts_ecs_for_named_fleets() { - let dir = std::env::temp_dir().join(format!("oab-normalize-runtime-named-{}", std::process::id())); + let dir = std::env::temp_dir().join(format!( + "oab-normalize-runtime-named-{}", + std::process::id() + )); let _ = std::fs::remove_dir_all(&dir); std::fs::create_dir_all(&dir).unwrap(); let path = dir.join("fleets.toml"); @@ -2829,7 +2931,10 @@ namespace = "prod" #[test] fn normalize_bindings_runtime_field_inserts_ecs_for_legacy_array_fleets() { - let dir = std::env::temp_dir().join(format!("oab-normalize-runtime-legacy-{}", std::process::id())); + let dir = std::env::temp_dir().join(format!( + "oab-normalize-runtime-legacy-{}", + std::process::id() + )); let _ = std::fs::remove_dir_all(&dir); std::fs::create_dir_all(&dir).unwrap(); let path = dir.join("fleets.toml"); @@ -2873,17 +2978,19 @@ namespace = "prod" ) .unwrap(); - let migrated = - migrate_legacy_k8s_bindings_at(&legacy_path, &fleets_path).expect("migrate"); + let migrated = migrate_legacy_k8s_bindings_at(&legacy_path, &fleets_path).expect("migrate"); assert!(migrated); assert!(!legacy_path.exists(), "legacy file should be renamed away"); assert!(dir.join("fleets-k8s.toml.migrated").exists()); let merged_text = std::fs::read_to_string(&fleets_path).unwrap(); - let merged: FleetBindings = toml::from_str(&merged_text).expect("merged fleets.toml still parses"); + let merged: FleetBindings = + toml::from_str(&merged_text).expect("merged fleets.toml still parses"); assert_eq!(merged.fleets.len(), 3); - let heph = merged.get("hephaestus").expect("migrated k8s fleet present"); + let heph = merged + .get("hephaestus") + .expect("migrated k8s fleet present"); assert_eq!(heph.runtime, FleetRuntime::K8s); assert_eq!(heph.context.as_deref(), Some("orbstack")); assert_eq!(heph.namespace.as_deref(), Some("openab-studio")); @@ -2916,7 +3023,10 @@ namespace = "prod" let text = "# my fleets\n\n[[fleet]]\nname = \"prod\"\nruntime = \"ecs\"\ncluster = \"oab\"\nprofile = \"orca-prod\"\n"; let parsed = save_bindings_text(&path, text).expect("save"); assert_eq!(parsed.fleets.len(), 1); - assert_eq!(parsed.for_cluster("oab").unwrap().profile.as_deref(), Some("orca-prod")); + assert_eq!( + parsed.for_cluster("oab").unwrap().profile.as_deref(), + Some("orca-prod") + ); // written verbatim — comment and layout preserved exactly assert_eq!(read_bindings_text(&path).unwrap(), text); let _ = std::fs::remove_dir_all(&dir); @@ -2968,7 +3078,10 @@ namespace = "prod" // studio_compose::Bundle::ZIP_FILENAME can't share a dependency edge // to enforce this with one constant (see provision_from_library) — // this is the cross-crate seam that catches drift instead. - assert_eq!(oabctl::studio_api::BUNDLE_ZIP_FILENAME, studio_compose::Bundle::ZIP_FILENAME); + assert_eq!( + oabctl::studio_api::BUNDLE_ZIP_FILENAME, + studio_compose::Bundle::ZIP_FILENAME + ); } #[test] @@ -3051,7 +3164,10 @@ aws_access_key_id = AKIA... fn k8s_service_account_from_principal_none_for_plain_username() { // A plain username (not a service account) or unset both mean "use // the namespace's default service account" — not an error. - assert_eq!(k8s_service_account_from_principal(Some("brett.chien")), None); + assert_eq!( + k8s_service_account_from_principal(Some("brett.chien")), + None + ); assert_eq!(k8s_service_account_from_principal(None), None); } @@ -3144,7 +3260,10 @@ aws_access_key_id = AKIA... &wizard_input(true, None), ) .unwrap_err(); - assert!(err.to_string().contains("predates /acp gateway support"), "{err}"); + assert!( + err.to_string().contains("predates /acp gateway support"), + "{err}" + ); } #[test] @@ -3168,7 +3287,10 @@ aws_access_key_id = AKIA... #[test] fn acp_compat_check_lets_a_custom_image_through_unverified() { - check_acp_image_compat("my-registry.example.com/custom:latest", &wizard_input(true, None)) - .expect("not an openab release tag — can't verify, don't block"); + check_acp_image_compat( + "my-registry.example.com/custom:latest", + &wizard_input(true, None), + ) + .expect("not an openab release tag — can't verify, don't block"); } } From 6fafbcdf4915b9d75727b8e278bbcb1ac8ba96d8 Mon Sep 17 00:00:00 2001 From: Devin Date: Wed, 30 Sep 2026 13:03:21 +0800 Subject: [PATCH 4/7] docs(adr): fold non-blocking review follow-ups into agent-lifecycle (#3) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fold the four wontfix-or-fold items from ADR-1's 8-axis review: - Stopping hard-loss edge: add the missing Stopping --> Stopped reclaim/hard-loss (no flush) edge to the state diagram and tighten principle 4 — a hard loss from *inside* Stopping skips 'state saved' and lands in the same absorbing Stopped, not a third outcome. - Section 9 lock-in/reversibility: new note enumerating what is expensive to reverse (the 6-state set, the identity_verified latch, per-instance identity + fencing epoch, the single-field dispatch predicate) vs cheap (projection rows, cause enums, attributes, deadline tuning). - R2 State.Paused: record that the AgentState surface labels are owned by the RuntimeDriver-contract ADR; in deployment-control-plane.md note the resolved naming consideration — the enum stays Paused, not Cordoned, because cordon is only one cause of ¬accepting_work (superseded fencing classifies the same way). - superseded: tighten to airtight instance-level phrasing — the attribute applies to the instance (agent is not paused); a replacement is a fresh lifecycle that classifies independently. --- docs/adr/agent-lifecycle.md | 39 ++++++++++++++++++++++------ docs/adr/deployment-control-plane.md | 4 +++ 2 files changed, 35 insertions(+), 8 deletions(-) diff --git a/docs/adr/agent-lifecycle.md b/docs/adr/agent-lifecycle.md index 642a934..454aeb7 100644 --- a/docs/adr/agent-lifecycle.md +++ b/docs/adr/agent-lifecycle.md @@ -67,7 +67,8 @@ stateDiagram-v2 Unhealthy --> Stopped : hard loss (OOM / crash / node death), no flush Running --> Stopping : stop / replace (desired=stopped) Paused --> Stopping : stop / replace - Stopping --> Stopped : state saved + Stopping --> Stopped : state saved (graceful) + Stopping --> Stopped : reclaim / hard loss (no flush) Running --> Stopped : reclaim (hard loss) Paused --> Stopped : reclaim (hard loss) Stopped --> [*] @@ -84,9 +85,12 @@ stateDiagram-v2 **Attributes, not states** (read alongside the state): `accepting_work` (Running vs Paused) — its authority is the **CP/director**, never the agent's -self-report; `superseded` / version-skew (a healthy instance whose desired -version has moved on) ⇒ `accepting_work=false`, so it classifies as **Paused** -and is never dispatched new work. *When and in what order* a superseded instance +self-report; `superseded` / version-skew is an **instance-level** attribute: +when a healthy, verified instance's desired version moves on, the CP sets +`accepting_work=false` *on that instance*, so it classifies as **Paused** and is +never dispatched new work while the attribute holds — the *agent* is not paused, +only that instance is; a replacement is a fresh instance with its own lifecycle +(typically `Starting`→`Running`). *When and in what order* a superseded instance is drained or replaced is a **fleet-level rollout** concern (e.g. make-before-break) — out of scope for this instance-level ADR; see the future rollout / RuntimeDriver ADR. Also: health `cause` = observed-bad vs @@ -111,8 +115,10 @@ unobservable; death `cause` enum; turn-level busy/idle. 4. **`reclaim` is two paths, not one.** A *planned* interruption (Spot/preempt notice — ECS ~120s SIGTERM, GKE ~30s + preStop) **compresses `Stopping`** into a short deadline. Only a *hard* loss (node death / SIGKILL / OOM) jumps - straight to `Stopped`. Durability never relies on the Stopping window — - **checkpoint while Running.** + straight to `Stopped` — **from any live state, `Stopping` included**: a hard + loss mid-flush skips `state saved` and lands in the same absorbing `Stopped`, + so "lost while Stopping" is not a third outcome. Durability never relies on + the Stopping window — **checkpoint while Running.** 5. **Runtime-independent.** Each driver projects native signals onto the 6 via the discriminators `(desiredStatus, accepting_work, health, identity_verified)`; the machine never changes per runtime. @@ -192,8 +198,25 @@ an ECS-only coincidence. and `restart:"no"` conditions above. - Detailed sub-states are **attributes** of the 6 (accepting_work, superseded, health-cause, death-cause enum, busy/idle), not new states. -- **Follow-ups:** a `RuntimeDriver` contract ADR (verbs apply / observe / scale - / cordon / …); an identity / lease / epoch spec ADR. +- **Lock-in / reversibility.** Expensive to reverse: the 6-state set itself — + the read-model, Studio, the dispatch predicate, and every driver's + conformance suite are all keyed on it, so adding or merging a state later is + a breaking change for every consumer; the `identity_verified` latch — each + driver must observe-and-remember "ever verified" per instance, which no + runtime exposes natively (k8s has no "ever Ready" field); per-instance + identity + the fencing epoch — unwinding them re-opens the trust model of + principles 1–2; and the single-field dispatch predicate — demoting `Paused` + to an attribute later silently restores the two-field predicate rejected in + §7, mis-dispatching every caller that forgets `&& accepting_work`. Cheap to + change: the §6 projection rows, the cause enums, new attributes, and + deadline/window tuning. +- **Follow-ups:** the `RuntimeDriver` contract ADR — + [ADR-2](./deployment-control-plane.md) (verbs apply / observe / scale / + cordon / …) — also owns the surface *labels* of the `AgentState` enum. `Paused` + is the state name while `cordon`/`resume` are the write-path verbs; whether + the shipped enum label stays `Paused` or aligns with the verbs is an open + naming consideration deferred to that ADR — the discriminator semantics + above are settled either way. Also: an identity / lease / epoch spec ADR. ## 10. More Information diff --git a/docs/adr/deployment-control-plane.md b/docs/adr/deployment-control-plane.md index 21284bd..268bb7f 100644 --- a/docs/adr/deployment-control-plane.md +++ b/docs/adr/deployment-control-plane.md @@ -123,6 +123,10 @@ One idempotent primitive: **`apply(Spec)`** — converge observed toward desired `accepting_work`, a Spec field). Without them the write model could not reach a state ADR-1 defines. `cordon` is *not* `stop`: it keeps the Instance alive and resumable. +- **Naming (ADR-1 review follow-up):** the enum label stays `Paused`, not + `Cordoned` — `cordon` is only *one* cause of `¬accepting_work` (a CP-set + `superseded` fence is another, and classifies the same way), so naming the + state after the verb would mis-scope it. - **dry-run / diff**: every write supports a preview that returns *what would change* without mutating — a safety valve, and important for agents (look before you leap). See §6 for `dry_run` as a first-class tool parameter. From 586928757e9c5371927c1f52132a90fbb8498702 Mon Sep 17 00:00:00 2001 From: Devin Date: Wed, 30 Sep 2026 13:04:58 +0800 Subject: [PATCH 5/7] test(agent-lifecycle): print stderr in fmt-guard failure output --- crates/agent-lifecycle/tests/rustfmt_clean.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/crates/agent-lifecycle/tests/rustfmt_clean.rs b/crates/agent-lifecycle/tests/rustfmt_clean.rs index cb298d1..bfceab6 100644 --- a/crates/agent-lifecycle/tests/rustfmt_clean.rs +++ b/crates/agent-lifecycle/tests/rustfmt_clean.rs @@ -13,7 +13,8 @@ fn workspace_is_rustfmt_clean() { .expect("spawn cargo fmt"); assert!( output.status.success(), - "cargo fmt --all -- --check failed:\n{}", - String::from_utf8_lossy(&output.stdout) + "cargo fmt --all -- --check failed:\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) ); } From 55ee935013d48e090b0a47e481738e5e73b4b674 Mon Sep 17 00:00:00 2001 From: Devin Date: Wed, 30 Sep 2026 13:04:58 +0800 Subject: [PATCH 6/7] docs(adr): point Paused-naming follow-up at ADR-2's recorded resolution --- docs/adr/agent-lifecycle.md | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/docs/adr/agent-lifecycle.md b/docs/adr/agent-lifecycle.md index 454aeb7..48b11a6 100644 --- a/docs/adr/agent-lifecycle.md +++ b/docs/adr/agent-lifecycle.md @@ -211,12 +211,14 @@ an ECS-only coincidence. change: the §6 projection rows, the cause enums, new attributes, and deadline/window tuning. - **Follow-ups:** the `RuntimeDriver` contract ADR — - [ADR-2](./deployment-control-plane.md) (verbs apply / observe / scale / - cordon / …) — also owns the surface *labels* of the `AgentState` enum. `Paused` - is the state name while `cordon`/`resume` are the write-path verbs; whether - the shipped enum label stays `Paused` or aligns with the verbs is an open - naming consideration deferred to that ADR — the discriminator semantics - above are settled either way. Also: an identity / lease / epoch spec ADR. + [ADR-2](./deployment-control-plane.md) (verbs apply / scale / cordon / + resume / stop / delete) — also owns the surface *labels* of the + `AgentState` enum: `Paused` is the state name while `cordon`/`resume` are + the write-path verbs. The naming alignment (state `Paused` vs verb + `cordon`) is recorded in ADR-2 §5 — the enum stays `Paused`, since cordon + is only one of several `¬accepting_work` causes; the discriminator + semantics above are settled either way. Also: an identity / lease / epoch + spec ADR. ## 10. More Information From cb913db3e6689354e5c5939d78d60171e9aca6ee Mon Sep 17 00:00:00 2001 From: Devin Date: Thu, 1 Oct 2026 06:58:30 +0800 Subject: [PATCH 7/7] docs(adr): point superseded rollout pointer at ADR-2, track issue #3 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review fixups on the follow-up fold-in: the "future RuntimeDriver ADR" reference predates ADR-2's acceptance — link it directly — and note issue #3 in the tracking header. Refs #3 Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- docs/adr/agent-lifecycle.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/docs/adr/agent-lifecycle.md b/docs/adr/agent-lifecycle.md index 48b11a6..c6d29ad 100644 --- a/docs/adr/agent-lifecycle.md +++ b/docs/adr/agent-lifecycle.md @@ -4,7 +4,7 @@ - **Date:** 2026-08-08 - **Author:** @brettchien - **Reviewers:** Mira (ECS), Jellyfish (control-plane), Falcon (MCP) — all LGTM -- **Tracking issues:** implementation openabdev/studio#2 +- **Tracking issues:** implementation openabdev/studio#2 · review follow-ups openabdev/studio#3 > **Y-statement.** In the context of running agents across heterogeneous > runtimes, facing the need for one glanceable, runtime-independent notion of @@ -92,8 +92,9 @@ never dispatched new work while the attribute holds — the *agent* is not pause only that instance is; a replacement is a fresh instance with its own lifecycle (typically `Starting`→`Running`). *When and in what order* a superseded instance is drained or replaced is a **fleet-level rollout** concern (e.g. -make-before-break) — out of scope for this instance-level ADR; see the future -rollout / RuntimeDriver ADR. Also: health `cause` = observed-bad vs +make-before-break) — out of scope for this instance-level ADR; the +RuntimeDriver contract is [ADR-2](./deployment-control-plane.md) and rollout +ordering remains future work. Also: health `cause` = observed-bad vs unobservable; death `cause` enum; turn-level busy/idle. ## 4. Principles