1use std::{
11 path::PathBuf,
12 sync::{Arc, Mutex, Weak},
13 time::{Duration, Instant},
14};
15
16use async_trait::async_trait;
17use scv_core::{
18 ApprovalGate, ProgressSink, Tool, ToolApprovals, ToolContext, ToolError, ToolOutput, ToolRisk,
19 ToolSpec,
20};
21use scv_protocol::{JobChange, JobStatus};
22use serde::Deserialize;
23use serde_json::{Map, Value, json};
24use tokio::sync::{mpsc, watch};
25use tokio_util::sync::CancellationToken;
26
27use crate::{
28 args::{Timeouts, bounded, parse_args, timeout_schema},
29 delegate::output::AgentReply,
30 sync::lock,
31};
32
33const MAX_FINISHED: usize = 16;
35const REPORT_REPLY_CHARS: usize = 6000;
37const REPORT_MAX_JOBS: usize = 4;
39const CANCEL_SETTLE: Duration = Duration::from_secs(10);
41const MAX_PENDING_CHANGES: usize = 64;
44const TASK_CHARS: usize = 80;
46
47pub struct BackgroundJobs {
50 limit: usize,
51 state: Mutex<JobsState>,
52 cancellation: CancellationToken,
53 finished: Option<mpsc::UnboundedSender<()>>,
55 approvals: Option<Arc<dyn ApprovalGate>>,
58}
59
60impl std::fmt::Debug for BackgroundJobs {
61 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
62 formatter
63 .debug_struct("BackgroundJobs")
64 .field("limit", &self.limit)
65 .finish_non_exhaustive()
66 }
67}
68
69#[derive(Default)]
70struct JobsState {
71 next: u64,
72 jobs: Vec<Job>,
73 changes: Vec<(String, JobChange)>,
75}
76
77impl JobsState {
78 fn record(&mut self, call_id: &str, change: JobChange) {
81 if call_id.is_empty() {
82 return;
83 }
84 if self.changes.len() >= MAX_PENDING_CHANGES {
85 self.changes.remove(0);
86 }
87 self.changes.push((call_id.to_owned(), change));
88 }
89}
90
91struct Job {
92 id: String,
93 tool: String,
94 task: String,
96 started: Instant,
97 progress: ProgressSink,
98 last_progress: Option<String>,
99 cancel: CancellationToken,
101 cancelled: bool,
103 outcome: Option<Outcome>,
104 reported: bool,
108 done: watch::Receiver<bool>,
109}
110
111struct Outcome {
112 output: ToolOutput,
113 elapsed: Duration,
114}
115
116#[derive(Debug, Clone)]
118pub struct JobReport {
119 pub job: String,
120 pub(crate) tool: String,
121 pub(crate) status: JobStatus,
122 pub(crate) session: Option<String>,
123 pub(crate) reply: String,
124}
125
126impl Drop for BackgroundJobs {
127 fn drop(&mut self) {
128 self.cancellation.cancel();
129 }
130}
131
132impl BackgroundJobs {
133 pub fn new(limit: usize, finished: Option<mpsc::UnboundedSender<()>>) -> Self {
135 Self {
136 limit,
137 state: Mutex::default(),
138 cancellation: CancellationToken::new(),
139 finished,
140 approvals: None,
141 }
142 }
143
144 #[must_use]
147 pub fn with_approvals(mut self, gate: Arc<dyn ApprovalGate>) -> Self {
148 self.approvals = Some(gate);
149 self
150 }
151
152 pub(crate) fn limit(&self) -> usize {
153 self.limit
154 }
155
156 fn state(&self) -> std::sync::MutexGuard<'_, JobsState> {
157 lock(&self.state)
158 }
159
160 fn start(
164 self: &Arc<Self>,
165 tool: Arc<dyn Tool>,
166 name: &str,
167 arguments: Value,
168 workspace: PathBuf,
169 call_id: &str,
170 ) -> Result<Value, ToolError> {
171 let task = task_line(
172 arguments
173 .get("prompt")
174 .and_then(Value::as_str)
175 .unwrap_or_default(),
176 );
177 let (id, progress, done_tx, cancellation) = {
178 let mut state = self.state();
179 let running = state
180 .jobs
181 .iter()
182 .filter(|job| job.outcome.is_none())
183 .count();
184 if running >= self.limit {
185 return Err(ToolError::limit(format!(
186 "{running} background jobs are already running, the limit \
187 (agent.max_background). Start this one after a job finishes, or \
188 stop one with agent_cancel if the user no longer needs it."
189 )));
190 }
191 state.next += 1;
192 let id = format!("job-{}", state.next);
193 let progress = ProgressSink::buffered();
194 let (done_tx, done) = watch::channel(false);
195 let cancel = self.cancellation.child_token();
196 state.record(
197 call_id,
198 JobChange {
199 job: id.clone(),
200 tool: name.to_owned(),
201 status: JobStatus::Running,
202 task: task.clone(),
203 },
204 );
205 state.jobs.push(Job {
206 id: id.clone(),
207 tool: name.to_owned(),
208 task,
209 started: Instant::now(),
210 progress: progress.clone(),
211 last_progress: None,
212 cancel: cancel.clone(),
213 cancelled: false,
214 outcome: None,
215 reported: false,
216 done,
217 });
218 (id, progress, done_tx, cancel)
219 };
220 let jobs = Arc::downgrade(self);
221 let job = id.clone();
222 let approvals = self.approvals.clone();
223 tokio::spawn(async move {
224 let started = Instant::now();
225 let mut context = ToolContext::new(workspace, cancellation);
226 if let Some(gate) = &approvals {
229 context.approvals = ToolApprovals::new(Arc::clone(gate), job.clone());
230 }
231 context.progress = progress;
232 let output = tool
233 .execute(arguments, context)
234 .await
235 .unwrap_or_else(ToolOutput::from);
236 finish(&jobs, &job, output, started.elapsed());
237 let _ = done_tx.send(true);
238 });
239 Ok(json!({
240 "job": id,
241 "tool": name,
242 "status": "running",
243 "background": true,
244 "note": "The agent is working in the background. SCV reports the result in a new \
245 turn when it finishes. agent_status shows its progress and agent_cancel \
246 stops it."
247 }))
248 }
249
250 async fn cancel(&self, job: &str, call_id: &str) -> Result<Value, ToolError> {
254 let mut done = {
255 let mut state = self.state();
256 let index = state
257 .jobs
258 .iter()
259 .position(|candidate| candidate.id == job)
260 .ok_or_else(|| unknown_job(job))?;
261 let entry = &mut state.jobs[index];
262 if entry.outcome.is_some() {
263 let (mut value, change) = entry.describe();
264 value["note"] = "The job had already finished.".into();
265 if let Some(change) = change {
266 state.record(call_id, change);
267 }
268 return Ok(value);
269 }
270 entry.cancelled = true;
271 entry.cancel.cancel();
272 let done = entry.done.clone();
273 if !entry.reported {
274 entry.reported = true;
275 let change = entry.change(JobStatus::Cancelled);
276 state.record(call_id, change);
277 }
278 done
279 };
280 let _ = tokio::time::timeout(CANCEL_SETTLE, done.wait_for(|finished| *finished)).await;
281 self.describe(Some(job), call_id)
282 }
283
284 async fn wait(
287 &self,
288 job: &str,
289 limit: Duration,
290 cancellation: &CancellationToken,
291 call_id: &str,
292 ) -> Result<Value, ToolError> {
293 let mut done = self
294 .state()
295 .jobs
296 .iter()
297 .find(|candidate| candidate.id == job)
298 .map(|candidate| candidate.done.clone())
299 .ok_or_else(|| unknown_job(job))?;
300 tokio::select! {
301 () = cancellation.cancelled() => return Err(ToolError::cancelled("wait cancelled")),
302 _ = tokio::time::timeout(limit, done.wait_for(|finished| *finished)) => {}
303 }
304 self.describe(Some(job), call_id)
305 }
306
307 fn describe(&self, job: Option<&str>, call_id: &str) -> Result<Value, ToolError> {
310 let mut state = self.state();
311 let mut changes = Vec::new();
312 let value = if let Some(job) = job {
313 let entry = state
314 .jobs
315 .iter_mut()
316 .find(|candidate| candidate.id == job)
317 .ok_or_else(|| unknown_job(job))?;
318 let (value, change) = entry.describe();
319 changes.extend(change);
320 value
321 } else {
322 let jobs: Vec<Value> = state
323 .jobs
324 .iter_mut()
325 .map(|entry| {
326 let (value, change) = entry.describe();
327 changes.extend(change);
328 value
329 })
330 .collect();
331 json!({ "jobs": jobs })
332 };
333 for change in changes {
334 state.record(call_id, change);
335 }
336 Ok(value)
337 }
338
339 pub fn take_changes(&self, call_id: &str) -> Vec<JobChange> {
342 let mut state = self.state();
343 let mut taken = Vec::new();
344 state.changes.retain(|(call, change)| {
345 if call == call_id {
346 taken.push(change.clone());
347 false
348 } else {
349 true
350 }
351 });
352 taken
353 }
354
355 pub fn take_unreported(&self) -> Vec<JobReport> {
358 let mut state = self.state();
359 state
360 .jobs
361 .iter_mut()
362 .filter(|job| job.outcome.is_some() && !job.reported)
363 .take(REPORT_MAX_JOBS)
364 .map(|job| {
365 job.reported = true;
366 let outcome = job.outcome.as_ref().expect("filtered on outcome");
367 let reply = AgentReply::read(&result_value(&outcome.output)).unwrap_or_default();
368 JobReport {
369 job: job.id.clone(),
370 tool: job.tool.clone(),
371 status: job_status(&outcome.output, &reply),
372 session: reply.session,
373 reply: bounded(
374 reply
375 .reply
376 .as_deref()
377 .unwrap_or(outcome.output.content.as_str()),
378 REPORT_REPLY_CHARS,
379 ),
380 }
381 })
382 .collect()
383 }
384
385 pub fn running(&self) -> usize {
387 self.state()
388 .jobs
389 .iter()
390 .filter(|job| job.outcome.is_none())
391 .count()
392 }
393}
394
395fn finish(jobs: &Weak<BackgroundJobs>, id: &str, output: ToolOutput, elapsed: Duration) {
396 let Some(jobs) = jobs.upgrade() else {
398 return;
399 };
400 {
401 let mut state = jobs.state();
402 if let Some(job) = state.jobs.iter_mut().find(|job| job.id == id) {
403 job.last_progress = job.progress.take().or(job.last_progress.take());
404 job.outcome = Some(Outcome { output, elapsed });
405 job.reported |= job.cancelled;
407 }
408 let finished = state
410 .jobs
411 .iter()
412 .filter(|job| job.outcome.is_some())
413 .count();
414 let mut excess = finished.saturating_sub(MAX_FINISHED);
415 state.jobs.retain(|job| {
416 if excess > 0 && job.outcome.is_some() && job.reported {
417 excess -= 1;
418 false
419 } else {
420 true
421 }
422 });
423 }
424 if let Some(finished) = &jobs.finished {
425 let _ = finished.send(());
426 }
427}
428
429impl Job {
430 fn change(&self, status: JobStatus) -> JobChange {
432 JobChange {
433 job: self.id.clone(),
434 tool: self.tool.clone(),
435 status,
436 task: self.task.clone(),
437 }
438 }
439
440 fn describe(&mut self) -> (Value, Option<JobChange>) {
443 if let Some(line) = self.progress.take() {
444 self.last_progress = Some(line);
445 }
446 let mut change = None;
447 let mut value = Map::new();
448 value.insert("job".into(), self.id.clone().into());
449 value.insert("tool".into(), self.tool.clone().into());
450 match &self.outcome {
451 None => {
452 value.insert("status".into(), "running".into());
453 value.insert(
454 "elapsed_seconds".into(),
455 self.started.elapsed().as_secs().into(),
456 );
457 if let Some(progress) = &self.last_progress {
458 value.insert("progress".into(), progress.clone().into());
459 }
460 }
461 Some(outcome) => {
462 let result = result_value(&outcome.output);
463 let status = if self.cancelled {
464 JobStatus::Cancelled
465 } else {
466 job_status(
467 &outcome.output,
468 &AgentReply::read(&result).unwrap_or_default(),
469 )
470 };
471 value.insert("status".into(), status.as_str().into());
472 value.insert("elapsed_seconds".into(), outcome.elapsed.as_secs().into());
473 value.insert("result".into(), result);
474 if !self.reported {
475 self.reported = true;
476 change = Some(self.change(status));
477 }
478 }
479 }
480 (Value::Object(value), change)
481 }
482}
483
484fn result_value(output: &ToolOutput) -> Value {
486 match serde_json::from_str::<Value>(&output.content) {
487 Ok(value @ Value::Object(_)) => value,
488 _ => json!({ "reply": output.content, "is_error": output.is_error() }),
489 }
490}
491
492fn job_status(output: &ToolOutput, reply: &AgentReply) -> JobStatus {
495 reply.status.unwrap_or(if output.is_error() {
496 JobStatus::Failed
497 } else {
498 JobStatus::Completed
499 })
500}
501
502fn task_line(prompt: &str) -> String {
505 let line = prompt
506 .lines()
507 .map(str::trim)
508 .find(|line| !line.is_empty())
509 .unwrap_or_default();
510 let mut task: String = line
511 .chars()
512 .filter(|character| !character.is_control())
513 .take(TASK_CHARS)
514 .collect();
515 if line.chars().count() > TASK_CHARS {
516 task.push('…');
517 }
518 task
519}
520
521fn unknown_job(job: &str) -> ToolError {
522 ToolError::invalid_arguments(format!(
523 "unknown background job {:?}; agent_status lists this session's jobs",
524 bounded(job, 64)
525 ))
526}
527
528pub fn report_prompt(reports: &[JobReport]) -> String {
530 let mut prompt = String::from(
531 "[SCV background report] Delegated work you started in the background has \
532 finished. The user did not send this message: tell them briefly what \
533 happened and the key result.\n",
534 );
535 for report in reports {
536 prompt.push_str(&format!(
537 "\n{} ({}{}): {}\n{}\n",
538 report.job,
539 report.tool,
540 report
541 .session
542 .as_deref()
543 .map_or_else(String::new, |session| format!(", conversation {session}")),
544 report.status,
545 report.reply.trim()
546 ));
547 }
548 prompt
549}
550
551pub(crate) struct BackgroundCapable {
553 pub(crate) inner: Arc<dyn Tool>,
554 pub(crate) jobs: Arc<BackgroundJobs>,
555}
556
557fn split_background(arguments: &Value) -> (Value, bool) {
559 let mut arguments = arguments.clone();
560 let background = arguments
561 .as_object_mut()
562 .and_then(|object| object.remove("background"))
563 .is_some_and(|value| value.as_bool() == Some(true));
564 (arguments, background)
565}
566
567#[async_trait]
568impl Tool for BackgroundCapable {
569 fn spec(&self) -> ToolSpec {
570 let mut spec = self.inner.spec();
571 if let Some(properties) = spec
572 .parameters
573 .get_mut("properties")
574 .and_then(Value::as_object_mut)
575 {
576 properties.insert(
577 "background".into(),
578 json!({
579 "type":"boolean",
580 "description":"Run in the background: the call returns a job handle at once, \
581 the user can keep talking to you while the agent works, and SCV reports \
582 the result in a new turn when it finishes. Use it for any substantial \
583 task; run in the foreground only for quick work whose result you need \
584 within this turn."
585 }),
586 );
587 }
588 spec.description.push_str(&format!(
589 " Set background to true for anything beyond a quick task: the call returns a job \
590 handle at once (at most {} running per session), SCV reports the result when the \
591 job finishes, agent_status shows progress, and agent_cancel stops it. A background \
592 job's own approval requests get only the answer this session would give without \
593 asking a person.",
594 self.jobs.limit()
595 ));
596 spec
597 }
598
599 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
600 self.inner.risk(&split_background(arguments).0)
601 }
602
603 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
604 let (arguments, background) = split_background(arguments);
605 let mut summary = self.inner.approval_summary(&arguments)?;
606 if background {
607 summary.push_str(
608 " Runs in the background: the call returns at once and the result is \
609 reported when the agent finishes.",
610 );
611 }
612 Ok(summary)
613 }
614
615 async fn execute(
616 &self,
617 arguments: Value,
618 context: ToolContext,
619 ) -> Result<ToolOutput, ToolError> {
620 let (arguments, background) = split_background(&arguments);
621 if !background {
622 return self.inner.execute(arguments, context).await;
623 }
624 self.inner.risk(&arguments)?;
626 let name = self.inner.spec().name;
627 let started = self.jobs.start(
628 Arc::clone(&self.inner),
629 &name,
630 arguments,
631 context.workspace,
632 &context.call_id,
633 )?;
634 Ok(ToolOutput::success(started.to_string()))
635 }
636}
637
638#[derive(Deserialize)]
639#[serde(deny_unknown_fields)]
640struct WaitArgs {
641 job: String,
642 timeout_seconds: Option<u64>,
643}
644
645#[derive(Deserialize)]
646#[serde(deny_unknown_fields)]
647struct StatusArgs {
648 job: Option<String>,
649}
650
651pub(crate) struct WaitTool {
653 pub(crate) jobs: Arc<BackgroundJobs>,
654 pub(crate) timeouts: Timeouts,
655}
656
657#[async_trait]
658impl Tool for WaitTool {
659 fn spec(&self) -> ToolSpec {
660 let mut timeout = timeout_schema(self.timeouts);
661 timeout["description"] = format!(
662 "Seconds to wait before returning the job still running. Defaults to {}; at most {}.",
663 self.timeouts.default.min(self.timeouts.max).as_secs(),
664 self.timeouts.max.as_secs()
665 )
666 .into();
667 ToolSpec {
668 name: "agent_wait".into(),
669 description: "Wait for a background agent job (the `job` handle an agent_* call with \
670 background: true returned) to finish, and return its result. Returns early with \
671 status running when the timeout passes. Waiting holds your turn open, so the \
672 user cannot reach you meanwhile; usually let SCV report the result instead."
673 .into(),
674 parameters: json!({
675 "type":"object",
676 "properties":{"job":{"type":"string"},"timeout_seconds":timeout},
677 "required":["job"],
678 "additionalProperties":false
679 }),
680 }
681 }
682
683 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
684 let args: WaitArgs = parse_args(arguments)?;
685 self.timeouts.resolve(args.timeout_seconds)?;
686 Ok(ToolRisk::ReadOnly)
687 }
688
689 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
690 let args: WaitArgs = parse_args(arguments)?;
691 Ok(format!(
692 "Wait for background job {}",
693 bounded(&args.job, 64)
694 ))
695 }
696
697 async fn execute(
698 &self,
699 arguments: Value,
700 context: ToolContext,
701 ) -> Result<ToolOutput, ToolError> {
702 let args: WaitArgs = parse_args(&arguments)?;
703 let limit = self.timeouts.resolve(args.timeout_seconds)?;
704 let value = self
705 .jobs
706 .wait(&args.job, limit, &context.cancellation, &context.call_id)
707 .await?;
708 Ok(ToolOutput::success(value.to_string()))
709 }
710}
711
712pub(crate) struct StatusTool {
714 pub(crate) jobs: Arc<BackgroundJobs>,
715}
716
717#[async_trait]
718impl Tool for StatusTool {
719 fn spec(&self) -> ToolSpec {
720 ToolSpec {
721 name: "agent_status".into(),
722 description: "Show this session's background agent jobs: running ones with their \
723 latest progress, finished ones with their result. Pass job for one job."
724 .into(),
725 parameters: json!({
726 "type":"object",
727 "properties":{"job":{"type":"string"}},
728 "additionalProperties":false
729 }),
730 }
731 }
732
733 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
734 let _: StatusArgs = parse_args(arguments)?;
735 Ok(ToolRisk::ReadOnly)
736 }
737
738 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
739 let args: StatusArgs = parse_args(arguments)?;
740 Ok(args.job.map_or_else(
741 || "List background jobs".into(),
742 |job| format!("Show background job {}", bounded(&job, 64)),
743 ))
744 }
745
746 async fn execute(
747 &self,
748 arguments: Value,
749 context: ToolContext,
750 ) -> Result<ToolOutput, ToolError> {
751 let args: StatusArgs = parse_args(&arguments)?;
752 let value = self.jobs.describe(args.job.as_deref(), &context.call_id)?;
753 Ok(ToolOutput::success(value.to_string()))
754 }
755}
756
757#[derive(Deserialize)]
758#[serde(deny_unknown_fields)]
759struct CancelArgs {
760 job: String,
761}
762
763pub(crate) struct CancelTool {
765 pub(crate) jobs: Arc<BackgroundJobs>,
766}
767
768#[async_trait]
769impl Tool for CancelTool {
770 fn spec(&self) -> ToolSpec {
771 ToolSpec {
772 name: "agent_cancel".into(),
773 description: "Stop a running background agent job (the `job` handle an agent_* call \
774 with background: true returned), for example when the user no longer wants \
775 it. The agent and every process it started are stopped; work it already wrote \
776 stays. Returns the job with status cancelled, and no report turn follows."
777 .into(),
778 parameters: json!({
779 "type":"object",
780 "properties":{"job":{"type":"string"}},
781 "required":["job"],
782 "additionalProperties":false
783 }),
784 }
785 }
786
787 fn risk(&self, arguments: &Value) -> Result<ToolRisk, ToolError> {
788 let _: CancelArgs = parse_args(arguments)?;
789 Ok(ToolRisk::Process)
790 }
791
792 fn approval_summary(&self, arguments: &Value) -> Result<String, ToolError> {
793 let args: CancelArgs = parse_args(arguments)?;
794 Ok(format!("Stop background job {}", bounded(&args.job, 64)))
795 }
796
797 async fn execute(
798 &self,
799 arguments: Value,
800 context: ToolContext,
801 ) -> Result<ToolOutput, ToolError> {
802 let args: CancelArgs = parse_args(&arguments)?;
803 let value = self.jobs.cancel(&args.job, &context.call_id).await?;
804 Ok(ToolOutput::success(value.to_string()))
805 }
806}
807
808#[cfg(test)]
809mod tests;