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}