| 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 |