Skip to main content

scv_tools/delegate/
background.rs

1//! Background delegations: an `agent` call with `background: true` returns a
2//! job handle at once while the agent keeps working; `agent_status` and
3//! `agent_wait` observe the job, `agent_cancel` stops it, and the session is
4//! told when it finishes so the server can report it in a turn of its own.
5//!
6//! Each call that starts a job, or shows the model a job's result, leaves a
7//! [`JobChange`] under its call ID, which the server hands to the session's
8//! clients with the call's `tool.completed` ([`BackgroundJobs::take_changes`]).
9//!
10//! A finished job stays unreported until a report turn about it succeeds. A
11//! failed one makes it due again after a delay, a bounded number of times;
12//! after that, or when another turn must not be tried, the server reports it
13//! to the client directly ([`BackgroundJobs::report_failed`]).
14//!
15//! A call that continues a conversation takes its place in that
16//! conversation's lane as it arrives. When a turn or another call is
17//! ahead, the call's busy policy decides: it steers the running turn, waits
18//! in the foreground, fails, or becomes a queued job that runs once its
19//! place comes up.
20
21mod lane;
22
23use std::{
24    sync::{Arc, Mutex, Weak},
25    time::{Duration, Instant},
26};
27
28use async_trait::async_trait;
29use scv_core::{
30    ApprovalGate, ProgressSink, Tool, ToolApprovals, ToolContext, ToolError, ToolOutput, ToolRisk,
31    ToolSpec,
32};
33use scv_protocol::{JobChange, JobReport, JobStatus, describe_reports};
34use serde::Deserialize;
35use serde_json::{Map, Value, json};
36use tokio::sync::{mpsc, watch};
37use tokio_util::sync::CancellationToken;
38
39use self::lane::{Lanes, Place};
40use crate::{
41    BusyBehavior,
42    args::{Timeouts, bounded, parse_args, timeout_schema},
43    delegate::{
44        agent::{AGENT_TOOL, AgentTool, Offered},
45        conversation,
46        output::AgentReply,
47        request::AgentArgs,
48    },
49    sync::lock,
50};
51
52/// Finished jobs a session keeps for `agent_status` beyond the running ones.
53const MAX_FINISHED: usize = 16;
54/// Characters of a finished job's reply quoted in a server-started report turn.
55const REPORT_REPLY_CHARS: usize = 6000;
56/// Jobs one report turn covers; any more wait for the next.
57const REPORT_MAX_JOBS: usize = 4;
58/// How long after each failed report turn the next one is due: after the
59/// first failure, then after the second. Once they run out, the job is
60/// reported directly.
61const REPORT_RETRY_DELAYS: [Duration; 2] = [Duration::from_secs(30), Duration::from_secs(120)];
62/// How long `agent_cancel` waits for a stopped job to settle.
63const CANCEL_SETTLE: Duration = Duration::from_secs(10);
64/// Job changes kept for calls whose `tool.completed` has not taken them,
65/// such as a call whose turn was aborted; the oldest go first.
66const MAX_PENDING_CHANGES: usize = 64;
67/// Characters of a delegated prompt's first line kept as its job's task.
68const TASK_CHARS: usize = 80;
69
70/// One session's background jobs. Dropping the store (with the session's
71/// tools) cancels every job still running.
72pub struct BackgroundJobs {
73    limit: usize,
74    state: Mutex<JobsState>,
75    cancellation: CancellationToken,
76    /// Woken whenever a job finishes, so the session can report it.
77    finished: Option<mpsc::UnboundedSender<()>>,
78    /// Decides a running job's nested approval requests, since no turn is
79    /// left to carry them to a person. Without it they are denied.
80    approvals: Option<Arc<dyn ApprovalGate>>,
81    /// The `agent` calls continuing each conversation, in arrival order.
82    lanes: Arc<Lanes>,
83}
84
85impl std::fmt::Debug for BackgroundJobs {
86    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
87        formatter
88            .debug_struct("BackgroundJobs")
89            .field("limit", &self.limit)
90            .finish_non_exhaustive()
91    }
92}
93
94#[derive(Default)]
95struct JobsState {
96    next: u64,
97    jobs: Vec<Job>,
98    /// Changes by the ID of the call that made them, oldest first.
99    changes: Vec<(String, JobChange)>,
100}
101
102impl JobsState {
103    /// Keep `change` for the call `call_id`, which a client learns of with
104    /// that call's `tool.completed`. A call without an ID has no event.
105    fn record(&mut self, call_id: &str, change: JobChange) {
106        if call_id.is_empty() {
107            return;
108        }
109        if self.changes.len() >= MAX_PENDING_CHANGES {
110            self.changes.remove(0);
111        }
112        self.changes.push((call_id.to_owned(), change));
113    }
114}
115
116struct Job {
117    id: String,
118    /// The tool that started it, `agent`.
119    tool: String,
120    /// The agent that runs it, such as `codex`.
121    agent: String,
122    /// The first line of the delegated prompt, shortened.
123    task: String,
124    started: Instant,
125    progress: ProgressSink,
126    last_progress: Option<String>,
127    /// Stops this job alone.
128    cancel: CancellationToken,
129    /// `agent_cancel` stopped it.
130    cancelled: bool,
131    /// A busy conversation's queued prompt, which `agent.max_queued_turns`
132    /// bounds instead of `agent.max_background`.
133    queued: bool,
134    outcome: Option<Outcome>,
135    /// The model has seen the result (through `agent_wait`, `agent_status`,
136    /// or a report turn that completed), the user stopped its report turn,
137    /// the server reported it directly, or the model asked for the stop with
138    /// `agent_cancel`, so it needs no report turn.
139    reported: bool,
140    /// A report turn about it runs.
141    reporting: bool,
142    /// Report turns about it that failed.
143    report_attempts: u32,
144    /// After a failed report turn: when the next one is due.
145    report_due: Option<tokio::time::Instant>,
146    done: watch::Receiver<bool>,
147}
148
149struct Outcome {
150    output: ToolOutput,
151    elapsed: Duration,
152}
153
154/// What a failed report turn leaves its jobs to.
155#[derive(Debug, PartialEq, Eq)]
156pub enum ReportFailure {
157    /// Another report turn about them is due after this long.
158    Retry(Duration),
159    /// No more report turns: the server reports them directly. `attempts`
160    /// counts the report turns that failed.
161    GiveUp {
162        attempts: u32,
163        reports: Vec<JobReport>,
164    },
165    /// The model saw each of them through a call during the turn.
166    Settled,
167}
168
169impl Drop for BackgroundJobs {
170    fn drop(&mut self) {
171        self.cancellation.cancel();
172    }
173}
174
175impl BackgroundJobs {
176    /// At most `limit` jobs run at once. `finished` is woken as jobs finish.
177    pub fn new(limit: usize, finished: Option<mpsc::UnboundedSender<()>>) -> Self {
178        Self {
179            limit,
180            state: Mutex::default(),
181            cancellation: CancellationToken::new(),
182            finished,
183            approvals: None,
184            lanes: Arc::default(),
185        }
186    }
187
188    /// Decide running jobs' nested approval requests with `gate`, which must
189    /// never grant more than the session's foreground would.
190    #[must_use]
191    pub fn with_approvals(mut self, gate: Arc<dyn ApprovalGate>) -> Self {
192        self.approvals = Some(gate);
193        self
194    }
195
196    pub(crate) fn limit(&self) -> usize {
197        self.limit
198    }
199
200    fn state(&self) -> std::sync::MutexGuard<'_, JobsState> {
201        lock(&self.state)
202    }
203
204    /// Start `tool` with `arguments` in the background for the call in
205    /// `context` (whose workspace and call ID it takes; the job has its own
206    /// cancellation), on `agent`, and return its job's
207    /// `{"job","agent","status":"running","background":true}` description.
208    fn start(
209        self: &Arc<Self>,
210        tool: Arc<dyn Tool>,
211        name: &str,
212        agent: &str,
213        arguments: Value,
214        context: &ToolContext,
215        turn: Turn,
216    ) -> Result<Value, ToolError> {
217        let task = task_line(
218            arguments
219                .get("prompt")
220                .and_then(Value::as_str)
221                .unwrap_or_default(),
222        );
223        let queued = matches!(turn, Turn::Queued(_));
224        let (id, progress, done_tx, cancellation) = {
225            let mut state = self.state();
226            let running = state
227                .jobs
228                .iter()
229                .filter(|job| job.outcome.is_none() && !job.queued)
230                .count();
231            if !queued && running >= self.limit {
232                return Err(ToolError::limit(format!(
233                    "{running} background jobs are already running, the limit \
234                     (agent.max_background). Start this one after a job finishes, or \
235                     stop one with agent_cancel if the user no longer needs it."
236                )));
237            }
238            state.next += 1;
239            let id = format!("job-{}", state.next);
240            let progress = ProgressSink::buffered();
241            let (done_tx, done) = watch::channel(false);
242            let cancel = self.cancellation.child_token();
243            state.record(
244                &context.call_id,
245                JobChange {
246                    job: id.clone(),
247                    tool: name.to_owned(),
248                    agent: agent.to_owned(),
249                    status: JobStatus::Running,
250                    task: task.clone(),
251                },
252            );
253            state.jobs.push(Job {
254                id: id.clone(),
255                tool: name.to_owned(),
256                agent: agent.to_owned(),
257                task,
258                started: Instant::now(),
259                progress: progress.clone(),
260                last_progress: None,
261                cancel: cancel.clone(),
262                cancelled: false,
263                queued,
264                outcome: None,
265                reported: false,
266                reporting: false,
267                report_attempts: 0,
268                report_due: None,
269                done,
270            });
271            (id, progress, done_tx, cancel)
272        };
273        let jobs = Arc::downgrade(self);
274        let job = id.clone();
275        let approvals = self.approvals.clone();
276        let workspace = context.workspace.clone();
277        tokio::spawn(async move {
278            let started = Instant::now();
279            let mut context = ToolContext::new(workspace, cancellation);
280            // No turn is left to ask a person, so nested approval requests
281            // get the session's unattended answer, or are denied.
282            if let Some(gate) = &approvals {
283                context.approvals = ToolApprovals::new(Arc::clone(gate), job.clone());
284            }
285            context.progress = progress;
286            let place = turn.place();
287            let output = match &place {
288                Some(place) if queued => {
289                    context.progress.report(&format!(
290                        "queued: waits for conversation {}'s running turn",
291                        place.handle()
292                    ));
293                    match place.wait_first(&context.cancellation).await {
294                        Ok(()) => {
295                            context.progress.take();
296                            context.progress.report(&format!(
297                                "started its turn in conversation {}",
298                                place.handle()
299                            ));
300                            tool.execute(arguments, context).await
301                        }
302                        Err(error) => Err(error),
303                    }
304                }
305                _ => tool.execute(arguments, context).await,
306            }
307            .unwrap_or_else(ToolOutput::from);
308            // Free the conversation before the result is visible, so a call
309            // continuing it after `agent_wait` finds it idle.
310            drop(place);
311            finish(&jobs, &job, output, started.elapsed());
312            let _ = done_tx.send(true);
313        });
314        let mut started = json!({
315            "job": id,
316            "agent": agent,
317            "status": "running",
318            "background": true,
319            "note": "The agent is working in the background. SCV reports the result in a new \
320                     turn when it finishes. agent_status shows its progress and agent_cancel \
321                     stops it."
322        });
323        if queued {
324            started["queued"] = true.into();
325            started["note"] = "The conversation is busy, so this prompt is queued: it runs in \
326                the background after the turns ahead of it, and SCV reports the result in a new \
327                turn when it finishes. agent_status shows where it stands and agent_cancel \
328                withdraws it."
329                .into();
330        }
331        Ok(started)
332    }
333
334    /// Stop `job` for the call `call_id` and wait briefly for it to settle,
335    /// then describe it. A stopped job needs no report turn: the model asked
336    /// for the stop.
337    async fn cancel(&self, job: &str, call_id: &str) -> Result<Value, ToolError> {
338        let mut done = {
339            let mut state = self.state();
340            let index = state
341                .jobs
342                .iter()
343                .position(|candidate| candidate.id == job)
344                .ok_or_else(|| unknown_job(job))?;
345            let entry = &mut state.jobs[index];
346            if entry.outcome.is_some() {
347                let (mut value, change) = entry.describe();
348                value["note"] = "The job had already finished.".into();
349                if let Some(change) = change {
350                    state.record(call_id, change);
351                }
352                return Ok(value);
353            }
354            entry.cancelled = true;
355            entry.cancel.cancel();
356            let done = entry.done.clone();
357            if !entry.reported {
358                entry.reported = true;
359                let change = entry.change(JobStatus::Cancelled);
360                state.record(call_id, change);
361            }
362            done
363        };
364        let _ = tokio::time::timeout(CANCEL_SETTLE, done.wait_for(|finished| *finished)).await;
365        self.describe(Some(job), call_id)
366    }
367
368    /// Wait up to `limit` for `job` and describe it for the call `call_id`;
369    /// a finished job is then marked seen.
370    async fn wait(
371        &self,
372        job: &str,
373        limit: Duration,
374        cancellation: &CancellationToken,
375        call_id: &str,
376    ) -> Result<Value, ToolError> {
377        let mut done = self
378            .state()
379            .jobs
380            .iter()
381            .find(|candidate| candidate.id == job)
382            .map(|candidate| candidate.done.clone())
383            .ok_or_else(|| unknown_job(job))?;
384        tokio::select! {
385            () = cancellation.cancelled() => return Err(ToolError::cancelled("wait cancelled")),
386            _ = tokio::time::timeout(limit, done.wait_for(|finished| *finished)) => {}
387        }
388        self.describe(Some(job), call_id)
389    }
390
391    /// Describe one job, or every job this session remembers, for the call
392    /// `call_id`; finished jobs described are marked seen.
393    fn describe(&self, job: Option<&str>, call_id: &str) -> Result<Value, ToolError> {
394        let mut state = self.state();
395        let mut changes = Vec::new();
396        let value = if let Some(job) = job {
397            let entry = state
398                .jobs
399                .iter_mut()
400                .find(|candidate| candidate.id == job)
401                .ok_or_else(|| unknown_job(job))?;
402            let (value, change) = entry.describe();
403            changes.extend(change);
404            value
405        } else {
406            let jobs: Vec<Value> = state
407                .jobs
408                .iter_mut()
409                .map(|entry| {
410                    let (value, change) = entry.describe();
411                    changes.extend(change);
412                    value
413                })
414                .collect();
415            json!({ "jobs": jobs })
416        };
417        for change in changes {
418            state.record(call_id, change);
419        }
420        Ok(value)
421    }
422
423    /// The jobs the call `call_id` started or showed the model the result of,
424    /// for its `tool.completed`.
425    pub fn take_changes(&self, call_id: &str) -> Vec<JobChange> {
426        let mut state = self.state();
427        let mut taken = Vec::new();
428        state.changes.retain(|(call, change)| {
429            if call == call_id {
430                taken.push(change.clone());
431                false
432            } else {
433                true
434            }
435        });
436        taken
437    }
438
439    /// Finished jobs the model has not seen yet whose report is due, at most
440    /// a few, for a report turn. Each stays unreported, and is not taken
441    /// again, until the turn ends: [`report_settled`](Self::report_settled)
442    /// or [`report_failed`](Self::report_failed).
443    pub fn take_unreported(&self) -> Vec<JobReport> {
444        let now = tokio::time::Instant::now();
445        let mut state = self.state();
446        state
447            .jobs
448            .iter_mut()
449            .filter(|job| job.awaits_report() && job.report_due.is_none_or(|due| due <= now))
450            .take(REPORT_MAX_JOBS)
451            .filter_map(|job| {
452                job.reporting = true;
453                job.report()
454            })
455            .collect()
456    }
457
458    /// The report turn about `jobs` ended without failing: it completed, so
459    /// the model has seen them, or the user cancelled it. Either settles them.
460    pub fn report_settled(&self, jobs: &[String]) {
461        let mut state = self.state();
462        for job in state.jobs.iter_mut().filter(|job| jobs.contains(&job.id)) {
463            job.reporting = false;
464            job.reported = true;
465            job.report_due = None;
466        }
467    }
468
469    /// The report turn about `jobs` failed, so the model has not seen them
470    /// (a failed turn leaves no history). With `retry`, another turn is due
471    /// after a delay, up to three turns in all; otherwise, or once they are
472    /// used up, the jobs are settled and returned for the server to report
473    /// directly. A turn that reported several jobs decides for all of them,
474    /// by the one that failed most.
475    pub fn report_failed(&self, jobs: &[String], retry: bool) -> ReportFailure {
476        let mut state = self.state();
477        let mut failed: Vec<&mut Job> = state
478            .jobs
479            .iter_mut()
480            .filter(|job| job.reporting && jobs.contains(&job.id))
481            .collect();
482        for job in &mut failed {
483            job.reporting = false;
484            job.report_attempts += 1;
485        }
486        // One the turn's model looked up through a call is settled already.
487        failed.retain(|job| !job.reported);
488        let Some(attempts) = failed.iter().map(|job| job.report_attempts).max() else {
489            return ReportFailure::Settled;
490        };
491        let delay = usize::try_from(attempts - 1)
492            .ok()
493            .and_then(|earlier| REPORT_RETRY_DELAYS.get(earlier));
494        if retry && let Some(&delay) = delay {
495            let due = tokio::time::Instant::now() + delay;
496            for job in failed {
497                job.report_due = Some(due);
498            }
499            return ReportFailure::Retry(delay);
500        }
501        let reports = failed
502            .into_iter()
503            .filter_map(|job| {
504                job.reported = true;
505                job.report_due = None;
506                job.report()
507            })
508            .collect();
509        ReportFailure::GiveUp { attempts, reports }
510    }
511
512    /// Make every report waiting out a failed turn due now, as once a turn
513    /// succeeds and the model is evidently reachable again. Whether any was
514    /// waiting.
515    pub fn retry_now(&self) -> bool {
516        let mut state = self.state();
517        let mut waiting = false;
518        for job in state.jobs.iter_mut().filter(|job| job.awaits_report()) {
519            waiting |= job.report_due.take().is_some();
520        }
521        waiting
522    }
523
524    /// When the earliest report waiting out a failed turn is due; it may
525    /// already have passed.
526    pub fn next_retry(&self) -> Option<tokio::time::Instant> {
527        self.state()
528            .jobs
529            .iter()
530            .filter(|job| job.awaits_report())
531            .filter_map(|job| job.report_due)
532            .min()
533    }
534
535    /// Jobs still running or whose result is not reported yet, including
536    /// those a report turn is reporting now.
537    pub fn pending(&self) -> usize {
538        self.state()
539            .jobs
540            .iter()
541            .filter(|job| job.outcome.is_none() || !job.reported)
542            .count()
543    }
544
545    /// Whether any job is still running.
546    pub fn running(&self) -> usize {
547        self.state()
548            .jobs
549            .iter()
550            .filter(|job| job.outcome.is_none())
551            .count()
552    }
553}
554
555/// How a job's call takes its conversation's turn.
556enum Turn {
557    /// A new conversation, which nothing else can reach yet.
558    New,
559    /// The conversation's next turn, which nothing was ahead of.
560    Next(Place),
561    /// A prompt queued behind the turns ahead of it.
562    Queued(Place),
563}
564
565impl Turn {
566    fn place(self) -> Option<Place> {
567        match self {
568            Self::New => None,
569            Self::Next(place) | Self::Queued(place) => Some(place),
570        }
571    }
572}
573
574fn finish(jobs: &Weak<BackgroundJobs>, id: &str, output: ToolOutput, elapsed: Duration) {
575    // The session ended first: its jobs were cancelled with it.
576    let Some(jobs) = jobs.upgrade() else {
577        return;
578    };
579    {
580        let mut state = jobs.state();
581        if let Some(job) = state.jobs.iter_mut().find(|job| job.id == id) {
582            job.last_progress = job.progress.take().or(job.last_progress.take());
583            job.outcome = Some(Outcome { output, elapsed });
584            // Stopped on request: the model already knows.
585            job.reported |= job.cancelled;
586        }
587        // Keep every running job and the newest finished ones.
588        let finished = state
589            .jobs
590            .iter()
591            .filter(|job| job.outcome.is_some())
592            .count();
593        let mut excess = finished.saturating_sub(MAX_FINISHED);
594        state.jobs.retain(|job| {
595            if excess > 0 && job.outcome.is_some() && job.reported {
596                excess -= 1;
597                false
598            } else {
599                true
600            }
601        });
602    }
603    if let Some(finished) = &jobs.finished {
604        let _ = finished.send(());
605    }
606}
607
608impl Job {
609    /// Finished, unseen by the model, and not in a report turn.
610    fn awaits_report(&self) -> bool {
611        self.outcome.is_some() && !self.reported && !self.reporting
612    }
613
614    /// The finished job's result, for a report.
615    fn report(&self) -> Option<JobReport> {
616        let outcome = self.outcome.as_ref()?;
617        let reply = AgentReply::read(&result_value(&outcome.output)).unwrap_or_default();
618        Some(JobReport {
619            job: self.id.clone(),
620            agent: self.agent.clone(),
621            task: self.task.clone(),
622            status: job_status(&outcome.output, &reply),
623            session: reply.session,
624            reply: bounded(
625                reply
626                    .reply
627                    .as_deref()
628                    .unwrap_or(outcome.output.content.as_str()),
629                REPORT_REPLY_CHARS,
630            ),
631        })
632    }
633
634    /// This job with `status`, as its clients learn of it.
635    fn change(&self, status: JobStatus) -> JobChange {
636        JobChange {
637            job: self.id.clone(),
638            tool: self.tool.clone(),
639            agent: self.agent.clone(),
640            status,
641            task: self.task.clone(),
642        }
643    }
644
645    /// The job as the model reads it, and its change when this is the first
646    /// time the model sees its result.
647    fn describe(&mut self) -> (Value, Option<JobChange>) {
648        if let Some(line) = self.progress.take() {
649            self.last_progress = Some(line);
650        }
651        let mut change = None;
652        let mut value = Map::new();
653        value.insert("job".into(), self.id.clone().into());
654        value.insert("agent".into(), self.agent.clone().into());
655        match &self.outcome {
656            None => {
657                value.insert("status".into(), "running".into());
658                value.insert(
659                    "elapsed_seconds".into(),
660                    self.started.elapsed().as_secs().into(),
661                );
662                if let Some(progress) = &self.last_progress {
663                    value.insert("progress".into(), progress.clone().into());
664                }
665            }
666            Some(outcome) => {
667                let result = result_value(&outcome.output);
668                let status = if self.cancelled {
669                    JobStatus::Cancelled
670                } else {
671                    job_status(
672                        &outcome.output,
673                        &AgentReply::read(&result).unwrap_or_default(),
674                    )
675                };
676                value.insert("status".into(), status.as_str().into());
677                value.insert("elapsed_seconds".into(), outcome.elapsed.as_secs().into());
678                value.insert("result".into(), result);
679                if !self.reported {
680                    self.reported = true;
681                    change = Some(self.change(status));
682                }
683            }
684        }
685        (Value::Object(value), change)
686    }
687}
688
689/// The agent tool's structured result, or its text when it was not JSON.
690fn result_value(output: &ToolOutput) -> Value {
691    match serde_json::from_str::<Value>(&output.content) {
692        Ok(value @ Value::Object(_)) => value,
693        _ => json!({ "reply": output.content, "is_error": output.is_error() }),
694    }
695}
696
697/// How a finished job ended: its agent's reported status, else whether its
698/// call failed.
699fn job_status(output: &ToolOutput, reply: &AgentReply) -> JobStatus {
700    reply.status.unwrap_or(if output.is_error() {
701        JobStatus::Failed
702    } else {
703        JobStatus::Completed
704    })
705}
706
707/// The first non-empty line of `prompt`, at most [`TASK_CHARS`] characters,
708/// which names a job to people.
709fn task_line(prompt: &str) -> String {
710    let line = prompt
711        .lines()
712        .map(str::trim)
713        .find(|line| !line.is_empty())
714        .unwrap_or_default();
715    let mut task: String = line
716        .chars()
717        .filter(|character| !character.is_control())
718        .take(TASK_CHARS)
719        .collect();
720    if line.chars().count() > TASK_CHARS {
721        task.push('…');
722    }
723    task
724}
725
726fn unknown_job(job: &str) -> ToolError {
727    ToolError::invalid_arguments(format!(
728        "unknown background job {:?}; agent_status lists this session's jobs",
729        bounded(job, 64)
730    ))
731}
732
733/// The server-started turn that reports finished jobs to the model.
734pub fn report_prompt(reports: &[JobReport]) -> String {
735    format!(
736        "[SCV background report] Delegated work you started in the background has \
737         finished. The user did not send this message: tell them briefly what \
738         happened and the key result.\n\n{}",
739        describe_reports(reports)
740    )
741}
742
743/// What the model is told of jobs the server reported directly, because
744/// their report turns failed with `error`: the user has their results.
745pub fn delivered_note(reports: &[JobReport], error: &str) -> String {
746    format!(
747        "[SCV background report, already delivered] Delegated work you started in \
748         the background has finished. Your report of it failed ({error}), so SCV \
749         sent the user the results below directly. The user did not send this \
750         message; repeat the results only if they ask.\n\n{}",
751        describe_reports(reports)
752    )
753}
754
755/// The `agent` tool, able to run a call in the background as well.
756pub(crate) struct BackgroundCapable {
757    pub(crate) inner: Arc<AgentTool>,
758    pub(crate) jobs: Arc<BackgroundJobs>,
759}
760
761/// The `agent` tool for a queued prompt, whose turn has come by the time it
762/// runs, so no busy policy applies to it again.
763struct QueuedTool(Arc<AgentTool>);
764
765#[async_trait]
766impl Tool for QueuedTool {
767    fn spec(&self) -> ToolSpec {
768        self.0.spec()
769    }
770    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
771        self.0.risk(arguments)
772    }
773    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
774        self.0.approval_summary(arguments)
775    }
776    async fn execute(
777        &self,
778        arguments: Value,
779        context: ToolContext,
780    ) -> Result<ToolOutput, ToolError> {
781        self.0.execute_when_idle(arguments, context).await
782    }
783}
784
785/// Split `background` off an agent call's arguments.
786fn split_background(arguments: &Value) -> (Value, bool) {
787    let mut arguments = arguments.clone();
788    let background = arguments
789        .as_object_mut()
790        .and_then(|object| object.remove("background"))
791        .is_some_and(|value| value.as_bool() == Some(true));
792    (arguments, background)
793}
794
795#[async_trait]
796impl Tool for BackgroundCapable {
797    fn spec(&self) -> ToolSpec {
798        let mut spec = self.inner.spec();
799        if let Some(properties) = spec
800            .parameters
801            .get_mut("properties")
802            .and_then(Value::as_object_mut)
803        {
804            properties.insert(
805                "background".into(),
806                json!({
807                    "type":"boolean",
808                    "description":"Run in the background: the call returns a job handle at once, \
809                        the user can keep talking to you while the agent works, and SCV reports \
810                        the result in a new turn when it finishes. Use it for any substantial \
811                        task; run in the foreground only for quick work whose result you need \
812                        within this turn."
813                }),
814            );
815        }
816        spec.description.push_str(&format!(
817            " Set background to true for anything beyond a quick task: the call returns a job \
818             handle at once (at most {} running per session), SCV reports the result when the \
819             job finishes, agent_status shows progress, and agent_cancel stops it. A background \
820             job's own approval requests get only the answer this session would give without \
821             asking a person.",
822            self.jobs.limit()
823        ));
824        spec
825    }
826
827    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
828        self.inner.risk(&split_background(arguments).0)
829    }
830
831    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
832        let (arguments, background) = split_background(arguments);
833        let mut summary = self.inner.approval_summary(&arguments)?;
834        if background {
835            summary.push_str(
836                " Runs in the background: the call returns at once and the result is \
837                 reported when the agent finishes.",
838            );
839        }
840        Ok(summary)
841    }
842
843    async fn execute(
844        &self,
845        arguments: Value,
846        context: ToolContext,
847    ) -> Result<ToolOutput, ToolError> {
848        let (arguments, background) = split_background(&arguments);
849        let (agent, routed) = self.inner.route(&arguments)?;
850        // Validate before returning a job handle, so a bad call fails now.
851        if background {
852            agent.backend.risk(&routed)?;
853        }
854        let Some(handle) = parse_args::<AgentArgs>(&routed)?
855            .session
856            .filter(|session| conversation::is_handle(session))
857        else {
858            return if background {
859                self.start(agent, arguments, &context, Turn::New)
860            } else {
861                self.inner.execute(arguments, context).await
862            };
863        };
864        // Taken before anything runs or waits, so calls keep their order.
865        let place = self.jobs.lanes.join(&handle);
866        if place.ahead == 0 && !agent.backend.busy(&routed)? {
867            return if background {
868                self.start(agent, arguments, &context, Turn::Next(place))
869            } else {
870                // Held until the turn ends, so a later call finds it busy.
871                let output = self.inner.execute(arguments, context).await;
872                drop(place);
873                output
874            };
875        }
876        let behavior = match agent.on_busy(&routed)? {
877            BusyBehavior::Steer => {
878                if let Some(output) = agent.backend.steer(routed, context.clone()).await? {
879                    return Ok(output);
880                }
881                agent.busy.fallback()
882            }
883            behavior => behavior,
884        };
885        // Calls ahead of this one other than the running turn.
886        let waiting = place.ahead.saturating_sub(1);
887        match behavior {
888            BusyBehavior::Fail => Err(ToolError::failed(if waiting == 0 {
889                format!("session busy: conversation {handle} is still running a turn")
890            } else {
891                format!(
892                    "session busy: conversation {handle} is still running a turn, and \
893                     {waiting} more prompts wait for it"
894                )
895            })),
896            // A background call cannot hold its caller, so it queues instead.
897            BusyBehavior::Wait if !background => {
898                place.wait_first(&context.cancellation).await?;
899                let output = self.inner.execute_when_idle(arguments, context).await;
900                drop(place);
901                output
902            }
903            _ if waiting >= agent.busy.max_queued_turns => Err(ToolError::limit(format!(
904                "session busy: conversation {handle} is running a turn and {waiting} prompts \
905                 already wait for it, of the {} agent.max_queued_turns allows. Send this one \
906                 once they have run, or stop one with agent_cancel.",
907                agent.busy.max_queued_turns
908            ))),
909            _ => self.start(agent, arguments, &context, Turn::Queued(place)),
910        }
911    }
912}
913
914impl BackgroundCapable {
915    /// Start the call on `agent` as a background job of this session.
916    fn start(
917        &self,
918        agent: &Offered,
919        arguments: Value,
920        context: &ToolContext,
921        turn: Turn,
922    ) -> Result<ToolOutput, ToolError> {
923        let tool: Arc<dyn Tool> = if matches!(turn, Turn::Queued(_)) {
924            Arc::new(QueuedTool(Arc::clone(&self.inner)))
925        } else {
926            Arc::clone(&self.inner) as Arc<dyn Tool>
927        };
928        let started = self
929            .jobs
930            .start(tool, AGENT_TOOL, &agent.name, arguments, context, turn)?;
931        Ok(ToolOutput::success(started.to_string()))
932    }
933}
934
935#[derive(Deserialize)]
936#[serde(deny_unknown_fields)]
937struct WaitArgs {
938    job: String,
939    timeout_seconds: Option<u64>,
940}
941
942#[derive(Deserialize)]
943#[serde(deny_unknown_fields)]
944struct StatusArgs {
945    job: Option<String>,
946}
947
948/// `agent_wait`: block until a background job finishes, or the timeout.
949pub(crate) struct WaitTool {
950    pub(crate) jobs: Arc<BackgroundJobs>,
951    pub(crate) timeouts: Timeouts,
952}
953
954#[async_trait]
955impl Tool for WaitTool {
956    fn spec(&self) -> ToolSpec {
957        let mut timeout = timeout_schema(self.timeouts);
958        timeout["description"] = format!(
959            "Seconds to wait before returning the job still running. Defaults to {}; at most {}.",
960            self.timeouts.default.min(self.timeouts.max).as_secs(),
961            self.timeouts.max.as_secs()
962        )
963        .into();
964        ToolSpec {
965            name: "agent_wait".into(),
966            description: "Wait for a background agent job (the `job` handle an agent call with \
967                background: true returned) to finish, and return its result. Returns early with \
968                status running when the timeout passes. Waiting holds your turn open, so the \
969                user cannot reach you meanwhile; usually let SCV report the result instead."
970                .into(),
971            parameters: json!({
972                "type":"object",
973                "properties":{"job":{"type":"string"},"timeout_seconds":timeout},
974                "required":["job"],
975                "additionalProperties":false
976            }),
977        }
978    }
979
980    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
981        let args: WaitArgs = parse_args(arguments)?;
982        self.timeouts.resolve(args.timeout_seconds)?;
983        Ok(ToolRisk::ReadOnly)
984    }
985
986    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
987        let args: WaitArgs = parse_args(arguments)?;
988        Ok(format!(
989            "Wait for background job {}",
990            bounded(&args.job, 64)
991        ))
992    }
993
994    async fn execute(
995        &self,
996        arguments: Value,
997        context: ToolContext,
998    ) -> Result<ToolOutput, ToolError> {
999        let args: WaitArgs = parse_args(&arguments)?;
1000        let limit = self.timeouts.resolve(args.timeout_seconds)?;
1001        let value = self
1002            .jobs
1003            .wait(&args.job, limit, &context.cancellation, &context.call_id)
1004            .await?;
1005        Ok(ToolOutput::success(value.to_string()))
1006    }
1007}
1008
1009/// `agent_status`: describe one background job or all of them.
1010pub(crate) struct StatusTool {
1011    pub(crate) jobs: Arc<BackgroundJobs>,
1012}
1013
1014#[async_trait]
1015impl Tool for StatusTool {
1016    fn spec(&self) -> ToolSpec {
1017        ToolSpec {
1018            name: "agent_status".into(),
1019            description: "Show this session's background agent jobs: running ones with their \
1020                latest progress, finished ones with their result. Pass job for one job."
1021                .into(),
1022            parameters: json!({
1023                "type":"object",
1024                "properties":{"job":{"type":"string"}},
1025                "additionalProperties":false
1026            }),
1027        }
1028    }
1029
1030    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
1031        let _: StatusArgs = parse_args(arguments)?;
1032        Ok(ToolRisk::ReadOnly)
1033    }
1034
1035    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
1036        let args: StatusArgs = parse_args(arguments)?;
1037        Ok(args.job.map_or_else(
1038            || "List background jobs".into(),
1039            |job| format!("Show background job {}", bounded(&job, 64)),
1040        ))
1041    }
1042
1043    async fn execute(
1044        &self,
1045        arguments: Value,
1046        context: ToolContext,
1047    ) -> Result<ToolOutput, ToolError> {
1048        let args: StatusArgs = parse_args(&arguments)?;
1049        let value = self.jobs.describe(args.job.as_deref(), &context.call_id)?;
1050        Ok(ToolOutput::success(value.to_string()))
1051    }
1052}
1053
1054#[derive(Deserialize)]
1055#[serde(deny_unknown_fields)]
1056struct CancelArgs {
1057    job: String,
1058}
1059
1060/// `agent_cancel`: stop one running background job.
1061pub(crate) struct CancelTool {
1062    pub(crate) jobs: Arc<BackgroundJobs>,
1063}
1064
1065#[async_trait]
1066impl Tool for CancelTool {
1067    fn spec(&self) -> ToolSpec {
1068        ToolSpec {
1069            name: "agent_cancel".into(),
1070            description: "Stop a running background agent job (the `job` handle an agent call \
1071                with background: true returned), for example when the user no longer wants \
1072                it. The agent and every process it started are stopped; work it already wrote \
1073                stays. Returns the job with status cancelled, and no report turn follows."
1074                .into(),
1075            parameters: json!({
1076                "type":"object",
1077                "properties":{"job":{"type":"string"}},
1078                "required":["job"],
1079                "additionalProperties":false
1080            }),
1081        }
1082    }
1083
1084    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
1085        let _: CancelArgs = parse_args(arguments)?;
1086        Ok(ToolRisk::Process)
1087    }
1088
1089    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
1090        let args: CancelArgs = parse_args(arguments)?;
1091        Ok(format!("Stop background job {}", bounded(&args.job, 64)))
1092    }
1093
1094    async fn execute(
1095        &self,
1096        arguments: Value,
1097        context: ToolContext,
1098    ) -> Result<ToolOutput, ToolError> {
1099        let args: CancelArgs = parse_args(&arguments)?;
1100        let value = self.jobs.cancel(&args.job, &context.call_id).await?;
1101        Ok(ToolOutput::success(value.to_string()))
1102    }
1103}
1104
1105#[cfg(test)]
1106mod tests;