Skip to main content

vv_agent/runtime/sub_task_manager/
sessions.rs

1use std::collections::BTreeMap;
2use std::panic::{catch_unwind, AssertUnwindSafe};
3use std::sync::Arc;
4use std::thread;
5
6use crate::runtime::sub_agent_sessions::{
7    register_sub_agent_session, sub_agent_session_registry, SubAgentSession,
8    SubAgentSessionListener,
9};
10use crate::tools::common::trim_portable_whitespace;
11use crate::types::{AgentStatus, SubTaskOutcome};
12use crate::workspace::WorkspaceBackend;
13
14use super::helpers::{normalize_failed_outcome, now_iso, panic_payload_to_string};
15use super::manager::SubTaskManager;
16use super::record::{ManagedSubAgentSession, ManagedSubTask};
17use super::types::{SubTaskLineage, SubTaskSessionAttachment, SubTaskTurnSnapshot};
18
19impl SubTaskManager {
20    pub fn attach_session(
21        &self,
22        task_id: impl Into<String>,
23        session_id: impl Into<String>,
24        agent_name: impl Into<String>,
25        task_title: impl Into<String>,
26        workspace_backend: Arc<dyn WorkspaceBackend>,
27        session: Arc<dyn SubAgentSession>,
28    ) {
29        self.attach_session_with_resolved(SubTaskSessionAttachment {
30            task_id: task_id.into(),
31            session_id: session_id.into(),
32            agent_name: agent_name.into(),
33            task_title: task_title.into(),
34            workspace_backend,
35            session,
36            resolved: BTreeMap::new(),
37        });
38    }
39
40    pub fn attach_session_with_resolved(&self, attachment: SubTaskSessionAttachment) {
41        self.attach_session_with_resolved_and_lineage(attachment, SubTaskLineage::default());
42    }
43
44    pub fn attach_session_with_resolved_and_lineage(
45        &self,
46        attachment: SubTaskSessionAttachment,
47        lineage: SubTaskLineage,
48    ) {
49        self.attach_session_inner(attachment, lineage, false);
50    }
51
52    pub(crate) fn attach_running_session_with_resolved_and_lineage(
53        &self,
54        attachment: SubTaskSessionAttachment,
55        lineage: SubTaskLineage,
56    ) {
57        self.attach_session_inner(attachment, lineage, true);
58    }
59
60    fn attach_session_inner(
61        &self,
62        attachment: SubTaskSessionAttachment,
63        lineage: SubTaskLineage,
64        running: bool,
65    ) {
66        let SubTaskSessionAttachment {
67            task_id,
68            session_id,
69            agent_name,
70            task_title,
71            workspace_backend,
72            session,
73            resolved,
74        } = attachment;
75        let listener_generation = {
76            let mut tasks = self.tasks.lock().expect("sub-task manager poisoned");
77            match tasks.get_mut(&task_id) {
78                Some(record) => {
79                    let initial_lineage = record
80                        .session
81                        .as_ref()
82                        .map(|attached| attached.initial_lineage.clone())
83                        .unwrap_or_else(|| SubTaskLineage {
84                            parent_run_id: lineage
85                                .parent_run_id
86                                .clone()
87                                .or_else(|| record.parent_run_id.clone()),
88                            parent_tool_call_id: lineage
89                                .parent_tool_call_id
90                                .clone()
91                                .or_else(|| record.parent_tool_call_id.clone()),
92                        });
93                    let session_changed = record
94                        .session
95                        .as_ref()
96                        .is_none_or(|attached| !Arc::ptr_eq(&attached.session, &session));
97                    if session_changed {
98                        record.session_generation =
99                            record.session_generation.saturating_add(1).max(1);
100                        record.manager_listener_generation = None;
101                    } else if record.session_generation == 0 {
102                        record.session_generation = 1;
103                    }
104                    record.session_id = session_id;
105                    record.agent_name = agent_name;
106                    if !task_title.is_empty() {
107                        record.task_title = task_title;
108                    }
109                    record.workspace_backend = Some(workspace_backend);
110                    record.session = Some(ManagedSubAgentSession {
111                        session: session.clone(),
112                        initial_lineage,
113                    });
114                    if !resolved.is_empty() {
115                        record.resolved = resolved;
116                    }
117                    if lineage.parent_run_id.is_some() {
118                        record.parent_run_id = lineage.parent_run_id;
119                    }
120                    if lineage.parent_tool_call_id.is_some() {
121                        record.parent_tool_call_id = lineage.parent_tool_call_id;
122                    }
123                    if running {
124                        record.running = true;
125                        record.outcome = None;
126                    }
127                    record.updated_at = now_iso();
128                    (record.manager_listener_generation != Some(record.session_generation))
129                        .then_some(record.session_generation)
130                }
131                None => {
132                    tasks.insert(
133                        task_id.clone(),
134                        ManagedSubTask {
135                            task_id: task_id.clone(),
136                            session_id,
137                            agent_name,
138                            task_title,
139                            workspace_backend: Some(workspace_backend),
140                            session: Some(ManagedSubAgentSession {
141                                session: session.clone(),
142                                initial_lineage: lineage.clone(),
143                            }),
144                            outcome: None,
145                            resolved,
146                            current_cycle_index: None,
147                            recent_activity: None,
148                            latest_cycle: None,
149                            latest_tool_call: None,
150                            parent_run_id: lineage.parent_run_id,
151                            parent_tool_call_id: lineage.parent_tool_call_id,
152                            running,
153                            worker_owned_run: false,
154                            handle: None,
155                            updated_at: now_iso(),
156                            session_generation: 1,
157                            manager_listener_generation: None,
158                        },
159                    );
160                    Some(1)
161                }
162            }
163        };
164
165        if let Some(session_generation) = listener_generation {
166            let tasks = Arc::downgrade(&self.tasks);
167            let listener_task_id = task_id.clone();
168            let listener: SubAgentSessionListener = Arc::new(move |event, payload| {
169                let Some(tasks) = tasks.upgrade() else {
170                    return;
171                };
172                SubTaskManager { tasks }.handle_session_event(
173                    &listener_task_id,
174                    session_generation,
175                    event,
176                    payload,
177                );
178            });
179            let _ = session.subscribe(listener);
180            let mut tasks = self.tasks.lock().expect("sub-task manager poisoned");
181            if let Some(record) = tasks.get_mut(&task_id) {
182                let session_is_current = record
183                    .session
184                    .as_ref()
185                    .is_some_and(|attached| Arc::ptr_eq(&attached.session, &session));
186                if session_is_current && record.session_generation == session_generation {
187                    record.manager_listener_generation = Some(session_generation);
188                }
189            }
190        }
191    }
192
193    pub fn continue_task(&self, task_id: &str, prompt: &str) -> Result<(), String> {
194        self.continue_task_inner(task_id, prompt, None)
195    }
196
197    pub(crate) fn continue_task_with_snapshot(
198        &self,
199        task_id: &str,
200        prompt: &str,
201        snapshot: SubTaskTurnSnapshot,
202    ) -> Result<(), String> {
203        self.continue_task_inner(task_id, prompt, Some(snapshot))
204    }
205
206    fn continue_task_inner(
207        &self,
208        task_id: &str,
209        prompt: &str,
210        snapshot: Option<SubTaskTurnSnapshot>,
211    ) -> Result<(), String> {
212        let prompt = trim_portable_whitespace(prompt);
213        if prompt.is_empty() {
214            return Err("Follow-up prompt cannot be empty.".to_string());
215        }
216
217        let (
218            session_id,
219            agent_name,
220            session,
221            resolved,
222            session_generation,
223            previous,
224            mut registration,
225        ) = {
226            let mut tasks = self.tasks.lock().expect("sub-task manager poisoned");
227            let Some(record) = tasks.get_mut(task_id) else {
228                return Err(format!("Sub-task {task_id} not found."));
229            };
230            if record.is_running() {
231                return Err(format!("Sub-task {task_id} is already running."));
232            }
233            if record
234                .outcome
235                .as_ref()
236                .is_some_and(|outcome| outcome.status == AgentStatus::MaxCycles)
237            {
238                return Err(format!(
239                    "Sub-task {task_id} reached max cycles and cannot continue."
240                ));
241            }
242            if record.session_id.trim().is_empty() {
243                return Err(format!("Sub-task {task_id} session is not attached."));
244            }
245            let Some(session) = record.session.clone() else {
246                return Err(format!("Sub-task {task_id} session is not attached."));
247            };
248            let continuation_lineage = snapshot
249                .as_ref()
250                .map(|snapshot| SubTaskLineage {
251                    parent_run_id: snapshot.parent_run_id.clone(),
252                    parent_tool_call_id: snapshot.parent_tool_call_id.clone(),
253                })
254                .unwrap_or_else(|| session.initial_lineage.clone());
255            let registration = ContinuationRegistration::register(
256                record.session_id.clone(),
257                session.session.clone(),
258            );
259            let previous = record.admit_continuation(prompt, &continuation_lineage, now_iso());
260            (
261                record.session_id.clone(),
262                record.agent_name.clone(),
263                session,
264                record.resolved.clone(),
265                record.session_generation,
266                previous,
267                registration,
268            )
269        };
270
271        if let Err(payload) = catch_unwind(AssertUnwindSafe(|| {
272            session.session.sanitize_for_resume();
273        })) {
274            let error = panic_payload_to_string(payload.as_ref());
275            let mut tasks = self.tasks.lock().expect("sub-task manager poisoned");
276            if let Some(record) = tasks.get_mut(task_id) {
277                record.rollback_continuation(previous);
278            }
279            return Err(format!(
280                "Sub-task {task_id} continuation setup failed: {error}"
281            ));
282        }
283
284        let tasks = self.tasks.clone();
285        let task_id_for_thread = task_id.to_string();
286        let prompt_for_thread = prompt.to_string();
287        let session_id_for_thread = session_id.clone();
288        let agent_name_for_thread = agent_name.clone();
289        let session_generation_for_thread = session_generation;
290        let turn_event_handler = snapshot
291            .as_ref()
292            .and_then(|snapshot| snapshot.event_handler.clone());
293        let session_for_thread = session.session.clone();
294        let mut task_records = self.tasks.lock().expect("sub-task manager poisoned");
295        let spawn_result = catch_unwind(AssertUnwindSafe(|| {
296            thread::Builder::new()
297                .name(format!("vv-agent-sub-task-{session_id_for_thread}"))
298                .spawn(move || {
299                    let _event_handler_scope =
300                        SubTaskTurnSnapshot::enter_event_handler_scope(turn_event_handler);
301                    let worker_result = catch_unwind(AssertUnwindSafe(|| match snapshot {
302                        Some(snapshot) => session_for_thread
303                            .continue_run_with_snapshot(&prompt_for_thread, snapshot),
304                        None => session_for_thread.continue_run(&prompt_for_thread),
305                    }));
306                    let outcome = normalize_failed_outcome(match worker_result {
307                        Ok(Ok(outcome)) => outcome,
308                        Ok(Err(error)) => SubTaskOutcome {
309                            task_id: task_id_for_thread.clone(),
310                            agent_name: agent_name_for_thread.clone(),
311                            status: AgentStatus::Failed,
312                            session_id: Some(session_id_for_thread.clone()),
313                            final_answer: None,
314                            wait_reason: None,
315                            error: Some(error),
316                            error_code: Some("sub_task_failed".to_string()),
317                            completion_reason: Some(crate::types::CompletionReason::Failed),
318                            completion_tool_name: None,
319                            partial_output: None,
320                            cycles: 0,
321                            todo_list: Vec::new(),
322                            resolved: resolved.clone(),
323                        },
324                        Err(payload) => SubTaskOutcome {
325                            task_id: task_id_for_thread.clone(),
326                            agent_name: agent_name_for_thread.clone(),
327                            status: AgentStatus::Failed,
328                            session_id: Some(session_id_for_thread.clone()),
329                            final_answer: None,
330                            wait_reason: None,
331                            error: Some(panic_payload_to_string(payload.as_ref())),
332                            error_code: Some("sub_task_failed".to_string()),
333                            completion_reason: Some(crate::types::CompletionReason::Failed),
334                            completion_tool_name: None,
335                            partial_output: None,
336                            cycles: 0,
337                            todo_list: Vec::new(),
338                            resolved,
339                        },
340                    });
341                    {
342                        let mut tasks = tasks
343                            .lock()
344                            .unwrap_or_else(|poisoned| poisoned.into_inner());
345                        if let Some(record) = tasks.get_mut(&task_id_for_thread) {
346                            let owns_session_generation = record.session_generation
347                                == session_generation_for_thread
348                                && record.session.as_ref().is_some_and(|attached| {
349                                    Arc::ptr_eq(&attached.session, &session_for_thread)
350                                });
351                            let mut outcome = outcome;
352                            if outcome.resolved.is_empty() && !record.resolved.is_empty() {
353                                outcome.resolved = record.resolved.clone();
354                            }
355                            if owns_session_generation {
356                                record.session_id = outcome
357                                    .session_id
358                                    .clone()
359                                    .unwrap_or_else(|| record.session_id.clone());
360                                record.agent_name = outcome.agent_name.clone();
361                                record.update_from_outcome(&outcome);
362                                record.outcome = Some(outcome);
363                                record.running = false;
364                                record.worker_owned_run = false;
365                                record.updated_at = now_iso();
366                            }
367                        }
368                    }
369                    sub_agent_session_registry()
370                        .unregister_if_matches(&session_id_for_thread, Some(&session_for_thread));
371                })
372        }));
373        let handle = match spawn_result {
374            Ok(Ok(handle)) => handle,
375            Ok(Err(error)) => {
376                if let Some(record) = task_records.get_mut(task_id) {
377                    record.rollback_continuation(previous);
378                }
379                return Err(format!(
380                    "Sub-task {task_id} continuation thread failed to spawn: {error}"
381                ));
382            }
383            Err(payload) => {
384                let error = panic_payload_to_string(payload.as_ref());
385                if let Some(record) = task_records.get_mut(task_id) {
386                    record.rollback_continuation(previous);
387                }
388                return Err(format!(
389                    "Sub-task {task_id} continuation thread failed to spawn: {error}"
390                ));
391            }
392        };
393
394        if let Some(record) = task_records.get_mut(task_id) {
395            record.handle = Some(handle);
396            record.updated_at = now_iso();
397        }
398        registration.commit();
399        Ok(())
400    }
401}
402
403struct ContinuationRegistration {
404    session_id: String,
405    session: Arc<dyn SubAgentSession>,
406    previous: Option<Arc<dyn SubAgentSession>>,
407    committed: bool,
408}
409
410impl ContinuationRegistration {
411    fn register(session_id: String, session: Arc<dyn SubAgentSession>) -> Self {
412        let previous = sub_agent_session_registry().get(&session_id);
413        register_sub_agent_session(session_id.clone(), session.clone());
414        Self {
415            session_id,
416            session,
417            previous,
418            committed: false,
419        }
420    }
421
422    fn commit(&mut self) {
423        self.committed = true;
424    }
425}
426
427impl Drop for ContinuationRegistration {
428    fn drop(&mut self) {
429        if self.committed {
430            return;
431        }
432        let removed = sub_agent_session_registry()
433            .unregister_if_matches(&self.session_id, Some(&self.session));
434        if removed {
435            if let Some(previous) = self.previous.clone() {
436                register_sub_agent_session(self.session_id.clone(), previous);
437            }
438        }
439    }
440}