use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::thread::JoinHandle;
use std::time::Duration;
use super::{Facts, facts, row};
use crate::app::Snapshot;
use crate::board::{BoardRow, Column};
use crate::boundary::Action;
use crate::boundary::dispatch::{self, Deps};
use crate::git_tree::{AgentState, GitTree};
use crate::opslog;
use crate::start::{BallSpec, Payload};
use crate::state::SnapshotCell;
use crate::ui_state::{Clock, UiState, content_hash};
pub struct PilotCtx {
pub deps: Deps,
pub cell: SnapshotCell,
pub clock: Arc<dyn Clock>,
pub ui_path: PathBuf,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Move {
Reap {
row: BoardRow,
claimant: String,
since: String,
},
Spawn { row: BoardRow },
}
impl PilotCtx {
pub fn period(&self) -> Duration {
crate::state::latest_snapshot(&self.cell).cadence.full_sweep
}
pub fn pass(&self) -> bool {
let snapshot = crate::state::latest_snapshot(&self.cell);
if snapshot.fleet.is_empty() {
return false;
}
let ts = self.clock.stamp();
let now: i64 = ts.parse().unwrap_or(0);
let mut ui = UiState::open(self.ui_path.clone());
let board = crate::board::build(&snapshot, &ui, now);
for fleet in &board.fleet {
if let Some(one) = plan(&snapshot, fleet, &board.rows, now) {
return self.fire(&snapshot, &mut ui, &ts, fleet, &one);
}
}
false
}
fn fire(
&self,
snapshot: &Arc<Snapshot>,
ui: &mut UiState,
ts: &str,
fleet: &Facts,
one: &Move,
) -> bool {
let deps = self.deps(snapshot, ts);
let entry = match one {
Move::Reap {
row,
claimant,
since,
} => {
let release = Action::Release {
project: row.project.clone(),
id: row.id.clone(),
name: claimant.clone(),
};
if !released(dispatch::dispatch(&deps, ui, ts, &release)) {
return false;
}
row::reaped(ts.to_owned(), &fleet.workspace, &row.id, claimant, since)
}
Move::Spawn { row } => {
let Some(conversation) = Self::birth(&deps, ui, ts, fleet, row) else {
return false;
};
row::spawned(ts.to_owned(), &fleet.workspace, &row.id, &conversation)
}
};
let _ = opslog::append(&deps.state_root, &entry);
true
}
fn birth(deps: &Deps, ui: &UiState, ts: &str, fleet: &Facts, row: &BoardRow) -> Option<String> {
let ball = deps
.snapshot
.balls_by_project
.get(&row.project)?
.iter()
.find(|b| b.id == row.id)?;
let payload = Payload::Ball {
project: row.project.clone(),
ball: BallSpec::Existing {
id: ball.id.clone(),
title: ball.title.clone(),
body: ball.body.clone(),
join: row.state,
},
};
let prepared = dispatch::prepare(deps, ts, &fleet.workspace, &payload).ok()?;
let goal = prepared.goal.clone();
dispatch::prompt(deps, ui, ts, &prepared, &goal).ok()
}
fn deps(&self, snapshot: &Arc<Snapshot>, ts: &str) -> Deps {
Deps {
snapshot: Arc::clone(snapshot),
mint_seed: content_hash(ts.as_bytes()),
..self.deps.clone()
}
}
}
fn released(reply: Result<crate::boundary::reply::Reply, String>) -> bool {
matches!(reply, Ok(crate::boundary::reply::Reply::Outcome(o)) if o.ok())
}
pub fn plan(snap: &Snapshot, fleet: &Facts, rows: &[BoardRow], now: i64) -> Option<Move> {
reap(snap, fleet, rows, now).or_else(|| spawn(fleet, rows))
}
fn reap(snap: &Snapshot, fleet: &Facts, rows: &[BoardRow], now: i64) -> Option<Move> {
let lease = i64::try_from(fleet.lease?.as_secs()).ok()?;
let tree = snap.trees.get(&fleet.workspace)?;
rows.iter().filter(|r| held_here(r, fleet)).find_map(|row| {
let claimant = row.claimant.clone()?;
let idle = quiet_for(tree, row, now)?;
(idle >= lease).then(move || Move::Reap {
row: row.clone(),
claimant,
since: format!("lease expired {} ago", facts::secs_label(idle - lease)),
})
})
}
fn spawn(fleet: &Facts, rows: &[BoardRow]) -> Option<Move> {
if !fleet.has_room() {
return None;
}
rows.iter()
.find(|r| r.column == Column::Ready && r.project == fleet.project)
.map(|row| Move::Spawn { row: row.clone() })
}
fn held_here(row: &BoardRow, fleet: &Facts) -> bool {
row.column == Column::Claimed && row.workspace.as_deref() == Some(&fleet.workspace)
}
fn quiet_for(tree: &GitTree, row: &BoardRow, now: i64) -> Option<i64> {
let mut newest: Option<i64> = None;
for agent in &tree.agents {
let root = crate::nav::convs::root_of(&tree.agents, &agent.agent_id)?;
if !row.drones.iter().any(|d| d.root_id == root) {
continue;
}
if matches!(agent.state, AgentState::Live | AgentState::InFlight) {
return None;
}
newest = Some(newest.unwrap_or(i64::MIN).max(agent.last_action_unix));
}
newest.map(|last| now.saturating_sub(last))
}
pub struct Pilot {
stop: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
}
impl Pilot {
pub fn spawn(ctx: PilotCtx) -> Self {
let stop = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&stop);
let handle = std::thread::spawn(move || {
while !flag.load(Ordering::Relaxed) {
ctx.pass();
std::thread::park_timeout(ctx.period());
}
});
Self {
stop,
handle: Some(handle),
}
}
}
impl Drop for Pilot {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(handle) = self.handle.take() {
handle.thread().unpark();
let _ = handle.join();
}
}
}
#[cfg(test)]
mod tests;