Skip to main content

scv_tools/
background.rs

1//! Background delegations: an `agent_*` call with `background: true` returns a
2//! job handle at once while the agent keeps working; `agent_wait` and
3//! `agent_status` observe the job, and the session is told when it finishes
4//! so the server can report it in a turn of its own.
5
6use std::{
7    path::PathBuf,
8    sync::{Arc, Mutex, PoisonError, Weak},
9    time::{Duration, Instant},
10};
11
12use async_trait::async_trait;
13use scv_core::{ProgressSink, Tool, ToolContext, ToolError, ToolOutput, ToolRisk, ToolSpec};
14use serde::Deserialize;
15use serde_json::{Map, Value, json};
16use tokio::sync::{mpsc, watch};
17use tokio_util::sync::CancellationToken;
18
19use crate::{Timeouts, bounded, parse_args, timeout_schema};
20
21/// Finished jobs a session keeps for `agent_status` beyond the running ones.
22const MAX_FINISHED: usize = 16;
23/// Characters of a finished job's reply quoted in a server-started report turn.
24const REPORT_REPLY_CHARS: usize = 6000;
25/// Jobs one report turn covers; any more wait for the next.
26const REPORT_MAX_JOBS: usize = 4;
27
28/// One session's background jobs. Dropping the store (with the session's
29/// tools) cancels every job still running.
30pub struct BackgroundJobs {
31    limit: usize,
32    state: Mutex<JobsState>,
33    cancellation: CancellationToken,
34    /// Woken whenever a job finishes, so the session can report it.
35    finished: Option<mpsc::UnboundedSender<()>>,
36}
37
38impl std::fmt::Debug for BackgroundJobs {
39    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
40        formatter
41            .debug_struct("BackgroundJobs")
42            .field("limit", &self.limit)
43            .finish_non_exhaustive()
44    }
45}
46
47#[derive(Default)]
48struct JobsState {
49    next: u64,
50    jobs: Vec<Job>,
51}
52
53struct Job {
54    id: String,
55    tool: String,
56    started: Instant,
57    progress: ProgressSink,
58    last_progress: Option<String>,
59    outcome: Option<Outcome>,
60    /// The model has seen the result (through `agent_wait`, `agent_status`,
61    /// or a report turn), so it needs no report turn.
62    reported: bool,
63    done: watch::Receiver<bool>,
64}
65
66struct Outcome {
67    output: ToolOutput,
68    elapsed: Duration,
69}
70
71/// A finished job not yet seen by the model, for a server-started report turn.
72#[derive(Debug, Clone)]
73pub struct JobReport {
74    pub job: String,
75    pub tool: String,
76    pub status: String,
77    pub session: Option<String>,
78    pub reply: String,
79}
80
81impl Drop for BackgroundJobs {
82    fn drop(&mut self) {
83        self.cancellation.cancel();
84    }
85}
86
87impl BackgroundJobs {
88    /// At most `limit` jobs run at once. `finished` is woken as jobs finish.
89    pub fn new(limit: usize, finished: Option<mpsc::UnboundedSender<()>>) -> Self {
90        Self {
91            limit,
92            state: Mutex::default(),
93            cancellation: CancellationToken::new(),
94            finished,
95        }
96    }
97
98    pub fn limit(&self) -> usize {
99        self.limit
100    }
101
102    fn state(&self) -> std::sync::MutexGuard<'_, JobsState> {
103        self.state.lock().unwrap_or_else(PoisonError::into_inner)
104    }
105
106    /// Start `tool` with `arguments` in the background and return its job's
107    /// `{"job","status":"running","background":true}` description.
108    fn start(
109        self: &Arc<Self>,
110        tool: Arc<dyn Tool>,
111        name: &str,
112        arguments: Value,
113        workspace: PathBuf,
114    ) -> Result<Value, ToolError> {
115        let (id, progress, done_tx, cancellation) = {
116            let mut state = self.state();
117            let running = state
118                .jobs
119                .iter()
120                .filter(|job| job.outcome.is_none())
121                .count();
122            if running >= self.limit {
123                return Err(ToolError(format!(
124                    "{running} background jobs are already running, the limit \
125                     (agent.max_background); wait for one with agent_wait first"
126                )));
127            }
128            state.next += 1;
129            let id = format!("job-{}", state.next);
130            let progress = ProgressSink::buffered();
131            let (done_tx, done) = watch::channel(false);
132            state.jobs.push(Job {
133                id: id.clone(),
134                tool: name.to_owned(),
135                started: Instant::now(),
136                progress: progress.clone(),
137                last_progress: None,
138                outcome: None,
139                reported: false,
140                done,
141            });
142            (id, progress, done_tx, self.cancellation.child_token())
143        };
144        let jobs = Arc::downgrade(self);
145        let job = id.clone();
146        tokio::spawn(async move {
147            let started = Instant::now();
148            let mut context = ToolContext::new(workspace, cancellation);
149            // A background job cannot ask anyone: nested approval requests
150            // are denied, so it runs on the agent's own permissions.
151            context.progress = progress;
152            let output = tool
153                .execute(arguments, context)
154                .await
155                .unwrap_or_else(|error| ToolOutput::failure(error.to_string()));
156            finish(&jobs, &job, output, started.elapsed());
157            let _ = done_tx.send(true);
158        });
159        Ok(json!({
160            "job": id,
161            "tool": name,
162            "status": "running",
163            "background": true,
164            "note": "The agent is working in the background. SCV reports the result when it \
165                     finishes; agent_wait blocks for it and agent_status shows its progress."
166        }))
167    }
168
169    /// Wait up to `limit` for `job` and describe it; a finished job is then
170    /// marked seen.
171    async fn wait(
172        &self,
173        job: &str,
174        limit: Duration,
175        cancellation: &CancellationToken,
176    ) -> Result<Value, ToolError> {
177        let mut done = self
178            .state()
179            .jobs
180            .iter()
181            .find(|candidate| candidate.id == job)
182            .map(|candidate| candidate.done.clone())
183            .ok_or_else(|| unknown_job(job))?;
184        tokio::select! {
185            _ = cancellation.cancelled() => return Err(ToolError("wait cancelled".into())),
186            _ = tokio::time::timeout(limit, done.wait_for(|finished| *finished)) => {}
187        }
188        self.describe(Some(job))
189    }
190
191    /// Describe one job, or every job this session remembers; finished jobs
192    /// described are marked seen.
193    fn describe(&self, job: Option<&str>) -> Result<Value, ToolError> {
194        let mut state = self.state();
195        if let Some(job) = job {
196            let entry = state
197                .jobs
198                .iter_mut()
199                .find(|candidate| candidate.id == job)
200                .ok_or_else(|| unknown_job(job))?;
201            return Ok(entry.describe());
202        }
203        let jobs: Vec<Value> = state.jobs.iter_mut().map(Job::describe).collect();
204        Ok(json!({ "jobs": jobs }))
205    }
206
207    /// Finished jobs the model has not seen yet, marked seen, for a report
208    /// turn. At most a few per call; the rest stay for the next.
209    pub fn take_unreported(&self) -> Vec<JobReport> {
210        let mut state = self.state();
211        state
212            .jobs
213            .iter_mut()
214            .filter(|job| job.outcome.is_some() && !job.reported)
215            .take(REPORT_MAX_JOBS)
216            .map(|job| {
217                job.reported = true;
218                let outcome = job.outcome.as_ref().expect("filtered on outcome");
219                let result = result_value(&outcome.output);
220                JobReport {
221                    job: job.id.clone(),
222                    tool: job.tool.clone(),
223                    status: job_status(&outcome.output, &result),
224                    session: result
225                        .get("session")
226                        .and_then(Value::as_str)
227                        .map(str::to_owned),
228                    reply: bounded(
229                        result
230                            .get("reply")
231                            .and_then(Value::as_str)
232                            .unwrap_or(outcome.output.content.as_str()),
233                        REPORT_REPLY_CHARS,
234                    ),
235                }
236            })
237            .collect()
238    }
239
240    /// Whether any job is still running.
241    pub fn running(&self) -> usize {
242        self.state()
243            .jobs
244            .iter()
245            .filter(|job| job.outcome.is_none())
246            .count()
247    }
248}
249
250fn finish(jobs: &Weak<BackgroundJobs>, id: &str, output: ToolOutput, elapsed: Duration) {
251    // The session ended first: its jobs were cancelled with it.
252    let Some(jobs) = jobs.upgrade() else {
253        return;
254    };
255    {
256        let mut state = jobs.state();
257        if let Some(job) = state.jobs.iter_mut().find(|job| job.id == id) {
258            job.last_progress = job.progress.take().or(job.last_progress.take());
259            job.outcome = Some(Outcome { output, elapsed });
260        }
261        // Keep every running job and the newest finished ones.
262        let finished = state
263            .jobs
264            .iter()
265            .filter(|job| job.outcome.is_some())
266            .count();
267        let mut excess = finished.saturating_sub(MAX_FINISHED);
268        state.jobs.retain(|job| {
269            if excess > 0 && job.outcome.is_some() && job.reported {
270                excess -= 1;
271                false
272            } else {
273                true
274            }
275        });
276    }
277    if let Some(finished) = &jobs.finished {
278        let _ = finished.send(());
279    }
280}
281
282impl Job {
283    fn describe(&mut self) -> Value {
284        if let Some(line) = self.progress.take() {
285            self.last_progress = Some(line);
286        }
287        let mut value = Map::new();
288        value.insert("job".into(), self.id.clone().into());
289        value.insert("tool".into(), self.tool.clone().into());
290        match &self.outcome {
291            None => {
292                value.insert("status".into(), "running".into());
293                value.insert(
294                    "elapsed_seconds".into(),
295                    self.started.elapsed().as_secs().into(),
296                );
297                if let Some(progress) = &self.last_progress {
298                    value.insert("progress".into(), progress.clone().into());
299                }
300            }
301            Some(outcome) => {
302                self.reported = true;
303                let result = result_value(&outcome.output);
304                value.insert("status".into(), job_status(&outcome.output, &result).into());
305                value.insert("elapsed_seconds".into(), outcome.elapsed.as_secs().into());
306                value.insert("result".into(), result);
307            }
308        }
309        Value::Object(value)
310    }
311}
312
313/// The agent tool's structured result, or its text when it was not JSON.
314fn result_value(output: &ToolOutput) -> Value {
315    match serde_json::from_str::<Value>(&output.content) {
316        Ok(value @ Value::Object(_)) => value,
317        _ => json!({ "reply": output.content, "is_error": output.is_error }),
318    }
319}
320
321fn job_status(output: &ToolOutput, result: &Value) -> String {
322    result
323        .get("status")
324        .and_then(Value::as_str)
325        .map(str::to_owned)
326        .unwrap_or_else(|| {
327            if output.is_error {
328                "failed".into()
329            } else {
330                "completed".into()
331            }
332        })
333}
334
335fn unknown_job(job: &str) -> ToolError {
336    ToolError(format!(
337        "unknown background job {:?}; agent_status lists this session's jobs",
338        bounded(job, 64)
339    ))
340}
341
342/// The server-started turn that reports finished jobs to the model.
343pub fn report_prompt(reports: &[JobReport]) -> String {
344    let mut prompt = String::from(
345        "[SCV background report] Delegated work you started in the background has \
346         finished. The user did not send this message: tell them briefly what \
347         happened and the key result.\n",
348    );
349    for report in reports {
350        prompt.push_str(&format!(
351            "\n{} ({}{}): {}\n{}\n",
352            report.job,
353            report.tool,
354            report
355                .session
356                .as_deref()
357                .map_or_else(String::new, |session| format!(", conversation {session}")),
358            report.status,
359            report.reply.trim()
360        ));
361    }
362    prompt
363}
364
365/// An agent tool that can also run in the background.
366pub(crate) struct BackgroundCapable {
367    pub(crate) inner: Arc<dyn Tool>,
368    pub(crate) jobs: Arc<BackgroundJobs>,
369}
370
371/// Split `background` off an agent call's arguments.
372fn split_background(arguments: &Value) -> (Value, bool) {
373    let mut arguments = arguments.clone();
374    let background = arguments
375        .as_object_mut()
376        .and_then(|object| object.remove("background"))
377        .is_some_and(|value| value.as_bool() == Some(true));
378    (arguments, background)
379}
380
381#[async_trait]
382impl Tool for BackgroundCapable {
383    fn spec(&self) -> ToolSpec {
384        let mut spec = self.inner.spec();
385        if let Some(properties) = spec
386            .parameters
387            .get_mut("properties")
388            .and_then(Value::as_object_mut)
389        {
390            properties.insert(
391                "background".into(),
392                json!({
393                    "type":"boolean",
394                    "description":"Run in the background: return a job handle at once and keep \
395                        working with the user while the agent runs. SCV reports the result when \
396                        it finishes. Use it for long work such as landing or releasing a change."
397                }),
398            );
399        }
400        spec.description.push_str(&format!(
401            " Set background to true for long work: the call returns a job handle at once \
402             (at most {} running per session), SCV reports the result when the job finishes, \
403             and agent_wait / agent_status observe it meanwhile.",
404            self.jobs.limit()
405        ));
406        spec
407    }
408
409    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
410        self.inner.risk(&split_background(arguments).0)
411    }
412
413    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
414        let (arguments, background) = split_background(arguments);
415        let mut summary = self.inner.approval_summary(&arguments)?;
416        if background {
417            summary.push_str(
418                " Runs in the background: the call returns at once and the result is \
419                 reported when the agent finishes.",
420            );
421        }
422        Ok(summary)
423    }
424
425    async fn execute(
426        &self,
427        arguments: Value,
428        context: ToolContext,
429    ) -> Result<ToolOutput, ToolError> {
430        let (arguments, background) = split_background(&arguments);
431        if !background {
432            return self.inner.execute(arguments, context).await;
433        }
434        // Validate before returning a job handle, so a bad call fails now.
435        self.inner.risk(&arguments)?;
436        let name = self.inner.spec().name;
437        let started =
438            self.jobs
439                .start(Arc::clone(&self.inner), &name, arguments, context.workspace)?;
440        Ok(ToolOutput::success(started.to_string()))
441    }
442}
443
444#[derive(Deserialize)]
445#[serde(deny_unknown_fields)]
446struct WaitArgs {
447    job: String,
448    timeout_seconds: Option<u64>,
449}
450
451#[derive(Deserialize)]
452#[serde(deny_unknown_fields)]
453struct StatusArgs {
454    job: Option<String>,
455}
456
457/// `agent_wait`: block until a background job finishes, or the timeout.
458pub(crate) struct WaitTool {
459    pub(crate) jobs: Arc<BackgroundJobs>,
460    pub(crate) timeouts: Timeouts,
461}
462
463#[async_trait]
464impl Tool for WaitTool {
465    fn spec(&self) -> ToolSpec {
466        let mut timeout = timeout_schema(self.timeouts);
467        timeout["description"] = format!(
468            "Seconds to wait before returning the job still running. Defaults to {}; at most {}.",
469            self.timeouts.default.min(self.timeouts.max).as_secs(),
470            self.timeouts.max.as_secs()
471        )
472        .into();
473        ToolSpec {
474            name: "agent_wait".into(),
475            description: "Wait for a background agent job (the `job` handle an agent_* call with \
476                background: true returned) to finish, and return its result. Returns early with \
477                status running when the timeout passes."
478                .into(),
479            parameters: json!({
480                "type":"object",
481                "properties":{"job":{"type":"string"},"timeout_seconds":timeout},
482                "required":["job"],
483                "additionalProperties":false
484            }),
485        }
486    }
487
488    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
489        let args: WaitArgs = parse_args(arguments)?;
490        self.timeouts.resolve(args.timeout_seconds)?;
491        Ok(ToolRisk::ReadOnly)
492    }
493
494    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
495        let args: WaitArgs = parse_args(arguments)?;
496        Ok(format!(
497            "Wait for background job {}",
498            bounded(&args.job, 64)
499        ))
500    }
501
502    async fn execute(
503        &self,
504        arguments: Value,
505        context: ToolContext,
506    ) -> Result<ToolOutput, ToolError> {
507        let args: WaitArgs = parse_args(&arguments)?;
508        let limit = self.timeouts.resolve(args.timeout_seconds)?;
509        let value = self
510            .jobs
511            .wait(&args.job, limit, &context.cancellation)
512            .await?;
513        Ok(ToolOutput::success(value.to_string()))
514    }
515}
516
517/// `agent_status`: describe one background job or all of them.
518pub(crate) struct StatusTool {
519    pub(crate) jobs: Arc<BackgroundJobs>,
520}
521
522#[async_trait]
523impl Tool for StatusTool {
524    fn spec(&self) -> ToolSpec {
525        ToolSpec {
526            name: "agent_status".into(),
527            description: "Show this session's background agent jobs: running ones with their \
528                latest progress, finished ones with their result. Pass job for one job."
529                .into(),
530            parameters: json!({
531                "type":"object",
532                "properties":{"job":{"type":"string"}},
533                "additionalProperties":false
534            }),
535        }
536    }
537
538    fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
539        let _: StatusArgs = parse_args(arguments)?;
540        Ok(ToolRisk::ReadOnly)
541    }
542
543    fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
544        let args: StatusArgs = parse_args(arguments)?;
545        Ok(args.job.map_or_else(
546            || "List background jobs".into(),
547            |job| format!("Show background job {}", bounded(&job, 64)),
548        ))
549    }
550
551    async fn execute(
552        &self,
553        arguments: Value,
554        _context: ToolContext,
555    ) -> Result<ToolOutput, ToolError> {
556        let args: StatusArgs = parse_args(&arguments)?;
557        let value = self.jobs.describe(args.job.as_deref())?;
558        Ok(ToolOutput::success(value.to_string()))
559    }
560}
561
562#[cfg(test)]
563mod tests {
564    use super::*;
565    use std::sync::atomic::{AtomicBool, Ordering};
566    use tokio::sync::Notify;
567
568    /// An agent that reports one progress line, then finishes when released
569    /// or stops when cancelled.
570    struct FakeAgent {
571        release: Arc<Notify>,
572        cancelled: Arc<AtomicBool>,
573    }
574
575    #[async_trait]
576    impl Tool for FakeAgent {
577        fn spec(&self) -> ToolSpec {
578            ToolSpec {
579                name: "agent_fake".into(),
580                description: "Fake agent.".into(),
581                parameters: json!({
582                    "type":"object",
583                    "properties":{"prompt":{"type":"string"}},
584                    "required":["prompt"],
585                    "additionalProperties":false
586                }),
587            }
588        }
589
590        fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
591            if arguments.get("prompt").and_then(Value::as_str).is_none() {
592                return Err(ToolError("prompt is required".into()));
593            }
594            Ok(ToolRisk::Delegate)
595        }
596
597        fn approval_summary(&self, _arguments: &Value) -> Result<String, ToolError> {
598            Ok("Launch the fake agent.".into())
599        }
600
601        async fn execute(
602            &self,
603            _arguments: Value,
604            context: ToolContext,
605        ) -> Result<ToolOutput, ToolError> {
606            context.progress.report("$ step one");
607            tokio::select! {
608                _ = self.release.notified() => Ok(ToolOutput::success(
609                    json!({"agent":"fake","status":"completed","reply":"all done","session":"fake-1"})
610                        .to_string(),
611                )),
612                _ = context.cancellation.cancelled() => {
613                    self.cancelled.store(true, Ordering::SeqCst);
614                    Err(ToolError("cancelled".into()))
615                }
616            }
617        }
618    }
619
620    struct Fixture {
621        tool: BackgroundCapable,
622        jobs: Arc<BackgroundJobs>,
623        release: Arc<Notify>,
624        cancelled: Arc<AtomicBool>,
625        finished: mpsc::UnboundedReceiver<()>,
626    }
627
628    fn fixture(limit: usize) -> Fixture {
629        let (finished_tx, finished) = mpsc::unbounded_channel();
630        let jobs = Arc::new(BackgroundJobs::new(limit, Some(finished_tx)));
631        let release = Arc::new(Notify::new());
632        let cancelled = Arc::new(AtomicBool::new(false));
633        let tool = BackgroundCapable {
634            inner: Arc::new(FakeAgent {
635                release: Arc::clone(&release),
636                cancelled: Arc::clone(&cancelled),
637            }),
638            jobs: Arc::clone(&jobs),
639        };
640        Fixture {
641            tool,
642            jobs,
643            release,
644            cancelled,
645            finished,
646        }
647    }
648
649    fn context() -> ToolContext {
650        ToolContext::new(std::env::temp_dir(), CancellationToken::new())
651    }
652
653    async fn start(tool: &BackgroundCapable) -> Value {
654        let output = tool
655            .execute(json!({"prompt":"work","background":true}), context())
656            .await
657            .unwrap();
658        serde_json::from_str(&output.content).unwrap()
659    }
660
661    #[tokio::test]
662    async fn background_calls_return_a_job_that_wait_and_status_observe() {
663        let mut fixture = fixture(2);
664        let started = start(&fixture.tool).await;
665        assert_eq!(
666            started,
667            json!({
668                "job":"job-1","tool":"agent_fake","status":"running","background":true,
669                "note":started["note"]
670            })
671        );
672        assert_eq!(
673            scv_protocol::background_job_update(&started.to_string()).started,
674            vec!["job-1".to_owned()]
675        );
676        // Running, with the job's latest progress.
677        let status = tokio::time::timeout(Duration::from_secs(5), async {
678            loop {
679                let status = fixture.jobs.describe(Some("job-1")).unwrap();
680                if status.get("progress").is_some() {
681                    break status;
682                }
683                tokio::time::sleep(Duration::from_millis(10)).await;
684            }
685        })
686        .await
687        .unwrap();
688        assert_eq!(status["status"], "running");
689        assert_eq!(status["progress"], "$ step one");
690        // A short wait returns it still running.
691        let waited = fixture
692            .jobs
693            .wait(
694                "job-1",
695                Duration::from_millis(50),
696                &CancellationToken::new(),
697            )
698            .await
699            .unwrap();
700        assert_eq!(waited["status"], "running");
701
702        fixture.release.notify_one();
703        let waited = fixture
704            .jobs
705            .wait("job-1", Duration::from_secs(5), &CancellationToken::new())
706            .await
707            .unwrap();
708        assert_eq!(waited["status"], "completed");
709        assert_eq!(waited["result"]["reply"], "all done");
710        assert_eq!(waited["result"]["session"], "fake-1");
711        assert_eq!(
712            scv_protocol::background_job_update(&waited.to_string()).settled,
713            vec!["job-1".to_owned()]
714        );
715        // The session was woken, but the model already saw the result.
716        fixture.finished.recv().await.unwrap();
717        assert!(fixture.jobs.take_unreported().is_empty());
718        let all = fixture.jobs.describe(None).unwrap();
719        assert_eq!(all["jobs"][0]["job"], "job-1");
720    }
721
722    #[tokio::test]
723    async fn a_finished_job_nobody_looked_at_is_reported_once() {
724        let mut fixture = fixture(2);
725        start(&fixture.tool).await;
726        fixture.release.notify_one();
727        tokio::time::timeout(Duration::from_secs(5), fixture.finished.recv())
728            .await
729            .unwrap()
730            .unwrap();
731        let reports = fixture.jobs.take_unreported();
732        assert_eq!(reports.len(), 1);
733        let report = &reports[0];
734        assert_eq!(
735            (
736                report.job.as_str(),
737                report.tool.as_str(),
738                report.status.as_str(),
739                report.session.as_deref(),
740                report.reply.as_str()
741            ),
742            (
743                "job-1",
744                "agent_fake",
745                "completed",
746                Some("fake-1"),
747                "all done"
748            )
749        );
750        let prompt = report_prompt(&reports);
751        assert!(prompt.starts_with("[SCV background report]"), "{prompt}");
752        assert!(
753            prompt.contains("job-1 (agent_fake, conversation fake-1): completed\nall done"),
754            "{prompt}"
755        );
756        assert!(fixture.jobs.take_unreported().is_empty(), "reported twice");
757    }
758
759    #[tokio::test]
760    async fn at_most_the_limit_runs_at_once() {
761        let mut fixture = fixture(1);
762        start(&fixture.tool).await;
763        let refused = fixture
764            .tool
765            .execute(json!({"prompt":"more","background":true}), context())
766            .await
767            .unwrap_err();
768        assert!(refused.0.contains("agent.max_background"), "{refused}");
769        fixture.release.notify_one();
770        fixture.finished.recv().await.unwrap();
771        assert_eq!(start(&fixture.tool).await["job"], "job-2");
772        assert_eq!(fixture.jobs.running(), 1);
773    }
774
775    #[tokio::test]
776    async fn dropping_the_session_store_cancels_running_jobs() {
777        let fixture = fixture(2);
778        start(&fixture.tool).await;
779        let cancelled = Arc::clone(&fixture.cancelled);
780        drop(fixture);
781        tokio::time::timeout(Duration::from_secs(5), async {
782            while !cancelled.load(Ordering::SeqCst) {
783                tokio::time::sleep(Duration::from_millis(10)).await;
784            }
785        })
786        .await
787        .expect("the job was cancelled with its session");
788    }
789
790    #[tokio::test]
791    async fn foreground_calls_pass_through_and_the_schema_offers_background() {
792        let fixture = fixture(2);
793        let spec = fixture.tool.spec();
794        assert_eq!(spec.name, "agent_fake");
795        assert_eq!(
796            spec.parameters["properties"]["background"]["type"],
797            "boolean"
798        );
799        assert!(
800            spec.description.contains("at most 2 running"),
801            "{}",
802            spec.description
803        );
804        let summary = fixture
805            .tool
806            .approval_summary(&json!({"prompt":"work","background":true}))
807            .unwrap();
808        assert!(summary.contains("Runs in the background"), "{summary}");
809        assert!(
810            !fixture
811                .tool
812                .approval_summary(&json!({"prompt":"work"}))
813                .unwrap()
814                .contains("background")
815        );
816        // A bad call fails at once instead of returning a job.
817        assert!(
818            fixture
819                .tool
820                .execute(json!({"background":true}), context())
821                .await
822                .is_err()
823        );
824        fixture.release.notify_one();
825        let output = fixture
826            .tool
827            .execute(json!({"prompt":"work","background":false}), context())
828            .await
829            .unwrap();
830        assert!(output.content.contains("all done"));
831        assert_eq!(fixture.jobs.running(), 0);
832        assert!(fixture.jobs.describe(Some("job-1")).is_err());
833    }
834
835    #[tokio::test]
836    async fn wait_and_status_tools_validate_and_are_read_only() {
837        let fixture = fixture(2);
838        let wait = WaitTool {
839            jobs: Arc::clone(&fixture.jobs),
840            timeouts: Timeouts {
841                default: Duration::from_secs(60),
842                max: Duration::from_secs(120),
843            },
844        };
845        assert_eq!(
846            wait.risk(&json!({"job":"job-1"})).unwrap(),
847            ToolRisk::ReadOnly
848        );
849        assert!(
850            wait.risk(&json!({"job":"job-1","timeout_seconds":121}))
851                .is_err()
852        );
853        let unknown = wait
854            .execute(json!({"job":"job-9"}), context())
855            .await
856            .unwrap_err();
857        assert!(unknown.0.contains("unknown background job"), "{unknown}");
858        let status = StatusTool {
859            jobs: Arc::clone(&fixture.jobs),
860        };
861        assert_eq!(status.risk(&json!({})).unwrap(), ToolRisk::ReadOnly);
862        let empty = status.execute(json!({}), context()).await.unwrap();
863        assert_eq!(empty.content, json!({"jobs":[]}).to_string());
864    }
865}