use crate::run;
use crate::workspace::Workspace;
use ostraka_core::record::Outcome;
use ostraka_runtime::progress::{Channel, Phase, Step};
use std::sync::mpsc::{self, Receiver, TryRecvError};
use std::thread::JoinHandle;
use std::time::Instant;
pub struct Finished {
pub run_id: String,
pub outcome: Option<Outcome>,
pub summary: String,
}
pub struct Session {
pub prompt: String,
pub steps: Vec<Step>,
pub phase: Option<Phase>,
pub elapsed_secs: u64,
pub finished: Option<Finished>,
pub failed: Option<String>,
pub stopping: bool,
stop: ostraka_adapter::interrupt::Stop,
started: Instant,
steps_in: Option<Receiver<Step>>,
done_in: Option<Receiver<Result<Finished, String>>>,
worker: Option<JoinHandle<()>>,
}
impl Session {
pub fn start(workspace: Workspace, args: run::Args) -> Self {
let prompt = args.prompt.clone();
let stop = ostraka_adapter::interrupt::Stop::new();
let worker_stop = stop.clone();
let (steps_out, steps_in) = mpsc::channel();
let (done_out, done_in) = mpsc::channel();
let worker = std::thread::spawn(move || {
let result = run::execute(
&workspace,
&args,
Some(Box::new(Channel(steps_out))),
&worker_stop,
)
.map(|report| Finished {
run_id: report.record.run_id.clone(),
outcome: report.record.outcome,
summary: match (&report.token, &report.refusal) {
(Some(_), _) => "approved — nothing merged".to_string(),
(None, Some(refusal)) => {
format!("rejected — {}", run::describe(refusal))
}
(None, None) => "rejected".to_string(),
},
})
.map_err(|e| e.to_string());
let _ = done_out.send(result);
});
Self {
stop,
prompt,
steps: Vec::new(),
phase: None,
elapsed_secs: 0,
finished: None,
failed: None,
stopping: false,
started: Instant::now(),
steps_in: Some(steps_in),
done_in: Some(done_in),
worker: Some(worker),
}
}
pub fn live(&self) -> bool {
self.finished.is_none() && self.failed.is_none()
}
pub fn drain(&mut self) {
let incoming: Vec<Step> = match &self.steps_in {
Some(rx) => rx.try_iter().collect(),
None => Vec::new(),
};
for step in incoming {
if let Step::Entered(phase) = step {
self.phase = Some(phase);
}
self.steps.push(step);
}
let ended = match &self.done_in {
Some(rx) => match rx.try_recv() {
Ok(result) => Some(result),
Err(TryRecvError::Empty) => None,
Err(TryRecvError::Disconnected) => {
Some(Err("the run stopped without saying why".to_string()))
}
},
None => None,
};
if let Some(result) = ended {
match result {
Ok(finished) => self.finished = Some(finished),
Err(why) => self.failed = Some(why),
}
self.phase = None;
self.stopping = false;
self.settle();
}
if self.live() {
self.elapsed_secs = self.started.elapsed().as_secs();
}
}
pub fn stop(&mut self) {
if self.live() {
self.stop.request();
self.stopping = true;
}
}
pub fn settle(&mut self) {
self.steps_in = None;
self.done_in = None;
if let Some(worker) = self.worker.take() {
let _ = worker.join();
}
}
#[cfg(test)]
pub fn recorded(prompt: &str, steps: Vec<Step>, finished: Option<Finished>) -> Self {
Self {
prompt: prompt.to_string(),
phase: steps.iter().rev().find_map(|s| match s {
Step::Entered(p) if finished.is_none() => Some(*p),
_ => None,
}),
steps,
elapsed_secs: 41,
finished,
failed: None,
stopping: false,
started: Instant::now(),
steps_in: None,
done_in: None,
worker: None,
stop: ostraka_adapter::interrupt::Stop::new(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use ostraka_core::record::Event;
#[test]
fn a_session_reads_the_phase_off_the_steps_it_is_given() {
let session = Session::recorded(
"do a thing",
vec![
Step::Entered(Phase::Authoring),
Step::Said {
phase: Phase::Authoring,
event: Event::Message {
text: "working".into(),
raw: None,
},
},
Step::Entered(Phase::Gating),
],
None,
);
assert_eq!(session.phase, Some(Phase::Gating));
assert!(session.live());
}
#[test]
fn a_finished_session_is_not_live_and_is_in_no_phase() {
let session = Session::recorded(
"do a thing",
vec![Step::Entered(Phase::Reviewing)],
Some(Finished {
run_id: "t1-20260907T000300Z".into(),
outcome: Some(Outcome::Approved),
summary: "approved — nothing merged".into(),
}),
);
assert!(!session.live());
assert_eq!(session.phase, None);
}
}