Skip to main content

vv_agent/runtime/
background_sessions.rs

1mod listeners;
2mod options;
3mod session;
4mod subscription;
5#[cfg(test)]
6mod tests;
7
8use std::collections::BTreeMap;
9use std::path::PathBuf;
10use std::sync::atomic::{AtomicU64, Ordering};
11use std::sync::{Mutex, OnceLock};
12use std::thread;
13use std::time::Duration;
14
15use serde_json::{json, Value};
16
17use crate::runtime::processes::start_captured_process_with_env;
18use crate::workspace::WorkspaceBackend;
19
20use listeners::notify_background_listeners;
21use session::BackgroundSession;
22
23pub use listeners::BackgroundSessionListener;
24pub use options::{BackgroundSessionAdoptOptions, BackgroundSessionStartOptions};
25pub use subscription::BackgroundSessionSubscription;
26
27static MANAGER: OnceLock<BackgroundSessionManager> = OnceLock::new();
28
29pub fn background_session_manager() -> &'static BackgroundSessionManager {
30    MANAGER.get_or_init(BackgroundSessionManager::default)
31}
32
33#[derive(Default)]
34pub struct BackgroundSessionManager {
35    sessions: Mutex<BTreeMap<String, BackgroundSession>>,
36    next_id: AtomicU64,
37    next_listener_id: AtomicU64,
38}
39
40impl BackgroundSessionManager {
41    pub fn start(
42        &self,
43        command: impl Into<String>,
44        cwd: impl Into<PathBuf>,
45        timeout_seconds: u64,
46        options: BackgroundSessionStartOptions,
47    ) -> Result<String, String> {
48        let command = command.into();
49        let cwd = cwd.into();
50        let prepared = super::shell::prepare_shell_execution(
51            &command,
52            options.auto_confirm,
53            options.stdin.as_deref(),
54            options.shell.as_deref(),
55            options.windows_shell_priority.as_deref(),
56        )?;
57        let started = start_captured_process_with_env(
58            &prepared.command,
59            &cwd,
60            prepared.stdin.as_deref(),
61            options.env.as_ref(),
62        )
63        .map_err(|error| error.to_string())?;
64        Ok(self.adopt_running_process(
65            command,
66            cwd,
67            timeout_seconds,
68            started.child,
69            started.output_path,
70            prepared.shell,
71        ))
72    }
73
74    pub fn adopt_running_process(
75        &self,
76        command: impl Into<String>,
77        cwd: impl Into<PathBuf>,
78        timeout_seconds: u64,
79        child: std::process::Child,
80        output_path: PathBuf,
81        shell: Option<String>,
82    ) -> String {
83        let mut options =
84            BackgroundSessionAdoptOptions::new(command, cwd, timeout_seconds, child, output_path);
85        options.shell = shell;
86        self.adopt_running_process_with_options(options)
87    }
88
89    pub fn adopt_running_process_with_options(
90        &self,
91        options: BackgroundSessionAdoptOptions,
92    ) -> String {
93        let id = self.next_id.fetch_add(1, Ordering::Relaxed) + 1;
94        let session_id = format!("bg_{id:012x}");
95        let session = BackgroundSession::from_adopt_options(session_id.clone(), options);
96        self.sessions
97            .lock()
98            .expect("background session manager poisoned")
99            .insert(session_id.clone(), session);
100        self.start_watch_thread(session_id.clone());
101        session_id
102    }
103
104    pub fn subscribe(
105        &'static self,
106        session_id: &str,
107        listener: BackgroundSessionListener,
108    ) -> BackgroundSessionSubscription {
109        let mut snapshot = None;
110        let listener_id = self.next_listener_id.fetch_add(1, Ordering::Relaxed) + 1;
111        {
112            let mut sessions = self
113                .sessions
114                .lock()
115                .expect("background session manager poisoned");
116            let Some(session) = sessions.get_mut(session_id) else {
117                return BackgroundSessionSubscription::noop();
118            };
119            if session.is_terminal() {
120                snapshot = Some(session.snapshot());
121            } else {
122                session.add_listener(listener_id, listener.clone());
123            }
124        }
125        if let Some(payload) = snapshot {
126            listener(&payload);
127            return BackgroundSessionSubscription::noop();
128        }
129        BackgroundSessionSubscription::new(session_id.to_string(), listener_id, self)
130    }
131
132    fn unsubscribe(&self, session_id: &str, listener_id: u64) {
133        let mut sessions = self
134            .sessions
135            .lock()
136            .expect("background session manager poisoned");
137        if let Some(session) = sessions.get_mut(session_id) {
138            session.remove_listener(listener_id);
139        }
140    }
141
142    fn start_watch_thread(&self, session_id: String) {
143        let thread_name = format!("vv-agent-bg-{session_id}");
144        let _ = thread::Builder::new()
145            .name(thread_name)
146            .spawn(move || loop {
147                thread::sleep(Duration::from_millis(200));
148                let payload = background_session_manager().check(&session_id);
149                let status = payload
150                    .get("status")
151                    .and_then(Value::as_str)
152                    .unwrap_or("missing");
153                if status != "running" {
154                    break;
155                }
156            });
157    }
158
159    pub fn check(&self, session_id: &str) -> Value {
160        let mut sessions = self
161            .sessions
162            .lock()
163            .expect("background session manager poisoned");
164        let Some(session) = sessions.get_mut(session_id) else {
165            return json!({
166                "status": "missing",
167                "session_id": session_id,
168                "error": "Background session not found",
169            });
170        };
171
172        if session.is_terminal() {
173            return session.snapshot();
174        }
175
176        let elapsed = session.elapsed();
177        if session.timed_out(elapsed) {
178            session.finalize_timeout();
179            let payload = session.snapshot();
180            let terminal_listeners = session.take_listeners();
181            drop(sessions);
182            notify_background_listeners(terminal_listeners, &payload);
183            return payload;
184        }
185
186        match session.try_wait() {
187            Ok(Some(exit_code)) => {
188                session.finalize_completed(exit_code);
189                let payload = session.snapshot();
190                let terminal_listeners = session.take_listeners();
191                drop(sessions);
192                notify_background_listeners(terminal_listeners, &payload);
193                payload
194            }
195            Ok(None) => session.running_snapshot(elapsed),
196            Err(error) => {
197                session.finalize_failed_with_output(-1, error.to_string());
198                let payload = session.snapshot();
199                let terminal_listeners = session.take_listeners();
200                drop(sessions);
201                notify_background_listeners(terminal_listeners, &payload);
202                payload
203            }
204        }
205    }
206
207    pub(crate) fn check_for_tool(
208        &self,
209        session_id: &str,
210        fallback_backend: std::sync::Arc<dyn WorkspaceBackend>,
211        fallback_task_id: &str,
212        fallback_tool_call_id: &str,
213    ) -> Value {
214        let payload = self.check(session_id);
215        if payload.get("status").and_then(Value::as_str) == Some("running") {
216            return payload;
217        }
218        let mut sessions = self
219            .sessions
220            .lock()
221            .expect("background session manager poisoned");
222        let Some(session) = sessions.get_mut(session_id) else {
223            return payload;
224        };
225        session.ensure_artifact(fallback_backend, fallback_task_id, fallback_tool_call_id);
226        session.snapshot()
227    }
228}