Skip to main content

roder_core/goals/
runtime.rs

1use 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    /// Keep goal mutations excluded until idle turn admission is committed.
33    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    /// Keep the admitted turn's configuration for automatic continuation. A
47    /// later explicit turn replaces it; user input and attachments are never
48    /// replayed, and this authority context is never persisted to disk.
49    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        // Delegated work consumes the ancestor goal's budget, even while the
194        // lead waits for its children. Charge each ancestor once.
195        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(&current).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}