Skip to main content

roder_core/review/
run.rs

1//! Runs a review as a detached, read-only sub-turn.
2//!
3//! Roder has no one-shot sub-agent with a per-call system prompt, so a review
4//! is a child thread plus a single [`Runtime::start_turn`] whose
5//! [`InstructionBundle::system`] is replaced by the rubric. The thread is
6//! created with a read-only `tool_allowlist`, which is enforced both when tools
7//! are advertised and when a call is dispatched — the global policy mode is
8//! never touched, so a review can never widen what the parent thread may do.
9
10use std::path::{Path, PathBuf};
11use std::sync::Arc;
12
13use anyhow::Context;
14use roder_api::events::{
15    EventEnvelope, ReviewCompleted, ReviewFailed, ReviewStarted, RoderEvent, ThreadId, TurnId,
16};
17use roder_api::inference::{InferenceEvent, InstructionBundle};
18use roder_api::review::{ReviewId, ReviewOutput, ReviewRequest, ReviewTarget};
19use roder_api::transcript::{AssistantMessage, TranscriptItem, UserMessage};
20use roder_api::version_control::{
21    VcsProvider, VcsProviderResolution, VcsProviderResolver, VcsResolveRequest,
22};
23use time::OffsetDateTime;
24use tokio::sync::broadcast;
25
26use crate::review::REVIEW_RUBRIC;
27use crate::review::parse::parse_review_output;
28use crate::review::prompt::resolve_review_request;
29use crate::review::render::{
30    render_review_output_text, synthetic_user_action, synthetic_user_action_interrupted,
31};
32use crate::runtime::{CreateThreadRequest, FINAL_ANSWER_PHASE, Runtime, StartTurnRequest};
33
34/// The only tools a review thread may call.
35///
36/// `shell` is included because `git diff`/`git log` is how a reviewer sees a
37/// change; every tool that writes files, edits code, applies patches, or spawns
38/// agents is deliberately absent. Names are matched exactly against the tool
39/// registry (`allowlist_permits`), so they must stay in sync with
40/// `roder-tools`.
41pub const REVIEW_TOOL_ALLOWLIST: &[&str] = &["read_file", "list_files", "grep", "glob", "shell"];
42
43/// How many completed reviews stay resolvable by id for a later publish.
44const REVIEW_HISTORY_LIMIT: usize = 64;
45
46const REVIEW_INTERRUPTED_MESSAGE: &str =
47    "Review was interrupted. Re-run `/review` and wait for it to complete.";
48
49#[derive(Debug, Clone)]
50pub struct StartReviewRequest {
51    /// Thread the review was requested from; the review runs in a child of it.
52    pub thread_id: ThreadId,
53    pub target: ReviewTarget,
54    pub workspace: String,
55    pub provider_override: Option<String>,
56    pub model_override: Option<String>,
57}
58
59#[derive(Debug, Clone, PartialEq, Eq)]
60pub struct ReviewHandle {
61    pub review_id: ReviewId,
62    pub review_thread_id: ThreadId,
63    pub turn_id: TurnId,
64}
65
66#[derive(Debug, Clone)]
67pub struct ReviewCompletion {
68    pub handle: ReviewHandle,
69    pub output: ReviewOutput,
70}
71
72/// Everything a publisher needs about a finished review, kept in memory so
73/// `review/publish` can act on a review id alone.
74#[derive(Debug, Clone)]
75pub struct ReviewRecord {
76    pub review_id: ReviewId,
77    pub parent_thread_id: ThreadId,
78    pub review_thread_id: ThreadId,
79    pub turn_id: TurnId,
80    pub request: ReviewRequest,
81    pub output: ReviewOutput,
82    pub workspace_root: PathBuf,
83}
84
85/// A review whose turn is running but whose findings have not been collected.
86struct StartedReview {
87    handle: ReviewHandle,
88    request: ReviewRequest,
89    parent_thread_id: ThreadId,
90    events: broadcast::Receiver<EventEnvelope>,
91}
92
93impl Runtime {
94    /// Starts a review and returns as soon as its turn is admitted. Findings
95    /// arrive as [`RoderEvent::ReviewCompleted`].
96    pub async fn start_review(
97        self: &Arc<Self>,
98        req: StartReviewRequest,
99    ) -> anyhow::Result<ReviewHandle> {
100        let started = self.begin_review(req).await?;
101        let handle = started.handle.clone();
102        let runtime = Arc::clone(self);
103        // `finish_review` reports its own failure as `ReviewFailed`, so the
104        // detached task has nothing left to do with the error.
105        tokio::spawn(async move {
106            let _ = runtime.finish_review(started).await;
107        });
108        Ok(handle)
109    }
110
111    /// Starts a review and awaits its findings.
112    pub async fn run_review(
113        self: &Arc<Self>,
114        req: StartReviewRequest,
115    ) -> anyhow::Result<ReviewCompletion> {
116        let started = self.begin_review(req).await?;
117        self.finish_review(started).await
118    }
119
120    /// Looks up a completed review by id.
121    pub async fn review_record(&self, review_id: &str) -> Option<ReviewRecord> {
122        self.reviews
123            .lock()
124            .await
125            .iter()
126            .rev()
127            .find(|record| record.review_id == review_id)
128            .cloned()
129    }
130
131    /// The most recently completed review, used when a publish request omits
132    /// the review id.
133    pub async fn latest_review_record(&self) -> Option<ReviewRecord> {
134        self.reviews.lock().await.last().cloned()
135    }
136
137    async fn begin_review(
138        self: &Arc<Self>,
139        req: StartReviewRequest,
140    ) -> anyhow::Result<StartedReview> {
141        if let Some(turn_id) = self.active_turn_for_thread(&req.thread_id).await {
142            anyhow::bail!(
143                "thread {} has an active turn ({turn_id}); wait for it to finish before starting a review",
144                req.thread_id
145            );
146        }
147
148        // `[review].model` reserves a stronger model for reviews without
149        // touching the model the requesting thread talks to.
150        let model_override = match req.model_override.clone() {
151            Some(model) => Some(model),
152            None => self.review_config().await.model,
153        };
154
155        let workspace_root = PathBuf::from(&req.workspace);
156        let vcs = self.review_vcs_provider(&workspace_root).await?;
157        let request = resolve_review_request(&req.target, &workspace_root, vcs.as_ref()).await?;
158
159        let review_thread = self
160            .create_thread_with(CreateThreadRequest {
161                title: Some(format!("review: {}", request.label)),
162                workspace: req.workspace.clone(),
163                workspace_id: None,
164                root_id: Some(req.thread_id.clone()),
165                provider: req.provider_override.clone(),
166                model: model_override.clone(),
167                selection_mode: None,
168                tool_allowlist: REVIEW_TOOL_ALLOWLIST
169                    .iter()
170                    .map(|name| (*name).to_string())
171                    .collect(),
172                developer_instructions: None,
173                external_tools: Vec::new(),
174                runner: None,
175            })
176            .await?;
177
178        // Subscribe before the turn starts, otherwise its first transcript
179        // items can land before the receiver exists.
180        let events = self.subscribe_events();
181
182        let turn_id = self
183            .start_turn(StartTurnRequest {
184                thread_id: review_thread.thread_id.clone(),
185                message: request.prompt.clone(),
186                images: Vec::new(),
187                provider_override: req.provider_override,
188                model_override,
189                reasoning_override: None,
190                workspace: req.workspace.clone(),
191                instructions: InstructionBundle {
192                    system: Some(REVIEW_RUBRIC.to_string()),
193                    developer: None,
194                    developer_context: None,
195                },
196                developer_context: None,
197                task_ledger_required: false,
198                service_tier_override: None,
199            })
200            .await?;
201
202        let handle = ReviewHandle {
203            review_id: uuid::Uuid::new_v4().to_string(),
204            review_thread_id: review_thread.thread_id,
205            turn_id,
206        };
207        self.emit(RoderEvent::ReviewStarted(ReviewStarted {
208            thread_id: req.thread_id.clone(),
209            turn_id: handle.turn_id.clone(),
210            review_id: handle.review_id.clone(),
211            review_thread_id: handle.review_thread_id.clone(),
212            target: req.target,
213            label: request.label.clone(),
214            timestamp: OffsetDateTime::now_utc(),
215        }))
216        .await;
217
218        Ok(StartedReview {
219            handle,
220            request,
221            parent_thread_id: req.thread_id,
222            events,
223        })
224    }
225
226    async fn finish_review(
227        self: &Arc<Self>,
228        started: StartedReview,
229    ) -> anyhow::Result<ReviewCompletion> {
230        let StartedReview {
231            handle,
232            request,
233            parent_thread_id,
234            mut events,
235        } = started;
236
237        let collected =
238            collect_final_answer(&mut events, &handle.review_thread_id, &handle.turn_id).await;
239
240        let text = match collected {
241            Ok(text) => text,
242            Err(error) => {
243                self.echo_review_into_parent(&parent_thread_id, &handle.turn_id, None)
244                    .await;
245                self.emit(RoderEvent::ReviewFailed(ReviewFailed {
246                    thread_id: parent_thread_id,
247                    turn_id: Some(handle.turn_id),
248                    review_id: handle.review_id,
249                    review_thread_id: Some(handle.review_thread_id),
250                    error: error.to_string(),
251                    timestamp: OffsetDateTime::now_utc(),
252                }))
253                .await;
254                return Err(error);
255            }
256        };
257
258        let output = parse_review_output(&text);
259        self.echo_review_into_parent(&parent_thread_id, &handle.turn_id, Some(&output))
260            .await;
261        self.remember_review(ReviewRecord {
262            review_id: handle.review_id.clone(),
263            parent_thread_id: parent_thread_id.clone(),
264            review_thread_id: handle.review_thread_id.clone(),
265            turn_id: handle.turn_id.clone(),
266            workspace_root: request.workspace_root.clone(),
267            request,
268            output: output.clone(),
269        })
270        .await;
271        self.emit(RoderEvent::ReviewCompleted(ReviewCompleted {
272            thread_id: parent_thread_id,
273            turn_id: handle.turn_id.clone(),
274            review_id: handle.review_id.clone(),
275            review_thread_id: handle.review_thread_id.clone(),
276            output: output.clone(),
277            timestamp: OffsetDateTime::now_utc(),
278        }))
279        .await;
280
281        Ok(ReviewCompletion { handle, output })
282    }
283
284    async fn review_vcs_provider(
285        &self,
286        workspace_root: &Path,
287    ) -> anyhow::Result<Arc<dyn VcsProvider>> {
288        let resolution = self
289            .registry
290            .version_control_resolver()
291            .resolve_provider(VcsResolveRequest {
292                workspace_root: workspace_root.to_path_buf(),
293                preferred_provider_id: None,
294            })
295            .await
296            .with_context(|| {
297                format!(
298                    "resolving a version control provider for {}",
299                    workspace_root.display()
300                )
301            })?;
302        match resolution {
303            VcsProviderResolution::Available { provider, .. } => Ok(provider),
304            VcsProviderResolution::Unavailable { workspace_root } => anyhow::bail!(
305                "no version control provider claims {}; a review needs a repository to diff",
306                workspace_root.display()
307            ),
308        }
309    }
310
311    /// Records the review in the requesting thread the way codex does: a
312    /// synthetic `<user_action>` turn plus the rendered summary, so the parent
313    /// agent can act on findings the user keeps.
314    async fn echo_review_into_parent(
315        &self,
316        thread_id: &ThreadId,
317        turn_id: &TurnId,
318        output: Option<&ReviewOutput>,
319    ) {
320        let (user_text, assistant_text) = match output {
321            Some(output) => (
322                synthetic_user_action(output),
323                render_review_output_text(output),
324            ),
325            None => (
326                synthetic_user_action_interrupted(),
327                REVIEW_INTERRUPTED_MESSAGE.to_string(),
328            ),
329        };
330        // Echoing is best effort: a persistence failure must not turn a
331        // successful review into a failed one.
332        let _ = self
333            .persist_turn_item(
334                thread_id,
335                turn_id,
336                &TranscriptItem::UserMessage(UserMessage::text(user_text)),
337            )
338            .await;
339        let _ = self
340            .persist_turn_item(
341                thread_id,
342                turn_id,
343                &TranscriptItem::AssistantMessage(AssistantMessage {
344                    text: assistant_text,
345                    phase: None,
346                }),
347            )
348            .await;
349    }
350
351    async fn remember_review(&self, record: ReviewRecord) {
352        let mut reviews = self.reviews.lock().await;
353        reviews.retain(|existing| existing.review_id != record.review_id);
354        reviews.push(record);
355        let overflow = reviews.len().saturating_sub(REVIEW_HISTORY_LIMIT);
356        if overflow > 0 {
357            reviews.drain(..overflow);
358        }
359    }
360}
361
362/// Drains the bus until the review turn ends, returning its final answer.
363///
364/// The persisted `final_answer` assistant item is authoritative; streamed
365/// message deltas are only a fallback for providers that end a turn without
366/// one.
367async fn collect_final_answer(
368    events: &mut broadcast::Receiver<EventEnvelope>,
369    thread_id: &ThreadId,
370    turn_id: &TurnId,
371) -> anyhow::Result<String> {
372    let mut final_answer: Option<String> = None;
373    let mut streamed = String::new();
374    loop {
375        let envelope = match events.recv().await {
376            Ok(envelope) => envelope,
377            Err(broadcast::error::RecvError::Lagged(_)) => continue,
378            Err(broadcast::error::RecvError::Closed) => {
379                anyhow::bail!("the runtime event bus closed before the review turn finished")
380            }
381        };
382        if envelope.thread_id.as_ref() != Some(thread_id)
383            || envelope.turn_id.as_ref() != Some(turn_id)
384        {
385            continue;
386        }
387        match &envelope.event {
388            RoderEvent::TranscriptItemAppended(event) => {
389                if let Some(TranscriptItem::AssistantMessage(message)) = &event.item
390                    && message.phase.as_deref() == Some(FINAL_ANSWER_PHASE)
391                {
392                    final_answer = Some(message.text.clone());
393                }
394            }
395            RoderEvent::InferenceEventReceived(event) => {
396                if let InferenceEvent::MessageDelta(delta) = &event.event {
397                    streamed.push_str(&delta.text);
398                }
399            }
400            RoderEvent::TurnCompleted(_) => break,
401            RoderEvent::TurnFailed(event) => {
402                anyhow::bail!("review turn failed: {}", event.error)
403            }
404            _ => {}
405        }
406    }
407    Ok(final_answer.unwrap_or(streamed))
408}
409
410#[cfg(test)]
411#[path = "run_tests.rs"]
412mod tests;