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 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 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 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}