use crate::author;
use crate::gate::{self, MergeToken, Refusal};
use crate::progress::{Phase, Watcher};
use crate::record::RunLog;
use crate::review;
use crate::route::Routing;
use crate::worktree;
use crate::{Error, Result};
use ostraka_adapter::{AdapterOutcome, VendorAdapter};
use ostraka_core::clock::now_rfc3339;
use ostraka_core::config::Config;
use ostraka_core::gate::{Approval, Verdict};
use ostraka_core::identity::ActorId;
use ostraka_core::record::{Event, Outcome, RunRecord};
use ostraka_core::task::TaskSpec;
use std::path::Path;
pub struct RunReport {
pub record: RunRecord,
pub token: Option<MergeToken>,
pub refusal: Option<Refusal>,
pub diff: String,
}
impl RunReport {
pub fn approved(&self) -> bool {
self.token.is_some()
}
}
pub struct Places<'a> {
pub repo: &'a Path,
pub worktrees: &'a Path,
pub records: &'a Path,
pub name: &'a str,
pub notes: Option<&'a Path>,
pub skills: Option<&'a Path>,
}
pub fn run_task(
places: &Places<'_>,
config: &Config,
routing: &Routing,
task: &TaskSpec,
reviewer_identity: &ActorId,
watcher: Option<Box<dyn Watcher>>,
) -> Result<RunReport> {
let repo = places.repo;
let run_id = format!("{}-{}", task.id, now_rfc3339().replace([':', '-'], ""));
let mut log = RunLog::create(places.records, &run_id)?.watched_by(watcher);
log.enter(Phase::Isolating);
let mut record = RunRecord {
run_id: run_id.clone(),
task_id: task.id.clone(),
prompt: task.prompt.clone(),
author: task.author.clone(),
adapter: routing.author.id().to_string(),
repository: places.name.to_string(),
started_at: now_rfc3339(),
finished_at: None,
checks: Vec::new(),
approval: None,
usage: Vec::new(),
outcome: None,
};
let wt = worktree::create(repo, places.worktrees, &run_id, &task.base_ref)?;
log.enter(Phase::Preparing);
match worktree::prepare(
repo,
wt.path(),
&config.worktree,
places.notes,
places.skills,
config.gate.timeout_secs.map(std::time::Duration::from_secs),
) {
Ok(steps) => {
for step in steps {
log.append(&Event::Message {
text: format!("prepared: {step}"),
raw: None,
})?;
}
}
Err(problem) => {
return finish(
log,
record,
Outcome::Failed,
None,
Some(Refusal::SetupFailed {
step: problem.step,
reason: problem.reason,
}),
String::new(),
);
}
}
log.enter(Phase::Authoring);
let authoring = TaskSpec {
prompt: author::author_prompt(
&task.prompt,
worktree::linked(wt.path(), "notes"),
worktree::linked(wt.path(), "skills"),
),
..task.clone()
};
let author = drive(routing.author.as_ref(), &authoring, wt.path(), &mut log)?;
record.usage.extend(author.usage.clone());
let touched = worktree::touched_paths(wt.path())?;
log.append(&Event::Finished {
exit_code: author.exit_code,
files_touched: touched.clone(),
})?;
if author.interrupted {
return finish(
log,
record,
Outcome::Rejected,
None,
Some(Refusal::Interrupted),
String::new(),
);
}
if author.timed_out {
return finish(
log,
record,
Outcome::Rejected,
None,
Some(Refusal::TimedOut {
after_secs: config.policy.timeout_secs.unwrap_or_default(),
}),
String::new(),
);
}
if author.exit_code != Some(0) {
let refusal = Refusal::AuthorFailed {
code: author
.exit_code
.map(|c| c.to_string())
.unwrap_or_else(|| "no exit code".to_string()),
diagnostics: author.diagnostics.clone(),
};
return finish(
log,
record,
Outcome::Rejected,
None,
Some(refusal),
String::new(),
);
}
if touched.is_empty() {
return finish(
log,
record,
Outcome::Rejected,
None,
Some(Refusal::NoChange),
String::new(),
);
}
if !config.policy.permits(&touched) {
return finish(
log,
record,
Outcome::Rejected,
None,
Some(Refusal::PolicyViolation {
reason: format!("policy forbids writing outside declared paths: {touched:?}"),
}),
String::new(),
);
}
log.enter(Phase::Gating);
let passed = match gate::run_checks(&config.gate, wt.path(), &mut |record| {
log_checked(&mut log, record)
}) {
Ok(p) => p,
Err(refusal) => {
if let Refusal::ChecksFailed { records, .. } = &refusal {
record.checks.clone_from(records);
}
return finish(
log,
record,
Outcome::Rejected,
None,
Some(refusal),
String::new(),
);
}
};
record.checks = passed.records().to_vec();
log.enter(Phase::Reviewing);
let diff = worktree::diff(wt.path())?;
let (verdict, reviewer_usage) =
collect_verdict(routing.reviewer.as_ref(), task, &diff, wt.path(), &mut log)?;
record.usage.extend(reviewer_usage);
let approval = Approval {
reviewer: reviewer_identity.clone(),
verdict: verdict.clone(),
};
record.approval = Some(approval.clone());
match gate::evaluate(
passed,
&task.author,
&approval,
config.gate.review.must_differ_from_author,
) {
Ok(token) => {
let message = format!(
"{}\n\nRun: {run_id}\nAuthored-by: {} ({})\nReviewed-by: {} ({})",
task.prompt,
task.author,
routing.author.id(),
reviewer_identity,
routing.reviewer.id(),
);
worktree::commit(wt.path(), &message, &task.author)?;
if let Err(e) = worktree::release(repo, &wt) {
log.append(&Event::Error {
message: format!("the worktree could not be removed: {e}"),
raw: None,
})?;
}
finish(log, record, Outcome::Approved, Some(token), None, diff)
}
Err(refusal) => finish(log, record, Outcome::Rejected, None, Some(refusal), diff),
}
}
fn log_checked(log: &mut RunLog, record: &ostraka_core::gate::CheckRecord) {
log.checked(record);
}
fn drive(
adapter: &dyn VendorAdapter,
task: &TaskSpec,
worktree: &Path,
log: &mut RunLog,
) -> Result<AdapterOutcome> {
let mut session = adapter.launch(task, worktree)?;
while let Some(event) = session.next_event() {
log.append(&event)?;
}
let outcome = session.finish();
if let Some(diagnostics) = &outcome.diagnostics {
log.append(&Event::Error {
message: format!("author exited abnormally: {diagnostics}"),
raw: None,
})?;
}
Ok(outcome)
}
type Reviewed = (Verdict, Option<ostraka_core::record::TokenUsage>);
fn collect_verdict(
reviewer: &dyn VendorAdapter,
task: &TaskSpec,
diff: &str,
worktree: &Path,
log: &mut RunLog,
) -> Result<Reviewed> {
let marker = review::verdict_marker(&task.id);
let review_task = TaskSpec {
id: format!("{}-review", task.id),
prompt: review::review_prompt(&task.prompt, diff, &marker),
adapter: reviewer.id().to_string(),
author: task.author.clone(),
base_ref: task.base_ref.clone(),
model: None,
};
let mut session = match reviewer.launch(&review_task, worktree) {
Ok(s) => s,
Err(e) => {
return Ok((
Verdict::Reject {
reason: format!("reviewer could not be launched: {e}"),
},
None,
));
}
};
let mut spoken = String::new();
while let Some(event) = session.next_event() {
if let Event::Message { text, .. } = &event {
spoken.push_str(text);
spoken.push('\n');
}
log.append(&event)?;
}
let outcome = session.finish();
if outcome.exit_code != Some(0) {
let code = outcome
.exit_code
.map(|c| c.to_string())
.unwrap_or_else(|| "no exit code".to_string());
let reason = match &outcome.diagnostics {
Some(d) => format!("reviewer could not run (exit {code}): {d}"),
None => format!("reviewer could not run (exit {code}), and said nothing"),
};
log.append(&Event::Error {
message: reason.clone(),
raw: None,
})?;
return Ok((Verdict::Reject { reason }, outcome.usage));
}
Ok((review::parse_verdict(&spoken, &marker), outcome.usage))
}
fn finish(
log: RunLog,
mut record: RunRecord,
outcome: Outcome,
token: Option<MergeToken>,
refusal: Option<Refusal>,
diff: String,
) -> Result<RunReport> {
record.finished_at = Some(now_rfc3339());
record.outcome = Some(outcome);
log.write_record(&record)?;
Ok(RunReport {
record,
token,
refusal,
diff,
})
}
pub fn replay(records_root: &Path, run_id: &str) -> Result<(RunRecord, Vec<Event>)> {
let dir = records_root.join("runs").join(run_id);
let record_text = std::fs::read_to_string(dir.join("record.json"))
.map_err(|e| Error::Other(format!("no run {run_id:?}: {e}")))?;
let record: RunRecord = serde_json::from_str(&record_text)
.map_err(|e| Error::Other(format!("run record is unreadable: {e}")))?;
let events = crate::record::read_events(&dir)?;
Ok((record, events))
}