1use 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
34pub const REVIEW_TOOL_ALLOWLIST: &[&str] = &["read_file", "list_files", "grep", "glob", "shell"];
42
43const 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 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#[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
85struct StartedReview {
87 handle: ReviewHandle,
88 request: ReviewRequest,
89 parent_thread_id: ThreadId,
90 events: broadcast::Receiver<EventEnvelope>,
91}
92
93impl Runtime {
94 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 tokio::spawn(async move {
106 let _ = runtime.finish_review(started).await;
107 });
108 Ok(handle)
109 }
110
111 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 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 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 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 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 })
199 .await?;
200
201 let handle = ReviewHandle {
202 review_id: uuid::Uuid::new_v4().to_string(),
203 review_thread_id: review_thread.thread_id,
204 turn_id,
205 };
206 self.emit(RoderEvent::ReviewStarted(ReviewStarted {
207 thread_id: req.thread_id.clone(),
208 turn_id: handle.turn_id.clone(),
209 review_id: handle.review_id.clone(),
210 review_thread_id: handle.review_thread_id.clone(),
211 target: req.target,
212 label: request.label.clone(),
213 timestamp: OffsetDateTime::now_utc(),
214 }))
215 .await;
216
217 Ok(StartedReview {
218 handle,
219 request,
220 parent_thread_id: req.thread_id,
221 events,
222 })
223 }
224
225 async fn finish_review(
226 self: &Arc<Self>,
227 started: StartedReview,
228 ) -> anyhow::Result<ReviewCompletion> {
229 let StartedReview {
230 handle,
231 request,
232 parent_thread_id,
233 mut events,
234 } = started;
235
236 let collected =
237 collect_final_answer(&mut events, &handle.review_thread_id, &handle.turn_id).await;
238
239 let text = match collected {
240 Ok(text) => text,
241 Err(error) => {
242 self.echo_review_into_parent(&parent_thread_id, &handle.turn_id, None)
243 .await;
244 self.emit(RoderEvent::ReviewFailed(ReviewFailed {
245 thread_id: parent_thread_id,
246 turn_id: Some(handle.turn_id),
247 review_id: handle.review_id,
248 review_thread_id: Some(handle.review_thread_id),
249 error: error.to_string(),
250 timestamp: OffsetDateTime::now_utc(),
251 }))
252 .await;
253 return Err(error);
254 }
255 };
256
257 let output = parse_review_output(&text);
258 self.echo_review_into_parent(&parent_thread_id, &handle.turn_id, Some(&output))
259 .await;
260 self.remember_review(ReviewRecord {
261 review_id: handle.review_id.clone(),
262 parent_thread_id: parent_thread_id.clone(),
263 review_thread_id: handle.review_thread_id.clone(),
264 turn_id: handle.turn_id.clone(),
265 workspace_root: request.workspace_root.clone(),
266 request,
267 output: output.clone(),
268 })
269 .await;
270 self.emit(RoderEvent::ReviewCompleted(ReviewCompleted {
271 thread_id: parent_thread_id,
272 turn_id: handle.turn_id.clone(),
273 review_id: handle.review_id.clone(),
274 review_thread_id: handle.review_thread_id.clone(),
275 output: output.clone(),
276 timestamp: OffsetDateTime::now_utc(),
277 }))
278 .await;
279
280 Ok(ReviewCompletion { handle, output })
281 }
282
283 async fn review_vcs_provider(
284 &self,
285 workspace_root: &Path,
286 ) -> anyhow::Result<Arc<dyn VcsProvider>> {
287 let resolution = self
288 .registry
289 .version_control_resolver()
290 .resolve_provider(VcsResolveRequest {
291 workspace_root: workspace_root.to_path_buf(),
292 preferred_provider_id: None,
293 })
294 .await
295 .with_context(|| {
296 format!(
297 "resolving a version control provider for {}",
298 workspace_root.display()
299 )
300 })?;
301 match resolution {
302 VcsProviderResolution::Available { provider, .. } => Ok(provider),
303 VcsProviderResolution::Unavailable { workspace_root } => anyhow::bail!(
304 "no version control provider claims {}; a review needs a repository to diff",
305 workspace_root.display()
306 ),
307 }
308 }
309
310 async fn echo_review_into_parent(
314 &self,
315 thread_id: &ThreadId,
316 turn_id: &TurnId,
317 output: Option<&ReviewOutput>,
318 ) {
319 let (user_text, assistant_text) = match output {
320 Some(output) => (
321 synthetic_user_action(output),
322 render_review_output_text(output),
323 ),
324 None => (
325 synthetic_user_action_interrupted(),
326 REVIEW_INTERRUPTED_MESSAGE.to_string(),
327 ),
328 };
329 let _ = self
332 .persist_turn_item(
333 thread_id,
334 turn_id,
335 &TranscriptItem::UserMessage(UserMessage::text(user_text)),
336 )
337 .await;
338 let _ = self
339 .persist_turn_item(
340 thread_id,
341 turn_id,
342 &TranscriptItem::AssistantMessage(AssistantMessage {
343 text: assistant_text,
344 phase: None,
345 }),
346 )
347 .await;
348 }
349
350 async fn remember_review(&self, record: ReviewRecord) {
351 let mut reviews = self.reviews.lock().await;
352 reviews.retain(|existing| existing.review_id != record.review_id);
353 reviews.push(record);
354 let overflow = reviews.len().saturating_sub(REVIEW_HISTORY_LIMIT);
355 if overflow > 0 {
356 reviews.drain(..overflow);
357 }
358 }
359}
360
361async fn collect_final_answer(
367 events: &mut broadcast::Receiver<EventEnvelope>,
368 thread_id: &ThreadId,
369 turn_id: &TurnId,
370) -> anyhow::Result<String> {
371 let mut final_answer: Option<String> = None;
372 let mut streamed = String::new();
373 loop {
374 let envelope = match events.recv().await {
375 Ok(envelope) => envelope,
376 Err(broadcast::error::RecvError::Lagged(_)) => continue,
377 Err(broadcast::error::RecvError::Closed) => {
378 anyhow::bail!("the runtime event bus closed before the review turn finished")
379 }
380 };
381 if envelope.thread_id.as_ref() != Some(thread_id)
382 || envelope.turn_id.as_ref() != Some(turn_id)
383 {
384 continue;
385 }
386 match &envelope.event {
387 RoderEvent::TranscriptItemAppended(event) => {
388 if let Some(TranscriptItem::AssistantMessage(message)) = &event.item
389 && message.phase.as_deref() == Some(FINAL_ANSWER_PHASE)
390 {
391 final_answer = Some(message.text.clone());
392 }
393 }
394 RoderEvent::InferenceEventReceived(event) => {
395 if let InferenceEvent::MessageDelta(delta) = &event.event {
396 streamed.push_str(&delta.text);
397 }
398 }
399 RoderEvent::TurnCompleted(_) => break,
400 RoderEvent::TurnFailed(event) => {
401 anyhow::bail!("review turn failed: {}", event.error)
402 }
403 _ => {}
404 }
405 }
406 Ok(final_answer.unwrap_or(streamed))
407}
408
409#[cfg(test)]
410#[path = "run_tests.rs"]
411mod tests;