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 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 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 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
362async 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;