Skip to main content

roder_core/runtime/
turn_start.rs

1use super::*;
2
3impl Runtime {
4    pub fn start_turn(
5        self: &Arc<Self>,
6        req: StartTurnRequest,
7    ) -> BoxFuture<'_, anyhow::Result<TurnId>> {
8        Box::pin(async move {
9            self.start_turn_admitted(req, false)
10                .await?
11                .ok_or_else(|| anyhow::anyhow!("explicit turn was not admitted"))
12        })
13    }
14
15    pub(crate) fn start_goal_turn_if_idle(
16        self: &Arc<Self>,
17        req: StartTurnRequest,
18    ) -> BoxFuture<'_, anyhow::Result<Option<TurnId>>> {
19        self.start_turn_admitted(req, true)
20    }
21
22    fn start_turn_admitted(
23        self: &Arc<Self>,
24        mut req: StartTurnRequest,
25        goal_only: bool,
26    ) -> BoxFuture<'_, anyhow::Result<Option<TurnId>>> {
27        Box::pin(async move {
28            let _thread_admission = self.thread_admission(&req.thread_id).await;
29            let turn_admission = self.turn_admission.lock().await;
30            let _goal_admission = if goal_only {
31                if self.has_active_turn_for_thread(&req.thread_id).await {
32                    return Ok(None);
33                }
34                let Some((guard, goal)) = self.goals.admit_continuation(&req.thread_id).await?
35                else {
36                    return Ok(None);
37                };
38                if let Some(previous) = self.goals.continuation_request(&req.thread_id).await {
39                    req = previous;
40                }
41                req.message = crate::goals::continuation_prompt(&goal);
42                Some(guard)
43            } else {
44                None
45            };
46            self.ensure_execution_authority()?;
47            anyhow::ensure!(
48                self.accepting_turns.load(Ordering::Acquire),
49                "runtime is quiescing and cannot accept new turns"
50            );
51            req.workspace = validate_thread_workspace(&req.workspace)?;
52            let team_member = self.teams.member_for_thread(&req.thread_id).await;
53            let cfg = self.config.read().await.clone();
54            let provider = req
55                .provider_override
56                .clone()
57                .unwrap_or_else(|| cfg.default_provider.clone());
58            self.engine_for(&provider)?;
59            let turn_id = uuid::Uuid::new_v4().to_string();
60            let mut initial_mailbox_ack = None;
61            if let Some((team_id, member)) = &team_member {
62                let pending = self
63                    .teams
64                    .reserve_pending_mailbox_messages(team_id, &member.id, &turn_id)
65                    .await?;
66                if !pending.is_empty()
67                    && let Some(team) = self.read_team(team_id).await
68                {
69                    let mailbox = format_mailbox_messages(&team, &pending);
70                    req.message = if req.message.trim().is_empty() {
71                        mailbox
72                    } else {
73                        format!("{mailbox}\n\n[Direct task input]\n{}", req.message)
74                    };
75                    initial_mailbox_ack = Some(MailboxDeliveryAck {
76                        team_id: team_id.clone(),
77                        message_ids: pending.iter().map(|message| message.id.clone()).collect(),
78                    });
79                }
80            }
81            let (abort_handle, abort_registration) = AbortHandle::new_pair();
82            let drain = Arc::new(TurnDrainHandle {
83                thread_id: req.thread_id.clone(),
84                interrupt_requested: AtomicBool::new(false),
85                interrupt_reason: Mutex::new(None),
86                completed: AtomicBool::new(false),
87                completed_notify: Notify::new(),
88            });
89            let (steer_changed, steering) = tokio::sync::watch::channel(0);
90            let active = ActiveTurnHandle {
91                thread_id: req.thread_id.clone(),
92                abort: abort_handle,
93                steers: Arc::new(Mutex::new(Vec::new())),
94                steer_changed,
95                drain,
96            };
97            self.active_turns
98                .write()
99                .await
100                .insert(turn_id.clone(), active);
101            self.record_turn_lifecycle(
102                req.thread_id.clone(),
103                turn_id.clone(),
104                TurnLifecycleState::Running,
105                TurnCleanupState::NotRequested,
106                None,
107            )
108            .await;
109            self.active_turn_contexts.write().await.insert(
110                turn_id.clone(),
111                InheritedTurnContext {
112                    workspace: req.workspace.clone(),
113                    instructions: req.instructions.clone(),
114                    developer_context: req.developer_context.clone(),
115                },
116            );
117            if let Some((team_id, member)) = team_member {
118                let updated = match self
119                    .teams
120                    .update_member(&team_id, &member.id, |member| {
121                        member.current_turn_id = Some(turn_id.clone());
122                        member.status = TeamMemberStatus::Running;
123                        member.final_message = None;
124                        member.terminal_error = None;
125                    })
126                    .await
127                {
128                    Ok(updated) => updated,
129                    Err(error) => {
130                        self.active_turns.write().await.remove(&turn_id);
131                        self.active_turn_contexts.write().await.remove(&turn_id);
132                        self.teams
133                            .release_mailbox_reservations_for_turn(&turn_id)
134                            .await;
135                        return Err(error);
136                    }
137                };
138                if let Some(member) = updated
139                    .members
140                    .into_iter()
141                    .find(|candidate| candidate.id == member.id)
142                {
143                    self.emit(RoderEvent::TeamMemberStatusChanged(
144                        TeamMemberStatusChanged {
145                            team_id,
146                            member_id: member.id,
147                            member_thread_id: member.thread_id,
148                            status: TeamMemberStatus::Running,
149                            timestamp: OffsetDateTime::now_utc(),
150                        },
151                    ))
152                    .await;
153                }
154            }
155            self.goals.remember_turn_options(&req).await;
156            let runtime = Arc::clone(self);
157            let turn_req = req;
158            let thread_id_for_task = turn_req.thread_id.clone();
159            let turn_id_for_task = turn_id.clone();
160            tokio::spawn(async move {
161                let result = Abortable::new(
162                    runtime.run_turn(
163                        turn_req,
164                        turn_id_for_task.clone(),
165                        initial_mailbox_ack,
166                        steering,
167                        goal_only,
168                    ),
169                    abort_registration,
170                )
171                .await;
172                /*
173                 * A failed sibling in a parallel tool batch drops in-flight external tool
174                 * futures (`try_join_all` in `route_tool_calls`), stranding their
175                 * `pending_external_tool_calls` entries. Sweep before reporting the turn
176                 * outcome so every `thread/toolExecutionRequested` gets a terminal
177                 * resolution; on clean completion the map holds nothing for this turn.
178                 */
179                runtime
180                    .cancel_pending_external_tool_calls_for_turn(&turn_id_for_task)
181                    .await;
182                let completed = matches!(&result, Ok(Ok(TurnRunOutcome::Completed)));
183                let stopped_status = match &result {
184                    Ok(Err(error)) => Some(crate::goals::status_after_error(error)),
185                    Ok(Ok(TurnRunOutcome::Stopped)) => {
186                        Some(roder_api::goals::ThreadGoalStatus::Blocked)
187                    }
188                    Err(_) => Some(roder_api::goals::ThreadGoalStatus::Paused),
189                    _ => None,
190                };
191                if let Err(error) = runtime
192                    .goals
193                    .finish_turn(&thread_id_for_task, &turn_id_for_task, stopped_status)
194                    .await
195                {
196                    eprintln!("failed to finalize goal accounting: {error}");
197                }
198                match &result {
199                    Ok(Err(err)) => {
200                        let (cleanup, ownership) =
201                            runtime.await_provider_turn_cleanup(&turn_id_for_task).await;
202                        // This wrapper owns terminal failure emission for returned errors.
203                        runtime
204                            .emit(RoderEvent::TurnFailed(TurnFailed {
205                                thread_id: thread_id_for_task.clone(),
206                                turn_id: turn_id_for_task.clone(),
207                                error: err.to_string(),
208                                error_kind: err
209                                    .downcast_ref::<roder_api::provider_error::ProviderFailure>()
210                                    .map(|failure| failure.kind.retry_cause().to_string()),
211                                usage: None,
212                                timestamp: OffsetDateTime::now_utc(),
213                            }))
214                            .await;
215                        runtime
216                            .record_turn_lifecycle_with_ownership(
217                                thread_id_for_task.clone(),
218                                turn_id_for_task.clone(),
219                                TurnLifecycleState::Failed,
220                                cleanup,
221                                Some(TurnLifecycleReason::ProviderFailure),
222                                ownership,
223                            )
224                            .await;
225                        let _ = runtime
226                            .complete_team_member_turn_with_result(
227                                &thread_id_for_task,
228                                &turn_id_for_task,
229                                TeamMemberStatus::Failed,
230                                None,
231                                Some(err.to_string()),
232                            )
233                            .await;
234                    }
235                    Ok(Ok(TurnRunOutcome::Stopped)) => {
236                        let (cleanup, ownership) =
237                            runtime.await_provider_turn_cleanup(&turn_id_for_task).await;
238                        runtime
239                            .record_turn_lifecycle_with_ownership(
240                                thread_id_for_task.clone(),
241                                turn_id_for_task.clone(),
242                                TurnLifecycleState::Failed,
243                                cleanup,
244                                Some(TurnLifecycleReason::ProviderFailure),
245                                ownership,
246                            )
247                            .await;
248                        let _ = runtime
249                            .complete_team_member_turn_with_result(
250                                &thread_id_for_task,
251                                &turn_id_for_task,
252                                TeamMemberStatus::Failed,
253                                None,
254                                Some("turn stopped before completion".to_string()),
255                            )
256                            .await;
257                    }
258                    Err(_) => {
259                        let reason = if let Some(handle) = runtime
260                            .turn_drains
261                            .read()
262                            .await
263                            .get(&turn_id_for_task)
264                            .cloned()
265                        {
266                            handle
267                                .interrupt_reason
268                                .lock()
269                                .await
270                                .unwrap_or(TurnLifecycleReason::RuntimeFailure)
271                        } else {
272                            TurnLifecycleReason::RuntimeFailure
273                        };
274                        let (cleanup, ownership) =
275                            runtime.await_provider_turn_cleanup(&turn_id_for_task).await;
276                        runtime
277                            .record_turn_lifecycle_with_ownership(
278                                thread_id_for_task.clone(),
279                                turn_id_for_task.clone(),
280                                TurnLifecycleState::Interrupted,
281                                cleanup,
282                                Some(reason),
283                                ownership,
284                            )
285                            .await;
286                        runtime
287                            .emit(RoderEvent::TurnInterrupted(TurnInterrupted {
288                                thread_id: thread_id_for_task.clone(),
289                                turn_id: turn_id_for_task.clone(),
290                                timestamp: OffsetDateTime::now_utc(),
291                            }))
292                            .await;
293                        let _ = runtime
294                            .complete_team_member_turn_with_result(
295                                &thread_id_for_task,
296                                &turn_id_for_task,
297                                TeamMemberStatus::Interrupted,
298                                None,
299                                None,
300                            )
301                            .await;
302                    }
303                    Ok(Ok(TurnRunOutcome::Completed)) => {}
304                }
305                runtime
306                    .teams
307                    .release_mailbox_reservations_for_turn(&turn_id_for_task)
308                    .await;
309                runtime.active_turns.write().await.remove(&turn_id_for_task);
310                if let Some(drain) = runtime.turn_drains.write().await.remove(&turn_id_for_task) {
311                    drain.completed.store(true, Ordering::Release);
312                    drain.completed_notify.notify_waiters();
313                }
314                runtime
315                    .active_turn_selections
316                    .write()
317                    .await
318                    .remove(&turn_id_for_task);
319                runtime
320                    .active_turn_contexts
321                    .write()
322                    .await
323                    .remove(&turn_id_for_task);
324                if !completed {
325                    // Completion paths above consume registered provider cleanup
326                    // handles. A setup failure before an engine can stream has no
327                    // such handle; this is a harmless final defensive sweep.
328                    let _ = runtime.provider_turn_cleanups.lock().map(|mut cleanups| {
329                        cleanups.remove(&turn_id_for_task);
330                    });
331                }
332                runtime.active_turns_changed.notify_waiters();
333                if completed {
334                    let _ = runtime
335                        .continue_active_goal_after_turn(thread_id_for_task)
336                        .await;
337                }
338            });
339            drop(turn_admission);
340            Ok(Some(turn_id))
341        })
342    }
343}