返回 CodeWhale
cloud_dispatch.rs
根目录 / crates / tui / src / cloud_dispatch.rs
1 //! Daytona cloud-agent dispatch contract.
2 //!
3 //! Local `cw` / Codewhale may propose sending a coding agent to Daytona so
4 //! heavy work can raise a branch and open a PR while the TUI stays responsive.
5 //! This module is the only first-class offload seam:
6 //!
7 //! - remotes are explicit forges: `github`, `cnb`, `gitee`
8 //! - CWC's convention is preserved: a remote *named* `github` is authoritative
9 //! GitHub; `origin` is classified by URL and is often the CNB mirror
10 //! - confirmation is required; nothing spends or pushes silently
11 //! - missing Daytona credentials fail closed; success is never faked
12 //!
13 //! The engine that drives a confirmed job end to end (sandbox → harness →
14 //! forge PR → teardown) lives in [`crate::dispatch_runner`]; this module owns
15 //! the persisted contract, the launcher seam, and the fail-closed gates.
16 //! Auto-decide heuristics remain leftover: Codewhale may propose, never
17 //! confirm itself.
18 //!
19 //! A sandbox create receipt is infrastructure identity, not Computer
20 //! entitlement. Metering requires [`crate::computer_meter`] admission plus a
21 //! provider-accepted active observation.
22
23 use std::fs;
24 use std::path::{Path, PathBuf};
25 use std::process::Command;
26 use std::time::{SystemTime, UNIX_EPOCH};
27
28 use anyhow::{Context, Result, anyhow, bail};
29 use codewhale_paths::codewhale_home;
30 use codewhale_secrets::Secrets;
31 use serde::{Deserialize, Serialize};
32
33 use crate::computer_meter::{
34 ComputerAdmission, ComputerMeterError, ComputerMeterReceipt, ProviderObservation,
35 issue_computer_meter_receipt,
36 };
37
38 const MAX_PROMPT_CHARS: usize = 4_000;
39 const MAX_REMOTE_BYTES: usize = 4 * 1024;
40 const JOB_KIND: &str = "cloud";
41 const DAYTONA_API_KEY_ENV: &str = "DAYTONA_API_KEY";
42 const DAYTONA_API_URL_ENV: &str = "DAYTONA_API_URL";
43 const CWC_DAYTONA_TOKEN_ENV: &str = "CWC_DAYTONA_TOKEN";
44 const CWC_DAYTONA_ENDPOINT_ENV: &str = "CWC_DAYTONA_ENDPOINT";
45 const KEYRING_SLOT: &str = "daytona";
46 const DEFAULT_DAYTONA_API: &str = "https://app.daytona.io/api";
47 /// Path inside the sandbox where the target repository is cloned.
48 pub const SANDBOX_WORKSPACE: &str = "/workspace";
49 const MAX_HARNESS_OUTPUT_CHARS: usize = 200_000;
50 const READY_POLL_ATTEMPTS: u32 = 40;
51 const READY_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_secs(3);
52 /// Sandbox label key carrying the cloud job id (see
53 /// [`LiveDaytonaLauncher::create_sandbox`]); the reconciler joins sandbox
54 /// labels back to job records through it.
55 pub const SANDBOX_JOB_LABEL: &str = "codewhale.job";
56 /// Sandbox label marking Codewhale dispatch sandboxes — the product tag the
57 /// label reconciler filters the provider's sandbox list on.
58 pub const SANDBOX_PRODUCT_LABEL: &str = "codewhale.product";
59 pub const SANDBOX_PRODUCT_VALUE: &str = "dispatch";
60 /// The only environment variable that carries the Codewhale account machine
61 /// token (`cwc_key_…`) into the sandbox, so the preinstalled `codewhale`
62 /// authenticates as the dispatching account. Same name the CLI's CI path
63 /// reads — one token, one meaning, every host.
64 pub const CLOUD_AGENT_TOKEN_ENV: &str = "CODEWHALE_API_KEY";
65 /// Operator override for the cloud-agent snapshot name.
66 pub const CLOUD_AGENT_SNAPSHOT_ENV: &str = "CODEWHALE_DISPATCH_SNAPSHOT";
67 /// Default Daytona snapshot name. The image definition is maintained in
68 /// `computer/snapshots/cloud-agent/Dockerfile`; its pinned Engine version and
69 /// credential/acceptance limitations are documented beside it.
70 pub const DEFAULT_CLOUD_AGENT_SNAPSHOT: &str = "codewhale-cloud-agent";
71 /// Active jobs older than this are stale. The declared harness budget for
72 /// one cloud-agent turn is an hour (`HARNESS_TIMEOUT_SECS`), so an active
73 /// record with no terminal state after 90 minutes (harness budget plus
74 /// control-plane slack) cannot be a healthy run — its runner is gone (TUI
75 /// quit, crash, killed process) and the record is the only witness. The
76 /// sweep fails such records and tears their sandboxes down.
77 pub const STALE_ACTIVE_JOB_SECS: u64 = 90 * 60;
78
79 /// Explicit PR forge. Never inferred from a generic "origin means GitHub" rule.
80 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
81 #[serde(rename_all = "lowercase")]
82 pub enum Forge {
83 Github,
84 Cnb,
85 Gitee,
86 }
87
88 impl Forge {
89 /// Stable CLI / TUI slug.
90 pub fn as_str(self) -> &'static str {
91 match self {
92 Self::Github => "github",
93 Self::Cnb => "cnb",
94 Self::Gitee => "gitee",
95 }
96 }
97
98 /// Parse a user-supplied forge slug.
99 pub fn parse(value: &str) -> Option<Self> {
100 match value.trim().to_ascii_lowercase().as_str() {
101 "github" | "gh" => Some(Self::Github),
102 "cnb" => Some(Self::Cnb),
103 "gitee" => Some(Self::Gitee),
104 _ => None,
105 }
106 }
107 }
108
109 /// One `git remote` row after fetch/push duplicates are collapsed.
110 #[derive(Debug, Clone, PartialEq, Eq)]
111 pub struct GitRemote {
112 pub name: String,
113 pub url: String,
114 }
115
116 /// A remote that has been classified as a supported forge.
117 #[derive(Debug, Clone, PartialEq, Eq)]
118 pub struct SelectedRemote {
119 pub forge: Forge,
120 pub name: String,
121 pub url: String,
122 }
123
124 /// Where a Daytona API key was found. Never carries the secret.
125 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
126 pub enum CredentialSource {
127 Env,
128 CwcEnv,
129 Keyring,
130 }
131
132 /// Presence of Daytona credentials. Absence is fail-closed.
133 #[derive(Debug, Clone, PartialEq, Eq)]
134 pub enum CredentialState {
135 Missing,
136 Present { source: CredentialSource },
137 }
138
139 /// First-class cloud job lifecycle. `kind` is always `cloud`.
140 ///
141 /// The runner path is `Proposed` (queued for an explicit confirm) →
142 /// `Launching` → `Running` (harness turn in the sandbox) → `OpeningPr` →
143 /// `Done`, with `Failed` / `Canceled` reachable from every active state and
144 /// `Refused` reserved for the fail-closed membership/credential gate.
145 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
146 #[serde(rename_all = "lowercase")]
147 pub enum CloudJobStatus {
148 Proposed,
149 Refused,
150 Launching,
151 Running,
152 #[serde(rename = "openingpr")]
153 OpeningPr,
154 Done,
155 Failed,
156 Canceled,
157 }
158
159 /// Durable cloud job record, listed on the same `/jobs` surface as Bash jobs.
160 ///
161 /// Fields added after the first landing carry `#[serde(default)]` so job
162 /// records written by earlier builds still load.
163 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
164 pub struct CloudJob {
165 pub id: String,
166 pub kind: String,
167 pub status: CloudJobStatus,
168 pub prompt: String,
169 pub forge: Forge,
170 pub remote_name: String,
171 pub remote_url: String,
172 pub branch: String,
173 pub confirmed: bool,
174 pub sandbox_id: Option<String>,
175 pub pr_url: Option<String>,
176 pub refusal: Option<String>,
177 pub note: String,
178 pub created_unix: u64,
179 /// Default branch of the agent's clone (the PR base), when known.
180 #[serde(default)]
181 pub base_branch: Option<String>,
182 /// Head commit the agent produced, when known.
183 #[serde(default)]
184 pub head_sha: Option<String>,
185 /// One-line truthful summary of what the agent did, when reported.
186 #[serde(default)]
187 pub agent_summary: Option<String>,
188 /// When the job reached a terminal state (`done`/`failed`/`canceled`).
189 #[serde(default)]
190 pub finished_unix: Option<u64>,
191 /// Intent record: a sandbox create POST is in flight for this job. Set
192 /// and persisted *before* the POST so that a create whose response is
193 /// slow (client timeout, process death) is still reconcilable by
194 /// sandbox label even though the sandbox id never arrived. Cleared once
195 /// the id is recorded.
196 #[serde(default)]
197 pub sandbox_pending: bool,
198 }
199
200 /// Validated plan that still requires an explicit confirm to spend or push.
201 #[derive(Debug, Clone, PartialEq, Eq)]
202 pub struct DispatchPlan {
203 pub prompt: String,
204 pub remote: SelectedRemote,
205 pub branch: String,
206 }
207
208 /// Result of proposing or confirming a dispatch.
209 #[derive(Debug, Clone, PartialEq, Eq)]
210 pub enum DispatchOutcome {
211 Proposal(CloudJob),
212 Refused(CloudJob),
213 Accepted(CloudJob),
214 }
215
216 /// Why a plan or launch cannot proceed.
217 #[derive(Debug, Clone, PartialEq, Eq)]
218 pub enum DispatchError {
219 EmptyPrompt,
220 PromptTooLong,
221 InvalidBranch,
222 UnknownForge,
223 AmbiguousRemote,
224 RemoteMissing(Forge),
225 NoSupportedRemote,
226 UnsafeRemote,
227 }
228
229 impl std::fmt::Display for DispatchError {
230 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
231 match self {
232 Self::EmptyPrompt => write!(f, "A cloud dispatch needs a non-empty task prompt."),
233 Self::PromptTooLong => write!(
234 f,
235 "A cloud dispatch prompt must be at most {MAX_PROMPT_CHARS} characters."
236 ),
237 Self::InvalidBranch => write!(
238 f,
239 "The requested branch is not a safe git ref (no leading '-', no '..', no shell metacharacters)."
240 ),
241 Self::UnknownForge => write!(
242 f,
243 "Remote must be one of github, cnb, or gitee. CWC treats the `github` remote as GitHub and `origin` as the CNB mirror when that URL is cnb.cool."
244 ),
245 Self::AmbiguousRemote => write!(
246 f,
247 "This workspace has more than one forge remote. Pass --remote github|cnb|gitee (CWC: `github` is GitHub, `origin` is often CNB)."
248 ),
249 Self::RemoteMissing(forge) => write!(
250 f,
251 "No {0} remote is configured. Add a `{0}` remote or a URL on {1}.",
252 forge.as_str(),
253 forge_host(*forge)
254 ),
255 Self::NoSupportedRemote => write!(
256 f,
257 "No GitHub, CNB, or Gitee remote was found. Remotes are classified by name (`github`/`cnb`/`gitee`) then by host."
258 ),
259 Self::UnsafeRemote => write!(
260 f,
261 "The selected remote URL is unsafe (leading-dash, embedded userinfo, or a non-forge host). Refusing to clone or show it."
262 ),
263 }
264 }
265 }
266
267 impl std::error::Error for DispatchError {}
268
269 /// File-backed store under `$CODEWHALE_HOME/cloud-jobs`.
270 #[derive(Debug, Clone)]
271 pub struct CloudJobStore {
272 root: PathBuf,
273 }
274
275 impl CloudJobStore {
276 /// Resolve the process Codewhale home.
277 pub fn from_env() -> Result<Self> {
278 let home = codewhale_home()
279 .map_err(|err| anyhow!(err.to_string()))?
280 .ok_or_else(|| anyhow!("CODEWHALE_HOME / user home is unavailable"))?;
281 Ok(Self::from_path(home.join("cloud-jobs")))
282 }
283
284 /// Test and injected-root constructor.
285 pub fn from_path(root: PathBuf) -> Self {
286 Self { root }
287 }
288
289 /// Persist a job atomically. Never writes credentials.
290 ///
291 /// The record is written through a unique temporary file (never a fixed
292 /// `.json.tmp` two writers would share) and renamed into place while the
293 /// job's lock is held, so it cannot interleave with a
294 /// [`Self::save_unless_canceled`] or a cancel.
295 pub fn save(&self, job: &CloudJob) -> Result<()> {
296 self.with_job_lock(&job.id, || self.write_record(job))
297 }
298
299 /// Cancel-authoritative save: refuse to overwrite a `canceled` record.
300 ///
301 /// The runner does read-modify-write phase saves while `/dispatch
302 /// cancel` (or `--cancel`) can flip the record concurrently; an
303 /// unconditional phase save landing after the cancel would resurrect a
304 /// dead run — including its branch push and PR. Returns `Ok(false)`
305 /// (leaving the cancellation exactly as the user left it) when the
306 /// persisted record is already `canceled`, `Ok(true)` after a normal
307 /// save.
308 ///
309 /// The load, check and save run under the job's advisory file lock, and
310 /// [`cancel_job`] flips the record under the same lock, so the check and
311 /// the write are one step with respect to a cancel in another process.
312 pub fn save_unless_canceled(&self, job: &CloudJob) -> Result<bool> {
313 self.with_job_lock(&job.id, || {
314 match self.load(&job.id) {
315 Ok(current) if current.status == CloudJobStatus::Canceled => return Ok(false),
316 Ok(_) => {}
317 Err(error)
318 if error
319 .downcast_ref::<std::io::Error>()
320 .is_some_and(|error| error.kind() == std::io::ErrorKind::NotFound) => {}
321 Err(error) => return Err(error.context("cannot verify the persisted job status")),
322 }
323 self.write_record(job)?;
324 Ok(true)
325 })
326 }
327
328 /// Run `f` holding the job's exclusive advisory lock
329 /// (`<root>/<id>.json.lock`). The lock is released when `f` returns.
330 ///
331 /// Limitation: the lock is advisory, so it serializes Codewhale writers
332 /// only; it is not reentrant, so `f` must not call a locking method.
333 fn with_job_lock<T>(&self, id: &str, f: impl FnOnce() -> Result<T>) -> Result<T> {
334 let path = self.job_path(id)?;
335 fs::create_dir_all(&self.root).context("failed to create cloud-jobs directory")?;
336 let lock_file = fs::OpenOptions::new()
337 .create(true)
338 .truncate(false)
339 .read(true)
340 .write(true)
341 .open(path.with_extension("json.lock"))
342 .context("failed to open the cloud job lock")?;
343 let mut lock = fd_lock::RwLock::new(lock_file);
344 let _guard = lock
345 .write()
346 .context("failed to lock the cloud job record")?;
347 f()
348 }
349
350 /// Write the record through a unique private temporary file. Callers
351 /// hold the job lock.
352 fn write_record(&self, job: &CloudJob) -> Result<()> {
353 let path = self.job_path(&job.id)?;
354 let body = serde_json::to_vec_pretty(job).context("failed to encode cloud job")?;
355 crate::utils::write_atomic(&path, &body).context("failed to commit the cloud job record")
356 }
357
358 /// Load one job by id.
359 pub fn load(&self, id: &str) -> Result<CloudJob> {
360 let path = self.job_path(id)?;
361 let body = fs::read(&path).with_context(|| format!("cloud job {id} was not found"))?;
362 serde_json::from_slice(&body).context("cloud job record is invalid JSON")
363 }
364
365 /// Newest-first listing.
366 pub fn list(&self) -> Result<Vec<CloudJob>> {
367 if !self.root.exists() {
368 return Ok(Vec::new());
369 }
370 let mut jobs = Vec::new();
371 for entry in fs::read_dir(&self.root).context("failed to read cloud-jobs")? {
372 let entry = entry?;
373 let name = entry.file_name();
374 let name = name.to_string_lossy();
375 if !name.starts_with("cloud_") || !name.ends_with(".json") {
376 continue;
377 }
378 let body = fs::read(entry.path())?;
379 if let Ok(job) = serde_json::from_slice::<CloudJob>(&body) {
380 jobs.push(job);
381 }
382 }
383 jobs.sort_by_key(|a| std::cmp::Reverse(a.created_unix));
384 Ok(jobs)
385 }
386
387 fn job_path(&self, id: &str) -> Result<PathBuf> {
388 if !valid_job_id(id) {
389 bail!("cloud job id must look like cloud_<hex>");
390 }
391 Ok(self.root.join(format!("{id}.json")))
392 }
393 }
394
395 /// Classify a git remote. Named `github` / `cnb` / `gitee` win over URL.
396 pub fn classify_remote(name: &str, url: &str) -> Option<Forge> {
397 match name.trim().to_ascii_lowercase().as_str() {
398 "github" => Some(Forge::Github),
399 "cnb" => Some(Forge::Cnb),
400 "gitee" => Some(Forge::Gitee),
401 _ => classify_url(url),
402 }
403 }
404
405 /// True when `branch` is a forge default that a non-force push could
406 /// fast-forward past review (`main` / `master` / `HEAD`).
407 pub fn is_forge_default_branch(branch: &str) -> bool {
408 matches!(
409 branch.trim().to_ascii_lowercase().as_str(),
410 "main" | "master" | "head"
411 )
412 }
413
414 /// Whether a git remote is safe to clone, show, or hand to a sandbox.
415 ///
416 /// Rejects leading-dash injection (`--upload-pack=…`), embedded userinfo,
417 /// and network remotes that do not classify as a supported forge. Local
418 /// path remotes (offline fixtures) are allowed when they do not start
419 /// with `-` and carry no userinfo.
420 pub fn safe_git_remote_url(raw: &str) -> bool {
421 validate_git_remote_url(raw).is_ok()
422 }
423
424 /// Classify and validate `job.remote_url` before any `git clone` or
425 /// sandbox clone. Returns the trimmed URL on success.
426 pub fn validate_git_remote_url(raw: &str) -> Result<String> {
427 let url = raw.trim();
428 if url.is_empty() || url.len() > MAX_REMOTE_BYTES {
429 bail!("remote url is empty or oversized");
430 }
431 if url.starts_with('-') {
432 bail!("remote url must not start with '-'");
433 }
434 if url.chars().any(char::is_control) {
435 bail!("remote url contains control characters");
436 }
437 if remote_has_userinfo(url) {
438 bail!("remote url must not embed userinfo");
439 }
440 if looks_like_network_git_url(url) && classify_url(url).is_none() {
441 bail!("remote url is not a supported forge");
442 }
443 Ok(url.to_string())
444 }
445
446 /// Display form of a remote: userinfo is never printed.
447 pub fn redact_remote_url(raw: &str) -> String {
448 redact_url_userinfo(raw)
449 }
450
451 fn remote_has_userinfo(url: &str) -> bool {
452 if let Ok(parsed) = reqwest::Url::parse(url) {
453 return !parsed.username().is_empty() || parsed.password().is_some();
454 }
455 // scp-style `user:token@host:path` (plain `git@host:path` is identity, not a secret).
456 if let Some((userinfo, _host)) = url.split_once('@') {
457 return userinfo.contains(':') && !userinfo.eq_ignore_ascii_case("git");
458 }
459 false
460 }
461
462 fn looks_like_network_git_url(url: &str) -> bool {
463 url.contains("://") || url.contains('@')
464 }
465
466 /// Strip `user:token@` (and URL userinfo) from free-form text for notes.
467 pub fn redact_url_userinfo(text: &str) -> String {
468 if let Ok(parsed) = reqwest::Url::parse(text)
469 && (!parsed.username().is_empty() || parsed.password().is_some())
470 {
471 let mut redacted = parsed;
472 let _ = redacted.set_username("");
473 let _ = redacted.set_password(None);
474 return redacted.to_string();
475 }
476 // scp-style or pasted `scheme://user:token@host`.
477 if let Some(scheme_end) = text.find("://") {
478 let rest = &text[scheme_end + 3..];
479 if let Some(at) = rest.find('@') {
480 let userinfo = &rest[..at];
481 if userinfo.contains(':') {
482 return format!("{}://[redacted]@{}", &text[..scheme_end], &rest[at + 1..]);
483 }
484 }
485 } else if let Some(at) = text.find('@') {
486 let userinfo = &text[..at];
487 if userinfo.contains(':') && !userinfo.eq_ignore_ascii_case("git") {
488 return format!("[redacted]@{}", &text[at + 1..]);
489 }
490 }
491 text.to_string()
492 }
493
494 /// Classify a clone URL by host. `origin` uses this path.
495 pub fn classify_url(url: &str) -> Option<Forge> {
496 let host = remote_host(url)?;
497 match host.as_str() {
498 "github.com" | "www.github.com" => Some(Forge::Github),
499 "cnb.cool" | "www.cnb.cool" => Some(Forge::Cnb),
500 "gitee.com" | "www.gitee.com" => Some(Forge::Gitee),
501 _ => None,
502 }
503 }
504
505 /// Parse `git remote -v` text into unique remotes (fetch URL preferred).
506 pub fn parse_remote_listing(text: &str) -> Vec<GitRemote> {
507 let mut remotes = Vec::new();
508 for line in text.lines() {
509 let line = line.trim();
510 if line.is_empty() || line.len() > MAX_REMOTE_BYTES {
511 continue;
512 }
513 let mut parts = line.split_whitespace();
514 let Some(name) = parts.next() else { continue };
515 let Some(url) = parts.next() else { continue };
516 if name.is_empty() || url.is_empty() || name.chars().any(char::is_control) {
517 continue;
518 }
519 if remotes.iter().any(|remote: &GitRemote| remote.name == name) {
520 continue;
521 }
522 remotes.push(GitRemote {
523 name: name.to_string(),
524 url: url.to_string(),
525 });
526 }
527 remotes
528 }
529
530 /// Read remotes from a workspace. Missing git is an empty list, not a panic.
531 pub fn discover_remotes(workspace: &Path) -> Vec<GitRemote> {
532 let output = Command::new("git")
533 .args(["-C", &workspace.to_string_lossy(), "remote", "-v"])
534 .output();
535 match output {
536 Ok(output) if output.status.success() => {
537 parse_remote_listing(&String::from_utf8_lossy(&output.stdout))
538 }
539 _ => Vec::new(),
540 }
541 }
542
543 /// Choose the forge remote. Multiple forges require an explicit request.
544 pub fn select_remote(
545 remotes: &[GitRemote],
546 requested: Option<Forge>,
547 ) -> Result<SelectedRemote, DispatchError> {
548 let classified: Vec<SelectedRemote> = remotes
549 .iter()
550 .filter_map(|remote| {
551 classify_remote(&remote.name, &remote.url).map(|forge| SelectedRemote {
552 forge,
553 name: remote.name.clone(),
554 url: remote.url.clone(),
555 })
556 })
557 .collect();
558
559 if let Some(forge) = requested {
560 return classified
561 .into_iter()
562 .find(|remote| remote.forge == forge)
563 .or_else(|| prefer_named(remotes, forge))
564 .ok_or(DispatchError::RemoteMissing(forge));
565 }
566
567 let mut unique = Vec::new();
568 for remote in classified {
569 if !unique
570 .iter()
571 .any(|seen: &SelectedRemote| seen.forge == remote.forge)
572 {
573 unique.push(remote);
574 }
575 }
576 match unique.len() {
577 0 => Err(DispatchError::NoSupportedRemote),
578 1 => Ok(unique.remove(0)),
579 _ => Err(DispatchError::AmbiguousRemote),
580 }
581 }
582
583 /// Build a plan. Does not spend, push, or write a job.
584 pub fn plan_dispatch(
585 remotes: &[GitRemote],
586 prompt: &str,
587 requested: Option<Forge>,
588 branch: Option<&str>,
589 ) -> Result<DispatchPlan, DispatchError> {
590 let prompt = validate_prompt(prompt)?;
591 let remote = select_remote(remotes, requested)?;
592 if !safe_git_remote_url(&remote.url) {
593 return Err(DispatchError::UnsafeRemote);
594 }
595 let branch = match branch.map(str::trim).filter(|value| !value.is_empty()) {
596 Some(value) => {
597 if !valid_branch(value) {
598 return Err(DispatchError::InvalidBranch);
599 }
600 value.to_string()
601 }
602 None => default_branch(),
603 };
604 Ok(DispatchPlan {
605 prompt,
606 remote,
607 branch,
608 })
609 }
610
611 /// Discover Daytona credentials without returning or logging the secret.
612 pub fn discover_credentials() -> CredentialState {
613 if env_present(DAYTONA_API_KEY_ENV) {
614 return CredentialState::Present {
615 source: CredentialSource::Env,
616 };
617 }
618 if env_present(CWC_DAYTONA_TOKEN_ENV) {
619 return CredentialState::Present {
620 source: CredentialSource::CwcEnv,
621 };
622 }
623 match Secrets::auto_detect().get(KEYRING_SLOT) {
624 Ok(Some(value)) if !value.trim().is_empty() => CredentialState::Present {
625 source: CredentialSource::Keyring,
626 },
627 _ => CredentialState::Missing,
628 }
629 }
630
631 /// Codewhale membership check. `/dispatch` cloud agents ship with the
632 /// account, so the fail-closed gate is sign-in — never a provider key.
633 fn membership_signed_in() -> bool {
634 let Ok(secrets) = codewhale_secrets::account::secure_account_session_secrets() else {
635 return false;
636 };
637 let store = codewhale_secrets::account::AccountSessionStore::new(
638 secrets,
639 None,
640 codewhale_secrets::account::DEFAULT_ACCOUNT_API_BASE,
641 );
642 matches!(store.load(), Ok(Some(_)))
643 }
644
645 /// Auto-decide leftover: Codewhale may propose, but never confirm itself.
646 pub fn should_auto_confirm(_plan: &DispatchPlan) -> bool {
647 false
648 }
649
650 /// Presence of the Codewhale account machine token that authenticates the
651 /// IN-SANDBOX agent. Mirrors [`CredentialState`]: a fact about availability,
652 /// never the value — the token is used in the create body and nothing else.
653 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
654 pub enum MachineTokenState {
655 Present,
656 Missing,
657 }
658
659 /// Discover the account machine token without returning or logging it.
660 /// The sandbox's `codewhale` resolves the account's configured model and
661 /// raises the PR as this identity, so a dispatch without it would spend on
662 /// a sandbox whose agent cannot authenticate — confirm refuses instead.
663 pub fn discover_machine_token() -> MachineTokenState {
664 if read_cloud_agent_token().is_some() {
665 MachineTokenState::Present
666 } else {
667 MachineTokenState::Missing
668 }
669 }
670
671 /// Truthful refusal for a confirmed dispatch with no account machine token.
672 /// Names the requirement and the fix, never a provider brand, never a secret.
673 pub fn missing_machine_token_message() -> String {
674 "cloud dispatch fails closed: the cloud agent runs Codewhale itself, so it \
675 needs a Codewhale account machine token to act as your account. Set \
676 CODEWHALE_API_KEY to a `cwc_key_…` machine key (Account → API keys in the \
677 web app) and confirm again."
678 .to_string()
679 }
680
681 /// Propose or launch. Confirmation and credentials are enforced here.
682 ///
683 /// With `confirm`, the job is persisted as `launching` and returned as
684 /// `Accepted`; the caller then starts the remote runner
685 /// ([`crate::dispatch_runner::run_confirmed_job`]) so this function stays
686 /// synchronous and offline-testable.
687 pub fn execute_dispatch(
688 store: &CloudJobStore,
689 plan: DispatchPlan,
690 confirm: bool,
691 credentials: &CredentialState,
692 machine_token: &MachineTokenState,
693 ) -> Result<DispatchOutcome> {
694 let mut job = CloudJob {
695 id: allocate_job_id(&plan),
696 kind: JOB_KIND.to_string(),
697 status: CloudJobStatus::Proposed,
698 prompt: plan.prompt.clone(),
699 forge: plan.remote.forge,
700 remote_name: plan.remote.name.clone(),
701 remote_url: plan.remote.url.clone(),
702 branch: plan.branch.clone(),
703 confirmed: false,
704 sandbox_id: None,
705 pr_url: None,
706 refusal: None,
707 note: proposal_note(&plan),
708 created_unix: unix_now(),
709 base_branch: None,
710 head_sha: None,
711 agent_summary: None,
712 finished_unix: None,
713 sandbox_pending: false,
714 };
715
716 if !confirm {
717 job.note = format!(
718 "{}. Confirm with `codewhale dispatch --confirm {}` or `/dispatch confirm {}`. Never silent spend, never silent push.",
719 job.note, job.id, job.id
720 );
721 store.save(&job)?;
722 return Ok(DispatchOutcome::Proposal(job));
723 }
724
725 job.confirmed = true;
726 if matches!(credentials, CredentialState::Missing) {
727 job.status = CloudJobStatus::Refused;
728 job.refusal = Some(missing_credentials_message());
729 job.note = missing_credentials_message();
730 job.finished_unix = Some(unix_now());
731 store.save(&job)?;
732 return Ok(DispatchOutcome::Refused(job));
733 }
734 if matches!(machine_token, MachineTokenState::Missing) {
735 job.status = CloudJobStatus::Refused;
736 job.refusal = Some(missing_machine_token_message());
737 job.note = missing_machine_token_message();
738 job.finished_unix = Some(unix_now());
739 store.save(&job)?;
740 return Ok(DispatchOutcome::Refused(job));
741 }
742
743 job.status = CloudJobStatus::Launching;
744 job.note = "Cloud agent confirmed; the sandbox is launching and the runner will raise the branch and open the PR. Watch `codewhale dispatch --show` or `/dispatch show`.".to_string();
745 store.save(&job)?;
746 Ok(DispatchOutcome::Accepted(job))
747 }
748
749 /// Confirm a previously proposed job — in place, under the SAME id.
750 ///
751 /// The record is mutated (`proposed` → `launching`, `confirmed = true`) and
752 /// saved as itself; a second confirm finds `launching` and refuses. Routing
753 /// through `execute_dispatch` instead would allocate a fresh job id (its ids
754 /// hash the plan plus `unix_now()` at second granularity), leaving the
755 /// proposal re-confirmable without limit — every confirm another sandbox
756 /// and another PR.
757 pub fn confirm_job(
758 store: &CloudJobStore,
759 id: &str,
760 credentials: &CredentialState,
761 machine_token: &MachineTokenState,
762 ) -> Result<DispatchOutcome> {
763 // One locked read-modify-write, as in `cancel_job`: two racing confirms
764 // cannot both pass the `proposed` check (two sandboxes, two PRs), and a
765 // cancel that lands mid-confirm cannot be overwritten back to `launching`.
766 store.with_job_lock(id, || {
767 let mut job = store.load(id)?;
768 if job.status != CloudJobStatus::Proposed {
769 bail!(
770 "Cloud job {} is {} and cannot be confirmed.",
771 job.id,
772 status_label(job.status)
773 );
774 }
775 job.confirmed = true;
776 if matches!(credentials, CredentialState::Missing) {
777 job.status = CloudJobStatus::Refused;
778 job.refusal = Some(missing_credentials_message());
779 job.note = missing_credentials_message();
780 job.finished_unix = Some(unix_now());
781 store.write_record(&job)?;
782 return Ok(DispatchOutcome::Refused(job));
783 }
784 if matches!(machine_token, MachineTokenState::Missing) {
785 job.status = CloudJobStatus::Refused;
786 job.refusal = Some(missing_machine_token_message());
787 job.note = missing_machine_token_message();
788 job.finished_unix = Some(unix_now());
789 store.write_record(&job)?;
790 return Ok(DispatchOutcome::Refused(job));
791 }
792
793 job.status = CloudJobStatus::Launching;
794 job.note = "Cloud agent confirmed; the sandbox is launching and the runner will raise the branch and open the PR. Watch `codewhale dispatch --show` or `/dispatch show`.".to_string();
795 store.write_record(&job)?;
796 Ok(DispatchOutcome::Accepted(job))
797 })
798 }
799
800 /// True while the job may still hold a live sandbox.
801 pub fn job_is_active(status: CloudJobStatus) -> bool {
802 matches!(
803 status,
804 CloudJobStatus::Launching | CloudJobStatus::Running | CloudJobStatus::OpeningPr
805 )
806 }
807
808 /// Cancel a job, tearing down a live sandbox when one exists.
809 ///
810 /// The record flips to `canceled` first (so a concurrent runner step sees it),
811 /// then teardown runs best-effort through the launcher; a teardown failure is
812 /// reported in the note, never silently dropped.
813 pub fn cancel_job(
814 store: &CloudJobStore,
815 id: &str,
816 launcher: &dyn DaytonaLauncher,
817 ) -> Result<CloudJob> {
818 // The flip is one locked read-modify-write, so a runner phase save cannot
819 // land between this load and this write (and lose either the cancel or
820 // the sandbox id the runner just recorded).
821 let flipped = store.with_job_lock(id, || {
822 let mut job = store.load(id)?;
823 // Terminal records stay as they are. `done` above all: its PR exists,
824 // and relabeling it `canceled` would misstate the outcome.
825 if matches!(
826 job.status,
827 CloudJobStatus::Canceled
828 | CloudJobStatus::Failed
829 | CloudJobStatus::Refused
830 | CloudJobStatus::Done
831 ) {
832 return Ok((job, false));
833 }
834 let sandbox_may_exist = job.sandbox_id.is_some() || job.sandbox_pending;
835 job.status = CloudJobStatus::Canceled;
836 job.finished_unix = Some(unix_now());
837 job.note = if sandbox_may_exist {
838 "Canceled locally; the cloud agent sandbox is being torn down.".to_string()
839 } else {
840 "Canceled locally before a sandbox was created.".to_string()
841 };
842 store.write_record(&job)?;
843 Ok((job, true))
844 })?;
845 let mut job = match flipped {
846 (job, false) => return Ok(job),
847 (job, true) => job,
848 };
849 if let Some(sandbox_id) = job.sandbox_id.clone() {
850 let receipt = SandboxReceipt {
851 sandbox_id,
852 toolbox_url: None,
853 };
854 match launcher.teardown(&receipt) {
855 Ok(()) => {
856 job.note = "Canceled locally; the cloud agent sandbox was torn down.".to_string();
857 }
858 Err(error) => {
859 job.note = format!(
860 "Canceled locally; sandbox teardown failed and may need a retry: {}",
861 sanitize_error(&error.to_string())
862 );
863 }
864 }
865 store.save(&job)?;
866 }
867 // A create whose POST landed but whose response never arrived leaves no
868 // recorded id (`sandbox_pending` with no `sandbox_id`). Best-effort
869 // label pass: delete any sandbox the provider still holds for this job.
870 if job.sandbox_pending && job.sandbox_id.is_none() {
871 match reconcile_job_sandboxes(launcher, &job.id) {
872 Ok(0) => {}
873 Ok(count) => {
874 job.note = format!(
875 "Canceled locally; {count} unrecorded cloud agent sandbox(es) labeled for this job were torn down."
876 );
877 let _ = store.save(&job);
878 }
879 Err(error) => {
880 job.note = format!(
881 "{} Unrecorded sandbox reconcile failed and may need a retry: {}",
882 job.note,
883 sanitize_error(&error.to_string())
884 );
885 let _ = store.save(&job);
886 }
887 }
888 }
889 Ok(job)
890 }
891
892 /// Outcome of one label-reconcile pass.
893 #[derive(Debug, Clone, PartialEq, Eq, Default)]
894 pub struct ReconcileReport {
895 /// Sandboxes deleted because their job is terminal or unknown.
896 pub deleted: Vec<String>,
897 /// Sandboxes left running because their job is still active.
898 pub live: usize,
899 }
900
901 /// Delete dispatch-labeled sandboxes whose job no longer needs them.
902 ///
903 /// Every sandbox [`LiveDaytonaLauncher::create_sandbox`] makes is labeled
904 /// with its job id and the dispatch product tag, so the provider's sandbox
905 /// list can be joined back to the store even when a record died mid-create.
906 /// A sandbox is deletable when its job label is missing or invalid, its job
907 /// record is absent from the store, or that job is terminal (`refused`,
908 /// `failed`, `canceled`, `done`). Sandboxes for active jobs are left alone —
909 /// their runner owns them. One sandbox failing to delete never stops the
910 /// rest; the report records what actually happened.
911 pub fn reconcile_sandboxes(
912 store: &CloudJobStore,
913 launcher: &dyn DaytonaLauncher,
914 ) -> Result<ReconcileReport> {
915 let mut report = ReconcileReport::default();
916 for sandbox in launcher.list_job_sandboxes()? {
917 let deletable = match sandbox
918 .job_id
919 .as_deref()
920 .filter(|job_id| valid_job_id(job_id))
921 {
922 None => true,
923 Some(job_id) => match store.load(job_id) {
924 Err(_) => true,
925 Ok(job) => !job_is_active(job.status),
926 },
927 };
928 if !deletable {
929 report.live += 1;
930 continue;
931 }
932 let receipt = SandboxReceipt {
933 sandbox_id: sandbox.sandbox_id.clone(),
934 toolbox_url: None,
935 };
936 if launcher.teardown(&receipt).is_ok() {
937 report.deleted.push(sandbox.sandbox_id);
938 }
939 }
940 Ok(report)
941 }
942
943 /// Best-effort: delete every dispatch sandbox labeled for one job id. Used
944 /// by cancel when the create POST may have landed without a receipt. Same
945 /// failure rule as [`reconcile_sandboxes`]: a sandbox that cannot be
946 /// deleted is skipped, not fatal.
947 pub fn reconcile_job_sandboxes(launcher: &dyn DaytonaLauncher, job_id: &str) -> Result<usize> {
948 let mut deleted = 0;
949 for sandbox in launcher.list_job_sandboxes()? {
950 if sandbox.job_id.as_deref() != Some(job_id) {
951 continue;
952 }
953 let receipt = SandboxReceipt {
954 sandbox_id: sandbox.sandbox_id.clone(),
955 toolbox_url: None,
956 };
957 if launcher.teardown(&receipt).is_ok() {
958 deleted += 1;
959 }
960 }
961 Ok(deleted)
962 }
963
964 /// Startup sweep: fail stale active jobs and tear their sandboxes down.
965 ///
966 /// The TUI spawns the runner detached, so quitting the TUI (or crashing)
967 /// leaves a `launching`/`running`/`openingpr` record whose runner is gone —
968 /// nothing else reconciles it, and any sandbox it made bills forever. On
969 /// startup, every active job older than [`STALE_ACTIVE_JOB_SECS`] is marked
970 /// `failed` with a truthful note, its recorded sandbox is torn down, and the
971 /// record is saved *before* teardown so a crash mid-sweep still leaves a
972 /// terminal record the label reconciler will clean up after. Returns the
973 /// swept records (empty when the store is unreadable — never fatal).
974 pub fn sweep_stale_jobs(
975 store: &CloudJobStore,
976 launcher: &dyn DaytonaLauncher,
977 now_unix: u64,
978 ) -> Vec<CloudJob> {
979 let Ok(jobs) = store.list() else {
980 return Vec::new();
981 };
982 let mut swept = Vec::new();
983 for mut job in jobs {
984 if !job_is_active(job.status) {
985 continue;
986 }
987 let age_secs = now_unix.saturating_sub(job.created_unix);
988 if age_secs < STALE_ACTIVE_JOB_SECS {
989 continue;
990 }
991 let sandbox_note = if job.sandbox_id.is_some() || job.sandbox_pending {
992 "Sandbox teardown was attempted; any sandbox left behind is deleted by the label reconciler on this startup."
993 } else {
994 "No sandbox was recorded for this job."
995 };
996 job.status = CloudJobStatus::Failed;
997 job.finished_unix = Some(now_unix);
998 job.refusal =
999 Some("stale: the runner stopped without recording a terminal state".to_string());
1000 job.note = format!(
1001 "Marked stale by the startup sweep: no terminal state for {} minutes and the declared harness budget is 60. {sandbox_note}",
1002 age_secs / 60,
1003 );
1004 // Persist the terminal record first: a crash mid-teardown must still
1005 // leave a failed job the label reconciler can clean up.
1006 if store.save(&job).is_ok() {
1007 if let Some(sandbox_id) = job.sandbox_id.clone() {
1008 let receipt = SandboxReceipt {
1009 sandbox_id,
1010 toolbox_url: None,
1011 };
1012 let _ = launcher.teardown(&receipt);
1013 }
1014 swept.push(job);
1015 }
1016 }
1017 swept
1018 }
1019
1020 /// Quit-path warning for the TUI: names live cloud jobs that quitting would
1021 /// leave behind. `None` when no job is `launching`/`running`/`openingpr` or
1022 /// the store cannot be read (a warning must never block the quit path).
1023 pub fn live_job_quit_warning(store: &CloudJobStore) -> Option<String> {
1024 let jobs = store.list().ok()?;
1025 let live: Vec<&str> = jobs
1026 .iter()
1027 .filter(|job| job_is_active(job.status))
1028 .map(|job| job.id.as_str())
1029 .take(3)
1030 .collect();
1031 if live.is_empty() {
1032 return None;
1033 }
1034 Some(format!(
1035 "Cloud job(s) {} still running: quitting now leaves the sandbox up until the next startup sweep reconciles it; /dispatch cancel <id> tears it down immediately.",
1036 live.join(", ")
1037 ))
1038 }
1039
1040 /// Human list used by `/jobs` (cloud kind) and `codewhale dispatch --list`.
1041 pub fn format_job_list(jobs: &[CloudJob]) -> String {
1042 if jobs.is_empty() {
1043 return "Cloud jobs (0)\nNo cloud-agent jobs yet. Use `codewhale dispatch <prompt>` or `/dispatch <prompt>`.".to_string();
1044 }
1045 let mut lines = vec![
1046 format!("Cloud jobs ({}) kind=cloud", jobs.len()),
1047 "----------------------------------------".to_string(),
1048 ];
1049 for job in jobs {
1050 lines.push(format!(
1051 "{} {:9} {} {} branch={}",
1052 job.id,
1053 status_label(job.status),
1054 job.forge.as_str(),
1055 job.remote_name,
1056 job.branch
1057 ));
1058 lines.push(format!(" prompt: {}", one_line(&job.prompt, 120)));
1059 if let Some(sandbox) = job.sandbox_id.as_ref() {
1060 lines.push(format!(" sandbox: {sandbox}"));
1061 }
1062 if let Some(pr) = job.pr_url.as_ref() {
1063 lines.push(format!(" pr: {pr}"));
1064 }
1065 if let Some(minutes) = runtime_minutes(job) {
1066 lines.push(format!(" runtime: {minutes}m"));
1067 }
1068 }
1069 lines.push(
1070 "Controls: /dispatch show <id>, /dispatch confirm <id>, /dispatch cancel <id>, /jobs list."
1071 .to_string(),
1072 );
1073 lines.join("\n")
1074 }
1075
1076 /// Whole minutes a terminal job was active, when internally known. This is
1077 /// local bookkeeping (created → finished), not a provider billing figure;
1078 /// sub-minute runs are omitted rather than rounded up.
1079 pub fn runtime_minutes(job: &CloudJob) -> Option<u64> {
1080 job.finished_unix
1081 .map(|end| end.saturating_sub(job.created_unix) / 60)
1082 .filter(|minutes| *minutes > 0)
1083 }
1084
1085 /// Job inspector used by `/dispatch show` and `/jobs show cloud_*`.
1086 pub fn format_job(job: &CloudJob) -> String {
1087 let mut lines = vec![
1088 format!("Cloud job {}", job.id),
1089 format!("Kind: {}", job.kind),
1090 format!("Status: {}", status_label(job.status)),
1091 format!("Forge: {}", job.forge.as_str()),
1092 format!(
1093 "Remote: {} {}",
1094 job.remote_name,
1095 redact_remote_url(&job.remote_url)
1096 ),
1097 format!("Branch: {}", job.branch),
1098 format!(
1099 "Base: {}",
1100 job.base_branch
1101 .as_deref()
1102 .unwrap_or("(detected at run time)")
1103 ),
1104 format!("Confirmed: {}", job.confirmed),
1105 format!("Sandbox: {}", job.sandbox_id.as_deref().unwrap_or("(none)")),
1106 format!("PR: {}", job.pr_url.as_deref().unwrap_or("(not opened)")),
1107 format!("Head: {}", job.head_sha.as_deref().unwrap_or("(pending)")),
1108 ];
1109 if let Some(minutes) = runtime_minutes(job) {
1110 lines.push(format!(
1111 "Runtime: {minutes}m (Codewhale bookkeeping, not a bill)"
1112 ));
1113 }
1114 if let Some(summary) = job.agent_summary.as_deref() {
1115 lines.push(format!("Agent: {}", one_line(summary, 200)));
1116 }
1117 lines.push(format!("Prompt: {}", job.prompt));
1118 lines.push(format!("Note: {}", job.note));
1119 if let Some(refusal) = job.refusal.as_ref() {
1120 lines.push(format!("Refusal: {refusal}"));
1121 }
1122 lines.join("\n")
1123 }
1124
1125 /// Status card for bare `/dispatch` and `codewhale dispatch --status`.
1126 ///
1127 /// `recent` is the newest slice of the job store; when the runner has
1128 /// receipts (sandbox id, PR URL, runtime) they are surfaced here verbatim.
1129 pub fn format_status(
1130 remotes: &[GitRemote],
1131 credentials: &CredentialState,
1132 recent: &[CloudJob],
1133 ) -> String {
1134 let mut lines = vec!["Codewhale cloud dispatch".to_string()];
1135 match credentials {
1136 CredentialState::Missing => {
1137 if membership_signed_in() {
1138 lines.push(
1139 "Cloud agents are not available for this account yet; cloud dispatch fails closed (no sandbox, no push, no PR)."
1140 .to_string(),
1141 );
1142 } else {
1143 lines.push(
1144 "Cloud agents are included with your Codewhale membership. Sign in with `codewhale login` to enable `/dispatch`; cloud dispatch fails closed until then (no sandbox, no push, no PR)."
1145 .to_string(),
1146 );
1147 }
1148 }
1149 CredentialState::Present { .. } => {
1150 lines.push(
1151 "Cloud agents: ready (account-linked). Confirmation is still required before spend or push."
1152 .to_string(),
1153 );
1154 }
1155 }
1156 if !recent.is_empty() {
1157 lines.push("Recent cloud jobs:".to_string());
1158 for job in recent.iter().take(5) {
1159 let mut line = format!(
1160 " {} {} {}",
1161 job.id,
1162 status_label(job.status),
1163 one_line(&job.prompt, 60)
1164 );
1165 if let Some(pr) = job.pr_url.as_deref() {
1166 line.push_str(&format!(" pr: {pr}"));
1167 } else if let Some(sandbox) = job.sandbox_id.as_deref() {
1168 line.push_str(&format!(" sandbox: {sandbox}"));
1169 }
1170 if let Some(minutes) = runtime_minutes(job) {
1171 line.push_str(&format!(" {minutes}m"));
1172 }
1173 lines.push(line);
1174 }
1175 }
1176 if remotes.is_empty() {
1177 lines.push("Remotes: none discovered.".to_string());
1178 } else {
1179 lines.push("Remotes:".to_string());
1180 for remote in remotes {
1181 let forge = classify_remote(&remote.name, &remote.url)
1182 .map(Forge::as_str)
1183 .unwrap_or("unsupported");
1184 lines.push(format!(
1185 " {} {forge} {}",
1186 remote.name,
1187 redact_remote_url(&remote.url)
1188 ));
1189 }
1190 }
1191 lines.push(
1192 "Offload: `codewhale dispatch \"<prompt>\" --remote github|cnb|gitee` then `--confirm`, or `/dispatch <prompt>` then `/dispatch confirm <id>`."
1193 .to_string(),
1194 );
1195 lines.join("\n")
1196 }
1197
1198 /// Resolve the API origin used by the live launcher. Never logs credentials.
1199 pub fn daytona_api_url() -> String {
1200 std::env::var(DAYTONA_API_URL_ENV)
1201 .ok()
1202 .or_else(|| std::env::var(CWC_DAYTONA_ENDPOINT_ENV).ok())
1203 .map(|value| value.trim().trim_end_matches('/').to_string())
1204 .filter(|value| !value.is_empty())
1205 .unwrap_or_else(|| DEFAULT_DAYTONA_API.to_string())
1206 }
1207
1208 /// Launch seam. Tests inject a recorder; production uses [`LiveDaytonaLauncher`].
1209 ///
1210 /// The methods after `create_sandbox` default to "unsupported" so partial
1211 /// fixtures keep compiling; the real runner (and the recording tests) drive
1212 /// the full protocol.
1213 pub trait DaytonaLauncher {
1214 /// Create the sandbox and return its id plus the (validated) toolbox URL.
1215 fn create_sandbox(&self, job: &CloudJob) -> Result<SandboxReceipt>;
1216 /// Block until the sandbox accepts toolbox calls (bounded poll).
1217 fn wait_ready(&self, _receipt: &SandboxReceipt) -> Result<()> {
1218 Ok(())
1219 }
1220 /// Clone the target forge repository inside the sandbox.
1221 fn clone_repository(&self, _receipt: &SandboxReceipt, _url: &str, _path: &str) -> Result<()> {
1222 bail!("this launcher does not support repository clones")
1223 }
1224 /// Run one harness command inside the sandbox and return bounded stdout.
1225 fn run_harness(&self, _receipt: &SandboxReceipt, _command: &HarnessCommand) -> Result<String> {
1226 bail!("this launcher does not support harness execution")
1227 }
1228 /// Collect the agent's work product from the sandbox.
1229 fn collect_patch(&self, _receipt: &SandboxReceipt) -> Result<PatchReceipt> {
1230 bail!("this launcher does not support patch collection")
1231 }
1232 /// Tear the sandbox down. Called on cancel, failure, and completion.
1233 fn teardown(&self, _receipt: &SandboxReceipt) -> Result<()> {
1234 bail!("this launcher does not support teardown")
1235 }
1236 /// List Codewhale-dispatch sandboxes with their job labels. Used by the
1237 /// reconciler (startup sweep and cancel); not part of the run protocol.
1238 fn list_job_sandboxes(&self) -> Result<Vec<LabeledSandbox>> {
1239 bail!("this launcher does not support sandbox listing")
1240 }
1241 }
1242
1243 /// One dispatch-labeled sandbox discovered by the reconciler.
1244 #[derive(Debug, Clone, PartialEq, Eq)]
1245 pub struct LabeledSandbox {
1246 pub sandbox_id: String,
1247 /// Job id from the `codewhale.job` label, when the label is present.
1248 pub job_id: Option<String>,
1249 }
1250
1251 /// Provider receipt for a created sandbox. `toolbox_url` is the validated
1252 /// per-sandbox toolbox origin returned by the create call, when present.
1253 /// `pr_url` is intentionally omitted from this slice.
1254 ///
1255 /// This is the infrastructure id of a created sandbox. It is not Computer
1256 /// entitlement and must not be treated as provider-accepted active seconds.
1257 #[derive(Debug, Clone, PartialEq, Eq)]
1258 pub struct SandboxReceipt {
1259 pub sandbox_id: String,
1260 pub toolbox_url: Option<String>,
1261 }
1262
1263 /// One command the runner asks the sandbox to execute. `argv` is exact: the
1264 /// recording tests pin it so the live protocol cannot drift silently.
1265 #[derive(Debug, Clone, PartialEq, Eq)]
1266 pub struct HarnessCommand {
1267 pub argv: Vec<String>,
1268 pub cwd: String,
1269 pub timeout_secs: u32,
1270 }
1271
1272 /// The agent's work product collected from the sandbox clone.
1273 #[derive(Debug, Clone, PartialEq, Eq)]
1274 pub struct PatchReceipt {
1275 /// Base branch of the clone (the PR base), e.g. `main`.
1276 pub base_branch: String,
1277 /// Head commit the agent produced.
1278 pub head_sha: String,
1279 /// One-line truthful summary of the agent's commit.
1280 pub summary: String,
1281 /// `git format-patch --stdout` payload (bounded) of the agent's commits.
1282 pub patch: String,
1283 }
1284
1285 /// Validate an outbound origin for credential-bearing HTTP calls.
1286 ///
1287 /// Rules:
1288 /// - `https` only for public hosts.
1289 /// - explicit loopback hosts (`localhost`, `127.0.0.1`, `::1`) are allowed
1290 /// only in debug builds, as the escape hatch for local smoke tests against
1291 /// a self-hosted sandbox service; release builds reject them outright.
1292 /// - the host must not be a private / link-local / reserved / multicast
1293 /// address or a `.local` / `.internal` name, and no userinfo may ride
1294 /// along.
1295 ///
1296 /// DNS-resolved rebinding is out of scope and documented as such.
1297 pub fn validate_outbound_origin(raw: &str) -> Result<reqwest::Url> {
1298 let trimmed = raw.trim();
1299 if trimmed.is_empty() || trimmed.len() > MAX_REMOTE_BYTES {
1300 bail!("outbound origin is empty or oversized");
1301 }
1302 let url = reqwest::Url::parse(trimmed).context("outbound origin is not a valid URL")?;
1303 if !matches!(url.scheme(), "http" | "https") {
1304 bail!("outbound origin must be http or https");
1305 }
1306 if !url.username().is_empty() || url.password().is_some() {
1307 bail!("outbound origin must not embed credentials");
1308 }
1309 let host = url
1310 .host_str()
1311 .context("outbound origin has no host")?
1312 .trim_end_matches('.')
1313 .to_ascii_lowercase();
1314 // `Url::host_str` keeps IPv6 brackets; strip them for the checks below.
1315 let host = host
1316 .strip_prefix('[')
1317 .and_then(|inner| inner.strip_suffix(']'))
1318 .map(str::to_string)
1319 .unwrap_or(host);
1320 let loopback_name = host == "localhost" || host == "127.0.0.1" || host == "::1";
1321 if loopback_name {
1322 if cfg!(debug_assertions) {
1323 return Ok(url);
1324 }
1325 bail!("loopback origins are not allowed in release builds");
1326 }
1327 if host.ends_with(".local") || host.ends_with(".internal") {
1328 bail!("outbound origin must be a public service host");
1329 }
1330 if let Ok(ip) = host.parse::<std::net::IpAddr>() {
1331 let blocked = match ip {
1332 std::net::IpAddr::V4(v4) => {
1333 let octets = v4.octets();
1334 v4.is_loopback()
1335 || v4.is_private()
1336 || v4.is_link_local()
1337 || v4.is_unspecified()
1338 || v4.is_broadcast()
1339 || v4.is_multicast()
1340 || v4.is_documentation()
1341 // 100.64.0.0/10 (carrier-grade NAT, `is_shared` is
1342 // not stable yet)
1343 || (octets[0] == 100 && (octets[1] & 0b1100_0000) == 0b0100_0000)
1344 }
1345 std::net::IpAddr::V6(v6) => {
1346 if let Some(v4) = v6.to_ipv4_mapped() {
1347 let octets = v4.octets();
1348 v4.is_loopback()
1349 || v4.is_private()
1350 || v4.is_link_local()
1351 || v4.is_unspecified()
1352 || v4.is_broadcast()
1353 || v4.is_multicast()
1354 || v4.is_documentation()
1355 || (octets[0] == 100 && (octets[1] & 0b1100_0000) == 0b0100_0000)
1356 } else {
1357 v6.is_loopback()
1358 || v6.is_unspecified()
1359 || v6.is_multicast()
1360 || (v6.segments()[0] & 0xfe00) == 0xfc00
1361 || (v6.segments()[0] & 0xffc0) == 0xfe80
1362 }
1363 }
1364 };
1365 if blocked {
1366 bail!("outbound origin must not target a loopback, private, or reserved address");
1367 }
1368 }
1369 if url.scheme() != "https" {
1370 bail!("outbound origin must use https");
1371 }
1372 Ok(url)
1373 }
1374
1375 /// Meter one closed interval on a dispatched cloud job.
1376 ///
1377 /// The job's sandbox id must match the provider observation. Wall-clock after
1378 /// create is not enough: the observation has to be provider-accepted active
1379 /// time bound to the immutable admission.
1380 pub fn meter_cloud_job(
1381 job: &CloudJob,
1382 admission: &ComputerAdmission,
1383 observation: ProviderObservation,
1384 ) -> Result<ComputerMeterReceipt, ComputerMeterError> {
1385 match job.sandbox_id.as_deref() {
1386 Some(sandbox_id) if sandbox_id == observation.provider_sandbox_id => {
1387 issue_computer_meter_receipt(admission, observation)
1388 }
1389 _ => Err(ComputerMeterError::AllocationMismatch {
1390 message: "Cloud job sandbox id does not match the provider allocation snapshot."
1391 .to_string(),
1392 }),
1393 }
1394 }
1395
1396 /// Real Daytona HTTP launcher. Fails closed on every step; never invents a
1397 /// PR URL; never logs or returns the API key.
1398 ///
1399 /// API shape (pinned against the published OpenAPI specs, see
1400 /// docs/DAYTONA_CLOUD_DISPATCH.md):
1401 /// - control plane `POST /sandbox`, `GET /sandbox/{id}`, `DELETE /sandbox/{id}`
1402 /// - toolbox `{toolboxProxyUrl}/{sandboxId}` with `POST /git/clone` and
1403 /// `POST /process/execute` (`{command, cwd, timeout}` → `{exitCode, result}`)
1404 pub struct LiveDaytonaLauncher;
1405
1406 impl LiveDaytonaLauncher {
1407 /// Total timeout for short control-plane calls (create/status/delete/
1408 /// list). A dispatched harness turn is NOT a short call — see
1409 /// [`Self::harness_client_budget_secs`].
1410 const CONTROL_PLANE_TIMEOUT_SECS: u64 = 120;
1411
1412 /// Slack added to a harness command's declared budget for the client
1413 /// that carries it: process start, clone drift, and response transfer
1414 /// are not part of the declared turn budget, but the client must still
1415 /// cut off eventually so a hung execute cannot hold a runner forever.
1416 const HARNESS_CLIENT_SLACK_SECS: u64 = 120;
1417
1418 fn blocking_client() -> Result<reqwest::blocking::Client> {
1419 // #6208: one client for the process. A `reqwest` client owns a
1420 // connection pool and a TLS configuration, so building one per call
1421 // paid a fresh TCP+TLS handshake on all nine call sites (one of them a
1422 // poll loop). `clone()` here is a refcount bump on that shared pool.
1423 //
1424 // The total timeout is attached per request instead, because a harness
1425 // turn carries a budget derived from its own command and must not
1426 // inherit the control-plane cap.
1427 static CLIENT: std::sync::OnceLock<Result<reqwest::blocking::Client, String>> =
1428 std::sync::OnceLock::new();
1429 CLIENT
1430 .get_or_init(|| {
1431 crate::tls::reqwest_blocking_client_builder()
1432 .connect_timeout(std::time::Duration::from_secs(8))
1433 .redirect(reqwest::redirect::Policy::none())
1434 .build()
1435 .map_err(|error| error.to_string())
1436 })
1437 .clone()
1438 .map_err(|message| anyhow!("failed to initialize the cloud agent client: {message}"))
1439 }
1440
1441 /// The total-timeout budget for a harness-carrying client, in seconds.
1442 /// Public to the crate so the runner's tests can pin it against the
1443 /// declared `HARNESS_TIMEOUT_SECS`.
1444 pub(crate) fn harness_client_budget_secs(command: &HarnessCommand) -> u64 {
1445 u64::from(command.timeout_secs).saturating_add(Self::HARNESS_CLIENT_SLACK_SECS)
1446 }
1447
1448 fn api_key() -> Result<String> {
1449 read_api_key().ok_or_else(|| anyhow!(missing_credentials_message()))
1450 }
1451
1452 /// Control-plane URL under the validated base.
1453 fn control_plane_url(path: &str) -> Result<reqwest::Url> {
1454 let base = validate_outbound_origin(&daytona_api_url())?;
1455 join_api_path(base, path).context("failed to build the cloud agent request URL")
1456 }
1457
1458 /// Toolbox base for one sandbox: `{toolboxProxyUrl}/{sandboxId}`.
1459 fn toolbox_base(receipt: &SandboxReceipt) -> Result<reqwest::Url> {
1460 if !valid_sandbox_id(&receipt.sandbox_id) {
1461 bail!("the sandbox id is not a usable path token");
1462 }
1463 let fallback = format!("{}/toolbox", DEFAULT_DAYTONA_API);
1464 let raw = receipt.toolbox_url.as_deref().unwrap_or(&fallback);
1465 let base = validate_outbound_origin(raw)?;
1466 join_api_path(base, &receipt.sandbox_id).context("failed to build the sandbox toolbox URL")
1467 }
1468
1469 /// Apply the dispatch labels via Daytona's dedicated labels endpoint.
1470 fn put_sandbox_labels(sandbox_id: &str, api_key: &str, job: &CloudJob) -> Result<()> {
1471 if !valid_sandbox_id(sandbox_id) {
1472 bail!("the sandbox id is not a usable path token");
1473 }
1474 let url = Self::control_plane_url(&format!("sandbox/{sandbox_id}/labels"))?;
1475 let body = serde_json::json!({
1476 "labels": {
1477 SANDBOX_JOB_LABEL: job.id,
1478 "codewhale.forge": job.forge.as_str(),
1479 SANDBOX_PRODUCT_LABEL: SANDBOX_PRODUCT_VALUE,
1480 }
1481 });
1482 let response = Self::send_json(reqwest::Method::PUT, &url, api_key, body)?;
1483 let status = response.status();
1484 if status.is_success() {
1485 Ok(())
1486 } else {
1487 bail!("cloud agent label apply failed (HTTP {status}).")
1488 }
1489 }
1490
1491 fn send_json(
1492 method: reqwest::Method,
1493 url: &reqwest::Url,
1494 api_key: &str,
1495 body: serde_json::Value,
1496 ) -> Result<reqwest::blocking::Response> {
1497 Self::send_json_on(
1498 &Self::blocking_client()?,
1499 Self::CONTROL_PLANE_TIMEOUT_SECS,
1500 method,
1501 url,
1502 api_key,
1503 body,
1504 )
1505 }
1506
1507 /// [`Self::send_json`] with an explicit total timeout, so a call whose
1508 /// declared budget differs from the control-plane cap (the harness turn)
1509 /// can carry its own without needing a client of its own.
1510 fn send_json_on(
1511 client: &reqwest::blocking::Client,
1512 total_secs: u64,
1513 method: reqwest::Method,
1514 url: &reqwest::Url,
1515 api_key: &str,
1516 body: serde_json::Value,
1517 ) -> Result<reqwest::blocking::Response> {
1518 client
1519 .request(method, url.clone())
1520 .timeout(std::time::Duration::from_secs(total_secs))
1521 .bearer_auth(api_key)
1522 .json(&body)
1523 .send()
1524 .context("could not reach the cloud agent service")
1525 }
1526 }
1527
1528 impl DaytonaLauncher for LiveDaytonaLauncher {
1529 fn create_sandbox(&self, job: &CloudJob) -> Result<SandboxReceipt> {
1530 let api_key = Self::api_key()?;
1531 let url = Self::control_plane_url("sandbox")?;
1532 let machine_token =
1533 read_cloud_agent_token().ok_or_else(|| anyhow!(missing_machine_token_message()))?;
1534 let body = create_sandbox_body(job, &machine_token, &cloud_agent_snapshot());
1535 let response = Self::send_json(reqwest::Method::POST, &url, &api_key, body)?;
1536 let status = response.status();
1537 let text = response.text().unwrap_or_default();
1538 if !status.is_success() {
1539 bail!("Cloud agent create failed (HTTP {status}).");
1540 }
1541 let parsed: serde_json::Value =
1542 serde_json::from_str(&text).context("the cloud agent service returned invalid JSON")?;
1543 let sandbox_id = parsed
1544 .get("id")
1545 .or_else(|| parsed.get("sandboxId"))
1546 .and_then(serde_json::Value::as_str)
1547 .unwrap_or("")
1548 .trim()
1549 .to_string();
1550 if !valid_sandbox_id(&sandbox_id) {
1551 // The provider says the sandbox exists (2xx) but gave us an id
1552 // we cannot safely interpolate into a path, so explicit teardown
1553 // is impossible. Best-effort delete with the raw string when it
1554 // is at least non-empty (the DELETE path itself validates and
1555 // will refuse dangerous shapes), and always name it in the
1556 // error so an operator can clean it up.
1557 let raw = parsed
1558 .get("id")
1559 .or_else(|| parsed.get("sandboxId"))
1560 .and_then(serde_json::Value::as_str)
1561 .unwrap_or("")
1562 .trim();
1563 if !raw.is_empty() && valid_sandbox_id(raw) {
1564 let _ = Self::send_json(
1565 reqwest::Method::DELETE,
1566 &Self::control_plane_url(&format!("sandbox/{raw}"))?,
1567 &api_key,
1568 serde_json::Value::Null,
1569 );
1570 }
1571 bail!(
1572 "Cloud agent create succeeded but returned no usable sandbox id (raw: \"{raw}\"). \
1573 A sandbox may need manual cleanup at the provider."
1574 );
1575 }
1576 let toolbox_url = parsed
1577 .get("toolboxProxyUrl")
1578 .and_then(serde_json::Value::as_str)
1579 .map(str::trim)
1580 .filter(|value| !value.is_empty() && value.len() <= MAX_REMOTE_BYTES)
1581 .and_then(|value| validate_outbound_origin(value).ok())
1582 .map(|url| url.to_string());
1583 // Daytona applies labels via a dedicated PUT, not the create body
1584 // (kept there for forward compatibility). Labels are load-bearing:
1585 // the orphan reconciler joins them back to job records, so a
1586 // sandbox without them is untraceable spend. If the PUT fails, the
1587 // already-created sandbox is torn down immediately and create
1588 // fails truthfully — no orphan, retryable — rather than returning
1589 // a receipt the reconciler can never find again.
1590 if let Err(label_error) = Self::put_sandbox_labels(&sandbox_id, &api_key, job) {
1591 let undo = Self::send_json(
1592 reqwest::Method::DELETE,
1593 &Self::control_plane_url(&format!("sandbox/{sandbox_id}"))?,
1594 &api_key,
1595 serde_json::Value::Null,
1596 )
1597 .and_then(|response| {
1598 let status = response.status();
1599 if status.is_success() || status.as_u16() == 404 {
1600 Ok(())
1601 } else {
1602 bail!("HTTP {status}")
1603 }
1604 });
1605 match undo {
1606 Ok(()) => bail!(
1607 "cloud agent created but its labels could not be applied ({}); \
1608 the sandbox was torn down — retry the job.",
1609 sanitize_error(&label_error.to_string())
1610 ),
1611 Err(undo_error) => bail!(
1612 "cloud agent {sandbox_id} was created but its labels could not be \
1613 applied ({}) and teardown also failed ({}); the sandbox needs \
1614 manual cleanup at the provider.",
1615 sanitize_error(&label_error.to_string()),
1616 sanitize_error(&undo_error.to_string())
1617 ),
1618 }
1619 }
1620 Ok(SandboxReceipt {
1621 sandbox_id,
1622 toolbox_url,
1623 })
1624 }
1625
1626 fn wait_ready(&self, receipt: &SandboxReceipt) -> Result<()> {
1627 let api_key = Self::api_key()?;
1628 if !valid_sandbox_id(&receipt.sandbox_id) {
1629 bail!("the sandbox id is not a usable path token");
1630 }
1631 let url = Self::control_plane_url(&format!("sandbox/{}", receipt.sandbox_id))?;
1632 for _ in 0..READY_POLL_ATTEMPTS {
1633 let response = Self::send_json(
1634 reqwest::Method::GET,
1635 &url,
1636 &api_key,
1637 serde_json::Value::Null,
1638 );
1639 match response {
1640 Ok(response) if response.status().is_success() => {
1641 let text = response.text().unwrap_or_default();
1642 if let Ok(parsed) = serde_json::from_str::<serde_json::Value>(&text)
1643 && let Some(state) = parsed.get("state").and_then(|v| v.as_str())
1644 {
1645 match state {
1646 "started" | "ready" => return Ok(()),
1647 "error" | "destroyed" | "archived" => {
1648 bail!("Cloud agent sandbox entered state '{state}'.");
1649 }
1650 _ => {}
1651 }
1652 } else {
1653 return Ok(());
1654 }
1655 }
1656 Err(error) => return Err(error),
1657 Ok(response) => {
1658 if response.status().as_u16() == 404 {
1659 bail!("Cloud agent sandbox disappeared before it was ready.");
1660 }
1661 }
1662 }
1663 std::thread::sleep(READY_POLL_INTERVAL);
1664 }
1665 bail!("Cloud agent sandbox was not ready in time.");
1666 }
1667
1668 fn clone_repository(&self, receipt: &SandboxReceipt, repo_url: &str, path: &str) -> Result<()> {
1669 let repo_url = validate_git_remote_url(repo_url)?;
1670 let api_key = Self::api_key()?;
1671 let url = Self::toolbox_base(receipt)?.join("git/clone")?;
1672 let body = serde_json::json!({ "url": repo_url, "path": path });
1673 let response = Self::send_json(reqwest::Method::POST, &url, &api_key, body)?;
1674 let status = response.status();
1675 if !status.is_success() {
1676 bail!("Cloud agent repository clone failed (HTTP {status}).");
1677 }
1678 Ok(())
1679 }
1680
1681 fn run_harness(&self, receipt: &SandboxReceipt, command: &HarnessCommand) -> Result<String> {
1682 let api_key = Self::api_key()?;
1683 let url = Self::toolbox_base(receipt)?.join("process/execute")?;
1684 // The toolbox executes one shell command string, so every argv
1685 // element is POSIX-single-quoted — a prompt cannot interpolate.
1686 let body = serde_json::json!({
1687 "command": shell_quote_join(&command.argv),
1688 "cwd": command.cwd,
1689 "timeout": command.timeout_secs,
1690 });
1691 // This call carries the declared turn budget (an hour for the agent
1692 // entry), so it asks for that budget plus slack per request — never the
1693 // 120s control-plane cap that used to bound it.
1694 let response = Self::send_json_on(
1695 &Self::blocking_client()?,
1696 Self::harness_client_budget_secs(command),
1697 reqwest::Method::POST,
1698 &url,
1699 &api_key,
1700 body,
1701 )?;
1702 let status = response.status();
1703 let text = response.text().unwrap_or_default();
1704 if !status.is_success() {
1705 bail!("Cloud agent harness execution failed (HTTP {status}).");
1706 }
1707 let parsed: serde_json::Value = serde_json::from_str(&text)
1708 .context("the sandbox returned an unreadable harness result")?;
1709 let exit_code = parsed
1710 .get("exitCode")
1711 .and_then(serde_json::Value::as_i64)
1712 .unwrap_or(0);
1713 let result = parsed
1714 .get("result")
1715 .and_then(serde_json::Value::as_str)
1716 .unwrap_or("");
1717 if exit_code != 0 {
1718 bail!(
1719 "Cloud agent harness exited with code {exit_code}: {}",
1720 sanitize_error(result)
1721 );
1722 }
1723 Ok(result.chars().take(MAX_HARNESS_OUTPUT_CHARS).collect())
1724 }
1725
1726 fn collect_patch(&self, receipt: &SandboxReceipt) -> Result<PatchReceipt> {
1727 let base_branch = self
1728 .run_harness(
1729 receipt,
1730 &HarnessCommand {
1731 argv: vec![
1732 "git".to_string(),
1733 "rev-parse".to_string(),
1734 "--abbrev-ref".to_string(),
1735 "origin/HEAD".to_string(),
1736 ],
1737 cwd: SANDBOX_WORKSPACE.to_string(),
1738 timeout_secs: 30,
1739 },
1740 )?
1741 .trim()
1742 .trim_start_matches("origin/")
1743 .to_string();
1744 let head_sha = self
1745 .run_harness(
1746 receipt,
1747 &HarnessCommand {
1748 argv: vec![
1749 "git".to_string(),
1750 "rev-parse".to_string(),
1751 "HEAD".to_string(),
1752 ],
1753 cwd: SANDBOX_WORKSPACE.to_string(),
1754 timeout_secs: 30,
1755 },
1756 )?
1757 .trim()
1758 .to_string();
1759 let summary = self
1760 .run_harness(
1761 receipt,
1762 &HarnessCommand {
1763 argv: vec![
1764 "git".to_string(),
1765 "log".to_string(),
1766 "-1".to_string(),
1767 "--format=%s".to_string(),
1768 ],
1769 cwd: SANDBOX_WORKSPACE.to_string(),
1770 timeout_secs: 30,
1771 },
1772 )?
1773 .trim()
1774 .to_string();
1775 let patch = self.run_harness(
1776 receipt,
1777 &HarnessCommand {
1778 argv: vec![
1779 "git".to_string(),
1780 "format-patch".to_string(),
1781 "origin/HEAD..HEAD".to_string(),
1782 "--stdout".to_string(),
1783 ],
1784 cwd: SANDBOX_WORKSPACE.to_string(),
1785 timeout_secs: 60,
1786 },
1787 )?;
1788 if base_branch.is_empty() || head_sha.len() < 7 {
1789 bail!("Cloud agent produced no branch head to raise.");
1790 }
1791 if patch.trim().is_empty() {
1792 bail!("Cloud agent produced an empty patch; refusing to open a PR.");
1793 }
1794 Ok(PatchReceipt {
1795 base_branch,
1796 head_sha,
1797 summary,
1798 patch,
1799 })
1800 }
1801
1802 fn teardown(&self, receipt: &SandboxReceipt) -> Result<()> {
1803 let api_key = Self::api_key()?;
1804 if !valid_sandbox_id(&receipt.sandbox_id) {
1805 bail!("the sandbox id is not a usable path token");
1806 }
1807 let url = Self::control_plane_url(&format!("sandbox/{}", receipt.sandbox_id))?;
1808 let response = Self::send_json(
1809 reqwest::Method::DELETE,
1810 &url,
1811 &api_key,
1812 serde_json::Value::Null,
1813 )?;
1814 // Daytona returns 204 on delete; treat 404 as already-gone success so
1815 // cancel/complete teardown is idempotent.
1816 let status = response.status();
1817 if status.is_success() || status.as_u16() == 404 {
1818 Ok(())
1819 } else {
1820 bail!("Cloud agent sandbox teardown failed (HTTP {status}).");
1821 }
1822 }
1823
1824 fn list_job_sandboxes(&self) -> Result<Vec<LabeledSandbox>> {
1825 let api_key = Self::api_key()?;
1826 // The provider's list call takes a JSON-encoded exact-match labels
1827 // filter (same OpenAPI family as create/get/delete above); filtering
1828 // on the product tag keeps the response to Codewhale dispatch
1829 // sandboxes only — never the user's own sandboxes on a shared key.
1830 let mut url = Self::control_plane_url("sandbox")?;
1831 url.query_pairs_mut().append_pair(
1832 "labels",
1833 &format!("{{\"{SANDBOX_PRODUCT_LABEL}\":\"{SANDBOX_PRODUCT_VALUE}\"}}"),
1834 );
1835 let response = Self::send_json(
1836 reqwest::Method::GET,
1837 &url,
1838 &api_key,
1839 serde_json::Value::Null,
1840 )?;
1841 let status = response.status();
1842 let text = response.text().unwrap_or_default();
1843 if !status.is_success() {
1844 bail!("Cloud agent sandbox listing failed (HTTP {status}).");
1845 }
1846 let parsed: serde_json::Value = serde_json::from_str(&text)
1847 .context("the cloud agent service returned an unreadable sandbox list")?;
1848 let rows = parsed.as_array().cloned().unwrap_or_default();
1849 let mut sandboxes = Vec::new();
1850 for row in rows {
1851 let sandbox_id = row
1852 .get("id")
1853 .and_then(serde_json::Value::as_str)
1854 .unwrap_or("")
1855 .trim()
1856 .to_string();
1857 if !valid_sandbox_id(&sandbox_id) {
1858 continue;
1859 }
1860 let job_id = row
1861 .get("labels")
1862 .and_then(|labels| labels.get(SANDBOX_JOB_LABEL))
1863 .and_then(serde_json::Value::as_str)
1864 .map(str::trim)
1865 .filter(|value| !value.is_empty())
1866 .map(str::to_string);
1867 sandboxes.push(LabeledSandbox { sandbox_id, job_id });
1868 }
1869 Ok(sandboxes)
1870 }
1871 }
1872
1873 /// Status / enablement copy shared by CLI and TUI fail-closed paths.
1874 /// Membership-first: the gate is sign-in, never a provider key.
1875 pub fn missing_credentials_message() -> String {
1876 if membership_signed_in() {
1877 "Cloud agents are not available for this account yet; cloud dispatch fails closed (no sandbox, no push, no PR).".to_string()
1878 } else {
1879 "Cloud agents are included with your Codewhale membership. Sign in with `codewhale login` to enable `/dispatch`; cloud dispatch fails closed until then (no sandbox, no push, no PR).".to_string()
1880 }
1881 }
1882
1883 fn prefer_named(remotes: &[GitRemote], forge: Forge) -> Option<SelectedRemote> {
1884 remotes
1885 .iter()
1886 .find(|remote| remote.name.eq_ignore_ascii_case(forge.as_str()))
1887 .map(|remote| SelectedRemote {
1888 forge,
1889 name: remote.name.clone(),
1890 url: remote.url.clone(),
1891 })
1892 }
1893
1894 fn validate_prompt(prompt: &str) -> Result<String, DispatchError> {
1895 let prompt = prompt.trim();
1896 if prompt.is_empty() {
1897 return Err(DispatchError::EmptyPrompt);
1898 }
1899 if prompt.chars().count() > MAX_PROMPT_CHARS {
1900 return Err(DispatchError::PromptTooLong);
1901 }
1902 Ok(prompt.to_string())
1903 }
1904
1905 fn valid_branch(branch: &str) -> bool {
1906 if is_forge_default_branch(branch) {
1907 return false;
1908 }
1909 if branch.is_empty()
1910 || branch.len() > 255
1911 || branch.starts_with('-')
1912 || branch.starts_with('/')
1913 || branch.ends_with('/')
1914 || branch.ends_with('.')
1915 || branch.contains("..")
1916 || branch.contains("//")
1917 || branch.contains("@{")
1918 {
1919 return false;
1920 }
1921 branch
1922 .bytes()
1923 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'/' | b'-'))
1924 }
1925
1926 fn valid_job_id(id: &str) -> bool {
1927 let Some(rest) = id.strip_prefix("cloud_") else {
1928 return false;
1929 };
1930 (8..=32).contains(&rest.len()) && rest.bytes().all(|byte| byte.is_ascii_hexdigit())
1931 }
1932
1933 fn allocate_job_id(plan: &DispatchPlan) -> String {
1934 use std::hash::{Hash, Hasher};
1935 let mut hasher = std::collections::hash_map::DefaultHasher::new();
1936 plan.prompt.hash(&mut hasher);
1937 plan.remote.forge.as_str().hash(&mut hasher);
1938 plan.branch.hash(&mut hasher);
1939 unix_now().hash(&mut hasher);
1940 format!("cloud_{:016x}", hasher.finish())
1941 }
1942
1943 fn default_branch() -> String {
1944 format!("codewhale/cloud-{}", unix_now())
1945 }
1946
1947 fn unix_now() -> u64 {
1948 SystemTime::now()
1949 .duration_since(UNIX_EPOCH)
1950 .map(|duration| duration.as_secs())
1951 .unwrap_or(0)
1952 }
1953
1954 /// Current unix time, shared with the runner for receipt timestamps.
1955 pub fn unix_timestamp() -> u64 {
1956 unix_now()
1957 }
1958
1959 fn env_present(name: &str) -> bool {
1960 std::env::var(name).is_ok_and(|value| !value.trim().is_empty())
1961 }
1962
1963 fn read_api_key() -> Option<String> {
1964 for name in [DAYTONA_API_KEY_ENV, CWC_DAYTONA_TOKEN_ENV] {
1965 if let Ok(value) = std::env::var(name) {
1966 let value = value.trim().to_string();
1967 if !value.is_empty() {
1968 return Some(value);
1969 }
1970 }
1971 }
1972 Secrets::auto_detect()
1973 .get(KEYRING_SLOT)
1974 .ok()
1975 .flatten()
1976 .map(|value| value.trim().to_string())
1977 .filter(|value| !value.is_empty())
1978 }
1979
1980 /// The account machine token for the in-sandbox agent. Read at create
1981 /// time only; used in the create body and never returned, logged, or
1982 /// persisted by the store. Shape-checked (`cwc_key_…`) so a misconfigured
1983 /// env var refuses at the confirm gate instead of paying for a sandbox
1984 /// whose agent can never authenticate.
1985 fn read_cloud_agent_token() -> Option<String> {
1986 let raw = std::env::var(CLOUD_AGENT_TOKEN_ENV).ok()?;
1987 machine_token_from_value(&raw)
1988 }
1989
1990 /// Pure core of [`read_cloud_agent_token`] and
1991 /// [`discover_machine_token`]: trim, shape-check, bound.
1992 pub(crate) fn machine_token_from_value(raw: &str) -> Option<String> {
1993 let value = raw.trim();
1994 let rest = value.strip_prefix("cwc_key_")?;
1995 // Documented shape: `cwc_key_` + 24 hex + separator + secret.
1996 if rest.len() < 25 {
1997 return None;
1998 }
1999 let (head, tail) = rest.split_at(24);
2000 if !head.bytes().all(|byte| byte.is_ascii_hexdigit()) {
2001 return None;
2002 }
2003 if !matches!(tail.as_bytes().first(), Some(b'_' | b'-')) {
2004 return None;
2005 }
2006 if !(MACHINE_TOKEN_MIN_BYTES..=MAX_MACHINE_TOKEN_BYTES).contains(&value.len()) {
2007 return None;
2008 }
2009 Some(value.to_string())
2010 }
2011
2012 /// `cwc_key_` (8) + at least the 24-hex id + one separator + a secret.
2013 const MACHINE_TOKEN_MIN_BYTES: usize = 40;
2014
2015 /// Bound on the injected machine token: real `cwc_key_…` keys are 76
2016 /// chars; the bound exists so a misconfigured env var cannot balloon the
2017 /// create body.
2018 const MAX_MACHINE_TOKEN_BYTES: usize = 512;
2019
2020 /// The cloud-agent snapshot to launch from. `CODEWHALE_DISPATCH_SNAPSHOT`
2021 /// overrides the default for operators; an invalid override falls back to
2022 /// the default rather than shipping an arbitrary string to the provider.
2023 fn cloud_agent_snapshot() -> String {
2024 std::env::var(CLOUD_AGENT_SNAPSHOT_ENV)
2025 .ok()
2026 .map(|value| value.trim().to_string())
2027 .filter(|value| valid_snapshot_name(value))
2028 .unwrap_or_else(|| DEFAULT_CLOUD_AGENT_SNAPSHOT.to_string())
2029 }
2030
2031 /// Snapshot names are provider path tokens: slug charset, bounded length.
2032 /// Same policy shape as [`valid_sandbox_id`].
2033 fn valid_snapshot_name(name: &str) -> bool {
2034 !name.is_empty()
2035 && name.len() <= 64
2036 && !name.starts_with(['.', '-'])
2037 && !name.contains("..")
2038 && name
2039 .chars()
2040 .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.')
2041 }
2042
2043 /// The cloud-agent sandbox create body (pure so tests pin the contract).
2044 ///
2045 /// Founder decision 2026-08-29: the sandbox launches from the Codewhale
2046 /// cloud-agent snapshot — Daytona snapshots are the only restorable
2047 /// create source (raw images are snapshot-build inputs) — with the CLI
2048 /// preinstalled, and the account machine token injected as
2049 /// `CODEWHALE_API_KEY` so the in-sandbox `codewhale exec` authenticates as
2050 /// the dispatching account and resolves the account's configured model.
2051 /// No provider API key ever widens into the sandbox (BYOK stays local);
2052 /// the sandbox speaks only with the Codewhale account.
2053 fn create_sandbox_body(job: &CloudJob, machine_token: &str, snapshot: &str) -> serde_json::Value {
2054 serde_json::json!({
2055 "name": format!("cw-{}", job.id.replace('_', "-")),
2056 "snapshot": snapshot,
2057 "env": {
2058 CLOUD_AGENT_TOKEN_ENV: machine_token,
2059 },
2060 "labels": {
2061 SANDBOX_JOB_LABEL: job.id,
2062 "codewhale.forge": job.forge.as_str(),
2063 SANDBOX_PRODUCT_LABEL: SANDBOX_PRODUCT_VALUE,
2064 }
2065 })
2066 }
2067
2068 /// Join an API path onto a base WITHOUT dropping the base's own path.
2069 /// `Url::join` with a relative path replaces the base's last segment
2070 /// (`https://host/api` + "sandbox" -> `https://host/sandbox`), which would
2071 /// silently drop the `/api` prefix every control-plane call needs — so the
2072 /// base is normalized to a trailing slash first.
2073 fn join_api_path(base: reqwest::Url, path: &str) -> Result<reqwest::Url> {
2074 let mut base = base;
2075 if !base.path().ends_with('/') {
2076 let path = format!("{}/", base.path());
2077 base.set_path(&path);
2078 }
2079 base.join(path.trim_start_matches('/'))
2080 .context("failed to join the cloud agent request path")
2081 }
2082
2083 fn remote_host(url: &str) -> Option<String> {
2084 let url = url.trim();
2085 if url.len() > MAX_REMOTE_BYTES {
2086 return None;
2087 }
2088 let host = if let Some(rest) = url
2089 .strip_prefix("https://")
2090 .or_else(|| url.strip_prefix("http://"))
2091 .or_else(|| url.strip_prefix("ssh://"))
2092 {
2093 let authority = rest.split('/').next()?;
2094 authority
2095 .rsplit_once('@')
2096 .map_or(authority, |(_, host)| host)
2097 } else if let Some((_, rest)) = url.split_once('@') {
2098 rest.split(':').next()?
2099 } else {
2100 return None;
2101 };
2102 let host = host.split(':').next()?.trim().to_ascii_lowercase();
2103 if host.is_empty() { None } else { Some(host) }
2104 }
2105
2106 fn forge_host(forge: Forge) -> &'static str {
2107 match forge {
2108 Forge::Github => "github.com",
2109 Forge::Cnb => "cnb.cool",
2110 Forge::Gitee => "gitee.com",
2111 }
2112 }
2113
2114 fn proposal_note(plan: &DispatchPlan) -> String {
2115 format!(
2116 "Proposed Codewhale cloud-agent offload to {} ({}) raising branch {}.",
2117 plan.remote.forge.as_str(),
2118 plan.remote.name,
2119 plan.branch
2120 )
2121 }
2122
2123 fn status_label(status: CloudJobStatus) -> &'static str {
2124 match status {
2125 CloudJobStatus::Proposed => "proposed",
2126 CloudJobStatus::Refused => "refused",
2127 CloudJobStatus::Launching => "launching",
2128 CloudJobStatus::Running => "running",
2129 CloudJobStatus::OpeningPr => "openingpr",
2130 CloudJobStatus::Done => "done",
2131 CloudJobStatus::Failed => "failed",
2132 CloudJobStatus::Canceled => "canceled",
2133 }
2134 }
2135
2136 fn one_line(value: &str, max: usize) -> String {
2137 let flat: String = value
2138 .chars()
2139 .map(|ch| if ch.is_control() { ' ' } else { ch })
2140 .collect();
2141 if flat.chars().count() <= max {
2142 flat
2143 } else {
2144 let mut out: String = flat.chars().take(max.saturating_sub(1)).collect();
2145 out.push('…');
2146 out
2147 }
2148 }
2149
2150 const SANITIZED_ERROR_MAX_CHARS: usize = 240;
2151
2152 /// Sanitized (single-line, secret-redacted, bounded) error text for job
2153 /// notes. Controls, newlines and tabs become one space so words never fuse;
2154 /// long text keeps its head and its tail, because harness and command output
2155 /// put the actual failure last.
2156 pub fn sanitize_error(message: &str) -> String {
2157 fn redact(text: &str) -> String {
2158 redact_machine_tokens(&redact_url_userinfo(text))
2159 }
2160 fn flatten(text: &str) -> String {
2161 text.split(|ch: char| ch.is_control() || ch.is_whitespace())
2162 .filter(|word| !word.is_empty())
2163 .collect::<Vec<_>>()
2164 .join(" ")
2165 }
2166 // A secret split by a control character (a newline, BEL, an ANSI
2167 // fragment) must still be caught, so redaction first sees the text with
2168 // every control removed, which rejoins it. Only when that finds nothing
2169 // do controls become spaces for readability; flattening only inserts
2170 // separators, so it cannot expose a token the rejoined text did not.
2171 let fused: String = message.chars().filter(|ch| !ch.is_control()).collect();
2172 let fused_redacted = redact(&fused);
2173 // Redact before cutting so a secret can never be split past the redactor.
2174 let redacted = if fused_redacted == fused {
2175 redact(&flatten(message))
2176 } else {
2177 flatten(&fused_redacted)
2178 };
2179 let count = redacted.chars().count();
2180 if count <= SANITIZED_ERROR_MAX_CHARS {
2181 return redacted;
2182 }
2183 let head_chars = SANITIZED_ERROR_MAX_CHARS / 3;
2184 let tail_chars = SANITIZED_ERROR_MAX_CHARS - head_chars - 1;
2185 let head: String = redacted.chars().take(head_chars).collect();
2186 let tail: String = redacted.chars().skip(count - tail_chars).collect();
2187 format!("{}…{}", head.trim_end(), tail.trim_start())
2188 }
2189
2190 /// Replace anything shaped like a Codewhale account machine token with its
2191 /// non-secret head + `[redacted]`. The sandbox environment carries
2192 /// `CODEWHALE_API_KEY`, so harness output and provider errors must never be
2193 /// able to echo a live token into a job record, note, or summary.
2194 pub fn redact_machine_tokens(text: &str) -> String {
2195 let mut out = String::with_capacity(text.len());
2196 let bytes = text.as_bytes();
2197 let mut i = 0;
2198 while i < bytes.len() {
2199 if text[i..].starts_with("cwc_key_") {
2200 let rest = &text[i + "cwc_key_".len()..];
2201 let end = rest
2202 .char_indices()
2203 .find(|(_, ch)| !ch.is_ascii_alphanumeric() && *ch != '_' && *ch != '-')
2204 .map(|(idx, _)| idx)
2205 .unwrap_or(rest.len());
2206 // The 24-hex id head is non-secret by design; the secret tail is not.
2207 if end >= 24 {
2208 // Non-secret head by design: `cwc_key_` + the 24-hex id.
2209 out.push_str("cwc_key_");
2210 out.push_str(&rest[..24]);
2211 out.push_str("_[redacted]");
2212 i += "cwc_key_".len() + end;
2213 continue;
2214 }
2215 }
2216 let ch = text[i..].chars().next().unwrap();
2217 out.push(ch);
2218 i += ch.len_utf8();
2219 }
2220 out
2221 }
2222
2223 /// Join argv into one POSIX shell command with every element single-quoted,
2224 /// so nothing (the prompt included) can interpolate when the sandbox toolbox
2225 /// executes the string.
2226 pub fn shell_quote_join(argv: &[String]) -> String {
2227 argv.iter()
2228 .map(|part| format!("'{}'", part.replace('\'', "'\\''")))
2229 .collect::<Vec<_>>()
2230 .join(" ")
2231 }
2232
2233 /// Sandbox ids must be plain URL-path-safe tokens before they are used in
2234 /// control-plane or toolbox paths.
2235 fn valid_sandbox_id(id: &str) -> bool {
2236 !id.is_empty()
2237 && id.len() <= 128
2238 && id
2239 .bytes()
2240 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_'))
2241 }
2242
2243 #[cfg(test)]
2244 mod tests {
2245 use super::*;
2246
2247 fn remotes(rows: &[(&str, &str)]) -> Vec<GitRemote> {
2248 rows.iter()
2249 .map(|(name, url)| GitRemote {
2250 name: (*name).to_string(),
2251 url: (*url).to_string(),
2252 })
2253 .collect()
2254 }
2255
2256 #[test]
2257 fn named_github_is_authoritative_even_when_origin_is_cnb() {
2258 assert_eq!(
2259 classify_remote("github", "https://cnb.cool/mirror/app.git"),
2260 Some(Forge::Github)
2261 );
2262 assert_eq!(
2263 classify_remote("origin", "https://cnb.cool/mirror/app.git"),
2264 Some(Forge::Cnb)
2265 );
2266 assert_eq!(
2267 classify_remote("origin", "https://github.com/codewhale-hq/CodeWhale.git"),
2268 Some(Forge::Github)
2269 );
2270 assert_eq!(
2271 classify_remote("gitee", "git@gitee.com:org/app.git"),
2272 Some(Forge::Gitee)
2273 );
2274 assert_eq!(
2275 classify_remote("upstream", "https://example.test/x.git"),
2276 None
2277 );
2278 }
2279
2280 #[test]
2281 fn ambiguous_github_and_cnb_require_an_explicit_remote() {
2282 let rows = remotes(&[
2283 ("github", "https://github.com/codewhale-hq/CodeWhale.git"),
2284 ("origin", "https://cnb.cool/codewhale.net/codewhale.git"),
2285 ]);
2286 assert_eq!(
2287 select_remote(&rows, None),
2288 Err(DispatchError::AmbiguousRemote)
2289 );
2290 let github = select_remote(&rows, Some(Forge::Github)).unwrap();
2291 assert_eq!(github.name, "github");
2292 let cnb = select_remote(&rows, Some(Forge::Cnb)).unwrap();
2293 assert_eq!(cnb.name, "origin");
2294 }
2295
2296 #[test]
2297 fn plan_rejects_empty_prompt_and_hostile_branch() {
2298 let rows = remotes(&[("github", "https://github.com/org/repo.git")]);
2299 assert_eq!(
2300 plan_dispatch(&rows, " ", None, None).unwrap_err(),
2301 DispatchError::EmptyPrompt
2302 );
2303 assert_eq!(
2304 plan_dispatch(&rows, "fix flake", None, Some("-bad;rm")).unwrap_err(),
2305 DispatchError::InvalidBranch
2306 );
2307 assert_eq!(
2308 plan_dispatch(&rows, "fix flake", None, Some("main")).unwrap_err(),
2309 DispatchError::InvalidBranch,
2310 "a non-force push must never target the forge default branch"
2311 );
2312 let plan = plan_dispatch(
2313 &rows,
2314 "fix flake",
2315 Some(Forge::Github),
2316 Some("codewhale/cloud-1"),
2317 )
2318 .unwrap();
2319 assert_eq!(plan.remote.forge, Forge::Github);
2320 assert_eq!(plan.branch, "codewhale/cloud-1");
2321 }
2322
2323 #[test]
2324 fn plan_rejects_leading_dash_and_userinfo_remotes() {
2325 assert_eq!(
2326 plan_dispatch(
2327 &remotes(&[("github", "--upload-pack=evil")]),
2328 "fix flake",
2329 Some(Forge::Github),
2330 None,
2331 )
2332 .unwrap_err(),
2333 DispatchError::UnsafeRemote
2334 );
2335 assert_eq!(
2336 plan_dispatch(
2337 &remotes(&[("github", "https://user:token@github.com/org/repo.git")]),
2338 "fix flake",
2339 Some(Forge::Github),
2340 None,
2341 )
2342 .unwrap_err(),
2343 DispatchError::UnsafeRemote
2344 );
2345 assert!(validate_git_remote_url("https://github.com/org/repo.git").is_ok());
2346 assert!(validate_git_remote_url("--upload-pack=evil").is_err());
2347 assert!(validate_git_remote_url("https://user:ghp_x@github.com/org/repo.git").is_err());
2348 }
2349
2350 #[test]
2351 fn unconfirmed_dispatch_writes_a_proposal_and_never_launches() {
2352 let temp = tempfile::tempdir().unwrap();
2353 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2354 let plan = plan_dispatch(
2355 &remotes(&[("github", "https://github.com/org/repo.git")]),
2356 "open a PR for the flake",
2357 Some(Forge::Github),
2358 Some("codewhale/cloud-test"),
2359 )
2360 .unwrap();
2361 let outcome = execute_dispatch(
2362 &store,
2363 plan,
2364 false,
2365 &CredentialState::Present {
2366 source: CredentialSource::Env,
2367 },
2368 &MachineTokenState::Present,
2369 )
2370 .unwrap();
2371 let DispatchOutcome::Proposal(job) = outcome else {
2372 panic!("expected proposal");
2373 };
2374 assert_eq!(job.status, CloudJobStatus::Proposed);
2375 assert!(!job.confirmed);
2376 assert!(job.sandbox_id.is_none());
2377 assert!(job.pr_url.is_none());
2378 assert_eq!(job.kind, "cloud");
2379 assert!(!should_auto_confirm(&DispatchPlan {
2380 prompt: job.prompt.clone(),
2381 remote: SelectedRemote {
2382 forge: job.forge,
2383 name: job.remote_name.clone(),
2384 url: job.remote_url.clone(),
2385 },
2386 branch: job.branch.clone(),
2387 }));
2388 }
2389
2390 #[test]
2391 fn confirmed_dispatch_fails_closed_without_credentials() {
2392 let temp = tempfile::tempdir().unwrap();
2393 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2394 let plan = plan_dispatch(
2395 &remotes(&[("origin", "https://cnb.cool/org/repo.git")]),
2396 "raise a CNB PR",
2397 Some(Forge::Cnb),
2398 Some("codewhale/cloud-cnb"),
2399 )
2400 .unwrap();
2401 let outcome = execute_dispatch(
2402 &store,
2403 plan,
2404 true,
2405 &CredentialState::Missing,
2406 &MachineTokenState::Present,
2407 )
2408 .unwrap();
2409 let DispatchOutcome::Refused(job) = outcome else {
2410 panic!("expected refuse");
2411 };
2412 assert_eq!(job.status, CloudJobStatus::Refused);
2413 assert!(job.confirmed);
2414 assert!(job.sandbox_id.is_none());
2415 assert!(job.pr_url.is_none());
2416 assert!(job.note.contains("cloud dispatch fails closed"));
2417 assert!(!job.note.contains("DAYTONA"));
2418 assert!(!job.note.contains("sk-"));
2419 }
2420
2421 #[test]
2422 fn confirmed_dispatch_without_machine_token_refuses_truthfully() {
2423 let temp = tempfile::tempdir().unwrap();
2424 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2425 let plan = plan_dispatch(
2426 &remotes(&[("origin", "https://cnb.cool/org/repo.git")]),
2427 "raise a CNB PR",
2428 Some(Forge::Cnb),
2429 Some("codewhale/cloud-cnb"),
2430 )
2431 .unwrap();
2432 // Daytona credentials ARE present: the machine token is the missing
2433 // fact, so the refusal must be the account-token one.
2434 let outcome = execute_dispatch(
2435 &store,
2436 plan,
2437 true,
2438 &CredentialState::Present {
2439 source: CredentialSource::Env,
2440 },
2441 &MachineTokenState::Missing,
2442 )
2443 .unwrap();
2444 let DispatchOutcome::Refused(job) = outcome else {
2445 panic!("expected refuse");
2446 };
2447 assert_eq!(job.status, CloudJobStatus::Refused);
2448 assert!(job.confirmed);
2449 assert!(job.sandbox_id.is_none());
2450 assert!(job.finished_unix.is_some(), "the refusal is terminal");
2451 assert!(job.note.contains("CODEWHALE_API_KEY"));
2452 assert!(job.note.contains("cwc_key_"));
2453 assert!(job.note.contains("cloud dispatch fails closed"));
2454 // No provider brand in user copy, no secret-shaped text.
2455 assert!(!job.note.contains("Daytona"));
2456 assert!(!job.note.contains("sk-"));
2457 }
2458
2459 #[test]
2460 fn confirm_without_machine_token_refuses_in_place() {
2461 let temp = tempfile::tempdir().unwrap();
2462 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2463 let plan = plan_dispatch(
2464 &remotes(&[("gitee", "https://gitee.com/org/repo.git")]),
2465 "gitee offload",
2466 Some(Forge::Gitee),
2467 Some("codewhale/cloud-confirm-token"),
2468 )
2469 .unwrap();
2470 let DispatchOutcome::Proposal(job) = execute_dispatch(
2471 &store,
2472 plan,
2473 false,
2474 &CredentialState::Present {
2475 source: CredentialSource::Env,
2476 },
2477 &MachineTokenState::Present,
2478 )
2479 .unwrap() else {
2480 panic!("expected proposal");
2481 };
2482 let id = job.id.clone();
2483
2484 let refused = match confirm_job(
2485 &store,
2486 &id,
2487 &CredentialState::Present {
2488 source: CredentialSource::Env,
2489 },
2490 &MachineTokenState::Missing,
2491 )
2492 .unwrap()
2493 {
2494 DispatchOutcome::Refused(job) => job,
2495 other => panic!("expected refuse, got {other:?}"),
2496 };
2497 // Refused in place: SAME id, one record, terminal.
2498 assert_eq!(refused.id, id);
2499 assert_eq!(refused.status, CloudJobStatus::Refused);
2500 assert!(refused.note.contains("CODEWHALE_API_KEY"));
2501 let listed = store.list().unwrap();
2502 assert_eq!(listed.len(), 1);
2503 assert_eq!(listed[0].id, id);
2504 // A second confirm cannot re-mint a run from the refusal.
2505 assert!(
2506 confirm_job(
2507 &store,
2508 &id,
2509 &CredentialState::Present {
2510 source: CredentialSource::Env,
2511 },
2512 &MachineTokenState::Present,
2513 )
2514 .is_err()
2515 );
2516 }
2517
2518 #[test]
2519 fn create_body_launches_from_the_cloud_agent_snapshot_with_the_account_token() {
2520 let job = CloudJob {
2521 id: "cloud_deadbeef".to_string(),
2522 kind: JOB_KIND.to_string(),
2523 status: CloudJobStatus::Launching,
2524 prompt: "pin the create contract".to_string(),
2525 forge: Forge::Github,
2526 remote_name: "origin".to_string(),
2527 remote_url: "https://github.com/org/repo.git".to_string(),
2528 branch: "codewhale/cloud-create-body".to_string(),
2529 confirmed: true,
2530 sandbox_id: None,
2531 pr_url: None,
2532 refusal: None,
2533 note: String::new(),
2534 created_unix: 1,
2535 base_branch: None,
2536 head_sha: None,
2537 agent_summary: None,
2538 finished_unix: None,
2539 sandbox_pending: false,
2540 };
2541 let token = "cwc_key_0123456789abcdef01234567_AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
2542 let body = create_sandbox_body(&job, token, DEFAULT_CLOUD_AGENT_SNAPSHOT);
2543 // The sandbox launches from the codewhale cloud-agent snapshot…
2544 assert_eq!(
2545 body.get("snapshot").and_then(serde_json::Value::as_str),
2546 Some(DEFAULT_CLOUD_AGENT_SNAPSHOT)
2547 );
2548 // …with exactly one injected env var: the account machine token.
2549 let env = body.get("env").expect("env present");
2550 assert_eq!(
2551 env.get(CLOUD_AGENT_TOKEN_ENV)
2552 .and_then(serde_json::Value::as_str),
2553 Some(token)
2554 );
2555 let env_len = env.as_object().map(|map| map.len()).unwrap_or(0);
2556 assert_eq!(env_len, 1, "no provider key widens into the sandbox");
2557 // Labels stay: the reconciler's join keys.
2558 let labels = body.get("labels").expect("labels present");
2559 assert_eq!(
2560 labels
2561 .get(SANDBOX_JOB_LABEL)
2562 .and_then(serde_json::Value::as_str),
2563 Some("cloud_deadbeef")
2564 );
2565 assert_eq!(
2566 labels
2567 .get(SANDBOX_PRODUCT_LABEL)
2568 .and_then(serde_json::Value::as_str),
2569 Some(SANDBOX_PRODUCT_VALUE)
2570 );
2571 // No provider API key ever ships into the sandbox.
2572 let serialized = body.to_string();
2573 assert!(!serialized.contains("DAYTONA_API_KEY"));
2574 assert!(!serialized.contains("sk-"));
2575 }
2576
2577 #[test]
2578 fn join_api_path_preserves_the_base_path_segment() {
2579 let base = reqwest::Url::parse("https://app.daytona.io/api").unwrap();
2580 let joined = join_api_path(base, "sandbox").unwrap();
2581 assert_eq!(
2582 joined.as_str(),
2583 "https://app.daytona.io/api/sandbox",
2584 "Url::join would drop the /api segment without the trailing-slash fix"
2585 );
2586 let base = reqwest::Url::parse("https://proxy.internal/control/").unwrap();
2587 let joined = join_api_path(base, "sandbox/abc/labels").unwrap();
2588 assert_eq!(
2589 joined.as_str(),
2590 "https://proxy.internal/control/sandbox/abc/labels"
2591 );
2592 }
2593
2594 #[test]
2595 fn machine_tokens_are_shape_checked_before_they_can_spend() {
2596 let good = "cwc_key_0123456789abcdef01234567_AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
2597 assert_eq!(machine_token_from_value(good).as_deref(), Some(good));
2598 // A bare env var (wrong secret family) never authorizes a sandbox.
2599 assert_eq!(machine_token_from_value("sk-not-a-machine-key"), None);
2600 assert_eq!(machine_token_from_value(""), None);
2601 assert_eq!(machine_token_from_value("cwc_key_short"), None);
2602 assert_eq!(
2603 machine_token_from_value(&format!("cwc_key_{}", "x".repeat(600))),
2604 None,
2605 "over-long values are refused, not truncated"
2606 );
2607 assert_eq!(
2608 machine_token_from_value(&format!(" {good} ")).as_deref(),
2609 Some(good),
2610 "surrounding whitespace is trimmed"
2611 );
2612 assert_eq!(
2613 machine_token_from_value("cwc_key_zzzzzzzzzzzzzzzzzzzzzzzz_nothex"),
2614 None,
2615 "the 24-char id must be hex"
2616 );
2617 assert_eq!(
2618 machine_token_from_value("cwc_key_0123456789abcdef01234567XXXX"),
2619 None,
2620 "documented shape requires a separator after the 24-hex id"
2621 );
2622 }
2623
2624 #[test]
2625 fn machine_tokens_never_survive_into_errors_or_summaries() {
2626 let token = "cwc_key_0123456789abcdef01234567_AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
2627 assert_eq!(
2628 redact_machine_tokens(&format!("auth failed for {token}")),
2629 "auth failed for cwc_key_0123456789abcdef01234567_[redacted]"
2630 );
2631 // Too short to be a real key head stays untouched (no false positives).
2632 assert_eq!(
2633 redact_machine_tokens("prefix cwc_key_abc suffix"),
2634 "prefix cwc_key_abc suffix"
2635 );
2636 // `dispatch_runner` tests that its summary path inherits this
2637 // redaction; asserting it here would pull a late module into the
2638 // runtime closure (scripts/split/module_graph.py).
2639 assert!(
2640 !sanitize_error("clone failed https://user:token@github.com/org/repo.git")
2641 .contains("token"),
2642 "userinfo must not survive sanitize_error"
2643 );
2644 }
2645
2646 #[test]
2647 fn sanitize_error_redacts_a_token_split_by_a_control_character() {
2648 let head = "0123456789abcdef01234567";
2649 for split in ['\u{7}', '\n', '\u{1b}'] {
2650 let message = format!("push failed: cwc_key_{head}{split}tailsecret0123456789 done");
2651 let sanitized = sanitize_error(&message);
2652 assert!(!sanitized.contains("tailsecret"), "{sanitized}");
2653 assert!(sanitized.contains("[redacted]"), "{sanitized}");
2654 }
2655 }
2656
2657 #[test]
2658 fn sanitize_error_flattens_whitespace_and_keeps_the_tail() {
2659 assert_eq!(
2660 sanitize_error("npm ERR!\tcode 1\n\n npm ERR!\u{7} missing\r\nscript"),
2661 "npm ERR! code 1 npm ERR! missing script"
2662 );
2663 let long = format!(
2664 "{} error: the real failure is here",
2665 "progress line\n".repeat(60)
2666 );
2667 let sanitized = sanitize_error(&long);
2668 assert!(
2669 sanitized.chars().count() <= SANITIZED_ERROR_MAX_CHARS,
2670 "{sanitized}"
2671 );
2672 assert!(
2673 sanitized.starts_with("progress line progress line"),
2674 "{sanitized}"
2675 );
2676 assert!(sanitized.contains('…'), "{sanitized}");
2677 assert!(
2678 sanitized.ends_with("error: the real failure is here"),
2679 "{sanitized}"
2680 );
2681 assert!(!sanitized.contains('\n'), "{sanitized}");
2682 }
2683
2684 #[test]
2685 fn snapshot_names_are_slug_charset_and_bounded() {
2686 assert!(valid_snapshot_name("codewhale-cloud-agent"));
2687 assert!(valid_snapshot_name("team.agent_2026"));
2688 assert!(!valid_snapshot_name(""));
2689 assert!(!valid_snapshot_name("has space"));
2690 assert!(!valid_snapshot_name("unicodé"));
2691 assert!(!valid_snapshot_name(".hidden"));
2692 assert!(!valid_snapshot_name("-leading-dash"));
2693 assert!(!valid_snapshot_name("dot..dot"));
2694 assert!(!valid_snapshot_name("a".repeat(65).as_str()));
2695 assert!(valid_snapshot_name("a".repeat(64).as_str()));
2696 }
2697
2698 #[test]
2699 fn confirmed_dispatch_queues_launching_without_touching_the_forge() {
2700 let temp = tempfile::tempdir().unwrap();
2701 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2702 let plan = plan_dispatch(
2703 &remotes(&[("gitee", "https://gitee.com/org/repo.git")]),
2704 "gitee offload",
2705 Some(Forge::Gitee),
2706 Some("codewhale/cloud-gitee"),
2707 )
2708 .unwrap();
2709 let outcome = execute_dispatch(
2710 &store,
2711 plan,
2712 true,
2713 &CredentialState::Present {
2714 source: CredentialSource::Keyring,
2715 },
2716 &MachineTokenState::Present,
2717 )
2718 .unwrap();
2719 let DispatchOutcome::Accepted(job) = outcome else {
2720 panic!("expected accept");
2721 };
2722 assert_eq!(job.status, CloudJobStatus::Launching);
2723 assert!(job.confirmed);
2724 assert!(job.sandbox_id.is_none());
2725 assert!(job.pr_url.is_none());
2726 assert!(job.note.contains("runner will raise the branch"));
2727 let listed = store.list().unwrap();
2728 assert_eq!(listed.len(), 1);
2729 assert_eq!(listed[0].id, job.id);
2730 // No sandbox exists yet, so cancel is a pure record flip.
2731 let canceled = cancel_job(&store, &job.id, &NoopLauncher).unwrap();
2732 assert_eq!(canceled.status, CloudJobStatus::Canceled);
2733 assert!(canceled.note.contains("before a sandbox"));
2734 }
2735
2736 #[test]
2737 fn confirm_job_confirms_in_place_under_the_same_id_and_is_not_reconfirmable() {
2738 let temp = tempfile::tempdir().unwrap();
2739 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2740 let plan = plan_dispatch(
2741 &remotes(&[("github", "https://github.com/org/repo.git")]),
2742 "open a PR for the flake",
2743 Some(Forge::Github),
2744 Some("codewhale/cloud-confirm"),
2745 )
2746 .unwrap();
2747 let DispatchOutcome::Proposal(job) = execute_dispatch(
2748 &store,
2749 plan,
2750 false,
2751 &CredentialState::Present {
2752 source: CredentialSource::Env,
2753 },
2754 &MachineTokenState::Present,
2755 )
2756 .unwrap() else {
2757 panic!("expected proposal");
2758 };
2759 let id = job.id.clone();
2760
2761 let confirmed = match confirm_job(
2762 &store,
2763 &id,
2764 &CredentialState::Present {
2765 source: CredentialSource::Env,
2766 },
2767 &MachineTokenState::Present,
2768 )
2769 .unwrap()
2770 {
2771 DispatchOutcome::Accepted(job) => job,
2772 other => panic!("expected accept, got {other:?}"),
2773 };
2774 // Same id, one record, launching: the proposal became the run.
2775 assert_eq!(confirmed.id, id);
2776 assert_eq!(confirmed.status, CloudJobStatus::Launching);
2777 assert!(confirmed.confirmed);
2778 let listed = store.list().unwrap();
2779 assert_eq!(listed.len(), 1, "confirm must not mint a second record");
2780 assert_eq!(listed[0].id, id);
2781
2782 // A second confirm is refused — the status gate, not the id, is the
2783 // guard (ids hash unix_now() at second granularity and can collide
2784 // across confirms).
2785 let again = confirm_job(
2786 &store,
2787 &id,
2788 &CredentialState::Present {
2789 source: CredentialSource::Env,
2790 },
2791 &MachineTokenState::Present,
2792 )
2793 .unwrap_err()
2794 .to_string();
2795 assert!(again.contains("cannot be confirmed"), "{again}");
2796 assert_eq!(store.list().unwrap().len(), 1);
2797 }
2798
2799 #[test]
2800 fn confirm_job_without_credentials_refuses_in_place_under_the_same_id() {
2801 let temp = tempfile::tempdir().unwrap();
2802 let store = CloudJobStore::from_path(temp.path().join("jobs"));
2803 let plan = plan_dispatch(
2804 &remotes(&[("github", "https://github.com/org/repo.git")]),
2805 "refuse me in place",
2806 Some(Forge::Github),
2807 Some("codewhale/cloud-refuse"),
2808 )
2809 .unwrap();
2810 let DispatchOutcome::Proposal(job) = execute_dispatch(
2811 &store,
2812 plan,
2813 false,
2814 &CredentialState::Present {
2815 source: CredentialSource::Env,
2816 },
2817 &MachineTokenState::Present,
2818 )
2819 .unwrap() else {
2820 panic!("expected proposal");
2821 };
2822 let id = job.id.clone();
2823 let refused = match confirm_job(
2824 &store,
2825 &id,
2826 &CredentialState::Missing,
2827 &MachineTokenState::Present,
2828 )
2829 .unwrap()
2830 {
2831 DispatchOutcome::Refused(job) => job,
2832 other => panic!("expected refuse, got {other:?}"),
2833 };
2834 assert_eq!(refused.id, id);
2835 assert_eq!(refused.status, CloudJobStatus::Refused);
2836 assert!(refused.confirmed);
2837 assert!(refused.finished_unix.is_some());
2838 assert_eq!(store.list().unwrap().len(), 1);
2839 // And it cannot be confirmed again either.
2840 assert!(
2841 confirm_job(
2842 &store,
2843 &id,
2844 &CredentialState::Missing,
2845 &MachineTokenState::Present
2846 )
2847 .is_err()
2848 );
2849 }
2850
2851 /// Launcher that can create but never tear down; used to pin the
2852 /// cancel-without-sandbox path without any network surface.
2853 struct NoopLauncher;
2854
2855 impl DaytonaLauncher for NoopLauncher {
2856 fn create_sandbox(&self, _job: &CloudJob) -> Result<SandboxReceipt> {
2857 bail!("no sandbox in this fixture")
2858 }
2859 }
2860
2861 #[test]
2862 fn old_job_records_without_runner_fields_still_load() {
2863 let temp = tempfile::tempdir().unwrap();
2864 let root = temp.path().join("jobs");
2865 std::fs::create_dir_all(&root).unwrap();
2866 // Exactly the JSON shape written by the first landing of the slice.
2867 let legacy = serde_json::json!({
2868 "id": "cloud_00000000000000ff",
2869 "kind": "cloud",
2870 "status": "running",
2871 "prompt": "legacy job",
2872 "forge": "github",
2873 "remote_name": "github",
2874 "remote_url": "https://github.com/org/repo.git",
2875 "branch": "codewhale/cloud-legacy",
2876 "confirmed": true,
2877 "sandbox_id": "sandbox_legacy",
2878 "pr_url": null,
2879 "refusal": null,
2880 "note": "legacy note",
2881 "created_unix": 1_000_u64
2882 });
2883 std::fs::write(
2884 root.join("cloud_00000000000000ff.json"),
2885 serde_json::to_vec_pretty(&legacy).unwrap(),
2886 )
2887 .unwrap();
2888 let store = CloudJobStore::from_path(root);
2889 let jobs = store.list().unwrap();
2890 assert_eq!(jobs.len(), 1);
2891 assert_eq!(jobs[0].sandbox_id.as_deref(), Some("sandbox_legacy"));
2892 assert!(jobs[0].base_branch.is_none());
2893 assert!(jobs[0].finished_unix.is_none());
2894 assert!(jobs[0].status == CloudJobStatus::Running);
2895 }
2896
2897 #[test]
2898 fn shell_quoting_and_sandbox_id_guards_hold() {
2899 // A hostile prompt stays one single-quoted argument: every embedded
2900 // quote becomes '\'' and no metacharacter can escape the quoting.
2901 let hostile = "fix it'; rm -rf /; echo '$(whoami)'";
2902 let joined = shell_quote_join(&[
2903 "codewhale".to_string(),
2904 "exec".to_string(),
2905 "--auto".to_string(),
2906 hostile.to_string(),
2907 ]);
2908 assert_eq!(
2909 joined,
2910 "'codewhale' 'exec' '--auto' 'fix it'\\''; rm -rf /; echo '\\''$(whoami)'\\'''"
2911 );
2912 // Behavioral pin: a real shell sees the prompt as ONE argument.
2913 let printed = std::process::Command::new("sh")
2914 .arg("-c")
2915 .arg(format!(
2916 "printf %s {}",
2917 shell_quote_join(&[hostile.to_string()])
2918 ))
2919 .output()
2920 .expect("sh is available in test environments");
2921 assert!(printed.status.success());
2922 assert_eq!(String::from_utf8_lossy(&printed.stdout), hostile);
2923 assert_eq!(shell_quote_join(&["a'b".to_string()]), "'a'\\''b'");
2924 assert!(valid_sandbox_id("sbx-123_abc"));
2925 for bad in ["", "../../evil", "sbx 1"] {
2926 assert!(!valid_sandbox_id(bad), "{bad:?} must be rejected");
2927 }
2928 }
2929
2930 #[test]
2931 fn net_guard_rejects_private_and_non_https_origins() {
2932 assert!(validate_outbound_origin("https://app.daytona.io/api").is_ok());
2933 assert!(validate_outbound_origin("https://gitee.com/api/v5").is_ok());
2934 assert!(validate_outbound_origin("ftp://example.com").is_err());
2935 assert!(validate_outbound_origin("https://user:pw@example.com/x").is_err());
2936 for blocked in [
2937 "http://example.com",
2938 "https://10.1.2.3/api",
2939 "https://192.168.1.10/api",
2940 "https://172.16.0.1/api",
2941 "https://169.254.169.254/latest/meta-data",
2942 "https://100.64.0.1/api",
2943 "https://0.0.0.0/api",
2944 "https://[fc00::1]/api",
2945 "https://[fe80::1]/api",
2946 "https://[::ffff:10.0.0.1]/api",
2947 "https://[::ffff:169.254.169.254]/latest/meta-data",
2948 "https://router.internal/api",
2949 "https://printer.local/api",
2950 ] {
2951 assert!(
2952 validate_outbound_origin(blocked).is_err(),
2953 "expected {blocked} to be rejected"
2954 );
2955 }
2956 // Loopback is the debug-only escape hatch for local smoke tests;
2957 // release builds reject it (pinned by the cfg! branch above).
2958 if cfg!(debug_assertions) {
2959 assert!(validate_outbound_origin("http://127.0.0.1:3986/api").is_ok());
2960 assert!(validate_outbound_origin("https://localhost/api").is_ok());
2961 }
2962 }
2963
2964 #[test]
2965 fn status_card_surfaces_runner_receipts_without_provider_branding() {
2966 let rows = remotes(&[("github", "https://github.com/org/repo.git")]);
2967 let jobs = vec![CloudJob {
2968 id: "cloud_00000000000000aa".to_string(),
2969 kind: "cloud".to_string(),
2970 status: CloudJobStatus::Done,
2971 prompt: "fix the flake".to_string(),
2972 forge: Forge::Github,
2973 remote_name: "github".to_string(),
2974 remote_url: "https://github.com/org/repo.git".to_string(),
2975 branch: "codewhale/cloud-1".to_string(),
2976 confirmed: true,
2977 sandbox_id: Some("sandbox_receipt_1".to_string()),
2978 pr_url: Some("https://github.com/org/repo/pull/7".to_string()),
2979 refusal: None,
2980 note: "done".to_string(),
2981 created_unix: 1_000,
2982 base_branch: Some("main".to_string()),
2983 head_sha: Some("abc1234def".to_string()),
2984 agent_summary: Some("Fixed the flake".to_string()),
2985 finished_unix: Some(1_960),
2986 sandbox_pending: false,
2987 }];
2988 let card = format_status(
2989 &rows,
2990 &CredentialState::Present {
2991 source: CredentialSource::Env,
2992 },
2993 &jobs,
2994 );
2995 assert!(card.contains("cloud_00000000000000aa"));
2996 assert!(card.contains("done"));
2997 assert!(card.contains("https://github.com/org/repo/pull/7"));
2998 assert!(card.contains("16m"));
2999 assert!(!card.contains("Daytona"));
3000 assert!(!card.contains("daytona"));
3001 let detail = format_job(&jobs[0]);
3002 assert!(detail.contains("Sandbox: sandbox_receipt_1"));
3003 assert!(detail.contains("PR: https://github.com/org/repo/pull/7"));
3004 assert!(detail.contains("Runtime: 16m"));
3005 assert!(detail.contains("Agent: Fixed the flake"));
3006 for banned in ["Daytona", "daytona"] {
3007 assert!(
3008 !detail.contains(banned),
3009 "the job card must not brand the sandbox operator: {banned}"
3010 );
3011 }
3012 }
3013
3014 #[test]
3015 fn parse_git_remote_listing_prefers_first_url_per_name() {
3016 let parsed = parse_remote_listing(
3017 "github\thttps://github.com/codewhale-hq/CodeWhale.git (fetch)\n\
3018 github\thttps://github.com/codewhale-hq/CodeWhale.git (push)\n\
3019 origin\thttps://cnb.cool/codewhale.net/codewhale.git (fetch)\n",
3020 );
3021 assert_eq!(parsed.len(), 2);
3022 assert_eq!(parsed[0].name, "github");
3023 assert_eq!(parsed[1].name, "origin");
3024 assert_eq!(
3025 classify_remote(&parsed[0].name, &parsed[0].url),
3026 Some(Forge::Github)
3027 );
3028 assert_eq!(
3029 classify_remote(&parsed[1].name, &parsed[1].url),
3030 Some(Forge::Cnb)
3031 );
3032 }
3033
3034 #[test]
3035 fn missing_credentials_copy_never_embeds_a_secret() {
3036 let message = missing_credentials_message();
3037 assert!(message.contains("cloud dispatch fails closed"));
3038 assert!(!message.contains("DAYTONA"));
3039 assert!(!message.contains("sk-"));
3040 assert!(!message.contains("Bearer"));
3041 }
3042
3043 /// Launcher fixture for sweep/reconcile tests: records teardowns and
3044 /// lists a configurable sandbox set; never touches a network.
3045 struct SweepLauncher {
3046 torn_down: std::sync::Mutex<Vec<String>>,
3047 listed: Vec<LabeledSandbox>,
3048 }
3049
3050 impl SweepLauncher {
3051 fn new(listed: Vec<LabeledSandbox>) -> Self {
3052 Self {
3053 torn_down: std::sync::Mutex::new(Vec::new()),
3054 listed,
3055 }
3056 }
3057
3058 fn torn_down(&self) -> Vec<String> {
3059 self.torn_down
3060 .lock()
3061 .map(|ids| ids.clone())
3062 .unwrap_or_default()
3063 }
3064 }
3065
3066 /// Records the persisted job status at the moment teardown runs, so
3067 /// tests can prove the terminal record was saved first.
3068 struct ObservingSweepLauncher {
3069 store: CloudJobStore,
3070 id: String,
3071 status_at_teardown: std::sync::Mutex<Option<CloudJobStatus>>,
3072 }
3073
3074 impl DaytonaLauncher for ObservingSweepLauncher {
3075 fn create_sandbox(&self, _job: &CloudJob) -> Result<SandboxReceipt> {
3076 bail!("no sandbox create in this fixture")
3077 }
3078 fn teardown(&self, _receipt: &SandboxReceipt) -> Result<()> {
3079 let status = self.store.load(&self.id).ok().map(|job| job.status);
3080 if let Ok(mut slot) = self.status_at_teardown.lock() {
3081 *slot = status;
3082 }
3083 Ok(())
3084 }
3085 fn list_job_sandboxes(&self) -> Result<Vec<LabeledSandbox>> {
3086 Ok(Vec::new())
3087 }
3088 }
3089
3090 impl DaytonaLauncher for SweepLauncher {
3091 fn create_sandbox(&self, _job: &CloudJob) -> Result<SandboxReceipt> {
3092 bail!("no sandbox create in this fixture")
3093 }
3094 fn teardown(&self, receipt: &SandboxReceipt) -> Result<()> {
3095 if let Ok(mut ids) = self.torn_down.lock() {
3096 ids.push(receipt.sandbox_id.clone());
3097 }
3098 Ok(())
3099 }
3100 fn list_job_sandboxes(&self) -> Result<Vec<LabeledSandbox>> {
3101 Ok(self.listed.clone())
3102 }
3103 }
3104
3105 fn stored_job(status: CloudJobStatus, created_unix: u64) -> CloudJob {
3106 CloudJob {
3107 id: "cloud_00000000000000e1".to_string(),
3108 kind: "cloud".to_string(),
3109 status,
3110 prompt: "fix the flake".to_string(),
3111 forge: Forge::Github,
3112 remote_name: "github".to_string(),
3113 remote_url: "https://github.com/org/repo.git".to_string(),
3114 branch: "codewhale/cloud-1".to_string(),
3115 confirmed: true,
3116 sandbox_id: None,
3117 pr_url: None,
3118 refusal: None,
3119 note: "n".to_string(),
3120 created_unix,
3121 base_branch: None,
3122 head_sha: None,
3123 agent_summary: None,
3124 finished_unix: None,
3125 sandbox_pending: false,
3126 }
3127 }
3128
3129 #[test]
3130 fn proposal_and_job_cards_never_carry_a_provider_brand() {
3131 // The proposal note is user copy (it rides `/dispatch show` and the
3132 // CLI card from the moment a job is proposed) — it names the
3133 // Codewhale cloud agent, never the sandbox operator.
3134 let plan = DispatchPlan {
3135 prompt: "offload me".to_string(),
3136 remote: SelectedRemote {
3137 forge: Forge::Github,
3138 name: "github".to_string(),
3139 url: "https://github.com/org/repo.git".to_string(),
3140 },
3141 branch: "codewhale/cloud-x".to_string(),
3142 };
3143 let note = proposal_note(&plan);
3144 assert!(note.contains("cloud-agent"), "{note}");
3145 let mut job = stored_job(CloudJobStatus::Proposed, 1);
3146 job.note = note.clone();
3147 let card = format_job(&job);
3148 for banned in ["Daytona", "daytona"] {
3149 assert!(
3150 !card.contains(banned),
3151 "the job card must not brand the sandbox operator: {banned}"
3152 );
3153 }
3154 assert!(card.contains("cloud-agent"));
3155 // And the list view (format_job_list) over the same record.
3156 let listing = format_job_list(&[job]);
3157 for banned in ["Daytona", "daytona"] {
3158 assert!(
3159 !listing.contains(banned),
3160 "listing must not brand: {banned}"
3161 );
3162 }
3163 }
3164
3165 #[test]
3166 fn harness_client_budget_covers_the_declared_turn_budget() {
3167 let hour = HarnessCommand {
3168 argv: vec!["codewhale".to_string()],
3169 cwd: SANDBOX_WORKSPACE.to_string(),
3170 timeout_secs: 3_600,
3171 };
3172 let budget = LiveDaytonaLauncher::harness_client_budget_secs(&hour);
3173 assert!(
3174 budget >= u64::from(hour.timeout_secs),
3175 "the client budget must cover the declared budget ({budget} < {})",
3176 hour.timeout_secs
3177 );
3178 assert!(
3179 budget > LiveDaytonaLauncher::CONTROL_PLANE_TIMEOUT_SECS,
3180 "an hour-long dispatched turn must not ride the 120s control-plane client"
3181 );
3182 // The budget scales with the declared timeout, not a fixed cap.
3183 let double = HarnessCommand {
3184 timeout_secs: 7_200,
3185 ..hour.clone()
3186 };
3187 assert!(LiveDaytonaLauncher::harness_client_budget_secs(&double) >= 7_200);
3188 // Short helper commands (collect_patch's git probes) keep a sane
3189 // bounded budget too.
3190 let probe = HarnessCommand {
3191 timeout_secs: 30,
3192 ..hour
3193 };
3194 assert!(LiveDaytonaLauncher::harness_client_budget_secs(&probe) >= 30);
3195 }
3196
3197 #[test]
3198 fn save_unless_canceled_refuses_to_resurrect_a_canceled_record() {
3199 let temp = tempfile::tempdir().unwrap();
3200 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3201 let mut job = stored_job(CloudJobStatus::Running, 10_000_000);
3202 store.save(&job).unwrap();
3203 // A phase save while the record is still active goes through.
3204 job.note = "phase save".to_string();
3205 assert!(store.save_unless_canceled(&job).unwrap());
3206 assert_eq!(store.load(&job.id).unwrap().note, "phase save");
3207 // The user cancels; the runner's next phase save must be refused and
3208 // leave the cancellation exactly as written.
3209 let mut canceled = store.load(&job.id).unwrap();
3210 canceled.status = CloudJobStatus::Canceled;
3211 canceled.note = "Canceled locally".to_string();
3212 canceled.finished_unix = Some(10_000_060);
3213 store.save(&canceled).unwrap();
3214 let mut stale_runner_copy = job.clone();
3215 stale_runner_copy.status = CloudJobStatus::OpeningPr;
3216 stale_runner_copy.note = "the runner's read-modify-write".to_string();
3217 assert!(!store.save_unless_canceled(&stale_runner_copy).unwrap());
3218 let persisted = store.load(&job.id).unwrap();
3219 assert_eq!(persisted.status, CloudJobStatus::Canceled);
3220 assert_eq!(persisted.note, "Canceled locally");
3221 assert_eq!(persisted.finished_unix, Some(10_000_060));
3222 }
3223
3224 #[test]
3225 fn a_phase_save_preserves_an_unreadable_job_record() {
3226 let temp = tempfile::tempdir().unwrap();
3227 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3228 let job = stored_job(CloudJobStatus::Running, 10_000_000);
3229 // A genuinely absent record can be created by the first phase save.
3230 assert!(store.save_unless_canceled(&job).unwrap());
3231 let path = store.job_path(&job.id).unwrap();
3232 let damaged = b"{\"status\":\"canceled\", interrupted write";
3233 fs::write(&path, damaged).unwrap();
3234 let error = store.save_unless_canceled(&job).unwrap_err();
3235 assert!(format!("{error:#}").contains("cannot verify the persisted job status"));
3236 assert_eq!(fs::read(&path).unwrap(), damaged);
3237
3238 // An I/O failure is uncertainty too, rather than permission to replace it.
3239 fs::remove_file(&path).unwrap();
3240 fs::create_dir(&path).unwrap();
3241 assert!(store.save_unless_canceled(&job).is_err());
3242 assert!(path.is_dir());
3243 }
3244
3245 #[test]
3246 fn a_phase_save_waits_for_a_cancel_holding_the_job_lock() {
3247 let temp = tempfile::tempdir().unwrap();
3248 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3249 let job = stored_job(CloudJobStatus::Running, 10_000_000);
3250 store.save(&job).unwrap();
3251 let mut runner_copy = job.clone();
3252 runner_copy.status = CloudJobStatus::OpeningPr;
3253
3254 // A cancel in progress holds the lock between its load and its write;
3255 // the runner's check-and-save must not slip into that window.
3256 let outcome = store
3257 .with_job_lock(&job.id, || {
3258 let runner_store = store.clone();
3259 let runner = std::thread::spawn(move || {
3260 runner_store.save_unless_canceled(&runner_copy).unwrap()
3261 });
3262 std::thread::sleep(std::time::Duration::from_millis(200));
3263 assert!(
3264 !runner.is_finished(),
3265 "the phase save ran inside the cancel's read-modify-write"
3266 );
3267 let mut canceled = store.load(&job.id)?;
3268 canceled.status = CloudJobStatus::Canceled;
3269 store.write_record(&canceled)?;
3270 Ok(runner)
3271 })
3272 .unwrap()
3273 .join()
3274 .unwrap();
3275 assert!(
3276 !outcome,
3277 "the phase save must see the cancel and stand down"
3278 );
3279 assert_eq!(
3280 store.load(&job.id).unwrap().status,
3281 CloudJobStatus::Canceled
3282 );
3283 assert!(
3284 !temp
3285 .path()
3286 .join("jobs")
3287 .join(format!("{}.json.tmp", job.id))
3288 .exists(),
3289 "records are written through unique temporaries"
3290 );
3291 }
3292
3293 #[test]
3294 fn a_confirm_waits_for_a_cancel_holding_the_job_lock() {
3295 let temp = tempfile::tempdir().unwrap();
3296 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3297 let mut job = stored_job(CloudJobStatus::Proposed, 10_000_000);
3298 job.confirmed = false;
3299 store.save(&job).unwrap();
3300
3301 // A confirm that loaded `proposed` before the cancel wrote must not
3302 // then overwrite the cancel with `launching` (and start a runner).
3303 let confirm = store
3304 .with_job_lock(&job.id, || {
3305 let confirm_store = store.clone();
3306 let id = job.id.clone();
3307 let confirm = std::thread::spawn(move || {
3308 confirm_job(
3309 &confirm_store,
3310 &id,
3311 &CredentialState::Present {
3312 source: CredentialSource::Env,
3313 },
3314 &MachineTokenState::Present,
3315 )
3316 });
3317 std::thread::sleep(std::time::Duration::from_millis(200));
3318 let mut canceled = store.load(&job.id)?;
3319 canceled.status = CloudJobStatus::Canceled;
3320 store.write_record(&canceled)?;
3321 Ok(confirm)
3322 })
3323 .unwrap()
3324 .join()
3325 .unwrap();
3326 assert!(confirm.is_err(), "a canceled proposal cannot be confirmed");
3327 assert_eq!(
3328 store.load(&job.id).unwrap().status,
3329 CloudJobStatus::Canceled
3330 );
3331 }
3332
3333 #[test]
3334 fn cancel_leaves_a_done_job_and_its_pr_alone() {
3335 let temp = tempfile::tempdir().unwrap();
3336 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3337 let mut done = stored_job(CloudJobStatus::Done, 10_000_000);
3338 done.pr_url = Some("https://github.com/org/repo/pull/7".to_string());
3339 store.save(&done).unwrap();
3340 let after = cancel_job(&store, &done.id, &NoopLauncher).unwrap();
3341 assert_eq!(after.status, CloudJobStatus::Done);
3342 assert_eq!(after.pr_url, done.pr_url);
3343 assert_eq!(store.load(&done.id).unwrap().status, CloudJobStatus::Done);
3344 }
3345
3346 #[test]
3347 fn sweep_fails_stale_active_jobs_and_tears_down_their_sandboxes() {
3348 let temp = tempfile::tempdir().unwrap();
3349 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3350 let now = 10_000_000_u64;
3351 // Stale: active well past the harness budget plus slack.
3352 let mut stale = stored_job(CloudJobStatus::Running, now - STALE_ACTIVE_JOB_SECS - 60);
3353 stale.sandbox_id = Some("sandbox_stale".to_string());
3354 store.save(&stale).unwrap();
3355 // Fresh: active but young; must be left exactly as it is.
3356 let mut fresh = stored_job(CloudJobStatus::Launching, now - 60);
3357 fresh.id = "cloud_00000000000000e2".to_string();
3358 fresh.sandbox_id = Some("sandbox_fresh".to_string());
3359 store.save(&fresh).unwrap();
3360 // Terminal: old but already done; never touched.
3361 let mut done = stored_job(CloudJobStatus::Done, now - STALE_ACTIVE_JOB_SECS * 2);
3362 done.id = "cloud_00000000000000e3".to_string();
3363 store.save(&done).unwrap();
3364
3365 let launcher = SweepLauncher::new(Vec::new());
3366 let swept = sweep_stale_jobs(&store, &launcher, now);
3367 assert_eq!(swept.len(), 1, "only the stale active job is swept");
3368 assert_eq!(swept[0].id, "cloud_00000000000000e1");
3369 let record = store.load("cloud_00000000000000e1").unwrap();
3370 assert_eq!(record.status, CloudJobStatus::Failed);
3371 assert_eq!(record.finished_unix, Some(now));
3372 assert!(record.note.contains("startup sweep"));
3373 assert!(record.note.contains("teardown was attempted"));
3374 assert_eq!(launcher.torn_down(), vec!["sandbox_stale".to_string()]);
3375 // The untouched records keep their state.
3376 assert_eq!(
3377 store.load("cloud_00000000000000e2").unwrap().status,
3378 CloudJobStatus::Launching
3379 );
3380 assert_eq!(
3381 store.load("cloud_00000000000000e3").unwrap().status,
3382 CloudJobStatus::Done
3383 );
3384 }
3385
3386 #[test]
3387 fn sweep_persists_the_terminal_record_before_teardown() {
3388 let temp = tempfile::tempdir().unwrap();
3389 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3390 let now = 10_000_000_u64;
3391 let mut stale = stored_job(CloudJobStatus::Running, now - STALE_ACTIVE_JOB_SECS - 60);
3392 stale.sandbox_id = Some("sandbox_stale".to_string());
3393 store.save(&stale).unwrap();
3394 let launcher = ObservingSweepLauncher {
3395 store: store.clone(),
3396 id: stale.id.clone(),
3397 status_at_teardown: std::sync::Mutex::new(None),
3398 };
3399 let swept = sweep_stale_jobs(&store, &launcher, now);
3400 assert_eq!(swept.len(), 1);
3401 let seen = *launcher.status_at_teardown.lock().unwrap();
3402 assert_eq!(
3403 seen,
3404 Some(CloudJobStatus::Failed),
3405 "teardown must observe an already-terminal record"
3406 );
3407 }
3408
3409 #[test]
3410 fn format_job_and_status_redact_remote_userinfo() {
3411 let mut job = stored_job(CloudJobStatus::Proposed, 10);
3412 job.remote_url = "https://user:token@github.com/org/repo.git".to_string();
3413 let card = format_job(&job);
3414 assert!(!card.contains("token"), "{card}");
3415 assert!(!card.contains("user:"), "{card}");
3416 assert!(card.contains("github.com/org/repo.git"), "{card}");
3417 let status = format_status(
3418 &remotes(&[("github", "https://user:token@github.com/org/repo.git")]),
3419 &CredentialState::Missing,
3420 &[],
3421 );
3422 assert!(!status.contains("token"), "{status}");
3423 }
3424
3425 #[test]
3426 fn reconcile_deletes_sandboxes_for_terminal_or_absent_jobs_and_keeps_active() {
3427 let temp = tempfile::tempdir().unwrap();
3428 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3429 let now = 10_000_000_u64;
3430 // Terminal job: its sandbox must go.
3431 let mut terminal = stored_job(CloudJobStatus::Canceled, now - 600);
3432 terminal.id = "cloud_00000000000000f1".to_string();
3433 terminal.sandbox_id = Some("sandbox_terminal".to_string());
3434 store.save(&terminal).unwrap();
3435 // Active job: its sandbox must stay.
3436 let mut active = stored_job(CloudJobStatus::Running, now - 60);
3437 active.id = "cloud_00000000000000f2".to_string();
3438 store.save(&active).unwrap();
3439
3440 let launcher = SweepLauncher::new(vec![
3441 LabeledSandbox {
3442 sandbox_id: "sandbox_terminal".to_string(),
3443 job_id: Some("cloud_00000000000000f1".to_string()),
3444 },
3445 LabeledSandbox {
3446 sandbox_id: "sandbox_active".to_string(),
3447 job_id: Some("cloud_00000000000000f2".to_string()),
3448 },
3449 // Labeled for a job that no longer exists in the store.
3450 LabeledSandbox {
3451 sandbox_id: "sandbox_ghost".to_string(),
3452 job_id: Some("cloud_0000000000000bad".to_string()),
3453 },
3454 // No usable job label at all.
3455 LabeledSandbox {
3456 sandbox_id: "sandbox_unlabeled".to_string(),
3457 job_id: None,
3458 },
3459 ]);
3460 let report = reconcile_sandboxes(&store, &launcher).unwrap();
3461 assert_eq!(report.deleted.len(), 3);
3462 assert!(report.deleted.contains(&"sandbox_terminal".to_string()));
3463 assert!(report.deleted.contains(&"sandbox_ghost".to_string()));
3464 assert!(report.deleted.contains(&"sandbox_unlabeled".to_string()));
3465 assert_eq!(report.live, 1);
3466 assert!(!launcher.torn_down().contains(&"sandbox_active".to_string()));
3467 }
3468
3469 #[test]
3470 fn cancel_deletes_an_unrecorded_sandbox_by_label() {
3471 let temp = tempfile::tempdir().unwrap();
3472 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3473 let mut pending = stored_job(CloudJobStatus::Running, 10_000_000);
3474 // The create POST landed but its response never arrived: the intent
3475 // is persisted and the id is unknowable — only the label can find it.
3476 pending.sandbox_pending = true;
3477 store.save(&pending).unwrap();
3478 let launcher = SweepLauncher::new(vec![
3479 LabeledSandbox {
3480 sandbox_id: "sandbox_lost".to_string(),
3481 job_id: Some(pending.id.clone()),
3482 },
3483 LabeledSandbox {
3484 sandbox_id: "sandbox_other".to_string(),
3485 job_id: Some("cloud_00000000000000f9".to_string()),
3486 },
3487 ]);
3488 let canceled = cancel_job(&store, &pending.id, &launcher).unwrap();
3489 assert_eq!(canceled.status, CloudJobStatus::Canceled);
3490 assert!(canceled.note.contains("unrecorded"));
3491 assert!(canceled.note.contains("torn down"));
3492 assert_eq!(
3493 launcher.torn_down(),
3494 vec!["sandbox_lost".to_string()],
3495 "cancel deletes only this job's labeled sandbox"
3496 );
3497 }
3498
3499 #[test]
3500 fn quit_warning_names_live_jobs_and_stays_quiet_otherwise() {
3501 let temp = tempfile::tempdir().unwrap();
3502 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3503 assert_eq!(live_job_quit_warning(&store), None);
3504 let mut live = stored_job(CloudJobStatus::OpeningPr, 10_000_000);
3505 live.id = "cloud_00000000000000aa".to_string();
3506 store.save(&live).unwrap();
3507 let warning = live_job_quit_warning(&store).expect("live job must warn");
3508 assert!(warning.contains("cloud_00000000000000aa"));
3509 assert!(warning.contains("/dispatch cancel"));
3510 assert!(!warning.contains("Daytona"));
3511 let mut done = stored_job(CloudJobStatus::Done, 10_000_000);
3512 done.id = "cloud_00000000000000ab".to_string();
3513 store.save(&done).unwrap();
3514 let mut dead = stored_job(CloudJobStatus::Failed, 10_000_000);
3515 dead.id = "cloud_00000000000000ac".to_string();
3516 store.save(&dead).unwrap();
3517 // Still exactly one live job after adding terminal siblings.
3518 let warning = live_job_quit_warning(&store).expect("live job must warn");
3519 assert!(warning.contains("cloud_00000000000000aa"));
3520 assert!(!warning.contains("cloud_00000000000000ab"));
3521 assert!(!warning.contains("cloud_00000000000000ac"));
3522 }
3523
3524 #[test]
3525 fn sandbox_create_is_not_computer_entitlement() {
3526 use crate::computer_meter::{
3527 ComputerAdmissionRequest, MeterBasis, bind_computer_admission,
3528 };
3529
3530 let temp = tempfile::tempdir().unwrap();
3531 let store = CloudJobStore::from_path(temp.path().join("jobs"));
3532 let plan = plan_dispatch(
3533 &remotes(&[("github", "https://github.com/org/repo.git")]),
3534 "meter honesty",
3535 Some(Forge::Github),
3536 Some("codewhale/cloud-meter"),
3537 )
3538 .unwrap();
3539 let outcome = execute_dispatch(
3540 &store,
3541 plan,
3542 true,
3543 &CredentialState::Present {
3544 source: CredentialSource::Keyring,
3545 },
3546 &MachineTokenState::Present,
3547 )
3548 .unwrap();
3549 let DispatchOutcome::Accepted(mut job) = outcome else {
3550 panic!("expected accept");
3551 };
3552 job.sandbox_id = Some("sbox_std8_a".to_string());
3553 let admission = bind_computer_admission(ComputerAdmissionRequest {
3554 admission_id: "adm_dispatch".to_string(),
3555 account_id: "acct_demo".to_string(),
3556 computer_id: "cmp_dispatch".to_string(),
3557 run_id: job.id.clone(),
3558 provider: "daytona".to_string(),
3559 profile_id: "standard-8".to_string(),
3560 funding_authority: "coding_membership_included".to_string(),
3561 quote_id: "quote_dispatch".to_string(),
3562 expires_at: "2026-08-31T18:00:00.000Z".to_string(),
3563 meter_revision: String::new(),
3564 catalog_revision: String::new(),
3565 })
3566 .unwrap();
3567 let idle = ProviderObservation {
3568 provider: "daytona".to_string(),
3569 provider_sandbox_id: "sbox_std8_a".to_string(),
3570 provider_event_ref: "daytona:sbox_std8_a:idle".to_string(),
3571 state: "running".to_string(),
3572 idle: true,
3573 provider_accepted: true,
3574 meter_basis: MeterBasis::WallClock,
3575 cpu: 2,
3576 memory_gib: 8,
3577 disk_gib: 8,
3578 started_at: "2026-08-31T12:00:00.000Z".to_string(),
3579 ended_at: "2026-08-31T13:00:00.000Z".to_string(),
3580 };
3581 assert_eq!(
3582 meter_cloud_job(&job, &admission, idle).unwrap_err().code(),
3583 "computer_meter_wall_clock_idle"
3584 );
3585 let accepted = ProviderObservation {
3586 provider: "daytona".to_string(),
3587 provider_sandbox_id: "sbox_std8_a".to_string(),
3588 provider_event_ref: "daytona:sbox_std8_a:running".to_string(),
3589 state: "running".to_string(),
3590 idle: false,
3591 provider_accepted: true,
3592 meter_basis: MeterBasis::ProviderAcceptedActive,
3593 cpu: 2,
3594 memory_gib: 8,
3595 disk_gib: 8,
3596 started_at: "2026-08-31T12:00:00.000Z".to_string(),
3597 ended_at: "2026-08-31T12:10:00.000Z".to_string(),
3598 };
3599 let receipt = meter_cloud_job(&job, &admission, accepted).unwrap();
3600 assert_eq!(receipt.accepted_seconds, 600);
3601 assert_eq!(receipt.standard_equivalent_seconds, 600);
3602 assert_eq!(receipt.admission_id, admission.admission_id);
3603 }
3604 }
3605
3605 lines RUST