1use crate::author;
8use crate::gate::{self, MergeToken, Refusal};
9use crate::progress::{Phase, Watcher};
10use crate::record::RunLog;
11use crate::review;
12use crate::route::Routing;
13use crate::worktree;
14use crate::{Error, Result};
15use ostraka_adapter::{AdapterOutcome, VendorAdapter};
16use ostraka_core::clock::now_rfc3339;
17use ostraka_core::config::Config;
18use ostraka_core::gate::{Approval, Verdict};
19use ostraka_core::identity::ActorId;
20use ostraka_core::record::{Event, Outcome, RunRecord};
21use ostraka_core::task::TaskSpec;
22use std::path::Path;
23
24pub struct RunReport {
26 pub record: RunRecord,
27 pub token: Option<MergeToken>,
29 pub refusal: Option<Refusal>,
30 pub diff: String,
31}
32
33impl RunReport {
34 pub fn approved(&self) -> bool {
35 self.token.is_some()
36 }
37}
38
39pub struct Places<'a> {
48 pub repo: &'a Path,
50 pub worktrees: &'a Path,
52 pub records: &'a Path,
54 pub name: &'a str,
56 pub notes: Option<&'a Path>,
59 pub skills: Option<&'a Path>,
63}
64
65pub fn run_task(
74 places: &Places<'_>,
75 config: &Config,
76 routing: &Routing,
77 task: &TaskSpec,
78 reviewer_identity: &ActorId,
79 watcher: Option<Box<dyn Watcher>>,
80) -> Result<RunReport> {
81 let repo = places.repo;
82 let run_id = format!("{}-{}", task.id, now_rfc3339().replace([':', '-'], ""));
83 let mut log = RunLog::create(places.records, &run_id)?.watched_by(watcher);
84 log.enter(Phase::Isolating);
85
86 let mut record = RunRecord {
87 run_id: run_id.clone(),
88 task_id: task.id.clone(),
89 prompt: task.prompt.clone(),
90 author: task.author.clone(),
91 adapter: routing.author.id().to_string(),
92 repository: places.name.to_string(),
93 started_at: now_rfc3339(),
94 finished_at: None,
95 checks: Vec::new(),
96 approval: None,
97 usage: Vec::new(),
98 outcome: None,
99 };
100
101 let wt = worktree::create(repo, places.worktrees, &run_id, &task.base_ref)?;
103
104 log.enter(Phase::Preparing);
108 match worktree::prepare(
109 repo,
110 wt.path(),
111 &config.worktree,
112 places.notes,
113 places.skills,
114 config.gate.timeout_secs.map(std::time::Duration::from_secs),
115 ) {
116 Ok(steps) => {
117 for step in steps {
118 log.append(&Event::Message {
119 text: format!("prepared: {step}"),
120 raw: None,
121 })?;
122 }
123 }
124 Err(problem) => {
125 return finish(
126 log,
127 record,
128 Outcome::Failed,
129 None,
130 Some(Refusal::SetupFailed {
131 step: problem.step,
132 reason: problem.reason,
133 }),
134 String::new(),
135 );
136 }
137 }
138
139 log.enter(Phase::Authoring);
142 let authoring = TaskSpec {
150 prompt: author::author_prompt(
151 &task.prompt,
152 worktree::linked(wt.path(), "notes"),
153 worktree::linked(wt.path(), "skills"),
154 ),
155 ..task.clone()
156 };
157 let author = drive(routing.author.as_ref(), &authoring, wt.path(), &mut log)?;
158 record.usage.extend(author.usage.clone());
159
160 let touched = worktree::touched_paths(wt.path())?;
162 log.append(&Event::Finished {
163 exit_code: author.exit_code,
164 files_touched: touched.clone(),
165 })?;
166
167 if author.interrupted {
175 return finish(
176 log,
177 record,
178 Outcome::Rejected,
179 None,
180 Some(Refusal::Interrupted),
181 String::new(),
182 );
183 }
184
185 if author.timed_out {
186 return finish(
187 log,
188 record,
189 Outcome::Rejected,
190 None,
191 Some(Refusal::TimedOut {
192 after_secs: config.policy.timeout_secs.unwrap_or_default(),
193 }),
194 String::new(),
195 );
196 }
197
198 if author.exit_code != Some(0) {
212 let refusal = Refusal::AuthorFailed {
213 code: author
214 .exit_code
215 .map(|c| c.to_string())
216 .unwrap_or_else(|| "no exit code".to_string()),
217 diagnostics: author.diagnostics.clone(),
218 };
219 return finish(
220 log,
221 record,
222 Outcome::Rejected,
223 None,
224 Some(refusal),
225 String::new(),
226 );
227 }
228
229 if touched.is_empty() {
230 return finish(
231 log,
232 record,
233 Outcome::Rejected,
234 None,
235 Some(Refusal::NoChange),
236 String::new(),
237 );
238 }
239
240 if !config.policy.permits(&touched) {
241 return finish(
242 log,
243 record,
244 Outcome::Rejected,
245 None,
246 Some(Refusal::PolicyViolation {
247 reason: format!("policy forbids writing outside declared paths: {touched:?}"),
248 }),
249 String::new(),
250 );
251 }
252
253 log.enter(Phase::Gating);
255 let passed = match gate::run_checks(&config.gate, wt.path(), &mut |record| {
256 log_checked(&mut log, record)
257 }) {
258 Ok(p) => p,
259 Err(refusal) => {
260 if let Refusal::ChecksFailed { records, .. } = &refusal {
264 record.checks.clone_from(records);
265 }
266 return finish(
267 log,
268 record,
269 Outcome::Rejected,
270 None,
271 Some(refusal),
272 String::new(),
273 );
274 }
275 };
276 record.checks = passed.records().to_vec();
277
278 log.enter(Phase::Reviewing);
280 let diff = worktree::diff(wt.path())?;
281 let (verdict, reviewer_usage) =
282 collect_verdict(routing.reviewer.as_ref(), task, &diff, wt.path(), &mut log)?;
283 record.usage.extend(reviewer_usage);
284 let approval = Approval {
285 reviewer: reviewer_identity.clone(),
286 verdict: verdict.clone(),
287 };
288 record.approval = Some(approval.clone());
289
290 match gate::evaluate(
292 passed,
293 &task.author,
294 &approval,
295 config.gate.review.must_differ_from_author,
296 ) {
297 Ok(token) => {
298 let message = format!(
301 "{}\n\nRun: {run_id}\nAuthored-by: {} ({})\nReviewed-by: {} ({})",
302 task.prompt,
303 task.author,
304 routing.author.id(),
305 reviewer_identity,
306 routing.reviewer.id(),
307 );
308 worktree::commit(wt.path(), &message, &task.author)?;
309 if let Err(e) = worktree::release(repo, &wt) {
318 log.append(&Event::Error {
319 message: format!("the worktree could not be removed: {e}"),
320 raw: None,
321 })?;
322 }
323 finish(log, record, Outcome::Approved, Some(token), None, diff)
324 }
325 Err(refusal) => finish(log, record, Outcome::Rejected, None, Some(refusal), diff),
326 }
327}
328
329fn log_checked(log: &mut RunLog, record: &ostraka_core::gate::CheckRecord) {
334 log.checked(record);
335}
336
337fn drive(
339 adapter: &dyn VendorAdapter,
340 task: &TaskSpec,
341 worktree: &Path,
342 log: &mut RunLog,
343) -> Result<AdapterOutcome> {
344 let mut session = adapter.launch(task, worktree)?;
345 while let Some(event) = session.next_event() {
346 log.append(&event)?;
347 }
348 let outcome = session.finish();
349 if let Some(diagnostics) = &outcome.diagnostics {
353 log.append(&Event::Error {
354 message: format!("author exited abnormally: {diagnostics}"),
355 raw: None,
356 })?;
357 }
358 Ok(outcome)
359}
360
361type Reviewed = (Verdict, Option<ostraka_core::record::TokenUsage>);
366
367fn collect_verdict(
368 reviewer: &dyn VendorAdapter,
369 task: &TaskSpec,
370 diff: &str,
371 worktree: &Path,
372 log: &mut RunLog,
373) -> Result<Reviewed> {
374 let marker = review::verdict_marker(&task.id);
378 let review_task = TaskSpec {
379 id: format!("{}-review", task.id),
380 prompt: review::review_prompt(&task.prompt, diff, &marker),
381 adapter: reviewer.id().to_string(),
382 author: task.author.clone(),
383 base_ref: task.base_ref.clone(),
384 model: None,
385 };
386
387 let mut session = match reviewer.launch(&review_task, worktree) {
388 Ok(s) => s,
389 Err(e) => {
390 return Ok((
391 Verdict::Reject {
392 reason: format!("reviewer could not be launched: {e}"),
393 },
394 None,
395 ));
396 }
397 };
398
399 let mut spoken = String::new();
400 while let Some(event) = session.next_event() {
401 if let Event::Message { text, .. } = &event {
402 spoken.push_str(text);
403 spoken.push('\n');
404 }
405 log.append(&event)?;
406 }
407 let outcome = session.finish();
408 if outcome.exit_code != Some(0) {
409 let code = outcome
410 .exit_code
411 .map(|c| c.to_string())
412 .unwrap_or_else(|| "no exit code".to_string());
413 let reason = match &outcome.diagnostics {
417 Some(d) => format!("reviewer could not run (exit {code}): {d}"),
418 None => format!("reviewer could not run (exit {code}), and said nothing"),
419 };
420 log.append(&Event::Error {
421 message: reason.clone(),
422 raw: None,
423 })?;
424 return Ok((Verdict::Reject { reason }, outcome.usage));
425 }
426
427 Ok((review::parse_verdict(&spoken, &marker), outcome.usage))
428}
429
430fn finish(
431 log: RunLog,
432 mut record: RunRecord,
433 outcome: Outcome,
434 token: Option<MergeToken>,
435 refusal: Option<Refusal>,
436 diff: String,
437) -> Result<RunReport> {
438 record.finished_at = Some(now_rfc3339());
439 record.outcome = Some(outcome);
440 log.write_record(&record)?;
441 Ok(RunReport {
442 record,
443 token,
444 refusal,
445 diff,
446 })
447}
448
449pub fn replay(records_root: &Path, run_id: &str) -> Result<(RunRecord, Vec<Event>)> {
451 let dir = records_root.join("runs").join(run_id);
452 let record_text = std::fs::read_to_string(dir.join("record.json"))
453 .map_err(|e| Error::Other(format!("no run {run_id:?}: {e}")))?;
454 let record: RunRecord = serde_json::from_str(&record_text)
455 .map_err(|e| Error::Other(format!("run record is unreadable: {e}")))?;
456 let events = crate::record::read_events(&dir)?;
457 Ok((record, events))
458}