返回 CodeWhale
operate.rs
根目录 / crates / tui / src / operate.rs
1 //! Operate: always-on fleet operation matching landed CWC `OperateRecord`
2 //! (`Hmbown/cwc` `20de981`, PR #284).
3 //!
4 //! One schema for `cw · operate` and CWC `/operate`. Burn rate is optional
5 //! (`null` = unbounded). The lead plans before workers. Pace throttles or
6 //! widens; it never stops the operation. Supervisor over nested instances —
7 //! no second `Engine::run_turn`.
8
9 use std::fs;
10 use std::path::{Component, Path, PathBuf};
11 use std::process::{Command, Stdio};
12
13 use anyhow::{Context, Result, bail};
14 use chrono::Utc;
15 use serde::{Deserialize, Serialize};
16 use uuid::Uuid;
17
18 use crate::automation_manager::{AutomationManager, AutomationRecord, AutomationStatus};
19
20 pub const CWC_OPERATE_SCHEMA_VERSION: u32 = 1;
21 pub const OPERATE_MAX_WRITERS: usize = 3;
22 /// `hold` admits no new writers past this budget (the 8% band is met; hold
23 /// the current width instead of widening).
24 pub const OPERATE_HOLD_WRITERS: usize = 2;
25 /// `throttle` (observed more than 8% over target) cuts worker concurrency to
26 /// one writer. Pace throttles; it never stops the operation.
27 pub const OPERATE_THROTTLE_WRITERS: usize = 1;
28 pub const OPERATE_KEEPALIVE_ID: &str = "cw-operate";
29 /// Follow-up lead runs recur hourly; the first lead-plan step is kicked to
30 /// the next scheduler tick instead of waiting for the first recurrence.
31 pub const OPERATE_KEEPALIVE_RRULE: &str = "FREQ=HOURLY;INTERVAL=1";
32 pub const AUTO_MERGE_CHECKER_ENV: &str = "CODEWHALE_AUTO_MERGE_CHECKER";
33 pub const DIRECTION_PATH_ENV: &str = "CODEWHALE_DIRECTION_PATH";
34 pub const CHECK_AUTO_MERGE_SCRIPT: &str = "scripts/check-auto-merge.py";
35 pub const AUTO_MERGE_SCRIPT: &str = "scripts/auto_merge.py";
36 pub const AUTO_MERGE_PR_SCRIPT: &str = "scripts/auto-merge-pr.py";
37
38 const PACE_BAND: f64 = 0.08;
39
40 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
41 #[serde(rename_all = "camelCase")]
42 pub struct OperateBurnRate {
43 pub kind: String,
44 pub amount_usd_per_hour: f64,
45 }
46
47 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
48 #[serde(rename_all = "snake_case")]
49 pub enum OperateStatus {
50 Planning,
51 Running,
52 #[serde(rename = "idle_blocked")]
53 IdleBlocked,
54 Cancelled,
55 }
56
57 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
58 #[serde(rename_all = "snake_case")]
59 pub enum OperateIdleReason {
60 MissingCredentials,
61 AwaitingLeadPlan,
62 DirectionEmpty,
63 HumanGated,
64 }
65
66 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
67 #[serde(rename_all = "snake_case")]
68 pub enum OperatePace {
69 Unbounded,
70 Hold,
71 Throttle,
72 Widen,
73 }
74
75 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
76 #[serde(rename_all = "camelCase")]
77 pub struct OperateRosterMember {
78 pub id: String,
79 pub display_name: String,
80 pub role: String,
81 /// Saved selection for the lead; empty when no model has been assigned.
82 pub model: String,
83 pub state: String,
84 }
85
86 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
87 #[serde(rename_all = "camelCase")]
88 pub struct OperatePlanSlice {
89 pub id: String,
90 pub title: String,
91 pub owner_id: String,
92 pub depends_on: Vec<String>,
93 pub est_cost_usd: f64,
94 pub start_offset_sec: u32,
95 pub duration_sec: u32,
96 }
97
98 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
99 pub struct OperateLeadPlan {
100 pub slices: Vec<OperatePlanSlice>,
101 }
102
103 /// Landed CWC `OperateRecord` (packages/contracts/src/operate.js @ 20de981).
104 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
105 #[serde(rename_all = "camelCase")]
106 pub struct Operation {
107 pub id: String,
108 pub schema_version: u32,
109 pub direction: String,
110 pub burn_rate: Option<OperateBurnRate>,
111 pub lead_operator: OperateRosterMember,
112 pub roster: Vec<OperateRosterMember>,
113 pub lead_plan: Option<OperateLeadPlan>,
114 pub status: OperateStatus,
115 pub idle_blocked_reason: Option<OperateIdleReason>,
116 pub pace: OperatePace,
117 pub writers_in_flight: usize,
118 pub workers_admitted: bool,
119 pub spent_usd: f64,
120 pub observed_burn_usd_per_hour: Option<f64>,
121 pub credentials_present: bool,
122 pub human_gated: bool,
123 pub human_gate: String,
124 pub created_at: String,
125 pub updated_at: String,
126 pub last_keep_alive_at: String,
127 #[serde(default)]
128 pub cancelled_at: String,
129 }
130
131 impl Operation {
132 /// Update display metadata from the saved selection without assigning a
133 /// route to planned workers or claiming that an executor is running.
134 pub(crate) fn set_lead_model(&mut self, model: &str) {
135 self.lead_operator.model = model.to_string();
136 for member in &mut self.roster {
137 if member.id == self.lead_operator.id {
138 member.model = model.to_string();
139 }
140 }
141 }
142
143 #[must_use]
144 pub fn new(direction: impl Into<String>, burn_usd_per_hour: Option<f64>) -> Self {
145 let now = Utc::now().to_rfc3339();
146 let lead = OperateRosterMember {
147 id: "lead".to_string(),
148 display_name: "Lead operator".to_string(),
149 role: "lead".to_string(),
150 model: String::new(),
151 state: "planning".to_string(),
152 };
153 let mut op = Self {
154 id: format!("op_{}", Uuid::new_v4()),
155 schema_version: CWC_OPERATE_SCHEMA_VERSION,
156 direction: normalize_direction(direction.into()),
157 burn_rate: normalize_burn_rate(burn_usd_per_hour),
158 lead_operator: lead.clone(),
159 roster: vec![lead],
160 lead_plan: None,
161 status: OperateStatus::Planning,
162 idle_blocked_reason: None,
163 pace: OperatePace::Unbounded,
164 writers_in_flight: 0,
165 workers_admitted: false,
166 spent_usd: 0.0,
167 observed_burn_usd_per_hour: None,
168 credentials_present: false,
169 human_gated: false,
170 human_gate: String::new(),
171 created_at: now.clone(),
172 updated_at: now.clone(),
173 last_keep_alive_at: now,
174 cancelled_at: String::new(),
175 };
176 op.project();
177 op
178 }
179
180 pub fn plan_from_direction(&mut self) {
181 self.lead_plan = slices_from_direction(&self.direction);
182 sync_plan_owners(self);
183 self.project();
184 }
185
186 pub fn project(&mut self) {
187 if self.status == OperateStatus::Cancelled {
188 self.idle_blocked_reason = None;
189 self.workers_admitted = false;
190 self.writers_in_flight = 0;
191 for member in &mut self.roster {
192 member.state = "idle".to_string();
193 }
194 return;
195 }
196 let (status, reason) = derive_status(self);
197 self.status = status;
198 self.idle_blocked_reason = reason;
199 self.workers_admitted = workers_admitted(self);
200 self.pace = derive_pace(self);
201 // Pace is a dispatch budget, not a label: the roster and the
202 // `writersInFlight` count the keepalive lead actually dispatches at
203 // come from `worker_dispatch_budget`, so over-target burn reduces
204 // real concurrency instead of only renaming it.
205 let budget = worker_dispatch_budget(self);
206 live_roster(self, budget);
207 self.writers_in_flight = if self.workers_admitted {
208 self.roster
209 .iter()
210 .filter(|member| member.role == "worker" && member.state == "in_flight")
211 .count()
212 } else {
213 0
214 };
215 if let Some(lead) = self.roster.iter().find(|member| member.role == "lead") {
216 self.lead_operator = lead.clone();
217 }
218 }
219 }
220
221 fn normalize_direction(value: String) -> String {
222 value.trim().chars().take(4000).collect()
223 }
224
225 fn normalize_burn_rate(amount: Option<f64>) -> Option<OperateBurnRate> {
226 parse_burn_amount(amount).ok().flatten()
227 }
228
229 /// CWC `normalizeOperateBurnRate`: number, `$/hr` object, or null.
230 pub fn parse_burn_rate(value: Option<&serde_json::Value>) -> Result<Option<OperateBurnRate>> {
231 let Some(value) = value else {
232 return Ok(None);
233 };
234 if value.is_null()
235 || value == &serde_json::Value::Bool(false)
236 || value.as_str().is_some_and(str::is_empty)
237 {
238 return Ok(None);
239 }
240 if let Some("unbounded") = value.get("kind").and_then(|kind| kind.as_str()) {
241 return Ok(None);
242 }
243 let amount = if value.is_number() || value.is_string() {
244 json_number(value)
245 } else {
246 json_number(
247 value
248 .get("amountUsdPerHour")
249 .or_else(|| value.get("usdPerHour"))
250 .or_else(|| value.get("amount"))
251 .unwrap_or(&serde_json::Value::Null),
252 )
253 };
254 if amount.is_none() {
255 anyhow::bail!("Burn rate is optional. When set, it must be a positive $/hr.");
256 }
257 parse_burn_amount(amount)
258 }
259
260 fn json_number(value: &serde_json::Value) -> Option<f64> {
261 value
262 .as_f64()
263 .or_else(|| value.as_i64().map(|n| n as f64))
264 .or_else(|| value.as_str().and_then(|s| s.parse().ok()))
265 }
266
267 fn parse_burn_amount(amount: Option<f64>) -> Result<Option<OperateBurnRate>> {
268 let Some(amount) = amount else {
269 return Ok(None);
270 };
271 if !amount.is_finite() || amount <= 0.0 {
272 anyhow::bail!("Burn rate is optional. When set, it must be a positive $/hr.");
273 }
274 if amount > 10_000.0 {
275 anyhow::bail!("Burn rate must be 10000 $/hr or less.");
276 }
277 let rounded = (amount * 100.0).round() / 100.0;
278 if rounded <= 0.0 {
279 // A sub-cent rate rounds to a $0/hr target, which the pace governor
280 // would treat as unbounded — reject it instead of silently dropping
281 // the requested cap.
282 anyhow::bail!("Burn rate must be at least $0.01/hr.");
283 }
284 Ok(Some(OperateBurnRate {
285 kind: "usd_per_hour".to_string(),
286 amount_usd_per_hour: rounded,
287 }))
288 }
289
290 fn derive_status(op: &Operation) -> (OperateStatus, Option<OperateIdleReason>) {
291 if op.status == OperateStatus::Cancelled {
292 return (OperateStatus::Cancelled, None);
293 }
294 if !op.credentials_present {
295 return (
296 OperateStatus::IdleBlocked,
297 Some(OperateIdleReason::MissingCredentials),
298 );
299 }
300 if op.direction.is_empty() {
301 return (
302 OperateStatus::IdleBlocked,
303 Some(OperateIdleReason::DirectionEmpty),
304 );
305 }
306 if op.human_gated {
307 return (
308 OperateStatus::IdleBlocked,
309 Some(OperateIdleReason::HumanGated),
310 );
311 }
312 if op
313 .lead_plan
314 .as_ref()
315 .is_none_or(|plan| plan.slices.is_empty())
316 {
317 return (
318 OperateStatus::IdleBlocked,
319 Some(OperateIdleReason::AwaitingLeadPlan),
320 );
321 }
322 (OperateStatus::Running, None)
323 }
324
325 fn workers_admitted(op: &Operation) -> bool {
326 op.status != OperateStatus::Cancelled
327 && op.credentials_present
328 && !op.direction.is_empty()
329 && op
330 .lead_plan
331 .as_ref()
332 .is_some_and(|plan| !plan.slices.is_empty())
333 && !op.human_gated
334 }
335
336 fn derive_pace(op: &Operation) -> OperatePace {
337 let Some(rate) = &op.burn_rate else {
338 return OperatePace::Unbounded;
339 };
340 let Some(observed) = op.observed_burn_usd_per_hour.filter(|value| *value > 0.0) else {
341 return OperatePace::Widen;
342 };
343 let target = rate.amount_usd_per_hour;
344 if target <= 0.0 {
345 return OperatePace::Unbounded;
346 }
347 let delta = (observed - target) / target;
348 if delta > PACE_BAND {
349 OperatePace::Throttle
350 } else if delta < -PACE_BAND {
351 OperatePace::Widen
352 } else {
353 OperatePace::Hold
354 }
355 }
356
357 /// Pace-driven worker-dispatch budget within the 8% band semantics.
358 ///
359 /// `widen`/`unbounded` opens the full writer width, `hold` admits no new
360 /// writers past the hold width, and `throttle` cuts concurrency to one
361 /// writer. While workers are admitted the budget is never zero — pace
362 /// throttles spend, it never stops the operation.
363 #[must_use]
364 pub fn worker_dispatch_budget(op: &Operation) -> usize {
365 if !op.workers_admitted {
366 return 0;
367 }
368 match op.pace {
369 OperatePace::Unbounded | OperatePace::Widen => OPERATE_MAX_WRITERS,
370 OperatePace::Hold => OPERATE_HOLD_WRITERS,
371 OperatePace::Throttle => OPERATE_THROTTLE_WRITERS,
372 }
373 }
374
375 fn live_roster(op: &mut Operation, budget: usize) {
376 let admitted = op.workers_admitted;
377 let cancelled = op.status == OperateStatus::Cancelled;
378 let has_plan = op
379 .lead_plan
380 .as_ref()
381 .is_some_and(|plan| !plan.slices.is_empty());
382 // Only the first `budget` workers in roster (plan) order dispatch; the
383 // rest stay idle until a slot frees or pace widens.
384 let mut worker_slot = 0usize;
385 for member in &mut op.roster {
386 if cancelled {
387 member.state = "idle".to_string();
388 } else if !op.credentials_present {
389 member.state = "blocked".to_string();
390 } else if member.role == "lead" && !has_plan {
391 member.state = "planning".to_string();
392 } else if !admitted {
393 member.state = if member.role == "lead" {
394 "planning".to_string()
395 } else {
396 "idle".to_string()
397 };
398 } else if member.role == "worker" {
399 worker_slot += 1;
400 member.state = if worker_slot <= budget {
401 "in_flight".to_string()
402 } else {
403 "idle".to_string()
404 };
405 } else {
406 member.state = "planning".to_string();
407 }
408 }
409 }
410
411 fn slices_from_direction(direction: &str) -> Option<OperateLeadPlan> {
412 let items = direction_items(direction);
413 if items.is_empty() {
414 return None;
415 }
416 let mut cursor = 0u32;
417 let slices = items
418 .into_iter()
419 .enumerate()
420 .map(|(index, title)| {
421 let duration_sec = 1800;
422 let slice = OperatePlanSlice {
423 id: format!("slice-{}", index + 1),
424 title,
425 owner_id: if index == 0 {
426 "lead".to_string()
427 } else {
428 format!("worker-{index}")
429 },
430 depends_on: if index == 0 {
431 Vec::new()
432 } else {
433 vec![format!("slice-{index}")]
434 },
435 est_cost_usd: 0.25,
436 start_offset_sec: cursor,
437 duration_sec,
438 };
439 cursor = cursor.saturating_add(duration_sec);
440 slice
441 })
442 .collect();
443 Some(OperateLeadPlan { slices })
444 }
445
446 fn direction_items(direction: &str) -> Vec<String> {
447 let mut items = Vec::new();
448 for raw in direction.lines() {
449 let line = raw.trim();
450 if line.is_empty() {
451 continue;
452 }
453 let item = line
454 .trim_start_matches(|c: char| {
455 c.is_ascii_digit() || c == '.' || c == ')' || c == '-' || c == '*' || c == '#'
456 })
457 .trim();
458 if !item.is_empty() {
459 items.push(item.to_string());
460 }
461 }
462 if items.is_empty() && !direction.trim().is_empty() {
463 items.push(direction.trim().to_string());
464 }
465 items
466 }
467
468 #[must_use]
469 pub fn render_plan_board(op: &Operation) -> String {
470 render_plan_board_locale(op, codewhale_localization::Locale::En)
471 }
472
473 /// Plan board with locale-aware chrome. Contract tokens (status / pace enum
474 /// values, slice ids, owner ids) stay verbatim; the surrounding prose comes
475 /// from the TUI locale packs.
476 #[must_use]
477 pub fn render_plan_board_locale(op: &Operation, locale: codewhale_localization::Locale) -> String {
478 use codewhale_localization::{MessageId, tr};
479 let tr_line = |id: MessageId| tr(locale, id).into_owned();
480
481 let mut out = String::new();
482 out.push_str(
483 &tr_line(MessageId::OperateBoardHeader)
484 .replace("{id}", &op.id)
485 .replace("{status}", &status_label(op))
486 .replace("{pace}", pace_label(op.pace))
487 .replace("{writers}", &op.writers_in_flight.to_string()),
488 );
489 out.push('\n');
490 match &op.burn_rate {
491 Some(rate) => out.push_str(
492 &tr_line(MessageId::OperateBoardBurnObserved)
493 .replace(
494 "{actual}",
495 &format!("{}", op.observed_burn_usd_per_hour.unwrap_or(0.0)),
496 )
497 .replace("{target}", &format!("{}", rate.amount_usd_per_hour)),
498 ),
499 None => out.push_str(&tr_line(MessageId::OperateBoardBurnNoCap)),
500 }
501 out.push('\n');
502 if op.direction.is_empty() {
503 out.push_str(&tr_line(MessageId::OperateBoardDirectionEmpty));
504 } else {
505 out.push_str(
506 &tr_line(MessageId::OperateBoardDirectionLine)
507 .replace("{line}", op.direction.lines().next().unwrap_or("")),
508 );
509 }
510 out.push('\n');
511 let Some(plan) = &op.lead_plan else {
512 out.push_str(&tr_line(MessageId::OperateBoardPlanMissing));
513 out.push('\n');
514 return out;
515 };
516 out.push_str(&tr_line(MessageId::OperateBoardPlanHeader));
517 out.push('\n');
518 for slice in &plan.slices {
519 out.push_str(&format!(
520 " {:<9} {:<8} {:>5} {:>5} {:>6.2} {:<7} {}\n",
521 slice.id,
522 slice.owner_id,
523 slice.start_offset_sec,
524 slice.duration_sec,
525 slice.est_cost_usd,
526 if slice.depends_on.is_empty() {
527 "-".to_string()
528 } else {
529 slice.depends_on.join(",")
530 },
531 slice.title
532 ));
533 }
534 out.push_str(&render_timeline(&plan.slices, locale));
535 out
536 }
537
538 fn render_timeline(slices: &[OperatePlanSlice], locale: codewhale_localization::Locale) -> String {
539 use codewhale_localization::{MessageId, tr};
540 let max_end = slices
541 .iter()
542 .map(|slice| slice.start_offset_sec.saturating_add(slice.duration_sec))
543 .max()
544 .unwrap_or(0)
545 .max(1);
546 let width = 24u32;
547 let mut out = format!("{}\n", tr(locale, MessageId::OperateBoardGantt));
548 for slice in slices {
549 let start = (slice.start_offset_sec.saturating_mul(width)) / max_end;
550 let end = (slice
551 .start_offset_sec
552 .saturating_add(slice.duration_sec)
553 .saturating_mul(width))
554 / max_end;
555 let end = end.max(start.saturating_add(1)).min(width);
556 let mut bar = vec!['.'; width as usize];
557 for idx in start..end {
558 if let Some(cell) = bar.get_mut(idx as usize) {
559 *cell = '#';
560 }
561 }
562 out.push_str(&format!(
563 " {:<9} {}\n",
564 slice.id,
565 bar.into_iter().collect::<String>()
566 ));
567 }
568 out
569 }
570
571 fn status_label(op: &Operation) -> String {
572 match (op.status, op.idle_blocked_reason) {
573 (OperateStatus::Cancelled, _) => "cancelled".to_string(),
574 (OperateStatus::Running, _) => "running".to_string(),
575 (OperateStatus::Planning, _) => "planning".to_string(),
576 (_, Some(OperateIdleReason::DirectionEmpty)) => "idle_blocked: direction_empty".to_string(),
577 (_, Some(OperateIdleReason::AwaitingLeadPlan)) => {
578 "idle_blocked: awaiting_lead_plan".to_string()
579 }
580 (_, Some(OperateIdleReason::MissingCredentials)) => {
581 "idle_blocked: missing_credentials".to_string()
582 }
583 (_, Some(OperateIdleReason::HumanGated)) => {
584 format!("idle_blocked: human_gated {}", op.human_gate)
585 }
586 (OperateStatus::IdleBlocked, None) => "idle_blocked".to_string(),
587 }
588 }
589
590 fn pace_label(pace: OperatePace) -> &'static str {
591 match pace {
592 OperatePace::Unbounded => "unbounded",
593 OperatePace::Hold => "hold",
594 OperatePace::Throttle => "throttle",
595 OperatePace::Widen => "widen",
596 }
597 }
598
599 /// Honor an explicit operator-provided path (env override) only when it is a
600 /// non-empty value without NUL bytes or `..` traversal segments that actually
601 /// names one regular file. Keeps env-provided values out of raw path
602 /// expressions (CodeQL "uncontrolled data in path" class).
603 fn explicit_file_path(raw: &str) -> Option<PathBuf> {
604 let trimmed = raw.trim();
605 if trimmed.is_empty() || trimmed.contains('\0') {
606 return None;
607 }
608 let path = PathBuf::from(trimmed);
609 if path
610 .components()
611 .any(|component| matches!(component, Component::ParentDir))
612 {
613 return None;
614 }
615 path.is_file().then_some(path)
616 }
617
618 /// Same hardening for an explicit directory (env override): no NUL bytes, no
619 /// `..` traversal segments. Existence is probed by the caller.
620 fn explicit_dir_path(raw: &str) -> Option<PathBuf> {
621 let trimmed = raw.trim();
622 if trimmed.is_empty() || trimmed.contains('\0') {
623 return None;
624 }
625 let path = PathBuf::from(trimmed);
626 if path
627 .components()
628 .any(|component| matches!(component, Component::ParentDir))
629 {
630 return None;
631 }
632 Some(path)
633 }
634
635 #[must_use]
636 pub fn discover_direction_path(workspace: &Path) -> Option<PathBuf> {
637 if let Ok(explicit) = std::env::var(DIRECTION_PATH_ENV)
638 && let Some(path) = explicit_file_path(&explicit)
639 {
640 return Some(path);
641 }
642 let local = workspace.join("DIRECTION.md");
643 if local.is_file() {
644 return Some(local);
645 }
646 materialize_ops_origin_main()
647 .ok()
648 .map(|root| root.join("DIRECTION.md"))
649 .filter(|path| path.is_file())
650 }
651
652 pub fn read_direction(workspace: &Path) -> Result<String> {
653 match discover_direction_path(workspace) {
654 Some(path) => fs::read_to_string(&path)
655 .with_context(|| format!("Failed to read direction {}", path.display())),
656 None => Ok(String::new()),
657 }
658 }
659
660 /// Resolve the saved keepalive pin, or the caller's effective session route
661 /// for a fresh/unpinned record. Auto stays a policy until Runtime admits a turn;
662 /// resolving its inventory here could run a paid classifier.
663 fn keepalive_route(
664 config: &crate::config::Config,
665 current: Option<&AutomationRecord>,
666 selection: Option<(&crate::config::ProviderIdentity, &str)>,
667 ) -> Result<(crate::config::ProviderIdentity, String, bool)> {
668 let (identity, model) = if let Some(record) = current
669 .filter(|record| record.model_provider.is_some() || record.model_provider_id.is_some())
670 {
671 let identity = config
672 .resolve_persisted_provider_identity(
673 record.model_provider.as_deref(),
674 record.model_provider_id.as_deref(),
675 )
676 .map_err(anyhow::Error::msg)?;
677 let model = record
678 .model
679 .as_deref()
680 .map(str::trim)
681 .filter(|model| !model.is_empty())
682 .context("Pinned Operate keepalive has no model; repair its saved route")?;
683 (identity, model.to_string())
684 } else if let Some((identity, model)) = selection {
685 let identity = config
686 .resolve_persisted_provider_identity(
687 Some(identity.provider.as_str()),
688 identity.persisted_id(),
689 )
690 .map_err(anyhow::Error::msg)?;
691 (identity, model.to_string())
692 } else {
693 (
694 config
695 .active_provider_identity()
696 .map_err(anyhow::Error::msg)?,
697 config.default_model(),
698 )
699 };
700 if model.trim().eq_ignore_ascii_case("auto") {
701 let mut scoped = config.clone();
702 scoped
703 .scope_to_provider_identity(&identity)
704 .map_err(anyhow::Error::msg)?;
705 let credentials = crate::config::has_api_key_for(&scoped, &identity);
706 return Ok((identity, "auto".to_string(), credentials));
707 }
708 let route =
709 crate::route_runtime::resolve_runtime_route_for_identity(config, &identity, Some(&model))
710 .map_err(anyhow::Error::msg)?;
711 let credentials = crate::config::has_api_key_for(&route.config, &route.identity);
712 Ok((route.identity, route.model, credentials))
713 }
714
715 /// Inspect the same saved route used by the scheduler. Credential presence is
716 /// a local readiness observation, not proof of provider execution.
717 pub(crate) fn keepalive_readiness(
718 manager: &AutomationManager,
719 config: &crate::config::Config,
720 selection: Option<(&crate::config::ProviderIdentity, &str)>,
721 ) -> Result<(String, bool)> {
722 let mut readiness = (String::new(), false);
723 manager.edit_automation(OPERATE_KEEPALIVE_ID, |current| {
724 if current
725 .as_ref()
726 .and_then(|record| record.execution_scope.as_deref())
727 .is_some_and(|scope| Some(scope) != manager.execution_scope())
728 {
729 bail!("Operate keepalive belongs to another Runtime execution scope");
730 }
731 let (_, model, credentials) = keepalive_route(config, current.as_ref(), selection)?;
732 readiness = (model, credentials);
733 Ok(None)
734 })?;
735 Ok(readiness)
736 }
737
738 #[derive(Debug, Clone, PartialEq, Eq)]
739 pub struct AutoMergeRequest<'a> {
740 pub pr: &'a str,
741 pub role: &'a str,
742 pub repo: &'a str,
743 }
744
745 #[derive(Debug, Clone, PartialEq, Eq)]
746 pub enum AutoMergeDecision {
747 Allow,
748 Deny { reason: String },
749 }
750
751 #[must_use]
752 pub fn check_auto_merge_args(repo: &str, pr: &str, agent: &str) -> Vec<String> {
753 vec![
754 CHECK_AUTO_MERGE_SCRIPT.to_string(),
755 "--repo".to_string(),
756 repo.to_string(),
757 "--pr".to_string(),
758 pr.to_string(),
759 "--agent".to_string(),
760 agent.to_string(),
761 ]
762 }
763
764 #[must_use]
765 pub fn auto_merge_pr_args(repo: &str, pr: &str, agent: &str) -> Vec<String> {
766 vec![
767 AUTO_MERGE_PR_SCRIPT.to_string(),
768 "--repo".to_string(),
769 repo.to_string(),
770 "--pr".to_string(),
771 pr.to_string(),
772 "--agent".to_string(),
773 agent.to_string(),
774 ]
775 }
776
777 /// Strictly validate an auto-merge request before any of it reaches argv.
778 ///
779 /// The checker is spawned without a shell, but these values still become
780 /// arguments to `python3` and then to `gh`, so they are held to the shapes
781 /// GitHub itself allows: `repo` is `owner/name`, `pr` is a positive decimal
782 /// number, and `agent` is a short role token. No value may start with `-`.
783 pub fn validate_auto_merge_request(request: &AutoMergeRequest<'_>) -> Result<(), String> {
784 fn is_owner(owner: &str) -> bool {
785 (1..=39).contains(&owner.len())
786 && !owner.starts_with('-')
787 && owner
788 .bytes()
789 .all(|b| b.is_ascii_alphanumeric() || b == b'-')
790 }
791 fn is_repo_name(name: &str) -> bool {
792 (1..=100).contains(&name.len())
793 && name != "."
794 && name != ".."
795 && !name.starts_with('-')
796 && name
797 .bytes()
798 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.'))
799 }
800 let repo_ok = request
801 .repo
802 .split_once('/')
803 .is_some_and(|(owner, name)| is_owner(owner) && is_repo_name(name));
804 if !repo_ok {
805 return Err("repo must be `owner/name` using GitHub name characters".to_string());
806 }
807 let pr_ok = (1..=10).contains(&request.pr.len())
808 && request.pr.bytes().all(|b| b.is_ascii_digit())
809 && request.pr.parse::<u64>().is_ok_and(|n| n > 0);
810 if !pr_ok {
811 return Err("pr must be a positive pull request number".to_string());
812 }
813 let agent_ok = (1..=64).contains(&request.role.len())
814 && !request.role.starts_with('-')
815 && request
816 .role
817 .bytes()
818 .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_'));
819 if !agent_ok {
820 return Err("agent must be 1-64 letters, digits, `-` or `_`".to_string());
821 }
822 Ok(())
823 }
824
825 pub fn evaluate_auto_merge(
826 request: AutoMergeRequest<'_>,
827 checker: Option<&Path>,
828 ) -> AutoMergeDecision {
829 if let Err(reason) = validate_auto_merge_request(&request) {
830 return AutoMergeDecision::Deny { reason };
831 }
832 let Some(checker) = checker else {
833 return AutoMergeDecision::Deny {
834 reason: "auto-merge checker missing; fail-closed".to_string(),
835 };
836 };
837 if !checker.exists() {
838 return AutoMergeDecision::Deny {
839 reason: "auto-merge checker missing; fail-closed".to_string(),
840 };
841 }
842 match Command::new("python3")
843 .arg(checker)
844 .arg("--repo")
845 .arg(request.repo)
846 .arg("--pr")
847 .arg(request.pr)
848 .arg("--agent")
849 .arg(request.role)
850 .stdin(Stdio::null())
851 .stdout(Stdio::null())
852 .stderr(Stdio::null())
853 .status()
854 {
855 Ok(status) if status.success() => AutoMergeDecision::Allow,
856 Ok(_) => AutoMergeDecision::Deny {
857 reason: "auto-merge checker refused".to_string(),
858 },
859 Err(error) => AutoMergeDecision::Deny {
860 reason: format!("auto-merge checker failed to start: {error}"),
861 },
862 }
863 }
864
865 #[must_use]
866 pub fn discover_auto_merge_checker(_workspace: &Path) -> Option<PathBuf> {
867 if let Ok(explicit) = std::env::var(AUTO_MERGE_CHECKER_ENV)
868 && let Some(path) = explicit_file_path(&explicit)
869 {
870 return Some(path);
871 }
872 materialize_ops_origin_main()
873 .ok()
874 .map(|root| root.join(CHECK_AUTO_MERGE_SCRIPT))
875 .filter(|path| path.is_file())
876 }
877
878 fn ops_git_candidates() -> Vec<PathBuf> {
879 let mut out = Vec::new();
880 for key in ["CODEWHALE_OPS_GIT", "CODEWHALE_OPS_ROOT"] {
881 if let Ok(path) = std::env::var(key)
882 && let Some(path) = explicit_dir_path(&path)
883 {
884 out.push(path);
885 }
886 }
887 out
888 }
889
890 fn git_origin_main_sha(repo: &Path) -> Option<String> {
891 let output = Command::new("git")
892 .arg("-C")
893 .arg(repo)
894 .args(["rev-parse", "origin/main"])
895 .stdin(Stdio::null())
896 .output()
897 .ok()?;
898 if !output.status.success() {
899 return None;
900 }
901 let sha = String::from_utf8_lossy(&output.stdout).trim().to_string();
902 // The sha becomes a path segment below; only plain hex of a plausible
903 // length may flow into it.
904 if !(7..=64).contains(&sha.len()) || !sha.chars().all(|c| c.is_ascii_hexdigit()) {
905 return None;
906 }
907 Some(sha)
908 }
909
910 pub fn materialize_ops_origin_main() -> Result<PathBuf> {
911 let repo = ops_git_candidates()
912 .into_iter()
913 .find(|path| git_origin_main_sha(path).is_some())
914 .context("no codewhale-ops git checkout with origin/main")?;
915 let sha = git_origin_main_sha(&repo).context("origin/main sha")?;
916 let dest = default_operate_dir()
917 .parent()
918 .unwrap_or(Path::new("."))
919 .join("ops-main")
920 .join(&sha[..12.min(sha.len())]);
921 let marker = dest.join(CHECK_AUTO_MERGE_SCRIPT);
922 if marker.is_file()
923 && dest.join("DIRECTION.md").is_file()
924 && dest.join(AUTO_MERGE_SCRIPT).is_file()
925 {
926 return Ok(dest);
927 }
928 fs::create_dir_all(&dest).with_context(|| format!("Failed to create {}", dest.display()))?;
929 let archive = Command::new("git")
930 .arg("-C")
931 .arg(&repo)
932 .args([
933 "archive",
934 "origin/main",
935 "--",
936 "DIRECTION.md",
937 CHECK_AUTO_MERGE_SCRIPT,
938 AUTO_MERGE_SCRIPT,
939 AUTO_MERGE_PR_SCRIPT,
940 "agent-workstreams/AUTO_MERGE.toml",
941 ])
942 .stdin(Stdio::null())
943 .output()
944 .context("git archive origin/main")?;
945 if !archive.status.success() {
946 anyhow::bail!(
947 "git archive origin/main failed: {}",
948 String::from_utf8_lossy(&archive.stderr).trim()
949 );
950 }
951 let mut child = Command::new("tar")
952 .arg("-x")
953 .arg("-C")
954 .arg(&dest)
955 .stdin(Stdio::piped())
956 .stdout(Stdio::null())
957 .stderr(Stdio::piped())
958 .spawn()
959 .context("tar extract ops origin/main")?;
960 if let Some(mut stdin) = child.stdin.take() {
961 use std::io::Write;
962 stdin.write_all(&archive.stdout)?;
963 }
964 let status = child.wait()?;
965 if !status.success() {
966 anyhow::bail!("failed to extract ops origin/main archive");
967 }
968 if !marker.is_file() {
969 anyhow::bail!("ops origin/main archive missing {CHECK_AUTO_MERGE_SCRIPT}");
970 }
971 Ok(dest)
972 }
973
974 pub fn default_operate_dir() -> PathBuf {
975 if let Ok(path) = std::env::var("CODEWHALE_OPERATE_DIR") {
976 let trimmed = path.trim();
977 if !trimmed.is_empty() {
978 return PathBuf::from(trimmed);
979 }
980 }
981 crate::automation_manager::default_automations_dir()
982 .parent()
983 .map(|parent| parent.join("operate"))
984 .unwrap_or_else(|| PathBuf::from("operate"))
985 }
986
987 pub struct OperationStore {
988 path: PathBuf,
989 lock_path: PathBuf,
990 }
991
992 impl OperationStore {
993 pub fn open(dir: impl Into<PathBuf>) -> Result<Self> {
994 let dir = dir.into();
995 fs::create_dir_all(&dir).with_context(|| format!("Failed to create {}", dir.display()))?;
996 let path = dir.join("current.json");
997 let lock_path = dir.join("current.json.lock");
998 Ok(Self { path, lock_path })
999 }
1000
1001 pub fn load(&self) -> Result<Option<Operation>> {
1002 // Match the skill-state discipline: pure readers take the shared
1003 // cross-process read lock only when one exists, so a read never
1004 // fabricates a lock file.
1005 if self.lock_path.exists() {
1006 let file = fs::File::open(&self.lock_path)
1007 .with_context(|| format!("Failed to open {}", self.lock_path.display()))?;
1008 let lock = fd_lock::RwLock::new(file);
1009 let _guard = lock
1010 .read()
1011 .with_context(|| format!("read-lock {}", self.path.display()))?;
1012 return self.load_unlocked();
1013 }
1014 self.load_unlocked()
1015 }
1016
1017 fn load_unlocked(&self) -> Result<Option<Operation>> {
1018 if !self.path.exists() {
1019 return Ok(None);
1020 }
1021 let raw = fs::read_to_string(&self.path)
1022 .with_context(|| format!("Failed to read {}", self.path.display()))?;
1023 let op: Operation = serde_json::from_str(&raw)
1024 .with_context(|| format!("Failed to parse {}", self.path.display()))?;
1025 Ok(Some(op))
1026 }
1027
1028 /// Save under the cross-process writer lock with an atomic temp+rename
1029 /// write, so concurrent Codewhale processes never interleave partial
1030 /// records.
1031 pub fn save(&self, op: &Operation) -> Result<()> {
1032 let file = self.open_lock_file()?;
1033 let mut lock = fd_lock::RwLock::new(file);
1034 let _guard = lock
1035 .write()
1036 .with_context(|| format!("write-lock {}", self.path.display()))?;
1037 self.save_unlocked(op)
1038 }
1039
1040 /// Read-merge-write under the cross-process writer lock: the latest
1041 /// on-disk record is reloaded *inside* the lock before `edit` runs, so a
1042 /// concurrent PATCH/keepalive/plan save can no longer be silently lost by
1043 /// a stale read. Returns `None` when no operation is recorded yet.
1044 pub fn mutate(
1045 &self,
1046 edit: impl FnOnce(&mut Operation) -> Result<()>,
1047 ) -> Result<Option<Operation>> {
1048 let file = self.open_lock_file()?;
1049 let mut lock = fd_lock::RwLock::new(file);
1050 let _guard = lock
1051 .write()
1052 .with_context(|| format!("write-lock {}", self.path.display()))?;
1053 let Some(mut op) = self.load_unlocked()? else {
1054 return Ok(None);
1055 };
1056 edit(&mut op)?;
1057 self.save_unlocked(&op)?;
1058 Ok(Some(op))
1059 }
1060
1061 fn open_lock_file(&self) -> Result<fs::File> {
1062 crate::session_manager::open_private_lock_file(&self.lock_path)
1063 .with_context(|| format!("Failed to open {}", self.lock_path.display()))
1064 }
1065
1066 fn save_unlocked(&self, op: &Operation) -> Result<()> {
1067 codewhale_config::persistence::atomic_write_json(&self.path, op)
1068 .with_context(|| format!("Failed to write {}", self.path.display()))
1069 }
1070 }
1071
1072 pub fn start_operation(
1073 store: &OperationStore,
1074 workspace: &Path,
1075 direction: Option<String>,
1076 burn_usd_per_hour: Option<f64>,
1077 credentials_present: bool,
1078 lead_model: &str,
1079 ) -> Result<Operation> {
1080 let mut direction = match direction {
1081 Some(text) if !text.trim().is_empty() => text,
1082 _ => read_direction(workspace)?,
1083 };
1084 if direction.trim().is_empty()
1085 && let Some(existing) = store.load()?
1086 {
1087 direction = existing.direction;
1088 }
1089 let mut op = Operation::new(direction, burn_usd_per_hour);
1090 op.set_lead_model(lead_model);
1091 op.credentials_present = credentials_present;
1092 op.project();
1093 store.save(&op)?;
1094 Ok(op)
1095 }
1096
1097 /// Re-entering Operate attaches to the recorded operation — same id, spend,
1098 /// lead plan, and roster — instead of minting a fresh record that silently
1099 /// resets progress. A cancelled record (or an empty store) starts a new
1100 /// operation. An explicitly provided direction is applied through the normal
1101 /// patch rules (which invalidate a superseded lead plan).
1102 pub fn attach_or_start_operation(
1103 store: &OperationStore,
1104 workspace: &Path,
1105 direction: Option<String>,
1106 burn_usd_per_hour: Option<f64>,
1107 credentials_present: bool,
1108 lead_model: &str,
1109 ) -> Result<Operation> {
1110 if store
1111 .load()?
1112 .is_some_and(|op| op.status != OperateStatus::Cancelled)
1113 {
1114 let attached = store.mutate(|op| {
1115 if let Some(text) = direction.as_ref().filter(|text| !text.trim().is_empty()) {
1116 apply_operate_patch(op, &serde_json::json!({ "direction": text }))?;
1117 }
1118 op.set_lead_model(lead_model);
1119 op.credentials_present = credentials_present;
1120 op.project();
1121 Ok(())
1122 })?;
1123 if let Some(op) = attached {
1124 return Ok(op);
1125 }
1126 }
1127 start_operation(
1128 store,
1129 workspace,
1130 direction,
1131 burn_usd_per_hour,
1132 credentials_present,
1133 lead_model,
1134 )
1135 }
1136
1137 pub fn apply_operate_patch(op: &mut Operation, patch: &serde_json::Value) -> Result<()> {
1138 if op.status == OperateStatus::Cancelled {
1139 anyhow::bail!("A cancelled Operation cannot be edited.");
1140 }
1141 if let Some(direction) = patch.get("direction") {
1142 let next = normalize_direction(direction.as_str().map(str::to_string).unwrap_or_default());
1143 if patch.get("leadPlan").is_none() && next != op.direction {
1144 // A changed direction supersedes the recorded lead plan: workers
1145 // must stop executing slices derived from the old direction. The
1146 // operation idles `awaiting_lead_plan` until the lead re-plans
1147 // (same patch may instead carry an explicit replacement plan).
1148 op.lead_plan = None;
1149 }
1150 op.direction = next;
1151 }
1152 if patch.get("burnRate").is_some() {
1153 op.burn_rate = parse_burn_rate(patch.get("burnRate"))?;
1154 }
1155 if let Some(plan) = patch.get("leadPlan") {
1156 op.lead_plan = if plan.is_null() {
1157 None
1158 } else {
1159 Some(serde_json::from_value(plan.clone()).context("leadPlan is invalid")?)
1160 };
1161 // Plan owners are the worker roster: PUT /v1/operate/plan and a
1162 // PATCHed plan must admit their owners or workers never dispatch.
1163 sync_plan_owners(op);
1164 }
1165 if let Some(flag) = patch.get("humanGated").and_then(serde_json::Value::as_bool) {
1166 op.human_gated = flag;
1167 }
1168 if let Some(gate) = patch.get("humanGate").and_then(serde_json::Value::as_str) {
1169 op.human_gate = gate.trim().chars().take(160).collect();
1170 if human_gate_for(&op.human_gate) {
1171 op.human_gated = true;
1172 }
1173 }
1174 if let Some(flag) = patch
1175 .get("credentialsPresent")
1176 .and_then(serde_json::Value::as_bool)
1177 {
1178 op.credentials_present = flag;
1179 }
1180 op.updated_at = Utc::now().to_rfc3339();
1181 op.project();
1182 Ok(())
1183 }
1184
1185 /// Add every non-lead plan owner to the roster (idempotent) so admitted
1186 /// slices have a worker to dispatch to.
1187 fn sync_plan_owners(op: &mut Operation) {
1188 if let Some(plan) = &op.lead_plan {
1189 for slice in &plan.slices {
1190 if slice.owner_id != "lead"
1191 && !slice.owner_id.is_empty()
1192 && !op.roster.iter().any(|member| member.id == slice.owner_id)
1193 {
1194 op.roster.push(OperateRosterMember {
1195 id: slice.owner_id.clone(),
1196 display_name: slice.owner_id.clone(),
1197 role: "worker".to_string(),
1198 model: String::new(),
1199 state: "idle".to_string(),
1200 });
1201 }
1202 }
1203 }
1204 }
1205
1206 pub fn cancel_operation(store: &OperationStore) -> Result<Option<Operation>> {
1207 store.mutate(|op| {
1208 let now = Utc::now().to_rfc3339();
1209 op.status = OperateStatus::Cancelled;
1210 op.cancelled_at = now.clone();
1211 op.updated_at = now.clone();
1212 op.last_keep_alive_at = now;
1213 op.project();
1214 Ok(())
1215 })
1216 }
1217
1218 pub fn keep_alive_observation(
1219 op: &mut Operation,
1220 observed_burn_usd_per_hour: Option<f64>,
1221 spent_usd: Option<f64>,
1222 credentials_present: Option<bool>,
1223 human_gated: Option<bool>,
1224 ) {
1225 let now = Utc::now().to_rfc3339();
1226 op.last_keep_alive_at = now.clone();
1227 op.updated_at = now;
1228 if let Some(burn) = observed_burn_usd_per_hour {
1229 op.observed_burn_usd_per_hour = Some(burn.max(0.0));
1230 }
1231 if let Some(spent) = spent_usd {
1232 op.spent_usd = spent.max(0.0);
1233 }
1234 if let Some(credentials) = credentials_present {
1235 op.credentials_present = credentials;
1236 }
1237 if let Some(gated) = human_gated {
1238 op.human_gated = gated;
1239 }
1240 op.project();
1241 }
1242
1243 /// Install (or refresh) the `cw-operate` keepalive bound to `workspace`.
1244 ///
1245 /// The record is built directly under the fixed id — no create-then-delete
1246 /// id swap that could orphan an active UUID-named automation. Reuse
1247 /// refreshes the prompt and workspace while preserving an explicit saved
1248 /// model/provider pair. Fresh and legacy unpinned records capture the caller's
1249 /// effective route together, including Auto intent and exact custom identity.
1250 ///
1251 /// `kick_now` schedules the first lead-plan step for the next scheduler tick
1252 /// (a fresh operation otherwise idles up to an hour awaiting its plan); the
1253 /// hourly recurrence covers follow-ups.
1254 pub(crate) fn upsert_keepalive(
1255 manager: &AutomationManager,
1256 workspace: &Path,
1257 kick_now: bool,
1258 config: &crate::config::Config,
1259 selection: Option<(&crate::config::ProviderIdentity, &str)>,
1260 ) -> Result<(String, bool)> {
1261 let now = Utc::now();
1262 let prompt = format!(
1263 "Keep Operate alive. Read the Operate record (current.json) and its direction, refresh the lead plan, and dispatch ready slices with at most `writersInFlight` concurrent workers — that budget already encodes pace (hold/throttle/widen), so honor it instead of widening on your own. Burn rate paces spend; it never stops the operation. Workspace: {}",
1264 workspace.display()
1265 );
1266 let mut readiness = (String::new(), false);
1267 manager.edit_automation(OPERATE_KEEPALIVE_ID, |current| {
1268 let scope = manager
1269 .execution_scope()
1270 .context("Operate execution ownership is unverified")?;
1271 if current
1272 .as_ref()
1273 .and_then(|record| record.execution_scope.as_deref())
1274 .is_some_and(|bound| bound != scope)
1275 {
1276 bail!("Operate keepalive belongs to another Runtime execution scope");
1277 }
1278 let (identity, model, ready) = keepalive_route(config, current.as_ref(), selection)?;
1279 readiness = (model.clone(), ready);
1280 let mut record = current.unwrap_or_else(|| AutomationRecord {
1281 schema_version: crate::automation_manager::CURRENT_AUTOMATION_SCHEMA_VERSION,
1282 execution_scope: Some(scope.to_string()),
1283 id: OPERATE_KEEPALIVE_ID.to_string(),
1284 name: "Operate keep-alive".to_string(),
1285 prompt: prompt.clone(),
1286 rrule: OPERATE_KEEPALIVE_RRULE.to_string(),
1287 cwds: Vec::new(),
1288 model: None,
1289 model_provider: None,
1290 model_provider_id: None,
1291 mode: Some("operate".to_string()),
1292 allow_shell: Some(true),
1293 trust_mode: Some(false),
1294 auto_approve: Some(false),
1295 delivery_mode: None,
1296 status: AutomationStatus::Active,
1297 created_at: now,
1298 updated_at: now,
1299 next_run_at: None,
1300 last_run_at: None,
1301 });
1302 record.schema_version = crate::automation_manager::CURRENT_AUTOMATION_SCHEMA_VERSION;
1303 if record.execution_scope.is_none() {
1304 record.execution_scope = Some(scope.to_string());
1305 }
1306 record.name = "Operate keep-alive".to_string();
1307 record.prompt = prompt;
1308 record.rrule = OPERATE_KEEPALIVE_RRULE.to_string();
1309 record.cwds = vec![workspace.to_path_buf()];
1310 record.model = Some(model);
1311 record.model_provider = Some(identity.provider.as_str().to_string());
1312 record.model_provider_id = identity.persisted_id().map(str::to_string);
1313 record.mode = Some("operate".to_string());
1314 record.allow_shell = Some(true);
1315 record.trust_mode = Some(false);
1316 record.auto_approve = Some(false);
1317 record.delivery_mode = None;
1318 record.status = AutomationStatus::Active;
1319 record.updated_at = now;
1320 if kick_now {
1321 record.next_run_at = Some(now);
1322 }
1323 Ok(Some(record))
1324 })?;
1325 Ok(readiness)
1326 }
1327
1328 /// Cancel tears the operation down *including* its keepalive: an unattended
1329 /// hourly lead run after cancel is pure cost. The automation is paused
1330 /// (never deleted) so its run history survives and a later start reactivates
1331 /// it. A missing keepalive is not an error; any other failure to establish
1332 /// its state is. A requested spending stop that cannot confirm the record is
1333 /// paused must never look like one that happened — an unreadable record may
1334 /// still be Active and scheduled.
1335 pub fn pause_keepalive(manager: &AutomationManager) -> Result<()> {
1336 let record = match manager.get_automation(OPERATE_KEEPALIVE_ID) {
1337 Ok(record) => record,
1338 Err(error) if is_not_found(&error) => return Ok(()),
1339 Err(error) => {
1340 return Err(error.context(
1341 "Operate keepalive state could not be established; it may still be scheduled",
1342 ));
1343 }
1344 };
1345 if matches!(record.status, AutomationStatus::Active) {
1346 manager.pause_automation(OPERATE_KEEPALIVE_ID)?;
1347 }
1348 Ok(())
1349 }
1350
1351 fn is_not_found(error: &anyhow::Error) -> bool {
1352 error.chain().any(|cause| {
1353 cause
1354 .downcast_ref::<std::io::Error>()
1355 .is_some_and(|io| io.kind() == std::io::ErrorKind::NotFound)
1356 })
1357 }
1358
1359 /// Pull the next keepalive lead run to the next scheduler tick (for example
1360 /// after a direction PATCH invalidated the plan). No-op when the keepalive is
1361 /// absent or paused (a paused keepalive belongs to a cancelled operation).
1362 pub fn kick_keepalive(manager: &AutomationManager) -> Result<bool> {
1363 let mut kicked = false;
1364 manager.edit_automation(OPERATE_KEEPALIVE_ID, |current| {
1365 let Some(mut record) = current else {
1366 return Ok(None);
1367 };
1368 if record.status == AutomationStatus::Active {
1369 let now = Utc::now();
1370 record.next_run_at = Some(now);
1371 record.updated_at = now;
1372 kicked = true;
1373 }
1374 Ok(Some(record))
1375 })?;
1376 Ok(kicked)
1377 }
1378
1379 #[must_use]
1380 pub fn human_gate_for(action: &str) -> bool {
1381 matches!(
1382 action,
1383 "deploy" | "billing" | "force-push" | "forbidden-pr" | "red-ci"
1384 )
1385 }
1386
1387 #[cfg(test)]
1388 mod tests {
1389 use super::*;
1390 use tempfile::TempDir;
1391
1392 fn route_fixture_config() -> crate::config::Config {
1393 toml::from_str(
1394 r#"
1395 provider = "route-a"
1396 [providers.route-a]
1397 kind = "openai-compatible"
1398 base_url = "https://route-a.example.test/v1"
1399 model = "same-model"
1400 auth_mode = "none"
1401 [providers.route-b]
1402 kind = "openai-compatible"
1403 base_url = "https://route-b.example.test/v1"
1404 model = "same-model"
1405 auth_mode = "api-key"
1406 api_key_env = "CW_OPERATE_MISSING_TEST_KEY"
1407 "#,
1408 )
1409 .expect("route fixture")
1410 }
1411
1412 #[test]
1413 fn keepalive_saved_route_controls_readiness_refresh_and_auto_intent() -> Result<()> {
1414 let _env = crate::test_support::lock_test_env();
1415 let _cli = crate::test_support::EnvVarGuard::remove("CODEWHALE_CLI_API_KEY");
1416 let _missing = crate::test_support::EnvVarGuard::remove("CW_OPERATE_MISSING_TEST_KEY");
1417 let root = TempDir::new()?;
1418 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
1419 let mut config = route_fixture_config();
1420 assert!(upsert_keepalive(&manager, root.path(), true, &config, None)?.1);
1421 let first = manager.get_automation(OPERATE_KEEPALIVE_ID)?;
1422 assert_eq!(first.model.as_deref(), Some("same-model"));
1423 assert_eq!(first.model_provider.as_deref(), Some("custom"));
1424 assert_eq!(first.model_provider_id.as_deref(), Some("route-a"));
1425 assert_eq!(first.auto_approve, Some(false));
1426
1427 // Explicitly edit the saved route. Parent route-a is credential-ready;
1428 // route-b is not, even though both expose the same model spelling.
1429 let mut pinned = first.clone();
1430 pinned.model_provider_id = Some("route-b".into());
1431 manager.save_automation(&pinned)?;
1432 assert!(!keepalive_readiness(&manager, &config, None)?.1);
1433 assert!(!upsert_keepalive(&manager, root.path(), false, &config, None)?.1);
1434 assert_eq!(
1435 manager
1436 .get_automation(OPERATE_KEEPALIVE_ID)?
1437 .model_provider_id,
1438 Some("route-b".into())
1439 );
1440
1441 pinned.model = Some("auto".into());
1442 manager.save_automation(&pinned)?;
1443 assert!(!upsert_keepalive(&manager, root.path(), false, &config, None)?.1);
1444 let auto = manager.get_automation(OPERATE_KEEPALIVE_ID)?;
1445 assert_eq!(auto.model.as_deref(), Some("auto"));
1446 assert_eq!(auto.model_provider_id.as_deref(), Some("route-b"));
1447 assert_eq!(auto.created_at, first.created_at);
1448
1449 config.providers.as_mut().unwrap().custom.remove("route-b");
1450 let before = serde_json::to_value(&auto)?;
1451 assert!(keepalive_readiness(&manager, &config, None).is_err());
1452 assert!(upsert_keepalive(&manager, root.path(), true, &config, None).is_err());
1453 assert_eq!(
1454 serde_json::to_value(manager.get_automation(OPERATE_KEEPALIVE_ID)?)?,
1455 before
1456 );
1457
1458 // Legacy unpinned GLM is a previous scheduler default, not an exact
1459 // provider choice. Migration replaces all route fields together.
1460 pinned.model = Some("GLM-5.3".into());
1461 pinned.model_provider = None;
1462 pinned.model_provider_id = None;
1463 manager.save_automation(&pinned)?;
1464 assert!(upsert_keepalive(&manager, root.path(), true, &config, None)?.1);
1465 let migrated = manager.get_automation(OPERATE_KEEPALIVE_ID)?;
1466 assert_eq!(migrated.model.as_deref(), Some("same-model"));
1467 assert_eq!(migrated.model_provider.as_deref(), Some("custom"));
1468 assert_eq!(migrated.model_provider_id.as_deref(), Some("route-a"));
1469 Ok(())
1470 }
1471
1472 #[test]
1473 fn keepalive_legacy_custom_records_the_table_id() -> Result<()> {
1474 let _env = crate::test_support::lock_test_env();
1475 let _cli = crate::test_support::EnvVarGuard::remove("CODEWHALE_CLI_API_KEY");
1476 let root = TempDir::new()?;
1477 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
1478 let config = crate::config::parse_config_base(
1479 r#"provider = "custom"
1480 [providers.custom]
1481 kind = "openai-compatible"
1482 base_url = "https://legacy.example.test/v1"
1483 model = "legacy-model"
1484 auth_mode = "none"
1485 "#,
1486 )?;
1487 upsert_keepalive(&manager, root.path(), false, &config, None)?;
1488 let record = manager.get_automation(OPERATE_KEEPALIVE_ID)?;
1489 assert_eq!(record.model.as_deref(), Some("legacy-model"));
1490 // The literal route is the `[providers.custom]` table since #6394.
1491 assert_eq!(record.model_provider.as_deref(), Some("custom"));
1492 assert_eq!(record.model_provider_id.as_deref(), Some("custom"));
1493 upsert_keepalive(&manager, root.path(), false, &config, None)?;
1494 assert_eq!(
1495 manager
1496 .get_automation(OPERATE_KEEPALIVE_ID)?
1497 .model_provider_id
1498 .as_deref(),
1499 Some("custom")
1500 );
1501 Ok(())
1502 }
1503
1504 struct RouteRecordingExecutor(
1505 std::sync::Arc<std::sync::Mutex<Vec<crate::runtime_threads::CreateThreadRequest>>>,
1506 );
1507
1508 #[async_trait::async_trait]
1509 impl crate::task_manager::TaskExecutor for RouteRecordingExecutor {
1510 async fn execute(
1511 &self,
1512 task: crate::task_manager::ExecutionTask,
1513 _events: tokio::sync::mpsc::Sender<crate::task_manager::TaskExecutionEvent>,
1514 _cancel: tokio_util::sync::CancellationToken,
1515 ) -> crate::task_manager::TaskExecutionResult {
1516 self.0.lock().unwrap().push(task.thread_request());
1517 crate::task_manager::TaskExecutionResult {
1518 status: crate::task_manager::TaskStatus::Completed,
1519 result_text: Some("route fixture completed".into()),
1520 error: None,
1521 terminal_reason: crate::task_manager::TaskTerminalReason::Completed,
1522 }
1523 }
1524 }
1525
1526 #[tokio::test]
1527 async fn keepalive_pin_reaches_task_and_survives_parent_change() -> Result<()> {
1528 let root = TempDir::new()?;
1529 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
1530 let mut config = route_fixture_config();
1531 upsert_keepalive(&manager, root.path(), false, &config, None)?;
1532 pause_keepalive(&manager)?;
1533 config.provider = Some("route-b".into());
1534 assert!(upsert_keepalive(&manager, root.path(), true, &config, None)?.1);
1535 let observations = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1536 let tasks = crate::task_manager::TaskManager::start_with_executor(
1537 crate::task_manager::TaskManagerConfig {
1538 data_dir: root.path().join("tasks"),
1539 worker_count: 1,
1540 default_workspace: root.path().to_path_buf(),
1541 default_model: "changed-default".into(),
1542 default_mode: "agent".into(),
1543 allow_shell: false,
1544 trust_mode: false,
1545 execution_limits: Default::default(),
1546 },
1547 std::sync::Arc::new(RouteRecordingExecutor(observations.clone())),
1548 )
1549 .await?;
1550 let shared = std::sync::Arc::new(tokio::sync::Mutex::new(manager));
1551 let run = crate::automation_manager::run_now_shared(&shared, OPERATE_KEEPALIVE_ID, &tasks)
1552 .await?;
1553 let id = run.task_id.as_deref().context("bound task")?;
1554 let task = crate::task_manager::wait_for_terminal_state(
1555 &tasks,
1556 id,
1557 std::time::Duration::from_secs(5),
1558 )
1559 .await?;
1560 assert_eq!(task.status, crate::task_manager::TaskStatus::Completed);
1561 assert_eq!(task.model, "same-model");
1562 assert_eq!(task.model_provider.as_deref(), Some("custom"));
1563 assert_eq!(task.model_provider_id.as_deref(), Some("route-a"));
1564 let observed = observations.lock().unwrap();
1565 assert_eq!(observed.len(), 1);
1566 assert_eq!(observed[0].model, Some(task.model.clone()));
1567 assert_eq!(observed[0].model_provider_id, task.model_provider_id);
1568 assert_eq!(observed[0].auto_approve, Some(false));
1569 drop(observed);
1570 let reopened = AutomationManager::open_for_test(root.path().join("automations"))?;
1571 assert!(keepalive_readiness(&reopened, &config, None)?.1);
1572 assert_eq!(
1573 reopened.list_runs(OPERATE_KEEPALIVE_ID, None)?[0]
1574 .task_id
1575 .as_deref(),
1576 Some(id)
1577 );
1578 // A subsequent explicit edit cannot rewrite an accepted task binding.
1579 let mut edited = reopened.get_automation(OPERATE_KEEPALIVE_ID)?;
1580 edited.model = Some("replacement-model".into());
1581 edited.model_provider_id = Some("route-b".into());
1582 reopened.save_automation(&edited)?;
1583 assert_eq!(
1584 tasks
1585 .read_bound_task(id)?
1586 .context("accepted task")?
1587 .model_provider_id,
1588 Some("route-a".into())
1589 );
1590 tasks.shutdown();
1591 Ok(())
1592 }
1593
1594 fn with_credentials(mut op: Operation) -> Operation {
1595 op.credentials_present = true;
1596 op.project();
1597 op
1598 }
1599
1600 #[test]
1601 fn unbounded_start_matches_cwc_contract() {
1602 let op = with_credentials(Operation::new("Keep shipping honest slices", None));
1603 assert_eq!(op.schema_version, 1);
1604 assert!(op.burn_rate.is_none());
1605 assert_eq!(op.pace, OperatePace::Unbounded);
1606 assert_eq!(op.status, OperateStatus::IdleBlocked);
1607 assert_eq!(
1608 op.idle_blocked_reason,
1609 Some(OperateIdleReason::AwaitingLeadPlan)
1610 );
1611 assert!(!op.workers_admitted);
1612 assert!(op.lead_operator.model.is_empty(), "no model was assigned");
1613 let json = serde_json::to_value(&op).expect("json");
1614 assert!(json.get("burnRate").unwrap().is_null());
1615 assert_eq!(json["leadOperator"]["model"], "");
1616 assert_eq!(json["schemaVersion"], 1);
1617 assert!(json.get("leadPlan").unwrap().is_null());
1618 assert_eq!(json["idleBlockedReason"], "awaiting_lead_plan");
1619 assert_eq!(json["workersAdmitted"], false);
1620 assert!(json["id"].as_str().unwrap().starts_with("op_"));
1621 }
1622
1623 #[test]
1624 fn operation_lead_label_follows_selection_and_planned_workers_stay_unassigned() -> Result<()> {
1625 let root = TempDir::new()?;
1626 let store = OperationStore::open(root.path())?;
1627 let operation = start_operation(
1628 &store,
1629 root.path(),
1630 Some("First bounded task\nSecond bounded task".into()),
1631 None,
1632 true,
1633 "selected-model",
1634 )?;
1635 assert_eq!(operation.lead_operator.model, "selected-model");
1636 let mut attached =
1637 attach_or_start_operation(&store, root.path(), None, None, true, "auto")?;
1638 assert_eq!(attached.id, operation.id);
1639 assert_eq!(attached.lead_operator.model, "auto");
1640 assert!(
1641 attached
1642 .roster
1643 .iter()
1644 .filter(|member| member.id == "lead")
1645 .all(|member| member.model == "auto")
1646 );
1647 attached.plan_from_direction();
1648 assert!(attached.roster.iter().any(|member| member.id != "lead"));
1649 assert!(
1650 attached
1651 .roster
1652 .iter()
1653 .filter(|member| member.id != "lead")
1654 .all(|member| member.model.is_empty())
1655 );
1656 assert!(
1657 store
1658 .load()?
1659 .context("saved operation")?
1660 .roster
1661 .iter()
1662 .filter(|member| member.id == "lead")
1663 .all(|member| member.model == "auto")
1664 );
1665 Ok(())
1666 }
1667
1668 /// The cross-process lock sidecar is owner-only like the operation file.
1669 #[cfg(unix)]
1670 #[test]
1671 fn operation_lock_file_is_owner_only() -> Result<()> {
1672 use std::os::unix::fs::PermissionsExt as _;
1673 let root = TempDir::new()?;
1674 let store = OperationStore::open(root.path())?;
1675 start_operation(&store, root.path(), None, None, true, "auto")?;
1676 let lock = root.path().join("current.json.lock");
1677 assert_eq!(fs::metadata(&lock)?.permissions().mode() & 0o777, 0o600);
1678 Ok(())
1679 }
1680
1681 #[test]
1682 fn create_without_credentials_fails_closed() {
1683 let op = Operation::new("Keep shipping honest slices", None);
1684 assert_eq!(op.status, OperateStatus::IdleBlocked);
1685 assert_eq!(
1686 op.idle_blocked_reason,
1687 Some(OperateIdleReason::MissingCredentials)
1688 );
1689 assert!(!op.workers_admitted);
1690 }
1691
1692 #[test]
1693 fn burn_rate_paces_and_never_stops() {
1694 let mut op = with_credentials(Operation::new(
1695 "Hold a $12/hr burn\nSecond slice\nThird slice",
1696 Some(12.0),
1697 ));
1698 op.plan_from_direction();
1699 assert_eq!(op.status, OperateStatus::Running);
1700 assert!(op.workers_admitted);
1701 keep_alive_observation(&mut op, Some(20.0), Some(80.0), None, None);
1702 assert_eq!(op.status, OperateStatus::Running);
1703 assert_eq!(op.pace, OperatePace::Throttle);
1704 // Throttle is a real dispatch cut, not a label: of the two planned
1705 // workers only one stays in flight while the operation keeps running.
1706 assert_eq!(op.writers_in_flight, OPERATE_THROTTLE_WRITERS);
1707 assert_eq!(
1708 op.roster
1709 .iter()
1710 .filter(|member| member.role == "worker" && member.state == "in_flight")
1711 .count(),
1712 OPERATE_THROTTLE_WRITERS
1713 );
1714 assert_eq!(
1715 op.roster
1716 .iter()
1717 .filter(|member| member.role == "worker" && member.state == "idle")
1718 .count(),
1719 1,
1720 "the worker past the throttle budget idles"
1721 );
1722 assert!(op.idle_blocked_reason.is_none());
1723 assert!(op.workers_admitted);
1724 let board = render_plan_board(&op);
1725 assert!(!board.contains("exhausted"));
1726 assert!(!board.contains("wallet"));
1727 assert_eq!(op.burn_rate.as_ref().unwrap().kind, "usd_per_hour");
1728 assert!((op.burn_rate.as_ref().unwrap().amount_usd_per_hour - 12.0).abs() < f64::EPSILON);
1729 }
1730
1731 #[test]
1732 fn hold_band_freezes_writer_width() {
1733 let mut op = with_credentials(Operation::new("one\ntwo\nthree", Some(12.0)));
1734 op.plan_from_direction();
1735 keep_alive_observation(&mut op, Some(12.0), None, None, None);
1736 assert_eq!(op.pace, OperatePace::Hold);
1737 assert_eq!(op.status, OperateStatus::Running);
1738 // Hold admits no new writers past the hold width; this plan has two
1739 // workers, so both stay in flight but the budget stops at two.
1740 assert_eq!(op.writers_in_flight, 2);
1741 let mut four = with_credentials(Operation::new("one\ntwo\nthree\nfour\nfive", Some(12.0)));
1742 four.plan_from_direction();
1743 keep_alive_observation(&mut four, Some(12.3), None, None, None);
1744 assert_eq!(four.pace, OperatePace::Hold);
1745 assert_eq!(
1746 four.writers_in_flight, OPERATE_HOLD_WRITERS,
1747 "hold never widens to the full writer width"
1748 );
1749 }
1750
1751 #[test]
1752 fn under_rate_widens() {
1753 let mut op = with_credentials(Operation::new("one\ntwo\nthree\nfour", Some(12.0)));
1754 op.plan_from_direction();
1755 keep_alive_observation(&mut op, Some(1.0), None, None, None);
1756 assert_eq!(op.pace, OperatePace::Widen);
1757 assert_eq!(op.status, OperateStatus::Running);
1758 assert_eq!(
1759 op.writers_in_flight, OPERATE_MAX_WRITERS,
1760 "widen opens the full writer width"
1761 );
1762 }
1763
1764 #[test]
1765 fn cancelled_field_matches_landed_cwc_contract() {
1766 // CWC `packages/contracts/src/operate.js` (`20de981`, PR #284) names
1767 // the field `cancelledAt` — `publicOperateRecord` always emits it,
1768 // as `""` before cancellation. There is no bare `cancelled` field.
1769 let dir = TempDir::new().expect("temp");
1770 let store = OperationStore::open(dir.path()).expect("store");
1771 let op = start_operation(
1772 &store,
1773 dir.path(),
1774 Some("Contract shape".into()),
1775 None,
1776 true,
1777 "selected-model",
1778 )
1779 .expect("start");
1780 let json = serde_json::to_value(&op).expect("json");
1781 assert_eq!(json["cancelledAt"], serde_json::json!(""));
1782 assert!(json.get("cancelled").is_none());
1783
1784 let cancelled = cancel_operation(&store).expect("cancel").expect("present");
1785 let json = serde_json::to_value(&cancelled).expect("json");
1786 assert!(!json["cancelledAt"].as_str().unwrap().is_empty());
1787 assert!(json.get("cancelled").is_none());
1788 }
1789
1790 #[test]
1791 fn direction_change_invalidates_stale_lead_plan() {
1792 let mut op = with_credentials(Operation::new("Old direction", None));
1793 op.plan_from_direction();
1794 assert_eq!(op.status, OperateStatus::Running);
1795
1796 apply_operate_patch(
1797 &mut op,
1798 &serde_json::json!({ "direction": "Brand new direction" }),
1799 )
1800 .expect("patch");
1801 assert_eq!(op.direction, "Brand new direction");
1802 assert!(
1803 op.lead_plan.is_none(),
1804 "a changed direction must not leave superseded slices executing"
1805 );
1806 assert_eq!(op.status, OperateStatus::IdleBlocked);
1807 assert_eq!(
1808 op.idle_blocked_reason,
1809 Some(OperateIdleReason::AwaitingLeadPlan)
1810 );
1811 assert!(!op.workers_admitted);
1812 assert_eq!(op.writers_in_flight, 0);
1813 }
1814
1815 #[test]
1816 fn same_direction_patch_keeps_plan() {
1817 let mut op = with_credentials(Operation::new("Steady", None));
1818 op.plan_from_direction();
1819 apply_operate_patch(&mut op, &serde_json::json!({ "direction": "Steady" })).expect("patch");
1820 assert!(op.lead_plan.is_some());
1821 assert_eq!(op.status, OperateStatus::Running);
1822 }
1823
1824 #[test]
1825 fn put_plan_admits_worker_owners() {
1826 let mut op = with_credentials(Operation::new("Slice it", None));
1827 let plan = serde_json::json!({
1828 "slices": [
1829 { "id": "slice-1", "title": "Scout", "ownerId": "lead",
1830 "dependsOn": [], "estCostUsd": 0.1, "startOffsetSec": 0, "durationSec": 600 },
1831 { "id": "slice-2", "title": "Build", "ownerId": "worker-7",
1832 "dependsOn": ["slice-1"], "estCostUsd": 0.2, "startOffsetSec": 600, "durationSec": 1200 }
1833 ]
1834 });
1835 apply_operate_patch(&mut op, &serde_json::json!({ "leadPlan": plan })).expect("patch");
1836 assert!(
1837 op.roster.iter().any(|member| member.id == "worker-7"),
1838 "plan owners must join the roster or workers never dispatch"
1839 );
1840 assert_eq!(op.status, OperateStatus::Running);
1841 assert!(op.workers_admitted);
1842 // Idempotent: re-applying the same plan must not duplicate the owner.
1843 let before = op.roster.len();
1844 apply_operate_patch(&mut op, &serde_json::json!({ "leadPlan": plan })).expect("patch");
1845 assert_eq!(op.roster.len(), before);
1846 }
1847
1848 #[test]
1849 fn attach_preserves_operation_record() {
1850 let dir = TempDir::new().expect("temp");
1851 let store = OperationStore::open(dir.path()).expect("store");
1852 let first = start_operation(
1853 &store,
1854 dir.path(),
1855 Some("Keep the lineage".into()),
1856 None,
1857 true,
1858 "selected-model",
1859 )
1860 .expect("start");
1861 let mut spent = first.clone();
1862 keep_alive_observation(&mut spent, Some(3.0), Some(42.0), None, None);
1863 store.save(&spent).expect("save spend");
1864
1865 let reentered =
1866 attach_or_start_operation(&store, dir.path(), None, None, true, "selected-model")
1867 .expect("attach");
1868 assert_eq!(reentered.id, first.id, "re-entry must attach, not reset");
1869 assert_eq!(reentered.spent_usd, 42.0);
1870 assert_eq!(reentered.observed_burn_usd_per_hour, Some(3.0));
1871
1872 // A cancelled record is terminal: re-entry starts a new operation.
1873 cancel_operation(&store).expect("cancel");
1874 let fresh =
1875 attach_or_start_operation(&store, dir.path(), None, None, true, "selected-model")
1876 .expect("restart");
1877 assert_ne!(fresh.id, first.id);
1878 // The fresh operation reuses the recorded direction and awaits its
1879 // own lead plan (CWC projects a plan-less record to idle_blocked).
1880 assert_eq!(fresh.direction, "Keep the lineage");
1881 assert_eq!(fresh.status, OperateStatus::IdleBlocked);
1882 assert_eq!(
1883 fresh.idle_blocked_reason,
1884 Some(OperateIdleReason::AwaitingLeadPlan)
1885 );
1886 assert_eq!(fresh.spent_usd, 0.0);
1887 }
1888
1889 #[test]
1890 fn mutate_reloads_latest_record_under_lock() {
1891 let dir = TempDir::new().expect("temp");
1892 let writer = OperationStore::open(dir.path()).expect("store a");
1893 let reader = OperationStore::open(dir.path()).expect("store b");
1894 start_operation(
1895 &writer,
1896 dir.path(),
1897 Some("First".into()),
1898 None,
1899 true,
1900 "selected-model",
1901 )
1902 .expect("start");
1903
1904 writer
1905 .mutate(|op| {
1906 apply_operate_patch(op, &serde_json::json!({ "direction": "Second" }))?;
1907 Ok(())
1908 })
1909 .expect("mutate a")
1910 .expect("present");
1911 reader
1912 .mutate(|op| {
1913 keep_alive_observation(op, Some(9.0), Some(5.0), None, None);
1914 Ok(())
1915 })
1916 .expect("mutate b")
1917 .expect("present");
1918
1919 // The second store reloaded under the lock, so the first store's
1920 // direction write survives alongside the keepalive observation.
1921 let merged = writer.load().expect("load").expect("present");
1922 assert_eq!(merged.direction, "Second");
1923 assert_eq!(merged.spent_usd, 5.0);
1924 assert_eq!(merged.observed_burn_usd_per_hour, Some(9.0));
1925 }
1926
1927 #[test]
1928 fn burn_rate_below_a_cent_is_rejected() {
1929 let err = parse_burn_rate(Some(&serde_json::json!(0.001))).expect_err("rejects");
1930 assert!(err.to_string().contains("at least $0.01/hr"));
1931 let zero_target = parse_burn_rate(Some(&serde_json::json!(0.004))).expect_err("rejects");
1932 assert!(zero_target.to_string().contains("at least $0.01/hr"));
1933 }
1934
1935 #[test]
1936 fn keepalive_reuse_refreshes_cwds_and_kicks_first_lead_run() {
1937 let dir = TempDir::new().expect("temp");
1938 let manager = AutomationManager::open_for_test(dir.path().to_path_buf()).expect("manager");
1939 let workspace_a = dir.path().join("workspace-a");
1940 let workspace_b = dir.path().join("workspace-b");
1941 fs::create_dir_all(&workspace_a).expect("dir a");
1942 fs::create_dir_all(&workspace_b).expect("dir b");
1943
1944 upsert_keepalive(&manager, &workspace_a, false, &route_fixture_config(), None)
1945 .expect("upsert a");
1946 let first = manager
1947 .get_automation(OPERATE_KEEPALIVE_ID)
1948 .expect("keepalive a");
1949 assert_eq!(first.cwds, vec![workspace_a.clone()]);
1950
1951 upsert_keepalive(&manager, &workspace_b, true, &route_fixture_config(), None)
1952 .expect("upsert b");
1953 let second = manager
1954 .get_automation(OPERATE_KEEPALIVE_ID)
1955 .expect("keepalive b");
1956 assert_eq!(
1957 second.cwds,
1958 vec![workspace_b],
1959 "reuse must retarget the workspace scheduled runs execute in"
1960 );
1961 assert_eq!(
1962 second.next_run_at.map(|at| at <= Utc::now()),
1963 Some(true),
1964 "kick schedules the first lead run for the next scheduler tick"
1965 );
1966 assert_eq!(second.model.as_deref(), Some("same-model"));
1967 assert_eq!(second.mode.as_deref(), Some("operate"));
1968 assert_eq!(second.rrule, OPERATE_KEEPALIVE_RRULE);
1969
1970 // Only one automation exists — no orphaned UUID-named twin.
1971 assert_eq!(
1972 manager
1973 .list_automations()
1974 .expect("list")
1975 .iter()
1976 .filter(|record| record.mode.as_deref() == Some("operate"))
1977 .count(),
1978 1
1979 );
1980 }
1981
1982 #[test]
1983 fn cancel_pauses_keepalive_so_no_cost_accrues() {
1984 let dir = TempDir::new().expect("temp");
1985 let manager = AutomationManager::open_for_test(dir.path().to_path_buf()).expect("manager");
1986 upsert_keepalive(&manager, dir.path(), false, &route_fixture_config(), None)
1987 .expect("upsert");
1988 pause_keepalive(&manager).expect("pause");
1989 let paused = manager
1990 .get_automation(OPERATE_KEEPALIVE_ID)
1991 .expect("keepalive");
1992 assert_eq!(paused.status, AutomationStatus::Paused);
1993 assert_eq!(paused.next_run_at, None, "nothing fires after cancel");
1994
1995 // Pausing is idempotent and a missing keepalive is not an error.
1996 pause_keepalive(&manager).expect("pause again");
1997 let empty =
1998 AutomationManager::open_for_test(dir.path().join("empty")).expect("empty manager");
1999 pause_keepalive(&empty).expect("missing keepalive is a no-op");
2000 assert!(!kick_keepalive(&empty).expect("kick missing"));
2001
2002 // A fresh start reactivates the keepalive.
2003 upsert_keepalive(&manager, dir.path(), false, &route_fixture_config(), None)
2004 .expect("reactivate");
2005 let active = manager
2006 .get_automation(OPERATE_KEEPALIVE_ID)
2007 .expect("keepalive");
2008 assert_eq!(active.status, AutomationStatus::Active);
2009
2010 // Kicks only touch an active keepalive.
2011 assert!(kick_keepalive(&manager).expect("kick"));
2012 pause_keepalive(&manager).expect("pause");
2013 assert!(!kick_keepalive(&manager).expect("kick paused"));
2014 }
2015
2016 #[test]
2017 fn pause_keepalive_refuses_to_report_a_stop_it_cannot_confirm() {
2018 let dir = TempDir::new().expect("temp");
2019 let manager = AutomationManager::open_for_test(dir.path().to_path_buf()).expect("manager");
2020 upsert_keepalive(&manager, dir.path(), false, &route_fixture_config(), None)
2021 .expect("upsert");
2022 // Fault fixture: the keepalive record exists but cannot be parsed, so
2023 // whether it is still Active (and spending) is unknown.
2024 let record = dir
2025 .path()
2026 .join("automations")
2027 .join(format!("{OPERATE_KEEPALIVE_ID}.json"));
2028 assert!(record.exists(), "fixture targets the real record path");
2029 fs::write(&record, b"{\"status\": \"active\", truncated").expect("corrupt record");
2030
2031 let error = pause_keepalive(&manager).expect_err("an unconfirmed stop must fail");
2032 assert!(
2033 format!("{error:#}").contains("could not be established"),
2034 "{error:#}"
2035 );
2036 }
2037
2038 #[test]
2039 fn explicit_env_paths_reject_traversal() {
2040 assert!(explicit_file_path(" /tmp/does-not-exist.md ").is_none());
2041 assert!(explicit_file_path("").is_none());
2042 assert!(explicit_file_path("/tmp/../etc/passwd").is_none());
2043 // Keep the temp dir alive for the whole assertion: dropping it first
2044 // would delete the file under the path.
2045 let dir = TempDir::new().expect("temp");
2046 let checker = dir.path().join("check.py");
2047 fs::write(&checker, "# marker").expect("write");
2048 let found =
2049 explicit_file_path(checker.to_str().expect("utf8")).expect("regular file accepted");
2050 assert_eq!(found, checker);
2051 assert!(explicit_dir_path("../escape").is_none());
2052 assert!(explicit_dir_path("ops/inside").is_some());
2053 }
2054
2055 #[test]
2056 fn missing_credentials_fail_closed() {
2057 let dir = TempDir::new().expect("temp");
2058 let store = OperationStore::open(dir.path()).expect("store");
2059 let op = start_operation(
2060 &store,
2061 dir.path(),
2062 Some("Do not spend silently".into()),
2063 None,
2064 false,
2065 "selected-model",
2066 )
2067 .expect("start");
2068 assert_eq!(op.status, OperateStatus::IdleBlocked);
2069 assert_eq!(
2070 op.idle_blocked_reason,
2071 Some(OperateIdleReason::MissingCredentials)
2072 );
2073 assert!(!op.workers_admitted);
2074 assert_eq!(op.writers_in_flight, 0);
2075 }
2076
2077 #[test]
2078 fn cancel_stays_cancelled_through_keep_alive() {
2079 let dir = TempDir::new().expect("temp");
2080 let store = OperationStore::open(dir.path()).expect("store");
2081 start_operation(
2082 &store,
2083 dir.path(),
2084 Some("Stop".into()),
2085 None,
2086 true,
2087 "selected-model",
2088 )
2089 .expect("start");
2090 let cancelled = cancel_operation(&store).expect("cancel").expect("present");
2091 let mut kept = cancelled;
2092 keep_alive_observation(&mut kept, Some(40.0), Some(999.0), None, None);
2093 assert_eq!(kept.status, OperateStatus::Cancelled);
2094 assert!(!kept.workers_admitted);
2095 }
2096
2097 #[test]
2098 fn lead_plan_is_the_gantt_model() {
2099 let mut op = with_credentials(Operation::new("Scout\nWrite", None));
2100 op.plan_from_direction();
2101 let plan = op.lead_plan.as_ref().expect("plan");
2102 assert_eq!(plan.slices.len(), 2);
2103 assert_eq!(plan.slices[0].owner_id, "lead");
2104 assert_eq!(plan.slices[1].depends_on, vec!["slice-1".to_string()]);
2105 assert_eq!(plan.slices[0].start_offset_sec, 0);
2106 assert!(plan.slices[0].duration_sec >= 1);
2107 let board = render_plan_board(&op);
2108 assert!(board.contains("gantt time →"), "{board}");
2109 assert!(board.contains("leadPlan"), "{board}");
2110 assert!(board.contains("No cap"), "{board}");
2111 }
2112
2113 #[test]
2114 fn empty_direction_is_idle_blocked() {
2115 let op = with_credentials(Operation::new("", None));
2116 assert_eq!(op.status, OperateStatus::IdleBlocked);
2117 assert_eq!(
2118 op.idle_blocked_reason,
2119 Some(OperateIdleReason::DirectionEmpty)
2120 );
2121 }
2122
2123 #[test]
2124 fn human_gates_do_not_include_merge() {
2125 assert!(human_gate_for("deploy"));
2126 assert!(human_gate_for("billing"));
2127 assert!(!human_gate_for("merge"));
2128 }
2129
2130 #[test]
2131 fn calls_landed_checker_flags() {
2132 assert_eq!(
2133 check_auto_merge_args("codewhale-hq/CodeWhale", "1234", "keel"),
2134 vec![
2135 "scripts/check-auto-merge.py",
2136 "--repo",
2137 "codewhale-hq/CodeWhale",
2138 "--pr",
2139 "1234",
2140 "--agent",
2141 "keel"
2142 ]
2143 );
2144 assert_eq!(
2145 auto_merge_pr_args("codewhale-hq/CodeWhale", "1234", "keel")[0],
2146 "scripts/auto-merge-pr.py"
2147 );
2148 let deny = evaluate_auto_merge(
2149 AutoMergeRequest {
2150 pr: "12",
2151 role: "keel",
2152 repo: "codewhale-hq/CodeWhale",
2153 },
2154 None,
2155 );
2156 assert!(matches!(deny, AutoMergeDecision::Deny { .. }));
2157 let _ = AUTO_MERGE_CHECKER_ENV;
2158 let _ = discover_auto_merge_checker(Path::new("/no-ops-here"));
2159 }
2160
2161 #[test]
2162 fn auto_merge_request_fields_are_validated_before_spawn() {
2163 let ok = |repo, pr, role| {
2164 validate_auto_merge_request(&AutoMergeRequest { pr, role, repo }).is_ok()
2165 };
2166 assert!(ok("codewhale-hq/CodeWhale", "1234", "keel"));
2167 assert!(ok("a-b/c.d_e-f", "1", "scout_2"));
2168 for (repo, pr, role) in [
2169 ("Hmbown", "1", "keel"),
2170 ("a/b/../../x", "1", "keel"),
2171 ("-x/y", "1", "keel"),
2172 ("x/-y", "1", "keel"),
2173 ("x/..", "1", "keel"),
2174 ("x y/z", "1", "keel"),
2175 ("x/y", "0", "keel"),
2176 ("x/y", "-1", "keel"),
2177 ("x/y", "1 2", "keel"),
2178 ("x/y", "12345678901", "keel"),
2179 ("x/y", "", "keel"),
2180 ("x/y", "1", ""),
2181 ("x/y", "1", "--fixture=/x"),
2182 ("x/y", "1", "keel ops"),
2183 ] {
2184 assert!(
2185 !ok(repo, pr, role),
2186 "{repo:?} {pr:?} {role:?} must be rejected"
2187 );
2188 }
2189 // A malformed request is denied even when a checker exists, so the
2190 // checker is never spawned with it.
2191 let dir = TempDir::new().expect("temp");
2192 let checker = dir.path().join("check-auto-merge.py");
2193 fs::write(&checker, "import sys\nsys.exit(0)\n").expect("write");
2194 assert!(matches!(
2195 evaluate_auto_merge(
2196 AutoMergeRequest {
2197 pr: "1",
2198 role: "--policy=x",
2199 repo: "x/y",
2200 },
2201 Some(&checker),
2202 ),
2203 AutoMergeDecision::Deny { .. }
2204 ));
2205 }
2206
2207 #[test]
2208 fn checker_exit_zero_allows() {
2209 let dir = TempDir::new().expect("temp");
2210 let checker = dir.path().join("check-auto-merge.py");
2211 fs::write(
2212 &checker,
2213 "#!/usr/bin/env python3\nimport argparse, sys\np=argparse.ArgumentParser()\np.add_argument('--repo')\np.add_argument('--pr')\np.add_argument('--agent', required=True)\np.parse_args()\nsys.exit(0)\n",
2214 )
2215 .expect("write");
2216 assert_eq!(
2217 evaluate_auto_merge(
2218 AutoMergeRequest {
2219 pr: "42",
2220 role: "keel",
2221 repo: "codewhale-hq/CodeWhale",
2222 },
2223 Some(&checker),
2224 ),
2225 AutoMergeDecision::Allow
2226 );
2227 }
2228
2229 #[test]
2230 fn plan_board_localizes_chrome_but_not_contract_tokens() {
2231 let mut op = with_credentials(Operation::new("Scout\nWrite", None));
2232 op.plan_from_direction();
2233 let english = render_plan_board_locale(&op, codewhale_localization::Locale::En);
2234 assert!(english.contains("gantt time →"), "{english}");
2235 assert!(english.contains("burn No cap"), "{english}");
2236 let japanese = render_plan_board_locale(&op, codewhale_localization::Locale::Ja);
2237 assert!(japanese.contains("ガント"), "{japanese}");
2238 // Contract tokens stay verbatim in every locale.
2239 assert!(japanese.contains("slice-1"), "{japanese}");
2240 assert!(japanese.contains(&op.id), "{japanese}");
2241 }
2242
2243 #[test]
2244 fn keepalive_automation_and_defaults() {
2245 let dir = TempDir::new().expect("temp");
2246 let manager = AutomationManager::open_for_test(dir.path().to_path_buf()).expect("manager");
2247 upsert_keepalive(&manager, dir.path(), false, &route_fixture_config(), None)
2248 .expect("upsert");
2249 let record = manager
2250 .get_automation(OPERATE_KEEPALIVE_ID)
2251 .expect("keepalive");
2252 assert_eq!(record.model.as_deref(), Some("same-model"));
2253 assert_eq!(record.mode.as_deref(), Some("operate"));
2254 assert_eq!(record.cwds, vec![dir.path().to_path_buf()]);
2255 assert_eq!(
2256 record.next_run_at, None,
2257 "without a kick the hourly recurrence owns the next run"
2258 );
2259 assert_eq!(OPERATE_MAX_WRITERS, 3);
2260 }
2261 }
2262
2262 lines RUST