返回 CodeWhale
dispatch.rs
根目录 / crates / cli / src / dispatch.rs
1 //! First-class Daytona cloud-agent offload: `codewhale dispatch`.
2
3 use std::io::{self, Write};
4 use std::path::PathBuf;
5
6 use anyhow::{Result, bail};
7 use clap::{Args, ValueEnum};
8 use codewhale_tui::cloud_dispatch::{
9 CloudJobStore, DispatchOutcome, Forge, LiveDaytonaLauncher, cancel_job, confirm_job,
10 discover_credentials, discover_machine_token, discover_remotes, execute_dispatch, format_job,
11 format_job_list, format_status, plan_dispatch,
12 };
13 use codewhale_tui::dispatch_runner::spawn_confirmed_runner;
14
15 #[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
16 enum ForgeArg {
17 Github,
18 Cnb,
19 Gitee,
20 }
21
22 impl From<ForgeArg> for Forge {
23 fn from(value: ForgeArg) -> Self {
24 match value {
25 ForgeArg::Github => Forge::Github,
26 ForgeArg::Cnb => Forge::Cnb,
27 ForgeArg::Gitee => Forge::Gitee,
28 }
29 }
30 }
31
32 #[derive(Debug, Args)]
33 pub(crate) struct DispatchArgs {
34 /// Task for the remote agent. Required unless listing or inspecting a job.
35 #[arg(value_name = "PROMPT")]
36 prompt: Vec<String>,
37 /// Forge that should receive the branch and PR: github, cnb, or gitee.
38 #[arg(long, value_enum)]
39 remote: Option<ForgeArg>,
40 /// Branch the remote agent will raise (default: codewhale/cloud-<unix>).
41 #[arg(long)]
42 branch: Option<String>,
43 /// Required to create Codewhale cloud-agent spend or push. Without this, only a proposal is written.
44 #[arg(long)]
45 confirm: bool,
46 /// Show remotes and whether Codewhale cloud-agent credentials are present (never prints secrets).
47 #[arg(long)]
48 status: bool,
49 /// List first-class cloud jobs (same kind shown by `/jobs`).
50 #[arg(long)]
51 list: bool,
52 /// Inspect one cloud job.
53 #[arg(long, value_name = "ID")]
54 show: Option<String>,
55 /// Cancel one cloud job.
56 #[arg(long, value_name = "ID")]
57 cancel: Option<String>,
58 /// Workspace whose git remotes are classified (default: current directory).
59 #[arg(long)]
60 cwd: Option<PathBuf>,
61 }
62
63 pub(crate) fn run(args: DispatchArgs) -> Result<()> {
64 let mut out = io::stdout().lock();
65 run_with(args, &mut out)
66 }
67
68 fn run_with<W: Write>(args: DispatchArgs, out: &mut W) -> Result<()> {
69 if [
70 args.status,
71 args.list,
72 args.show.is_some(),
73 args.cancel.is_some(),
74 !args.prompt.is_empty(),
75 ]
76 .iter()
77 .filter(|flag| **flag)
78 .count()
79 > 1
80 {
81 bail!("Use one of: a prompt, --status, --list, --show <id>, or --cancel <id>.");
82 }
83
84 let workspace = args
85 .cwd
86 .clone()
87 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
88 let remotes = discover_remotes(&workspace);
89 let credentials = discover_credentials();
90 let store = CloudJobStore::from_env()?;
91
92 if args.status
93 || (args.prompt.is_empty() && args.show.is_none() && args.cancel.is_none() && !args.list)
94 {
95 writeln!(
96 out,
97 "{}",
98 format_status(&remotes, &credentials, &recent_jobs(&store))
99 )?;
100 return Ok(());
101 }
102 if args.list {
103 writeln!(out, "{}", format_job_list(&store.list()?))?;
104 return Ok(());
105 }
106 if let Some(id) = args.show.as_deref() {
107 writeln!(out, "{}", format_job(&store.load(id)?))?;
108 return Ok(());
109 }
110 if let Some(id) = args.cancel.as_deref() {
111 writeln!(
112 out,
113 "{}",
114 format_job(&cancel_job(&store, id, &LiveDaytonaLauncher)?)
115 )?;
116 return Ok(());
117 }
118
119 let prompt = args.prompt.join(" ");
120 if prompt.starts_with("cloud_") && args.confirm && prompt.split_whitespace().count() == 1 {
121 let outcome = confirm_job(
122 &store,
123 prompt.trim(),
124 &credentials,
125 &discover_machine_token(),
126 )?;
127 let runner = spawn_accepted(&store, &outcome);
128 return write_then_join(out, &store, prompt.trim(), outcome, runner);
129 }
130
131 let plan = plan_dispatch(
132 &remotes,
133 &prompt,
134 args.remote.map(Forge::from),
135 args.branch.as_deref(),
136 )?;
137 let outcome = execute_dispatch(
138 &store,
139 plan,
140 args.confirm,
141 &credentials,
142 &discover_machine_token(),
143 )?;
144 let runner = spawn_accepted(&store, &outcome);
145 let job_id = outcome_job_id(&outcome).unwrap_or_default();
146 write_then_join(out, &store, &job_id, outcome, runner)
147 }
148
149 /// Print the outcome card, then stay attached to the runner — joining it even
150 /// when the print failed. A closed stdout (`| head`, a dead terminal) used to
151 /// return before the join, ending the process and the runner thread with it,
152 /// which orphaned the paid sandbox this wait exists to supervise. The print
153 /// error still wins the exit status.
154 fn write_then_join<W: Write>(
155 out: &mut W,
156 store: &CloudJobStore,
157 id: &str,
158 outcome: DispatchOutcome,
159 runner: Option<std::thread::JoinHandle<()>>,
160 ) -> Result<()> {
161 let written = write_outcome(out, outcome);
162 let joined = join_runner(out, store, id, runner);
163 written.and(joined)
164 }
165
166 /// The CLI stays attached to a confirmed run: the card prints immediately,
167 /// then the process waits for the runner so a paid sandbox is never
168 /// orphaned by an early exit. Ctrl-C exits the wait; the job stays recorded
169 /// and `--cancel` tears the sandbox down.
170 fn join_runner<W: Write>(
171 out: &mut W,
172 store: &CloudJobStore,
173 id: &str,
174 runner: Option<std::thread::JoinHandle<()>>,
175 ) -> Result<()> {
176 if let Some(runner) = runner {
177 runner
178 .join()
179 .map_err(|_| anyhow::anyhow!("the cloud agent runner panicked"))?;
180 if !id.is_empty()
181 && let Ok(job) = store.load(id)
182 {
183 writeln!(out, "{}", format_job(&job))?;
184 }
185 }
186 Ok(())
187 }
188
189 fn outcome_job_id(outcome: &DispatchOutcome) -> Option<String> {
190 match outcome {
191 DispatchOutcome::Proposal(job)
192 | DispatchOutcome::Refused(job)
193 | DispatchOutcome::Accepted(job) => Some(job.id.clone()),
194 }
195 }
196
197 /// Newest jobs for the status card's receipts section (best effort — an
198 /// unreadable store must not hide the card).
199 fn recent_jobs(store: &CloudJobStore) -> Vec<codewhale_tui::cloud_dispatch::CloudJob> {
200 store
201 .list()
202 .unwrap_or_default()
203 .into_iter()
204 .take(5)
205 .collect()
206 }
207
208 /// Start the background runner for a just-accepted confirm. The sandbox,
209 /// harness turn, branch push, PR open, and teardown all happen there.
210 fn spawn_accepted(
211 store: &CloudJobStore,
212 outcome: &DispatchOutcome,
213 ) -> Option<std::thread::JoinHandle<()>> {
214 match outcome {
215 DispatchOutcome::Accepted(job) => spawn_confirmed_runner(store.clone(), job.id.clone()),
216 _ => None,
217 }
218 }
219
220 fn write_outcome<W: Write>(out: &mut W, outcome: DispatchOutcome) -> Result<()> {
221 match outcome {
222 DispatchOutcome::Proposal(job) | DispatchOutcome::Accepted(job) => {
223 writeln!(out, "{}", format_job(&job))?;
224 Ok(())
225 }
226 DispatchOutcome::Refused(job) => {
227 writeln!(out, "{}", format_job(&job))?;
228 bail!("{}", job.note);
229 }
230 }
231 }
232
233 #[cfg(test)]
234 mod tests {
235 use super::*;
236 use crate::{Cli, Commands};
237 use clap::Parser;
238
239 fn args(argv: &[&str]) -> DispatchArgs {
240 let cli = Cli::try_parse_from(argv).unwrap();
241 let Some(Commands::Dispatch(args)) = cli.command else {
242 panic!("expected dispatch command");
243 };
244 args
245 }
246
247 #[test]
248 fn parses_the_obvious_dispatch_command() {
249 let parsed = args(&[
250 "codewhale",
251 "dispatch",
252 "fix",
253 "the",
254 "flake",
255 "--remote",
256 "github",
257 ]);
258 assert_eq!(parsed.prompt, ["fix", "the", "flake"]);
259 assert_eq!(parsed.remote, Some(ForgeArg::Github));
260 assert!(!parsed.confirm);
261 assert!(args(&["codewhale", "cloud-agent", "--status"]).status);
262 assert!(Cli::try_parse_from(["codewhale", "dispatch", "--remote", "gitlab"]).is_err());
263 }
264
265 #[test]
266 fn refused_confirmation_is_a_nonzero_error() {
267 use codewhale_tui::cloud_dispatch::{CloudJob, CloudJobStatus};
268 let job = CloudJob {
269 id: "cloud_00000000000000dd".to_string(),
270 kind: "cloud".to_string(),
271 status: CloudJobStatus::Refused,
272 prompt: "fix".to_string(),
273 forge: Forge::Github,
274 remote_name: "github".to_string(),
275 remote_url: "https://github.com/org/repo.git".to_string(),
276 branch: "codewhale/cloud-x".to_string(),
277 confirmed: true,
278 sandbox_id: None,
279 pr_url: None,
280 refusal: Some("no credentials".to_string()),
281 note: "Refused: cloud agents are not available.".to_string(),
282 created_unix: 1,
283 base_branch: None,
284 head_sha: None,
285 agent_summary: None,
286 finished_unix: Some(1),
287 sandbox_pending: false,
288 };
289 let error = write_outcome(&mut Vec::new(), DispatchOutcome::Refused(job)).unwrap_err();
290 assert!(
291 error.to_string().contains("Refused"),
292 "refused confirmations must not exit 0: {error}"
293 );
294 }
295
296 #[test]
297 fn status_is_fail_closed_and_never_prints_secrets() {
298 let temp = tempfile::tempdir().unwrap();
299 let mut output = Vec::new();
300 run_with(
301 DispatchArgs {
302 prompt: Vec::new(),
303 remote: None,
304 branch: None,
305 confirm: false,
306 status: true,
307 list: false,
308 show: None,
309 cancel: None,
310 cwd: Some(temp.path().to_path_buf()),
311 },
312 &mut output,
313 )
314 .unwrap();
315 let text = String::from_utf8(output).unwrap();
316 assert!(text.contains("Codewhale cloud dispatch"));
317 assert!(!text.contains("Daytona"));
318 assert!(!text.contains("sk-"));
319 assert!(!text.contains("Bearer"));
320 }
321
322 /// Audit R02-09: a stdout failure must not skip the runner join.
323 #[test]
324 fn a_failed_card_write_still_joins_the_runner() {
325 use std::sync::Arc;
326 use std::sync::atomic::{AtomicBool, Ordering};
327
328 struct ClosedPipe;
329 impl Write for ClosedPipe {
330 fn write(&mut self, _: &[u8]) -> io::Result<usize> {
331 Err(io::Error::from(io::ErrorKind::BrokenPipe))
332 }
333 fn flush(&mut self) -> io::Result<()> {
334 Err(io::Error::from(io::ErrorKind::BrokenPipe))
335 }
336 }
337
338 let finished = Arc::new(AtomicBool::new(false));
339 let runner = {
340 let finished = Arc::clone(&finished);
341 std::thread::spawn(move || {
342 std::thread::sleep(std::time::Duration::from_millis(100));
343 finished.store(true, Ordering::SeqCst);
344 })
345 };
346 let job: codewhale_tui::cloud_dispatch::CloudJob =
347 serde_json::from_value(serde_json::json!({
348 "id": "cloud_test", "kind": "cloud", "status": "running",
349 "prompt": "p", "forge": "github", "remote_name": "origin",
350 "remote_url": "https://example.invalid/r.git", "branch": "b",
351 "confirmed": true, "sandbox_id": null, "pr_url": null,
352 "refusal": null, "note": "", "created_unix": 0
353 }))
354 .expect("job fixture");
355 let dir = tempfile::tempdir().unwrap();
356 let store = CloudJobStore::from_path(dir.path().to_path_buf());
357
358 let result = write_then_join(
359 &mut ClosedPipe,
360 &store,
361 "cloud_test",
362 DispatchOutcome::Accepted(job),
363 Some(runner),
364 );
365
366 assert!(result.is_err(), "the print failure is still reported");
367 assert!(
368 finished.load(Ordering::SeqCst),
369 "returned before the runner finished"
370 );
371 }
372
373 #[test]
374 fn rendered_help_carries_no_provider_brand() {
375 use clap::CommandFactory;
376 let help = Cli::command()
377 .find_subcommand_mut("dispatch")
378 .expect("dispatch subcommand exists")
379 .render_help()
380 .to_string();
381 // The reworded cloud-agent copy must actually land in --help…
382 assert!(help.contains("cloud-agent"), "{help}");
383 assert!(help.contains("--confirm"), "{help}");
384 // …and no provider brand may leak into it.
385 for banned in ["Daytona", "daytona"] {
386 assert!(
387 !help.contains(banned),
388 "--help must not brand the operator: {banned}"
389 );
390 }
391 }
392 }
393
393 lines RUST