roder_core/goals/
runtime.rs1use super::*;
2
3impl RuntimeGoalController {
4 pub(crate) async fn continuation_request(
5 &self,
6 thread_id: &ThreadId,
7 ) -> Option<StartTurnRequest> {
8 self.cache
9 .lock()
10 .await
11 .continuation_requests
12 .get(thread_id)
13 .cloned()
14 }
15
16 pub(crate) async fn inherit_thread_goal_snapshot(
17 &self,
18 source: &ThreadId,
19 target: &ThreadId,
20 ) -> anyhow::Result<Option<ThreadGoal>> {
21 let _guard = self.mutation.lock().await;
22 self.flush_thread_progress(source).await?;
23 let Some(mut goal) = self.load_goal(source).await? else {
24 return Ok(None);
25 };
26 validate_thread_goal_objective(&goal.objective)?;
27 goal.thread_id = target.clone();
28 self.store_goal(goal.clone()).await?;
29 Ok(Some(goal))
30 }
31
32 pub(crate) async fn admit_continuation(
34 &self,
35 thread_id: &ThreadId,
36 ) -> anyhow::Result<Option<(tokio::sync::OwnedMutexGuard<()>, ThreadGoal)>> {
37 let guard = self.mutation.clone().lock_owned().await;
38 self.flush_thread_progress(thread_id).await?;
39 Ok(self
40 .load_goal(thread_id)
41 .await?
42 .filter(|goal| goal.status.is_active())
43 .map(|goal| (guard, goal)))
44 }
45
46 pub(crate) async fn remember_turn_options(&self, request: &StartTurnRequest) {
50 self.cache.lock().await.continuation_requests.insert(
51 request.thread_id.clone(),
52 StartTurnRequest {
53 thread_id: request.thread_id.clone(),
54 message: String::new(),
55 images: Vec::new(),
56 provider_override: request.provider_override.clone(),
57 model_override: request.model_override.clone(),
58 reasoning_override: request.reasoning_override.clone(),
59 workspace: request.workspace.clone(),
60 instructions: request.instructions.clone(),
61 developer_context: request.developer_context.clone(),
62 task_ledger_required: request.task_ledger_required,
63 service_tier_override: request.service_tier_override.clone(),
64 },
65 );
66 }
67}
68
69impl Runtime {
70 pub async fn thread_goal_get(
71 &self,
72 thread_id: &ThreadId,
73 ) -> anyhow::Result<Option<ThreadGoal>> {
74 self.goals.get_thread_goal(thread_id).await
75 }
76
77 pub async fn thread_goal_set(
78 &self,
79 thread_id: &ThreadId,
80 patch: ThreadGoalPatch,
81 ) -> anyhow::Result<Option<ThreadGoal>> {
82 self.goals.set_thread_goal(thread_id, patch).await
83 }
84
85 pub async fn thread_goal_clear(&self, thread_id: &ThreadId) -> anyhow::Result<bool> {
86 let cleared = self.goals.clear_thread_goal(thread_id).await?;
87 if cleared && let Some(turn_id) = self.active_turn_for_thread(thread_id).await {
88 let _ = self.steer_turn(thread_id.clone(), turn_id,
89 "The user cleared the thread goal. Stop autonomous goal work; respond to any remaining explicit user instructions.".into(), Vec::new()).await;
90 }
91 Ok(cleared)
92 }
93
94 pub async fn apply_external_goal_set_effects(
95 self: &Arc<Self>,
96 previous_goal: Option<ThreadGoal>,
97 goal: Option<ThreadGoal>,
98 ) -> anyhow::Result<Option<ThreadId>> {
99 let Some(goal) = goal else {
100 return Ok(None);
101 };
102 if goal.status != ThreadGoalStatus::Active {
103 if let Some(turn_id) = self.active_turn_for_thread(&goal.thread_id).await {
104 let message = if goal.status == ThreadGoalStatus::BudgetLimited {
105 super::prompts::budget_limit_prompt(&goal)
106 } else {
107 format!(
108 "The user set the goal status to {}. Stop autonomous goal work and report this status.",
109 goal.status.as_str()
110 )
111 };
112 self.steer_turn(goal.thread_id.clone(), turn_id.clone(), message, Vec::new())
113 .await?;
114 return Ok(Some(turn_id));
115 }
116 return Ok(None);
117 }
118
119 let objective_changed = previous_goal
120 .as_ref()
121 .is_none_or(|previous| previous.objective != goal.objective);
122 let resumed = previous_goal
123 .as_ref()
124 .is_some_and(|previous| !previous.status.is_active());
125 if (objective_changed || resumed)
126 && let Some(turn_id) = self.active_turn_for_thread(&goal.thread_id).await
127 {
128 self.steer_turn(
129 goal.thread_id.clone(),
130 turn_id.clone(),
131 if objective_changed { objective_updated_prompt(&goal) } else {
132 format!("The user resumed the goal. The earlier pause request is revoked. Start a fresh blocked audit and continue pursuing the full objective.\n\n{}", continuation_prompt(&goal))
133 },
134 Vec::new(),
135 )
136 .await?;
137 return Ok(Some(turn_id));
138 }
139
140 self.continue_active_goal_if_idle(goal.thread_id.clone())
141 .await
142 }
143
144 pub async fn continue_active_goal_if_idle(
145 self: &Arc<Self>,
146 thread_id: ThreadId,
147 ) -> anyhow::Result<Option<ThreadId>> {
148 if self.has_active_turn_for_thread(&thread_id).await {
149 return Ok(None);
150 }
151 let Some(goal) = self.goals.active_goal(&thread_id).await? else {
152 return Ok(None);
153 };
154 let previous = self.goals.continuation_request(&thread_id).await;
155 let mut request = match previous {
156 Some(request) => request,
157 None => StartTurnRequest {
158 thread_id: thread_id.clone(),
159 message: String::new(),
160 images: Vec::new(),
161 provider_override: None,
162 model_override: None,
163 reasoning_override: None,
164 workspace: self.workspace_for_thread(&thread_id).await?,
165 instructions: crate::default_instructions(),
166 developer_context: None,
167 task_ledger_required: false,
168 service_tier_override: None,
169 },
170 };
171 request.message = continuation_prompt(&goal);
172 self.start_goal_turn_if_idle(request).await
173 }
174
175 pub(crate) async fn continue_active_goal_after_turn(
176 self: &Arc<Self>,
177 thread_id: ThreadId,
178 ) -> anyhow::Result<Option<ThreadId>> {
179 self.continue_active_goal_if_idle(thread_id).await
180 }
181}
182
183impl Runtime {
184 pub(crate) async fn record_goal_token_usage(
185 &self,
186 thread_id: &ThreadId,
187 turn_id: &str,
188 tokens: i64,
189 ) -> anyhow::Result<()> {
190 self.goals
191 .record_turn_usage(thread_id, turn_id, tokens)
192 .await?;
193 let mut current = thread_id.clone();
196 let mut visited = std::collections::HashSet::from([current.clone()]);
197 while let Some((team_id, member)) = self.teams.member_for_thread(¤t).await {
198 let parent = match member.parent_thread_id {
199 Some(parent) => parent,
200 None => match self.teams.get(&team_id).await {
201 Some(team) if team.lead_thread_id != current => team.lead_thread_id,
202 _ => break,
203 },
204 };
205 if !visited.insert(parent.clone()) {
206 break;
207 }
208 self.goals.record_descendant_usage(&parent, tokens).await?;
209 current = parent;
210 }
211 Ok(())
212 }
213}