返回 CodeWhale
dispatch_runner.rs
根目录 / crates / tui / src / dispatch_runner.rs
1 //! Cloud-dispatch remote runner: confirmed job → sandbox → forge PR.
2 //!
3 //! This module is the engine behind `confirm_job` /
4 //! [`crate::cloud_dispatch::execute_dispatch`] with `confirm`: it drives one
5 //! cloud job from `launching` through a Codewhale-operated sandbox to a pull
6 //! request on the target forge (`github` | `cnb` | `gitee`), then tears the
7 //! sandbox down on completion, failure, or cancellation. The persisted
8 //! contract, credential discovery, and fail-closed membership gate stay in
9 //! [`crate::cloud_dispatch`]; nothing here re-implements them.
10 //!
11 //! Invariants:
12 //!
13 //! - one harness: the sandbox runs the same `codewhale exec --auto`
14 //! one-shot entry every local non-interactive caller uses (one
15 //! `Engine::run_turn` inside); the runner itself runs no model turn.
16 //! - credentials never widen: the sandbox credential stays inside
17 //! [`crate::cloud_dispatch::LiveDaytonaLauncher`]; forge tokens are read
18 //! from Codewhale service slots only at PR-open time and are never
19 //! printed, logged, or persisted into job records.
20 //! - fail closed: every phase that cannot honestly complete records a
21 //! `failed` (or keeps `canceled`) job with a truthful note; a PR URL is
22 //! never invented.
23 //! - teardown always runs: cancel, failure, and success all attempt sandbox
24 //! teardown; the job note records whether it succeeded.
25 //! - orphans reconcile: the TUI detaches this runner, so quitting the TUI
26 //! can orphan a live sandbox. The job record persists a create intent
27 //! before the POST, every sandbox is labeled with its job id, and
28 //! [`startup_reconcile`] fails stale active jobs and deletes labeled
29 //! sandboxes whose job no longer needs them.
30
31 use std::path::Path;
32 use std::process::Command;
33 use std::sync::Mutex;
34
35 use anyhow::{Context, Result, anyhow, bail};
36
37 use crate::cloud_dispatch::{
38 self, CloudJob, CloudJobStatus, CloudJobStore, DaytonaLauncher, Forge, HarnessCommand,
39 PatchReceipt, SANDBOX_WORKSPACE, SandboxReceipt, sanitize_error, unix_timestamp,
40 validate_outbound_origin,
41 };
42 use crate::dependencies::ExternalTool;
43
44 /// Ceiling for PR titles (forges truncate longer titles).
45 const MAX_TITLE_CHARS: usize = 96;
46 /// Ceiling for the PR body, kept well under forge limits.
47 const MAX_BODY_CHARS: usize = 6_000;
48 /// Harness timeout for one cloud-agent turn.
49 const HARNESS_TIMEOUT_SECS: u32 = 3_600;
50
51 /// The pull request the forge opened (or that `gh` reported).
52 #[derive(Debug, Clone, PartialEq, Eq)]
53 pub struct PrOpened {
54 pub url: String,
55 /// SHA actually applied and pushed — not the sandbox `patch.head_sha`.
56 pub head_sha: String,
57 }
58
59 /// Forge seam: raise the agent's branch and open the PR. Tests inject a
60 /// recorder; production uses [`LiveForgePr`].
61 pub trait ForgePr {
62 fn open(&self, job: &CloudJob, patch: &PatchReceipt) -> Result<PrOpened>;
63 }
64
65 /// Run one confirmed job end to end.
66 ///
67 /// Requires the job to be `launching` (or `running`, for a resumed runner):
68 /// a `proposed` job is refused — confirmation is the caller's explicit act,
69 /// never the runner's. Every phase persists its transition so `/dispatch
70 /// show` and `/jobs` stream real progress, and a `canceled` record at any
71 /// checkpoint stops the run and tears the sandbox down.
72 pub fn run_confirmed_job(
73 store: &CloudJobStore,
74 id: &str,
75 launcher: &dyn DaytonaLauncher,
76 forge: &dyn ForgePr,
77 ) -> Result<CloudJob> {
78 let mut job = store.load(id)?;
79 if !matches!(
80 job.status,
81 CloudJobStatus::Launching | CloudJobStatus::Running
82 ) {
83 bail!(
84 "Cloud job {id} is {} and cannot be run; confirm it first with `/dispatch confirm {id}`.",
85 status_word(job.status)
86 );
87 }
88 match drive(store, &mut job, launcher, forge) {
89 Ok(()) => store.load(id),
90 Err(error) => {
91 let message = sanitize_error(&error.to_string());
92 // Re-load before writing any failure: a job the user canceled
93 // stays canceled — a later run error must not overwrite the
94 // user's terminal word with `failed`. The error is appended to
95 // the note instead. And if a cancel lands inside the load→save
96 // span of the failure write, the disk record (already
97 // `canceled`, with the user's note) wins and is left alone.
98 let current = match store.load(id) {
99 Ok(mut record) if record.status == CloudJobStatus::Canceled => {
100 record.finished_unix =
101 Some(record.finished_unix.unwrap_or_else(unix_timestamp));
102 record.note = format!("{}. Run error after cancel: {message}", record.note);
103 let _ = store.save(&record);
104 record
105 }
106 loaded => {
107 let mut failed = loaded.unwrap_or_else(|_| job.clone());
108 failed.status = CloudJobStatus::Failed;
109 failed.refusal = Some(message.clone());
110 failed.finished_unix = Some(unix_timestamp());
111 failed.note = format!("Cloud agent run failed closed. {message}");
112 match store.save_unless_canceled(&failed) {
113 Ok(true) => failed,
114 _ => store.load(id).unwrap_or(failed),
115 }
116 }
117 };
118 teardown_best_effort(launcher, &current);
119 Err(error)
120 }
121 }
122 }
123
124 /// Background runner used by the CLI and TUI confirm paths. The store is
125 /// the source of truth; the thread result is intentionally dropped
126 /// (failures are recorded inside the job record). The CLI joins the handle
127 /// before exiting so a confirmed run is never orphaned mid-flight; the TUI
128 /// detaches it (the job record keeps the truth across restarts, and
129 /// `/dispatch cancel` tears a live sandbox down at any time).
130 pub fn spawn_confirmed_runner(
131 store: CloudJobStore,
132 id: String,
133 ) -> Option<std::thread::JoinHandle<()>> {
134 std::thread::Builder::new()
135 .name(format!("cw-dispatch-{id}"))
136 .spawn(move || {
137 let launcher = cloud_dispatch::LiveDaytonaLauncher;
138 let forge = LiveForgePr;
139 let _ = run_confirmed_job(&store, &id, &launcher, &forge);
140 })
141 .ok()
142 }
143
144 /// Best-effort startup reconciliation for the detached runner.
145 ///
146 /// The TUI spawns [`spawn_confirmed_runner`] detached, so quitting the TUI
147 /// or crashing mid-run orphans the record (and possibly a billing sandbox)
148 /// with nothing to reconcile it. Two passes, in order:
149 ///
150 /// 1. [`cloud_dispatch::sweep_stale_jobs`] — active records older than the
151 /// declared harness budget plus slack are failed and their recorded
152 /// sandboxes torn down;
153 /// 2. [`cloud_dispatch::reconcile_sandboxes`] — any dispatch-labeled
154 /// sandbox whose job is terminal or absent from the store is deleted by
155 /// label, covering creates whose id was never recorded.
156 ///
157 /// Never fatal and never blocks the caller's critical path beyond the
158 /// launcher's own bounded HTTP budget; returns a human receipt for the log
159 /// (empty when there was nothing to do).
160 pub fn startup_reconcile(store: &CloudJobStore, launcher: &dyn DaytonaLauncher) -> String {
161 let swept = cloud_dispatch::sweep_stale_jobs(store, launcher, cloud_dispatch::unix_timestamp());
162 let mut lines = Vec::new();
163 for job in &swept {
164 lines.push(format!(
165 "cloud dispatch startup sweep: job {} marked stale (failed) and its sandbox teardown attempted",
166 job.id
167 ));
168 }
169 match cloud_dispatch::reconcile_sandboxes(store, launcher) {
170 Ok(report) if !report.deleted.is_empty() => {
171 lines.push(format!(
172 "cloud dispatch label reconcile: deleted orphaned sandbox(es) {}",
173 report.deleted.join(", ")
174 ));
175 }
176 Ok(_) => {}
177 Err(error) => lines.push(format!(
178 "cloud dispatch label reconcile skipped: {}",
179 sanitize_error(&error.to_string())
180 )),
181 }
182 lines.join("\n")
183 }
184
185 fn drive(
186 store: &CloudJobStore,
187 job: &mut CloudJob,
188 launcher: &dyn DaytonaLauncher,
189 forge: &dyn ForgePr,
190 ) -> Result<()> {
191 // Launching → Running: create the sandbox. The intent record goes down
192 // BEFORE the POST: if the create response is slow and the client gives
193 // up (or the process dies), the sandbox may still come into being with
194 // no recorded id — `sandbox_pending` is what cancel and the label
195 // reconciler use to find and delete it by label.
196 //
197 // Every phase save below is cancel-authoritative
198 // (`save_unless_canceled`): a cancel that lands while a phase is in
199 // flight wins over the runner's read-modify-write, so a canceled job
200 // can never be resurrected into a later phase — above all never into
201 // the branch push / PR open.
202 job.sandbox_pending = true;
203 if !store.save_unless_canceled(job)? {
204 // Canceled before any sandbox exists: nothing to tear down, and the
205 // canceled record on disk stays the truth.
206 *job = store.load(&job.id)?;
207 return Ok(());
208 }
209 let receipt = launcher.create_sandbox(job)?;
210 job.status = CloudJobStatus::Running;
211 job.sandbox_pending = false;
212 job.sandbox_id = Some(receipt.sandbox_id.clone());
213 job.note = format!(
214 "Sandbox {} created; the Codewhale cloud agent turn is running.",
215 receipt.sandbox_id
216 );
217 if !store.save_unless_canceled(job)? {
218 return finish_canceled(store, job, launcher, &receipt);
219 }
220
221 launcher.wait_ready(&receipt)?;
222 let clone_url = cloud_dispatch::validate_git_remote_url(&job.remote_url)?;
223 launcher.clone_repository(&receipt, &clone_url, SANDBOX_WORKSPACE)?;
224 if cancel_requested(store, job)? {
225 return finish_canceled(store, job, launcher, &receipt);
226 }
227
228 // One agent turn through the standard one-shot harness entry.
229 let output = launcher.run_harness(&receipt, &harness_command(job))?;
230 job.agent_summary = Some(summary_line(&output));
231 if cancel_requested(store, job)? {
232 return finish_canceled(store, job, launcher, &receipt);
233 }
234
235 // Running → OpeningPr: collect the agent's work product.
236 let patch = launcher.collect_patch(&receipt)?;
237 job.status = CloudJobStatus::OpeningPr;
238 job.base_branch = Some(patch.base_branch.clone());
239 job.head_sha = Some(patch.head_sha.clone()); // replaced with the pushed sha after `forge.open`
240 job.note = format!(
241 "Agent turn complete ({}); raising branch {} and opening the PR on {}.",
242 patch.summary,
243 job.branch,
244 job.forge.as_str()
245 );
246 if !store.save_unless_canceled(job)? {
247 return finish_canceled(store, job, launcher, &receipt);
248 }
249
250 // OpeningPr → Done: push the branch and open the PR. The cancel check
251 // is the last gate before money-adjacent side effects on the forge.
252 //
253 // Known limitation: the push and PR open are remote side effects with no
254 // transaction to join, so a cancel that lands after this check cannot
255 // stop them. The outcome is then stated, not hidden: the record stays
256 // `canceled` (the user's word) and carries the PR URL the forge returned.
257 // A `done` job is terminal; a cancel after it changes nothing.
258 if cancel_requested(store, job)? {
259 return finish_canceled(store, job, launcher, &receipt);
260 }
261 let opened = forge.open(job, &patch)?;
262 job.status = CloudJobStatus::Done;
263 job.pr_url = Some(opened.url.clone());
264 job.head_sha = Some(opened.head_sha.clone());
265 job.finished_unix = Some(unix_timestamp());
266 job.note = format!(
267 "Cloud agent finished; PR opened at {}. {}",
268 opened.url,
269 teardown_note(launcher, &receipt)
270 );
271 if !store.save_unless_canceled(job)? {
272 // A cancel landed while the PR was opening. The forge accepted it,
273 // so the PR exists: keep its URL and say exactly that rather than
274 // claiming success or silently dropping the receipt.
275 let mut canceled = store.load(&job.id)?;
276 canceled.pr_url = job.pr_url.clone();
277 canceled.agent_summary = job.agent_summary.clone();
278 canceled.finished_unix = Some(canceled.finished_unix.unwrap_or_else(unix_timestamp));
279 canceled.note = format!(
280 "Canceled after the PR had already opened at {}; close it on the forge if it is unwanted. {}",
281 opened.url,
282 teardown_note(launcher, &receipt)
283 );
284 *job = canceled.clone();
285 return store.save(&canceled).map(|_| ());
286 }
287 Ok(())
288 }
289
290 /// True when the user canceled the job mid-run.
291 fn cancel_requested(store: &CloudJobStore, job: &CloudJob) -> Result<bool> {
292 Ok(store.load(&job.id)?.status == CloudJobStatus::Canceled)
293 }
294
295 fn finish_canceled(
296 store: &CloudJobStore,
297 job: &mut CloudJob,
298 launcher: &dyn DaytonaLauncher,
299 receipt: &SandboxReceipt,
300 ) -> Result<()> {
301 let mut canceled = store.load(&job.id)?;
302 canceled.agent_summary = job.agent_summary.clone();
303 canceled.finished_unix = Some(canceled.finished_unix.unwrap_or_else(unix_timestamp));
304 canceled.note = format!("Canceled mid-run. {}", teardown_note(launcher, receipt));
305 *job = canceled.clone();
306 store.save(&canceled)
307 }
308
309 fn teardown_best_effort(launcher: &dyn DaytonaLauncher, job: &CloudJob) {
310 if let Some(sandbox_id) = job.sandbox_id.clone() {
311 let receipt = SandboxReceipt {
312 sandbox_id,
313 toolbox_url: None,
314 };
315 let _ = launcher.teardown(&receipt);
316 }
317 }
318
319 fn teardown_note(launcher: &dyn DaytonaLauncher, receipt: &SandboxReceipt) -> String {
320 match launcher.teardown(receipt) {
321 Ok(()) => "The sandbox was torn down.".to_string(),
322 Err(error) => format!(
323 "Sandbox teardown failed and may need a retry: {}",
324 sanitize_error(&error.to_string())
325 ),
326 }
327 }
328
329 /// The exact harness invocation the sandbox runs: the one-shot
330 /// `codewhale exec --auto` entry — the same single-`Engine::run_turn` path
331 /// local non-interactive callers use, never a second engine.
332 pub fn harness_command(job: &CloudJob) -> HarnessCommand {
333 HarnessCommand {
334 argv: vec![
335 "codewhale".to_string(),
336 "exec".to_string(),
337 "--auto".to_string(),
338 job.prompt.clone(),
339 ],
340 cwd: SANDBOX_WORKSPACE.to_string(),
341 timeout_secs: HARNESS_TIMEOUT_SECS,
342 }
343 }
344
345 /// First non-empty line of harness output, bounded for notes and PR bodies.
346 pub fn summary_line(output: &str) -> String {
347 // Redacted first: the sandbox env carries the account machine token, and
348 // harness output must not be able to echo it into the job record.
349 crate::cloud_dispatch::redact_machine_tokens(output)
350 .lines()
351 .map(str::trim)
352 .find(|line| !line.is_empty())
353 .map(|line| line.chars().take(200).collect())
354 .unwrap_or_else(|| "cloud agent turn completed".to_string())
355 }
356
357 /// Owner/repo slug for a forge remote URL, or `None` when the URL does not
358 /// match the job's forge. Both https and `git@host:owner/repo.git` shapes
359 /// are accepted; a trailing `.git` is stripped.
360 pub fn forge_slug(forge: Forge, remote_url: &str) -> Option<String> {
361 if cloud_dispatch::classify_url(remote_url) != Some(forge) {
362 return None;
363 }
364 // `https://host/owner/repo.git` or `git@host:owner/repo.git`
365 let repo_path = if let Some((_, rest)) = remote_url.split_once("://") {
366 // Strip the host: the slug is the two path segments after it.
367 rest.split_once('/')?.1
368 } else {
369 remote_url.split_once(':')?.1
370 };
371 let segments: Vec<&str> = repo_path.split('/').collect();
372 if segments.len() < 2 {
373 return None;
374 }
375 let repo = segments[segments.len() - 1].trim_end_matches(".git");
376 let owner = segments[segments.len() - 2];
377 if owner.is_empty() || repo.is_empty() || repo == owner {
378 return None;
379 }
380 Some(format!("{owner}/{repo}"))
381 }
382
383 /// Truthful PR title: names the agent and its own summary.
384 pub fn compose_pr_title(job: &CloudJob, patch: &PatchReceipt) -> String {
385 let summary = if patch.summary.trim().is_empty() {
386 one_line(&job.prompt, 72)
387 } else {
388 one_line(&patch.summary, 72)
389 };
390 one_line(&format!("codewhale cloud: {summary}"), MAX_TITLE_CHARS)
391 }
392
393 /// Truthful PR body: what the agent did, the receipts Codewhale has, and an
394 /// explicit No-Issue line (no tracked issue; the cloud job is the record).
395 /// Sandbox-provider names never appear — the operator is Codewhale.
396 pub fn compose_pr_body(job: &CloudJob, patch: &PatchReceipt) -> String {
397 compose_pr_body_for_head(job, patch, &patch.head_sha)
398 }
399
400 /// PR body whose `Head:` is the sha actually applied/pushed.
401 pub fn compose_pr_body_for_head(job: &CloudJob, patch: &PatchReceipt, head_sha: &str) -> String {
402 let summary = if patch.summary.trim().is_empty() {
403 "(the agent's commit list is the record)".to_string()
404 } else {
405 patch.summary.trim().to_string()
406 };
407 let body = format!(
408 "Automated change by a Codewhale cloud agent.\n\n\
409 ## What the agent did\n{summary}\n\n\
410 ## Task\n{}\n\n\
411 ## Receipts\n\
412 - Cloud job: {}\n\
413 - Sandbox: {}\n\
414 - Branch: `{}` (base `{}`)\n\
415 - Head: `{}`\n\n\
416 No-Issue: cloud dispatch {} (no tracked issue; receipts above)",
417 job.prompt,
418 job.id,
419 job.sandbox_id.as_deref().unwrap_or("(pending)"),
420 job.branch,
421 patch.base_branch,
422 head_sha,
423 job.id,
424 );
425 one_line(&body, MAX_BODY_CHARS)
426 }
427
428 /// `gh pr create` argv for the GitHub path. The body rides in a file so the
429 /// prompt text never appears in `ps` output.
430 pub fn gh_pr_create_argv(
431 slug: &str,
432 base: &str,
433 head: &str,
434 title: &str,
435 body_file: &str,
436 ) -> Vec<String> {
437 vec![
438 "pr".to_string(),
439 "create".to_string(),
440 "--repo".to_string(),
441 slug.to_string(),
442 "--base".to_string(),
443 base.to_string(),
444 "--head".to_string(),
445 head.to_string(),
446 "--title".to_string(),
447 title.to_string(),
448 "--body-file".to_string(),
449 body_file.to_string(),
450 ]
451 }
452
453 /// Gitee v5 pull-request endpoint for a slug.
454 pub fn gitee_pr_url(slug: &str) -> String {
455 format!("https://gitee.com/api/v5/repos/{slug}/pulls")
456 }
457
458 /// CNB pull-request endpoint for a slug.
459 pub fn cnb_pr_url(slug: &str) -> String {
460 format!("https://api.cnb.cool/{slug}/-/pulls")
461 }
462
463 /// Production forge opener: local git for the branch push, then the forge's
464 /// own API surface for the PR (gh for GitHub; REST for CNB and Gitee with
465 /// service-slot tokens). Fails closed on missing tooling or tokens — the
466 /// branch push is not rolled back, and no PR URL is ever invented.
467 pub struct LiveForgePr;
468
469 impl ForgePr for LiveForgePr {
470 fn open(&self, job: &CloudJob, patch: &PatchReceipt) -> Result<PrOpened> {
471 let slug = forge_slug(job.forge, &job.remote_url).ok_or_else(|| {
472 anyhow!(
473 "the {} remote {} does not resolve to an owner/repo slug",
474 job.forge.as_str(),
475 job.remote_url
476 )
477 })?;
478 let title = compose_pr_title(job, patch);
479 let dir = tempfile::tempdir().context("could not stage the cloud agent branch")?;
480 let head_sha = prepare_branch(job, patch, dir.path())?;
481 push_branch(&dir.path().join("repo"), &job.remote_url, &job.branch)?;
482 let body = compose_pr_body_for_head(job, patch, &head_sha);
483 let url = match job.forge {
484 Forge::Github => open_pr_github(&slug, job, patch, &title, &body)?,
485 Forge::Gitee => open_pr_gitee(&slug, job, patch, &title, &body)?,
486 Forge::Cnb => open_pr_cnb(&slug, job, patch, &title, &body)?,
487 };
488 Ok(PrOpened { url, head_sha })
489 }
490 }
491
492 /// Shallow-clone the target repository, apply the agent's patch on a branch,
493 /// and return the head sha. Plain git subprocess work; never forceful.
494 fn prepare_branch(job: &CloudJob, patch: &PatchReceipt, dir: &Path) -> Result<String> {
495 let remote_url = cloud_dispatch::validate_git_remote_url(&job.remote_url)?;
496 git(
497 None,
498 &[
499 "clone",
500 "--quiet",
501 "--depth",
502 "50",
503 "--",
504 &remote_url,
505 &dir.join("repo").to_string_lossy(),
506 ],
507 )
508 .context("could not clone the target repository for the cloud agent branch")?;
509 let repo = dir.join("repo");
510 let patch_path = dir.join("agent.patch");
511 std::fs::write(&patch_path, &patch.patch).context("could not stage the agent patch")?;
512 git(
513 Some(&repo),
514 &["config", "user.name", "Codewhale Cloud Agent"],
515 )
516 .context("could not set the agent identity")?;
517 git(
518 Some(&repo),
519 &["config", "user.email", "cloud-agent@codewhale.invalid"],
520 )
521 .context("could not set the agent identity")?;
522 git(Some(&repo), &["checkout", "--quiet", "-b", &job.branch])
523 .context("could not create the cloud agent branch")?;
524 git(
525 Some(&repo),
526 &["am", "--quiet", "--3way", &patch_path.to_string_lossy()],
527 )
528 .context("the agent patch did not apply cleanly onto the target branch")?;
529 git(Some(&repo), &["rev-parse", "HEAD"]).map(|out| out.trim().to_string())
530 }
531
532 /// Push the prepared branch. Plain push only — `--force` is never passed, so
533 /// an existing branch that is not a fast-forward fails closed instead of
534 /// rewriting the target's history.
535 fn push_branch(repo: &Path, remote_url: &str, branch: &str) -> Result<()> {
536 let remote_url = cloud_dispatch::validate_git_remote_url(remote_url)?;
537 if cloud_dispatch::is_forge_default_branch(branch) {
538 bail!("refusing to push onto the forge default branch {branch}");
539 }
540 if looks_like_network_remote(&remote_url) && cloud_dispatch::classify_url(&remote_url).is_none()
541 {
542 bail!("refusing to push to a non-forge remote");
543 }
544 if remote_branch_exists(&remote_url, branch)? {
545 bail!("refusing to update existing remote branch {branch}");
546 }
547 git(
548 Some(repo),
549 &[
550 "push",
551 "--quiet",
552 "--",
553 &remote_url,
554 &format!("HEAD:refs/heads/{branch}"),
555 ],
556 )
557 .map(|_| ())
558 .context(
559 "could not push the cloud agent branch (it may need credentials or the branch may have moved)",
560 )
561 }
562
563 fn looks_like_network_remote(url: &str) -> bool {
564 url.contains("://") || url.contains('@')
565 }
566
567 fn remote_branch_exists(remote_url: &str, branch: &str) -> Result<bool> {
568 let listing =
569 git(None, &["ls-remote", "--heads", "--", remote_url, branch]).unwrap_or_default();
570 Ok(listing
571 .lines()
572 .any(|line| line.contains(&format!("refs/heads/{branch}"))))
573 }
574
575 fn git(cwd: Option<&Path>, args: &[&str]) -> Result<String> {
576 let mut command = Command::new("git");
577 if let Some(cwd) = cwd {
578 command.current_dir(cwd);
579 }
580 let output = command
581 .args(args)
582 .output()
583 .context("failed to start git for the cloud agent branch")?;
584 if !output.status.success() {
585 bail!(
586 "git {} failed: {}",
587 args.first().unwrap_or(&""),
588 sanitize_error(&String::from_utf8_lossy(&output.stderr))
589 );
590 }
591 Ok(String::from_utf8_lossy(&output.stdout).to_string())
592 }
593
594 fn open_pr_github(
595 slug: &str,
596 job: &CloudJob,
597 patch: &PatchReceipt,
598 title: &str,
599 body: &str,
600 ) -> Result<String> {
601 let body_dir = tempfile::tempdir().context("could not stage the PR body")?;
602 let body_file = body_dir.path().join("body.md");
603 std::fs::write(&body_file, body).context("could not write the PR body")?;
604 let mut command = crate::dependencies::Gh::command()
605 .ok_or_else(|| anyhow!("the GitHub pull request needs the gh CLI on PATH"))?;
606 command.args(gh_pr_create_argv(
607 slug,
608 &patch.base_branch,
609 &job.branch,
610 title,
611 &body_file.to_string_lossy(),
612 ));
613 let output = command
614 .output()
615 .context("failed to start gh for the pull request")?;
616 if !output.status.success() {
617 bail!(
618 "gh pr create failed: {}",
619 sanitize_error(&String::from_utf8_lossy(&output.stderr))
620 );
621 }
622 let url = String::from_utf8_lossy(&output.stdout).trim().to_string();
623 if !url.starts_with("https://") {
624 bail!("gh did not report a pull request URL; refusing to invent one.");
625 }
626 Ok(url)
627 }
628
629 fn open_pr_gitee(
630 slug: &str,
631 job: &CloudJob,
632 patch: &PatchReceipt,
633 title: &str,
634 body: &str,
635 ) -> Result<String> {
636 let token = read_service_token("gitee").ok_or_else(|| {
637 anyhow!("a Gitee access token is not configured in the Codewhale service slot; the branch was pushed but no pull request was opened")
638 })?;
639 let url = validate_outbound_origin(&gitee_pr_url(slug))?;
640 let response = crate::tls::reqwest_blocking_client_builder()
641 .connect_timeout(std::time::Duration::from_secs(8))
642 .timeout(std::time::Duration::from_secs(30))
643 .redirect(reqwest::redirect::Policy::none())
644 .build()
645 .context("could not initialize the Gitee client")?
646 .post(url)
647 .form(&[
648 ("access_token", token.as_str()),
649 ("title", title),
650 ("head", job.branch.as_str()),
651 ("base", patch.base_branch.as_str()),
652 ("body", body),
653 ])
654 .send()
655 .context("could not reach Gitee")?;
656 let status = response.status();
657 let text = response.text().unwrap_or_default();
658 if !status.is_success() {
659 bail!("Gitee pull request create failed (HTTP {status}).");
660 }
661 let parsed: serde_json::Value =
662 serde_json::from_str(&text).context("Gitee returned invalid JSON")?;
663 parsed
664 .get("html_url")
665 .and_then(serde_json::Value::as_str)
666 .map(str::trim)
667 .filter(|url| url.starts_with("https://"))
668 .map(|url| url.to_string())
669 .ok_or_else(|| anyhow!("Gitee did not report a pull request URL; refusing to invent one."))
670 }
671
672 fn open_pr_cnb(
673 slug: &str,
674 job: &CloudJob,
675 patch: &PatchReceipt,
676 title: &str,
677 body: &str,
678 ) -> Result<String> {
679 let token = read_service_token("cnb").ok_or_else(|| {
680 anyhow!("a CNB access token is not configured in the Codewhale service slot; the branch was pushed but no pull request was opened")
681 })?;
682 let url = validate_outbound_origin(&cnb_pr_url(slug))?;
683 let response = crate::tls::reqwest_blocking_client_builder()
684 .connect_timeout(std::time::Duration::from_secs(8))
685 .timeout(std::time::Duration::from_secs(30))
686 .redirect(reqwest::redirect::Policy::none())
687 .build()
688 .context("could not initialize the CNB client")?
689 .post(url)
690 .bearer_auth(&token)
691 .json(&serde_json::json!({
692 "title": title,
693 "head": job.branch,
694 "base": patch.base_branch,
695 "body": body,
696 }))
697 .send()
698 .context("could not reach CNB")?;
699 let status = response.status();
700 let text = response.text().unwrap_or_default();
701 if !status.is_success() {
702 bail!("CNB pull request create failed (HTTP {status}).");
703 }
704 let parsed: serde_json::Value =
705 serde_json::from_str(&text).context("CNB returned invalid JSON")?;
706 let number = parsed
707 .get("number")
708 .and_then(serde_json::Value::as_i64)
709 .filter(|number| *number > 0)
710 .ok_or_else(|| {
711 anyhow!("CNB did not report a pull request number; refusing to invent a URL.")
712 })?;
713 Ok(format!("https://cnb.cool/{slug}/-/pulls/{number}"))
714 }
715
716 /// Read a forge token from the Codewhale service slot. Never logged.
717 fn read_service_token(slot: &str) -> Option<String> {
718 codewhale_secrets::Secrets::auto_detect()
719 .get(slot)
720 .ok()
721 .flatten()
722 .map(|value| value.trim().to_string())
723 .filter(|value| !value.is_empty())
724 }
725
726 fn one_line(value: &str, max: usize) -> String {
727 let flat: String = value
728 .chars()
729 .map(|ch| {
730 if ch.is_control() && ch != '\n' {
731 ' '
732 } else {
733 ch
734 }
735 })
736 .collect();
737 if flat.chars().count() <= max {
738 flat
739 } else {
740 let mut out: String = flat.chars().take(max.saturating_sub(1)).collect();
741 out.push('…');
742 out
743 }
744 }
745
746 fn status_word(status: CloudJobStatus) -> &'static str {
747 match status {
748 CloudJobStatus::Proposed => "proposed",
749 CloudJobStatus::Refused => "refused",
750 CloudJobStatus::Launching => "launching",
751 CloudJobStatus::Running => "running",
752 CloudJobStatus::OpeningPr => "openingpr",
753 CloudJobStatus::Done => "done",
754 CloudJobStatus::Failed => "failed",
755 CloudJobStatus::Canceled => "canceled",
756 }
757 }
758
759 /// Callback fired at each launcher phase boundary (used by tests to
760 /// simulate mid-run cancellation).
761 pub type PhaseHook = Box<dyn Fn(&str) + Send + Sync>;
762
763 /// Recording launcher for offline tests. Records every phase by name and
764 /// replays canned results; `hook` fires at each phase boundary so tests can
765 /// simulate mid-run cancellation.
766 pub struct RecordingLauncher {
767 sandbox_id: String,
768 patch: PatchReceipt,
769 calls: Mutex<Vec<String>>,
770 pub hook: Option<PhaseHook>,
771 /// Sandboxes reported by `list_job_sandboxes` (the reconciler seam).
772 pub listed: Mutex<Vec<crate::cloud_dispatch::LabeledSandbox>>,
773 }
774
775 impl RecordingLauncher {
776 pub fn new(sandbox_id: &str, patch: PatchReceipt) -> Self {
777 Self {
778 sandbox_id: sandbox_id.to_string(),
779 patch,
780 calls: Mutex::new(Vec::new()),
781 hook: None,
782 listed: Mutex::new(Vec::new()),
783 }
784 }
785
786 pub fn calls(&self) -> Vec<String> {
787 self.calls
788 .lock()
789 .map(|calls| calls.clone())
790 .unwrap_or_default()
791 }
792
793 fn record(&self, phase: &str) {
794 if let Ok(mut calls) = self.calls.lock() {
795 calls.push(phase.to_string());
796 }
797 if let Some(hook) = self.hook.as_ref() {
798 hook(phase);
799 }
800 }
801 }
802
803 impl DaytonaLauncher for RecordingLauncher {
804 fn create_sandbox(&self, _job: &CloudJob) -> Result<SandboxReceipt> {
805 self.record("create");
806 Ok(SandboxReceipt {
807 sandbox_id: self.sandbox_id.clone(),
808 toolbox_url: Some("https://toolbox.example.test".to_string()),
809 })
810 }
811
812 fn wait_ready(&self, _receipt: &SandboxReceipt) -> Result<()> {
813 self.record("wait_ready");
814 Ok(())
815 }
816
817 fn clone_repository(&self, _receipt: &SandboxReceipt, url: &str, path: &str) -> Result<()> {
818 self.record("clone");
819 assert_eq!(path, SANDBOX_WORKSPACE);
820 assert!(
821 url.starts_with("https://"),
822 "live clone must be an https URL"
823 );
824 Ok(())
825 }
826
827 fn run_harness(&self, _receipt: &SandboxReceipt, _command: &HarnessCommand) -> Result<String> {
828 self.record("harness");
829 Ok("Fixed the flaky test and re-ran the suite.\nall green".to_string())
830 }
831
832 fn collect_patch(&self, _receipt: &SandboxReceipt) -> Result<PatchReceipt> {
833 self.record("collect");
834 Ok(self.patch.clone())
835 }
836
837 fn teardown(&self, _receipt: &SandboxReceipt) -> Result<()> {
838 self.record("teardown");
839 Ok(())
840 }
841
842 fn list_job_sandboxes(&self) -> Result<Vec<crate::cloud_dispatch::LabeledSandbox>> {
843 self.record("list");
844 Ok(self
845 .listed
846 .lock()
847 .map(|listed| listed.clone())
848 .unwrap_or_default())
849 }
850 }
851
852 /// Recording forge opener for offline tests.
853 pub struct RecordingForgePr {
854 pub url: String,
855 opened: Mutex<Vec<String>>,
856 }
857
858 impl RecordingForgePr {
859 pub fn new(url: &str) -> Self {
860 Self {
861 url: url.to_string(),
862 opened: Mutex::new(Vec::new()),
863 }
864 }
865
866 pub fn opened(&self) -> Vec<String> {
867 self.opened
868 .lock()
869 .map(|opened| opened.clone())
870 .unwrap_or_default()
871 }
872 }
873
874 impl ForgePr for RecordingForgePr {
875 fn open(&self, job: &CloudJob, patch: &PatchReceipt) -> Result<PrOpened> {
876 if let Ok(mut opened) = self.opened.lock() {
877 opened.push(format!("{}:{}", job.id, patch.head_sha));
878 }
879 Ok(PrOpened {
880 url: self.url.clone(),
881 head_sha: patch.head_sha.clone(),
882 })
883 }
884 }
885
886 #[cfg(test)]
887 mod tests {
888 use super::*;
889 use crate::cloud_dispatch::{
890 CloudJobStatus, CredentialSource, CredentialState, DispatchOutcome, GitRemote,
891 MachineTokenState, execute_dispatch, plan_dispatch,
892 };
893 use std::sync::Arc;
894
895 #[test]
896 fn summary_line_redacts_machine_tokens() {
897 let token = "cwc_key_0123456789abcdef01234567_AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
898 let summary = summary_line(&format!("done with {token}"));
899 assert!(!summary.contains("AAAAAAAA"));
900 assert!(summary.contains("[redacted]"));
901 }
902
903 fn fixture_patch() -> PatchReceipt {
904 PatchReceipt {
905 base_branch: "main".to_string(),
906 head_sha: "abc123def4567".to_string(),
907 summary: "Fix the flaky dispatch test".to_string(),
908 patch: "From abc123 Mon Sep 17 00:00:00 2001\nSubject: [PATCH] Fix the flake\n"
909 .to_string(),
910 }
911 }
912
913 fn confirmed_job(store: &CloudJobStore) -> CloudJob {
914 let plan = plan_dispatch(
915 &[GitRemote {
916 name: "github".to_string(),
917 url: "https://github.com/org/repo.git".to_string(),
918 }],
919 "open a PR that fixes the flake",
920 Some(Forge::Github),
921 Some("codewhale/cloud-runner-test"),
922 )
923 .unwrap();
924 match execute_dispatch(
925 store,
926 plan,
927 true,
928 &CredentialState::Present {
929 source: CredentialSource::Env,
930 },
931 &MachineTokenState::Present,
932 )
933 .unwrap()
934 {
935 DispatchOutcome::Accepted(job) => job,
936 other => panic!("expected accept, got {other:?}"),
937 }
938 }
939
940 #[test]
941 fn recording_lifecycle_reaches_done_with_receipts_and_teardown() {
942 let temp = tempfile::tempdir().unwrap();
943 let store = CloudJobStore::from_path(temp.path().join("jobs"));
944 let job = confirmed_job(&store);
945 let launcher = RecordingLauncher::new("sandbox_runner_1", fixture_patch());
946 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
947 let finished = run_confirmed_job(&store, &job.id, &launcher, &forge).unwrap();
948
949 assert_eq!(finished.status, CloudJobStatus::Done);
950 assert_eq!(
951 finished.pr_url.as_deref(),
952 Some("https://github.com/org/repo/pull/9")
953 );
954 assert_eq!(finished.sandbox_id.as_deref(), Some("sandbox_runner_1"));
955 assert_eq!(finished.base_branch.as_deref(), Some("main"));
956 assert_eq!(finished.head_sha.as_deref(), Some("abc123def4567"));
957 assert_eq!(
958 finished.agent_summary.as_deref(),
959 Some("Fixed the flaky test and re-ran the suite.")
960 );
961 assert!(finished.finished_unix.is_some());
962 assert!(finished.note.contains("PR opened at"));
963 assert!(finished.note.contains("torn down"));
964 // Full protocol order, teardown last.
965 assert_eq!(
966 launcher.calls(),
967 vec![
968 "create",
969 "wait_ready",
970 "clone",
971 "harness",
972 "collect",
973 "teardown"
974 ]
975 );
976 assert_eq!(forge.opened(), vec![format!("{}:abc123def4567", job.id)]);
977 // The persisted record streams the same truth.
978 assert_eq!(store.load(&job.id).unwrap().status, CloudJobStatus::Done);
979 }
980
981 #[test]
982 fn cancel_during_running_tears_down_and_never_opens_the_pr() {
983 let temp = tempfile::tempdir().unwrap();
984 let root = temp.path().join("jobs");
985 let store = CloudJobStore::from_path(root.clone());
986 let job = confirmed_job(&store);
987 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
988 // Cancel lands while the harness turn is in flight.
989 let cancel_root = root.clone();
990 let cancel_id = job.id.clone();
991 let canceled_seen = Arc::new(std::sync::atomic::AtomicBool::new(false));
992 let seen = canceled_seen.clone();
993 let mut launcher = RecordingLauncher::new("sandbox_runner_2", fixture_patch());
994 launcher.hook = Some(Box::new(move |phase| {
995 if phase == "harness" && !seen.swap(true, std::sync::atomic::Ordering::SeqCst) {
996 let store = CloudJobStore::from_path(cancel_root.clone());
997 let mut current = store.load(&cancel_id).unwrap();
998 current.status = CloudJobStatus::Canceled;
999 store.save(&current).unwrap();
1000 }
1001 }));
1002 let finished = run_confirmed_job(&store, &job.id, &launcher, &forge).unwrap();
1003
1004 assert_eq!(finished.status, CloudJobStatus::Canceled);
1005 assert!(finished.pr_url.is_none());
1006 assert!(finished.note.contains("Canceled mid-run"));
1007 assert!(finished.note.contains("torn down"));
1008 // The run stopped before collect and the forge never fired; teardown ran.
1009 assert_eq!(
1010 launcher.calls(),
1011 vec!["create", "wait_ready", "clone", "harness", "teardown"]
1012 );
1013 assert!(forge.opened().is_empty());
1014 assert!(canceled_seen.load(std::sync::atomic::Ordering::SeqCst));
1015 }
1016
1017 #[test]
1018 fn cancel_between_clone_and_harness_never_starts_the_turn() {
1019 let temp = tempfile::tempdir().unwrap();
1020 let root = temp.path().join("jobs");
1021 let store = CloudJobStore::from_path(root.clone());
1022 let job = confirmed_job(&store);
1023 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
1024 let mut launcher = RecordingLauncher::new("sandbox_runner_clone_cancel", fixture_patch());
1025 launcher.hook = Some(Box::new({
1026 let cancel = raw_cancel_from_hook(root, job.id.clone());
1027 move |phase| {
1028 if phase == "clone" {
1029 cancel(phase);
1030 }
1031 }
1032 }));
1033 let finished = run_confirmed_job(&store, &job.id, &launcher, &forge).unwrap();
1034 assert_eq!(finished.status, CloudJobStatus::Canceled);
1035 assert!(finished.pr_url.is_none());
1036 let calls = launcher.calls();
1037 assert!(
1038 calls.contains(&"clone".to_string()),
1039 "clone must have run: {calls:?}"
1040 );
1041 assert!(
1042 !calls.contains(&"harness".to_string()),
1043 "harness must not run after a post-clone cancel: {calls:?}"
1044 );
1045 assert!(forge.opened().is_empty());
1046 }
1047
1048 /// Cancels the job from inside a phase hook with a raw status flip (no)
1049 /// `finished_unix`), mirroring the reviewer's reproduction: cancel_job
1050 /// saves `canceled` while a phase is in flight.
1051 fn raw_cancel_from_hook(root: std::path::PathBuf, id: String) -> impl Fn(&str) + Send + Sync {
1052 move |_phase: &str| {
1053 let store = CloudJobStore::from_path(root.clone());
1054 let mut current = store.load(&id).unwrap();
1055 current.status = CloudJobStatus::Canceled;
1056 store.save(&current).unwrap();
1057 }
1058 }
1059
1060 #[test]
1061 fn cancel_during_create_wins_over_the_post_create_save() {
1062 let temp = tempfile::tempdir().unwrap();
1063 let root = temp.path().join("jobs");
1064 let store = CloudJobStore::from_path(root.clone());
1065 let job = confirmed_job(&store);
1066 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
1067 // The cancel lands while the create POST is in flight — before the
1068 // runner ever holds a receipt. The post-create phase save must
1069 // refuse to resurrect the run.
1070 let mut launcher = RecordingLauncher::new("sandbox_cancel_1", fixture_patch());
1071 launcher.hook = Some(Box::new({
1072 let cancel = raw_cancel_from_hook(root.clone(), job.id.clone());
1073 move |phase| {
1074 if phase == "create" {
1075 cancel(phase);
1076 }
1077 }
1078 }));
1079 let finished = run_confirmed_job(&store, &job.id, &launcher, &forge).unwrap();
1080
1081 assert_eq!(finished.status, CloudJobStatus::Canceled);
1082 assert!(finished.pr_url.is_none());
1083 assert!(finished.finished_unix.is_some());
1084 assert!(finished.note.contains("Canceled mid-run"));
1085 assert!(finished.note.contains("torn down"));
1086 // The run never reached readiness, the forge never fired, teardown ran.
1087 assert_eq!(launcher.calls(), vec!["create", "teardown"]);
1088 assert!(forge.opened().is_empty());
1089 let persisted = store.load(&job.id).unwrap();
1090 assert_eq!(persisted.status, CloudJobStatus::Canceled);
1091 assert!(persisted.finished_unix.is_some());
1092 }
1093
1094 #[test]
1095 fn cancel_before_the_sandbox_intent_save_never_creates_a_sandbox() {
1096 let temp = tempfile::tempdir().unwrap();
1097 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1098 let job = confirmed_job(&store);
1099 // The runner loaded the job; the cancel lands before its first save.
1100 let mut runner_copy = store.load(&job.id).unwrap();
1101 crate::cloud_dispatch::cancel_job(
1102 &store,
1103 &job.id,
1104 &RecordingLauncher::new("unused", fixture_patch()),
1105 )
1106 .unwrap();
1107 let launcher = RecordingLauncher::new("sandbox_never", fixture_patch());
1108 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
1109
1110 drive(&store, &mut runner_copy, &launcher, &forge).unwrap();
1111
1112 assert!(launcher.calls().is_empty(), "{:?}", launcher.calls());
1113 assert!(forge.opened().is_empty());
1114 assert_eq!(runner_copy.status, CloudJobStatus::Canceled);
1115 assert_eq!(
1116 store.load(&job.id).unwrap().status,
1117 CloudJobStatus::Canceled
1118 );
1119 }
1120
1121 #[test]
1122 fn cancel_during_collect_blocks_the_branch_raise_and_the_pr() {
1123 let temp = tempfile::tempdir().unwrap();
1124 let root = temp.path().join("jobs");
1125 let store = CloudJobStore::from_path(root.clone());
1126 let job = confirmed_job(&store);
1127 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
1128 // The cancel lands while the patch is being collected; the
1129 // OpeningPr phase save must refuse, so the branch is never raised
1130 // and the PR never opened.
1131 let mut launcher = RecordingLauncher::new("sandbox_cancel_2", fixture_patch());
1132 launcher.hook = Some(Box::new({
1133 let cancel = raw_cancel_from_hook(root.clone(), job.id.clone());
1134 move |phase| {
1135 if phase == "collect" {
1136 cancel(phase);
1137 }
1138 }
1139 }));
1140 let finished = run_confirmed_job(&store, &job.id, &launcher, &forge).unwrap();
1141
1142 assert_eq!(finished.status, CloudJobStatus::Canceled);
1143 assert!(finished.pr_url.is_none());
1144 assert!(finished.base_branch.is_none());
1145 assert!(finished.finished_unix.is_some());
1146 assert!(finished.note.contains("Canceled mid-run"));
1147 assert_eq!(
1148 launcher.calls(),
1149 vec![
1150 "create",
1151 "wait_ready",
1152 "clone",
1153 "harness",
1154 "collect",
1155 "teardown"
1156 ]
1157 );
1158 assert!(forge.opened().is_empty());
1159 }
1160
1161 /// Launcher whose harness step fails; earlier phases record normally.
1162 struct HarnessFailsLauncher {
1163 calls: Mutex<Vec<String>>,
1164 hook: Option<PhaseHook>,
1165 }
1166
1167 impl HarnessFailsLauncher {
1168 fn note(&self, phase: &str) {
1169 if let Ok(mut calls) = self.calls.lock() {
1170 calls.push(phase.to_string());
1171 }
1172 if let Some(hook) = self.hook.as_ref() {
1173 hook(phase);
1174 }
1175 }
1176 }
1177
1178 impl DaytonaLauncher for HarnessFailsLauncher {
1179 fn create_sandbox(&self, _job: &CloudJob) -> Result<SandboxReceipt> {
1180 self.note("create");
1181 Ok(SandboxReceipt {
1182 sandbox_id: "sandbox_err_1".to_string(),
1183 toolbox_url: None,
1184 })
1185 }
1186 fn wait_ready(&self, _receipt: &SandboxReceipt) -> Result<()> {
1187 self.note("wait_ready");
1188 Ok(())
1189 }
1190 fn clone_repository(&self, _receipt: &SandboxReceipt, url: &str, path: &str) -> Result<()> {
1191 self.note("clone");
1192 assert_eq!(path, SANDBOX_WORKSPACE);
1193 assert!(url.starts_with("https://"));
1194 Ok(())
1195 }
1196 fn run_harness(
1197 &self,
1198 _receipt: &SandboxReceipt,
1199 _command: &HarnessCommand,
1200 ) -> Result<String> {
1201 self.note("harness");
1202 bail!("harness exploded")
1203 }
1204 fn teardown(&self, _receipt: &SandboxReceipt) -> Result<()> {
1205 self.note("teardown");
1206 Ok(())
1207 }
1208 }
1209
1210 #[test]
1211 fn run_error_after_cancel_keeps_the_canceled_record() {
1212 let temp = tempfile::tempdir().unwrap();
1213 let root = temp.path().join("jobs");
1214 let store = CloudJobStore::from_path(root.clone());
1215 let job = confirmed_job(&store);
1216 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/9");
1217 // The cancel lands as the harness explodes: the failure arm must not
1218 // overwrite the user's `canceled` with `failed` — the error rides
1219 // along in the note.
1220 let mut launcher = HarnessFailsLauncher {
1221 calls: Mutex::new(Vec::new()),
1222 hook: None,
1223 };
1224 launcher.hook = Some(Box::new({
1225 let cancel = raw_cancel_from_hook(root.clone(), job.id.clone());
1226 move |phase| {
1227 if phase == "harness" {
1228 cancel(phase);
1229 }
1230 }
1231 }));
1232 let error = run_confirmed_job(&store, &job.id, &launcher, &forge)
1233 .unwrap_err()
1234 .to_string();
1235 assert!(error.contains("harness exploded"), "{error}");
1236
1237 let persisted = store.load(&job.id).unwrap();
1238 assert_eq!(
1239 persisted.status,
1240 CloudJobStatus::Canceled,
1241 "a user cancel survives a later run error"
1242 );
1243 assert!(persisted.note.contains("Run error after cancel"));
1244 assert!(persisted.note.contains("harness exploded"));
1245 assert!(persisted.refusal.is_none());
1246 assert!(persisted.pr_url.is_none());
1247 assert!(persisted.finished_unix.is_some());
1248 assert!(forge.opened().is_empty());
1249 let calls = launcher.calls.lock().unwrap().clone();
1250 assert_eq!(
1251 calls,
1252 vec!["create", "wait_ready", "clone", "harness", "teardown"]
1253 );
1254 }
1255
1256 #[test]
1257 fn the_declared_harness_budget_fits_the_harness_client_budget() {
1258 let temp = tempfile::tempdir().unwrap();
1259 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1260 let job = confirmed_job(&store);
1261 let command = harness_command(&job);
1262 assert_eq!(command.timeout_secs, HARNESS_TIMEOUT_SECS);
1263 assert!(
1264 cloud_dispatch::LiveDaytonaLauncher::harness_client_budget_secs(&command)
1265 >= u64::from(HARNESS_TIMEOUT_SECS),
1266 "the client that carries the harness turn must cover the declared hour"
1267 );
1268 }
1269
1270 #[test]
1271 fn sandbox_intent_is_persisted_before_the_create_post() {
1272 let temp = tempfile::tempdir().unwrap();
1273 let root = temp.path().join("jobs");
1274 let store = CloudJobStore::from_path(root.clone());
1275 let job = confirmed_job(&store);
1276 // Observed from inside the create phase: the intent must already be
1277 // on disk, so a create whose response never arrives is still
1278 // reconcilable by label.
1279 let intent_root = root.clone();
1280 let intent_id = job.id.clone();
1281 let seen_pending = Arc::new(std::sync::atomic::AtomicBool::new(false));
1282 let seen = seen_pending.clone();
1283 let mut launcher = RecordingLauncher::new("sandbox_intent_1", fixture_patch());
1284 launcher.hook = Some(Box::new(move |phase| {
1285 if phase == "create" {
1286 let store = CloudJobStore::from_path(intent_root.clone());
1287 let current = store.load(&intent_id).unwrap();
1288 assert!(current.sandbox_pending, "intent must precede the POST");
1289 seen.store(true, std::sync::atomic::Ordering::SeqCst);
1290 }
1291 }));
1292 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/11");
1293 let finished = run_confirmed_job(&store, &job.id, &launcher, &forge).unwrap();
1294
1295 assert!(seen_pending.load(std::sync::atomic::Ordering::SeqCst));
1296 assert_eq!(finished.status, CloudJobStatus::Done);
1297 assert!(!finished.sandbox_pending, "intent clears once the id lands");
1298 let persisted = store.load(&job.id).unwrap();
1299 assert!(!persisted.sandbox_pending);
1300 assert_eq!(persisted.sandbox_id.as_deref(), Some("sandbox_intent_1"));
1301 }
1302
1303 #[test]
1304 fn startup_reconcile_sweeps_stale_jobs_and_deletes_orphan_sandboxes() {
1305 let temp = tempfile::tempdir().unwrap();
1306 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1307 let job = confirmed_job(&store);
1308 // Age the record past the stale threshold and park it mid-run, the
1309 // state a quit/crash would leave behind.
1310 let mut stale = store.load(&job.id).unwrap();
1311 stale.status = CloudJobStatus::Running;
1312 stale.sandbox_id = Some("sandbox_orphan".to_string());
1313 stale.created_unix = stale
1314 .created_unix
1315 .saturating_sub(crate::cloud_dispatch::STALE_ACTIVE_JOB_SECS + 120);
1316 store.save(&stale).unwrap();
1317 // A sandbox labeled for a job that is not in the store at all.
1318 let launcher = RecordingLauncher::new("unused", fixture_patch());
1319 *launcher.listed.lock().unwrap() = vec![crate::cloud_dispatch::LabeledSandbox {
1320 sandbox_id: "sandbox_ghost".to_string(),
1321 job_id: Some("cloud_0000000000000bad".to_string()),
1322 }];
1323
1324 let receipt = startup_reconcile(&store, &launcher);
1325 assert!(
1326 receipt.contains(&job.id),
1327 "receipt names the swept job: {receipt}"
1328 );
1329 assert!(
1330 receipt.contains("sandbox_ghost"),
1331 "receipt names deletions: {receipt}"
1332 );
1333 assert_eq!(store.load(&job.id).unwrap().status, CloudJobStatus::Failed);
1334 assert!(launcher.calls().contains(&"teardown".to_string()));
1335 assert!(launcher.calls().contains(&"list".to_string()));
1336 }
1337
1338 #[test]
1339 fn only_confirmed_jobs_can_run() {
1340 let temp = tempfile::tempdir().unwrap();
1341 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1342 let plan = plan_dispatch(
1343 &[GitRemote {
1344 name: "github".to_string(),
1345 url: "https://github.com/org/repo.git".to_string(),
1346 }],
1347 "unconfirmed work",
1348 Some(Forge::Github),
1349 Some("codewhale/cloud-unconfirmed"),
1350 )
1351 .unwrap();
1352 match execute_dispatch(
1353 &store,
1354 plan,
1355 false,
1356 &CredentialState::Present {
1357 source: CredentialSource::Env,
1358 },
1359 &MachineTokenState::Present,
1360 )
1361 .unwrap()
1362 {
1363 DispatchOutcome::Proposal(job) => {
1364 let launcher = RecordingLauncher::new("never", fixture_patch());
1365 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/1");
1366 let error = run_confirmed_job(&store, &job.id, &launcher, &forge)
1367 .unwrap_err()
1368 .to_string();
1369 assert!(error.contains("confirm it first"), "{error}");
1370 assert!(launcher.calls().is_empty());
1371 }
1372 other => panic!("expected proposal, got {other:?}"),
1373 }
1374 }
1375
1376 #[test]
1377 fn launch_failure_fails_closed_with_sanitized_note() {
1378 let temp = tempfile::tempdir().unwrap();
1379 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1380 let job = confirmed_job(&store);
1381 struct FailingLauncher;
1382 impl DaytonaLauncher for FailingLauncher {
1383 fn create_sandbox(&self, _job: &CloudJob) -> Result<SandboxReceipt> {
1384 bail!("create exploded\u{1}")
1385 }
1386 }
1387 let forge = RecordingForgePr::new("https://github.com/org/repo/pull/2");
1388 let error = run_confirmed_job(&store, &job.id, &FailingLauncher, &forge).unwrap_err();
1389 assert!(error.to_string().contains("create exploded"));
1390 let failed = store.load(&job.id).unwrap();
1391 assert_eq!(failed.status, CloudJobStatus::Failed);
1392 assert!(failed.pr_url.is_none());
1393 assert!(failed.note.contains("failed closed"));
1394 assert!(!failed.note.contains('\u{1}'));
1395 assert!(forge.opened().is_empty());
1396 }
1397
1398 #[test]
1399 fn pr_title_and_body_are_truthful_unbranded_and_carry_receipts() {
1400 let temp = tempfile::tempdir().unwrap();
1401 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1402 let job = confirmed_job(&store);
1403 let patch = fixture_patch();
1404 let title = compose_pr_title(&job, &patch);
1405 assert_eq!(title, "codewhale cloud: Fix the flaky dispatch test");
1406 let body = compose_pr_body(&job, &patch);
1407 assert!(body.contains("Codewhale cloud agent"));
1408 assert!(body.contains("What the agent did"));
1409 assert!(body.contains("Fix the flaky dispatch test"));
1410 assert!(body.contains("open a PR that fixes the flake"));
1411 assert!(body.contains(&job.id));
1412 assert!(body.contains("(pending)"));
1413 assert!(body.contains("codewhale/cloud-runner-test"));
1414 assert!(body.contains("`main`"));
1415 assert!(body.contains("abc123def4567"));
1416 assert!(body.contains("No-Issue: cloud dispatch"));
1417 for banned in ["Daytona", "daytona"] {
1418 assert!(
1419 !body.contains(banned),
1420 "body must not brand the sandbox: {banned}"
1421 );
1422 }
1423 }
1424
1425 #[test]
1426 fn harness_command_is_the_standard_one_shot_entry() {
1427 let temp = tempfile::tempdir().unwrap();
1428 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1429 let job = confirmed_job(&store);
1430 let command = harness_command(&job);
1431 assert_eq!(
1432 command.argv,
1433 vec![
1434 "codewhale".to_string(),
1435 "exec".to_string(),
1436 "--auto".to_string(),
1437 job.prompt.clone(),
1438 ]
1439 );
1440 assert_eq!(command.cwd, SANDBOX_WORKSPACE);
1441 }
1442
1443 #[test]
1444 fn gh_argv_and_forge_endpoints_pin_the_pr_shapes() {
1445 let argv = gh_pr_create_argv(
1446 "org/repo",
1447 "main",
1448 "codewhale/cloud-1",
1449 "title",
1450 "/tmp/body.md",
1451 );
1452 assert_eq!(
1453 argv,
1454 vec![
1455 "pr",
1456 "create",
1457 "--repo",
1458 "org/repo",
1459 "--base",
1460 "main",
1461 "--head",
1462 "codewhale/cloud-1",
1463 "--title",
1464 "title",
1465 "--body-file",
1466 "/tmp/body.md",
1467 ]
1468 );
1469 assert_eq!(
1470 gitee_pr_url("org/repo"),
1471 "https://gitee.com/api/v5/repos/org/repo/pulls"
1472 );
1473 assert_eq!(
1474 cnb_pr_url("org/repo"),
1475 "https://api.cnb.cool/org/repo/-/pulls"
1476 );
1477 }
1478
1479 #[test]
1480 fn forge_slug_parses_https_and_ssh_and_rejects_foreign_hosts() {
1481 assert_eq!(
1482 forge_slug(
1483 Forge::Github,
1484 "https://github.com/codewhale-hq/CodeWhale.git"
1485 ),
1486 Some("codewhale-hq/CodeWhale".to_string())
1487 );
1488 assert_eq!(
1489 forge_slug(Forge::Github, "git@github.com:codewhale-hq/CodeWhale.git"),
1490 Some("codewhale-hq/CodeWhale".to_string())
1491 );
1492 assert_eq!(
1493 forge_slug(Forge::Cnb, "https://cnb.cool/codewhale.net/codewhale.git"),
1494 Some("codewhale.net/codewhale".to_string())
1495 );
1496 assert_eq!(
1497 forge_slug(Forge::Gitee, "https://gitee.com/org/repo.git"),
1498 Some("org/repo".to_string())
1499 );
1500 assert_eq!(
1501 forge_slug(Forge::Github, "https://gitee.com/org/repo.git"),
1502 None
1503 );
1504 assert_eq!(
1505 forge_slug(Forge::Cnb, "https://example.test/org/repo.git"),
1506 None
1507 );
1508 assert_eq!(
1509 forge_slug(Forge::Github, "https://github.com/only-repo"),
1510 None
1511 );
1512 }
1513
1514 #[test]
1515 fn push_branch_refuses_non_forge_remotes() {
1516 let error =
1517 push_branch(Path::new("."), "https://example.test/org/repo.git", "b").unwrap_err();
1518 let text = error.to_string();
1519 assert!(
1520 text.contains("non-forge") || text.contains("not a supported forge"),
1521 "{text}"
1522 );
1523 }
1524
1525 #[test]
1526 fn push_branch_refuses_forge_default_branch_names() {
1527 let error =
1528 push_branch(Path::new("."), "https://github.com/org/repo.git", "main").unwrap_err();
1529 assert!(error.to_string().contains("default branch"), "{}", error);
1530 }
1531
1532 #[test]
1533 fn prepare_branch_rejects_leading_dash_remote_before_clone() {
1534 let temp = tempfile::tempdir().unwrap();
1535 let mut job = confirmed_job(&CloudJobStore::from_path(temp.path().join("jobs")));
1536 job.remote_url = "--upload-pack=evil".to_string();
1537 let error = prepare_branch(&job, &fixture_patch(), temp.path()).unwrap_err();
1538 assert!(
1539 error.to_string().contains("must not start with '-'"),
1540 "{error}"
1541 );
1542 }
1543
1544 #[test]
1545 fn compose_pr_body_head_is_the_pushed_sha_not_the_sandbox_sha() {
1546 let job = CloudJob {
1547 id: "cloud_00000000000000cc".to_string(),
1548 kind: "cloud".to_string(),
1549 status: CloudJobStatus::OpeningPr,
1550 prompt: "fix".to_string(),
1551 forge: Forge::Github,
1552 remote_name: "github".to_string(),
1553 remote_url: "https://github.com/org/repo.git".to_string(),
1554 branch: "codewhale/cloud-b".to_string(),
1555 confirmed: true,
1556 sandbox_id: Some("sandbox".to_string()),
1557 pr_url: None,
1558 refusal: None,
1559 note: "n".to_string(),
1560 created_unix: 1,
1561 base_branch: None,
1562 head_sha: None,
1563 agent_summary: None,
1564 finished_unix: None,
1565 sandbox_pending: false,
1566 };
1567 let patch = fixture_patch();
1568 let body = compose_pr_body_for_head(&job, &patch, "pushedsha0000000000000000000000000001");
1569 assert!(body.contains("Head: `pushedsha0000000000000000000000000001`"));
1570 assert!(
1571 !body.contains(&format!("Head: `{}`", patch.head_sha)),
1572 "sandbox sha must not be the receipt Head"
1573 );
1574 }
1575
1576 /// Local-fixture integration: real git, no network — branch preparation
1577 /// applies a real format-patch, and a diverged push fails without force.
1578 #[test]
1579 fn prepare_branch_applies_a_real_patch_locally() {
1580 let temp = tempfile::tempdir().unwrap();
1581 let origin = temp.path().join("origin.git");
1582 git(
1583 None,
1584 &[
1585 "init",
1586 "--bare",
1587 "--quiet",
1588 "--initial-branch=main",
1589 &origin.to_string_lossy(),
1590 ],
1591 )
1592 .unwrap();
1593 let seed = temp.path().join("seed");
1594 git(
1595 None,
1596 &[
1597 "clone",
1598 "--quiet",
1599 &origin.to_string_lossy(),
1600 &seed.to_string_lossy(),
1601 ],
1602 )
1603 .unwrap();
1604 git(Some(&seed), &["config", "user.name", "Seeder"]).unwrap();
1605 git(Some(&seed), &["config", "user.email", "seed@example.test"]).unwrap();
1606 std::fs::write(seed.join("file.txt"), "base\n").unwrap();
1607 git(Some(&seed), &["add", "."]).unwrap();
1608 git(Some(&seed), &["commit", "--quiet", "-m", "base"]).unwrap();
1609 git(
1610 Some(&seed),
1611 &[
1612 "push",
1613 "--quiet",
1614 &origin.to_string_lossy(),
1615 "HEAD:refs/heads/main",
1616 ],
1617 )
1618 .unwrap();
1619
1620 // The agent's work product: a real patch.
1621 std::fs::write(seed.join("file.txt"), "base\nagent change\n").unwrap();
1622 git(Some(&seed), &["add", "."]).unwrap();
1623 git(Some(&seed), &["commit", "--quiet", "-m", "agent work"]).unwrap();
1624 let patch_text = git(Some(&seed), &["format-patch", "HEAD~1", "--stdout"]).unwrap();
1625
1626 let store = CloudJobStore::from_path(temp.path().join("jobs"));
1627 let mut job = confirmed_job(&store);
1628 // Point at the local fixture so preparation is offline.
1629 job.remote_url = origin.to_string_lossy().to_string();
1630 let patch = PatchReceipt {
1631 base_branch: "main".to_string(),
1632 head_sha: "fixture".to_string(),
1633 summary: "agent work".to_string(),
1634 patch: patch_text,
1635 };
1636 let dir = temp.path().join("agent");
1637 std::fs::create_dir_all(&dir).unwrap();
1638 let head = prepare_branch(&job, &patch, &dir).unwrap();
1639 assert_eq!(
1640 head.len(),
1641 40,
1642 "prepare_branch returns the applied head sha"
1643 );
1644
1645 // A diverged push onto the same branch must fail: no force, ever.
1646 let repo = dir.join("repo");
1647 std::fs::write(repo.join("other.txt"), "y\n").unwrap();
1648 git(Some(&repo), &["add", "."]).unwrap();
1649 git(Some(&repo), &["commit", "--quiet", "-m", "diverges"]).unwrap();
1650 // Seed the remote branch at the "diverges" commit.
1651 git(
1652 Some(&repo),
1653 &[
1654 "push",
1655 "--quiet",
1656 &origin.to_string_lossy(),
1657 "HEAD:refs/heads/codewhale/cloud-runner-test",
1658 ],
1659 )
1660 .unwrap();
1661 // Rewind and build a sibling commit: same parent, different content.
1662 git(Some(&repo), &["reset", "--quiet", "--hard", "HEAD~1"]).unwrap();
1663 std::fs::write(repo.join("other2.txt"), "z\n").unwrap();
1664 git(Some(&repo), &["add", "."]).unwrap();
1665 git(
1666 Some(&repo),
1667 &["commit", "--quiet", "-m", "diverges differently"],
1668 )
1669 .unwrap();
1670 let diverged = git(
1671 Some(&repo),
1672 &[
1673 "push",
1674 "--quiet",
1675 &origin.to_string_lossy(),
1676 "HEAD:refs/heads/codewhale/cloud-runner-test",
1677 ],
1678 );
1679 assert!(diverged.is_err(), "a diverged push must fail without force");
1680 }
1681
1682 #[test]
1683 fn compose_pr_body_bounds_oversized_prompts() {
1684 let job = CloudJob {
1685 id: "cloud_00000000000000bb".to_string(),
1686 kind: "cloud".to_string(),
1687 status: CloudJobStatus::OpeningPr,
1688 prompt: "p".repeat(9_000),
1689 forge: Forge::Github,
1690 remote_name: "github".to_string(),
1691 remote_url: "https://github.com/org/repo.git".to_string(),
1692 branch: "codewhale/cloud-b".to_string(),
1693 confirmed: true,
1694 sandbox_id: Some("sandbox".to_string()),
1695 pr_url: None,
1696 refusal: None,
1697 note: "n".to_string(),
1698 created_unix: 1,
1699 base_branch: None,
1700 head_sha: None,
1701 agent_summary: None,
1702 finished_unix: None,
1703 sandbox_pending: false,
1704 };
1705 let body = compose_pr_body(&job, &fixture_patch());
1706 assert!(body.chars().count() <= MAX_BODY_CHARS);
1707 }
1708
1709 #[test]
1710 fn one_line_flattens_control_characters() {
1711 assert_eq!(one_line("a\nb", 10), "a\nb");
1712 assert_eq!(one_line("a\u{1}b", 10), "a b");
1713 assert_eq!(one_line("abcdef", 3), "ab…");
1714 }
1715 }
1716
1716 lines RUST