use std::collections::{HashMap, HashSet};
use std::fmt;
use std::future::Future;
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use time::OffsetDateTime;
use tokio::process::Command;
use crate::engine::wave_config::read_wave_config;
use crate::lfd::id::LfdId;
use crate::lfd::types::{
Run, RunStatus, Session, SessionStatus, SessionUse, Wave, WaveStatus, LIVE_SESSION_STATUSES,
WAVE_SERVER_ENDPOINT_ENV, WAVE_SERVER_PID_ENV, WAVE_SERVER_SOURCE,
};
use crate::lfdb::{SharedStore, StoreResult};
use crate::wave::journal::WorkerOutcome;
use crate::wave::runtime::WaveRuntime;
fn is_active_run_status(status: RunStatus) -> bool {
matches!(
status,
RunStatus::Pending | RunStatus::Running | RunStatus::Waiting
)
}
async fn tmux_session_exists(session_name: &str) -> anyhow::Result<bool> {
let status = Command::new("tmux")
.args(["has-session", "-t", session_name])
.status()
.await
.map_err(|err| anyhow::anyhow!("tmux session probe failed: {err}"))?;
Ok(status.success())
}
pub const POLL_CADENCE: Duration = Duration::from_secs(10);
#[derive(Debug, Clone)]
pub struct RegistryConfig {
pub store: SharedStore,
pub wave: Wave,
pub cwd: String,
pub pid: u32,
pub force: bool,
}
#[derive(Debug)]
pub enum RegisterOutcome {
Registered(Box<Registration>),
Refused {
message: String,
},
}
#[derive(Debug, Clone)]
pub struct Registration {
store: SharedStore,
session: Session,
done: Arc<AtomicBool>,
}
impl Registration {
pub fn session_id(&self) -> &LfdId {
&self.session.id
}
pub async fn deregister(&self) {
if self.done.swap(true, Ordering::SeqCst) {
return;
}
let mut session = self.session.clone();
if !session.complete(0) {
return;
}
if let Err(err) = self.store.update_control_session(&session).await {
tracing::warn!(error = %err, "wave server deregistration failed; the next boot's pid probe will close the row");
}
}
pub fn deregister_blocking(&self) {
let registration = self.clone();
block_on(async move { registration.deregister().await });
}
}
pub async fn ensure_wave_row(
store: &SharedStore,
main_repo: &Path,
name: &str,
) -> StoreResult<Wave> {
if let Some(wave) = store.get_wave_by_name(name).await? {
return Ok(wave);
}
let mut wave = Wave::new(
LfdId::new(),
name.to_string(),
main_repo.display().to_string(),
);
if let Some(config) = read_wave_config(main_repo, name) {
if let Some(flow) = config.primary_flow {
wave.primary_flow = flow;
}
if let Some(goal) = config.goal.filter(|goal| !goal.trim().is_empty()) {
wave.goal = goal;
}
}
wave.paused = read_wave_config(main_repo, name)
.and_then(|config| config.paused)
.unwrap_or(false);
store.create_wave(&wave).await?;
tracing::info!(
wave = name,
wave_id = %wave.id,
"wave was not in the session registry; created its row"
);
Ok(wave)
}
pub async fn register(config: &RegistryConfig, endpoint: &str) -> StoreResult<RegisterOutcome> {
let wave_id = config.wave.id();
if let Some(mut existing) = live_brain_after_probe(&config.store, wave_id).await? {
if !config.force {
return Ok(RegisterOutcome::Refused {
message: format!(
"wave '{}' already has a live wave-agent session {} (source '{}')",
config.wave.name(),
existing.id,
existing.source
),
});
}
if let Some(pid) = wave_server_pid(&existing) {
if process_alive(pid).await {
let _ = Command::new("kill").arg(pid.to_string()).status().await;
}
} else if existing.is_tmux_backed() && !existing.tmux_name.is_empty() {
let _ = Command::new("tmux")
.args(["kill-session", "-t", &existing.tmux_name])
.status()
.await;
}
if existing.cancel() {
config.store.update_control_session(&existing).await?;
}
}
let session_id = LfdId::new();
let now = OffsetDateTime::now_utc();
let session = Session {
id: session_id,
wave_id: wave_id.clone(),
run_id: None,
parent_session_id: None,
session_use: SessionUse::WaveAgent,
step: "mind".to_string(),
agent: "lf".to_string(),
cwd: config.cwd.clone(),
argv: vec![
"lf".to_string(),
"wave".to_string(),
config.wave.name().clone(),
],
env: std::collections::BTreeMap::from([
(WAVE_SERVER_ENDPOINT_ENV.to_string(), endpoint.to_string()),
(WAVE_SERVER_PID_ENV.to_string(), config.pid.to_string()),
]),
source: WAVE_SERVER_SOURCE.to_string(),
tmux_name: String::new(),
status: SessionStatus::Running,
attached_at: None,
started_at: Some(now),
completed_at: None,
created_at: now,
completion_token: None,
};
config.store.register_session(&session).await?;
Ok(RegisterOutcome::Registered(Box::new(Registration {
store: config.store.clone(),
session,
done: Arc::new(AtomicBool::new(false)),
})))
}
pub async fn live_brain_after_probe(
store: &SharedStore,
wave_id: &LfdId,
) -> StoreResult<Option<Session>> {
let sessions = store
.list_control_sessions(Some(wave_id), Some(LIVE_SESSION_STATUSES))
.await?;
let mut live = None;
for mut session in sessions {
if session.session_use != SessionUse::WaveAgent {
continue;
}
if session.source == WAVE_SERVER_SOURCE {
let alive = match wave_server_pid(&session) {
Some(pid) => process_alive(pid).await,
None => false,
};
if !alive {
if session.complete(1) {
store.update_control_session(&session).await?;
tracing::info!(session_id = %session.id, "closed crashed wave server session");
}
continue;
}
}
if live.is_none() {
live = Some(session);
}
}
Ok(live)
}
fn wave_server_pid(session: &Session) -> Option<u32> {
if session.source != WAVE_SERVER_SOURCE {
return None;
}
session
.env
.get(WAVE_SERVER_PID_ENV)
.and_then(|pid| pid.parse().ok())
}
pub(crate) async fn wave_server_endpoint(
store: &SharedStore,
wave_id: &LfdId,
) -> anyhow::Result<Option<String>> {
let Some(session) = store.live_wave_agent_session(wave_id).await? else {
return Ok(None);
};
Ok(session
.env
.get(WAVE_SERVER_ENDPOINT_ENV)
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty()))
}
pub(crate) async fn process_alive(pid: u32) -> bool {
Command::new("kill")
.args(["-0", &pid.to_string()])
.status()
.await
.is_ok_and(|status| status.success())
}
fn block_on(future: impl Future<Output = ()> + Send + 'static) {
if tokio::runtime::Handle::try_current().is_ok() {
let _ = std::thread::spawn(move || block_on_new_runtime(future)).join();
return;
}
block_on_new_runtime(future);
}
fn block_on_new_runtime(future: impl Future<Output = ()>) {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime always builds")
.block_on(future);
}
pub struct StoreObserver {
runtime: Arc<WaveRuntime>,
store: SharedStore,
wave_id: LfdId,
terminal_cutoff: std::sync::Mutex<Option<i64>>,
seen: std::sync::Mutex<HashSet<String>>,
#[cfg(test)]
liveness_override: std::sync::Mutex<Option<LivenessProbe>>,
}
#[cfg(test)]
type LivenessProbe = Box<dyn Fn(&Session) -> bool + Send + Sync>;
impl fmt::Debug for StoreObserver {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("StoreObserver")
.field("wave_id", &self.wave_id)
.finish()
}
}
const TERMINAL_CUTOFF_MARGIN: time::Duration = time::Duration::seconds(60);
impl StoreObserver {
pub fn new(runtime: Arc<WaveRuntime>, store: SharedStore, wave_id: LfdId) -> Self {
Self {
runtime,
store,
wave_id,
terminal_cutoff: std::sync::Mutex::new(None),
seen: std::sync::Mutex::new(HashSet::new()),
#[cfg(test)]
liveness_override: std::sync::Mutex::new(None),
}
}
pub async fn run(self: Arc<Self>, cadence: Duration) {
loop {
self.poll_once().await;
tokio::time::sleep(cadence).await;
}
}
fn known(&self, run_id: &str) -> bool {
if self
.seen
.lock()
.expect("observer seen cache poisoned")
.contains(run_id)
{
return true;
}
if self.runtime.worker_known(run_id) {
self.mark_seen(run_id);
return true;
}
false
}
fn mark_seen(&self, run_id: &str) {
self.seen
.lock()
.expect("observer seen cache poisoned")
.insert(run_id.to_string());
}
async fn worker_alive(&self, session: &Session) -> bool {
#[cfg(test)]
if let Some(probe) = self
.liveness_override
.lock()
.expect("observer probe override poisoned")
.as_ref()
{
return probe(session);
}
tmux_session_exists(&session.tmux_name)
.await
.unwrap_or(false)
}
#[cfg(test)]
fn set_liveness_probe(&self, probe: impl Fn(&Session) -> bool + Send + Sync + 'static) {
*self
.liveness_override
.lock()
.expect("observer probe override poisoned") = Some(Box::new(probe));
}
pub async fn poll_once(&self) {
self.runtime.sweep_dead_children();
let poll_started = OffsetDateTime::now_utc();
let cutoff = *self
.terminal_cutoff
.lock()
.expect("observer cutoff poisoned");
let sessions = match cutoff {
None => {
self.store
.list_control_sessions(Some(&self.wave_id), None)
.await
}
Some(since) => {
self.store
.list_recent_control_sessions(&self.wave_id, since)
.await
}
};
let sessions = match sessions {
Ok(sessions) => sessions,
Err(err) => {
tracing::debug!(error = %err, "wave observer session read failed");
return;
}
};
*self
.terminal_cutoff
.lock()
.expect("observer cutoff poisoned") =
Some((poll_started - TERMINAL_CUTOFF_MARGIN).unix_timestamp());
let mut workers: Vec<(Session, String)> = sessions
.into_iter()
.filter(|session| session.session_use == SessionUse::Worker)
.filter_map(|session| {
let run_id = session.run_id.as_ref().map(ToString::to_string)?;
Some((session, run_id))
})
.collect();
let mut gone: HashSet<String> = HashSet::new();
for (session, run_id) in workers.iter_mut() {
if session.status != SessionStatus::Running
|| !session.is_tmux_backed()
|| session.tmux_name.is_empty()
|| self.worker_alive(session).await
{
continue;
}
if session.complete(1) {
if let Err(err) = self.store.update_control_session(session).await {
tracing::debug!(
error = %err,
session_id = %session.id,
"failed to close dead worker session; will retry next poll"
);
continue;
}
tracing::info!(
session_id = %session.id,
run_id,
"closed dead worker session (process gone)"
);
}
gone.insert(run_id.clone());
}
let in_flight: HashSet<String> = self
.runtime
.in_flight_workers()
.into_iter()
.map(|worker| worker.run_id)
.collect();
let needs_runs = workers.iter().any(|(session, run_id)| {
!self.known(run_id) || (session.status.is_terminal() && in_flight.contains(run_id))
});
if !needs_runs {
return;
}
let runs: HashMap<String, Run> = match self.store.list_runs(Some(&self.wave_id), None).await
{
Ok(runs) => runs
.into_iter()
.map(|run| (run.id.to_string(), run))
.collect(),
Err(err) => {
tracing::debug!(error = %err, "wave observer run read failed");
return;
}
};
for (session, run_id) in workers {
if !self.known(&run_id) {
let Some(run) = runs.get(&run_id) else {
tracing::warn!(run_id, "observed worker session with no run row; skipped");
continue;
};
if self.runtime.journal_run_observed(
&run_id,
session.id.as_str(),
&run.flow,
run.task.as_deref().unwrap_or(""),
) {
tracing::info!(
run_id,
session_id = %session.id,
flow = run.flow,
"worker dispatched"
);
}
self.mark_seen(&run_id);
}
if !session.status.is_terminal() {
continue;
}
let outcome = if session.status == SessionStatus::Succeeded {
WorkerOutcome::Completed
} else {
WorkerOutcome::Failed
};
let summary = if gone.contains(&run_id) {
"process gone".to_string()
} else {
let mut summary = format!("session {}", session.status.as_str());
if let Some(run) = runs.get(&run_id) {
if let Some(pr) = &run.pr {
summary.push_str(&format!("; pr {}", pr.url));
}
if let Some(error) = &run.error {
summary.push_str(&format!("; error: {error}"));
}
}
summary
};
if self
.runtime
.journal_run_completed(&run_id, outcome, &summary)
{
tracing::info!(run_id, outcome = outcome.name(), summary, "worker finished");
}
self.close_run_and_maybe_idle_wave(&runs, &run_id, outcome)
.await;
}
}
async fn close_run_and_maybe_idle_wave(
&self,
runs: &HashMap<String, Run>,
run_id: &str,
outcome: WorkerOutcome,
) {
let Some(run) = runs.get(run_id) else {
return;
};
if is_active_run_status(run.status) {
let mut run = run.clone();
run.status = match outcome {
WorkerOutcome::Completed => RunStatus::Completed,
WorkerOutcome::Failed => RunStatus::Failed,
};
run.ended_at = Some(OffsetDateTime::now_utc());
if let Err(err) = self.store.update_run(&run).await {
tracing::debug!(run_id, error = %err, "failed to close run row; retry next poll");
return;
}
}
match self.store.count_active_runs(&self.wave_id).await {
Ok(0) => self.reset_wave_repos_to_idle().await,
Ok(_) => {}
Err(err) => {
tracing::debug!(error = %err, "active-run count failed; wave status unchanged");
}
}
}
async fn reset_wave_repos_to_idle(&self) {
let Ok(Some(mut wave)) = self.store.get_wave(&self.wave_id).await else {
return;
};
let mut changed = false;
for repo in &mut wave.repos {
if repo.status == WaveStatus::Running {
repo.status = WaveStatus::Idle;
changed = true;
}
}
if !changed {
return;
}
if let Err(err) = self.store.update_wave(&wave).await {
tracing::debug!(wave_id = %self.wave_id, error = %err, "failed to reset wave status to idle");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::BTreeMap;
use crate::lfd::types::{
PullRequest, RepoWork, RunStackStatus, RunStatus, WaveStatus, TMUX_TERMINAL_SOURCE,
};
use crate::lfdb::{open_store, StorageConfig};
use crate::wave::journal::{journal_path, EventKind, Journal};
async fn temp_store(tmp: &std::path::Path) -> SharedStore {
Arc::new(
open_store(&StorageConfig::sqlite(tmp.join("lfd.db")))
.await
.expect("open sqlite store"),
)
}
fn make_wave(name: &str) -> Wave {
Wave {
id: LfdId::new(),
name: name.to_string(),
primary_flow: "ship-roadmap".to_string(),
goal: "ship-roadmap".to_string(),
metrics: Vec::new(),
repos: vec![RepoWork {
repo: "/tmp/repo".to_string(),
worktree: String::new(),
branch: String::new(),
status: WaveStatus::Idle,
iteration: 0,
cycle_start_iteration: 0,
position: 0,
}],
direction: Vec::new(),
area: Vec::new(),
paused: false,
created_at: Some(OffsetDateTime::now_utc()),
workers: 2,
parent_wave_id: None,
}
}
fn registry_config(store: SharedStore, wave: Wave, force: bool) -> RegistryConfig {
RegistryConfig {
store,
wave,
cwd: "/tmp/repo.ship".to_string(),
pid: std::process::id(),
force,
}
}
fn server_session(wave: &Wave, pid: u32) -> Session {
Session {
id: LfdId::new(),
wave_id: wave.id().clone(),
run_id: None,
parent_session_id: None,
session_use: SessionUse::WaveAgent,
step: "mind".to_string(),
agent: "lf".to_string(),
cwd: "/tmp/repo.ship".to_string(),
argv: vec!["lf".to_string(), "wave".to_string(), wave.name().clone()],
env: BTreeMap::from([
(
WAVE_SERVER_ENDPOINT_ENV.to_string(),
"127.0.0.1:9".to_string(),
),
(WAVE_SERVER_PID_ENV.to_string(), pid.to_string()),
]),
source: WAVE_SERVER_SOURCE.to_string(),
tmux_name: String::new(),
status: SessionStatus::Running,
attached_at: None,
started_at: Some(OffsetDateTime::now_utc()),
completed_at: None,
created_at: OffsetDateTime::now_utc(),
completion_token: None,
}
}
fn worker_session(wave: &Wave, run_id: &LfdId) -> Session {
Session {
id: LfdId::new(),
wave_id: wave.id().clone(),
run_id: Some(run_id.clone()),
parent_session_id: None,
session_use: SessionUse::Worker,
step: "dispatch:implement".to_string(),
agent: "lf".to_string(),
cwd: "/tmp/repo.ship".to_string(),
argv: Vec::new(),
env: BTreeMap::new(),
source: TMUX_TERMINAL_SOURCE.to_string(),
tmux_name: "lf-x".to_string(),
status: SessionStatus::Running,
attached_at: None,
started_at: Some(OffsetDateTime::now_utc()),
completed_at: None,
created_at: OffsetDateTime::now_utc(),
completion_token: None,
}
}
fn make_run(wave: &Wave, flow: &str, task: &str) -> Run {
Run {
id: LfdId::new(),
wave_id: wave.id().clone(),
repo: "/tmp/repo".to_string(),
flow: flow.to_string(),
task: Some(task.to_string()),
direction: Vec::new(),
area: Vec::new(),
iteration: 0,
step_index: 0,
status: RunStatus::Running,
worktree: "/tmp/repo.ship".to_string(),
branch: "ship-branch".to_string(),
started_at: Some(OffsetDateTime::now_utc()),
ended_at: None,
error: None,
flow_parents: Vec::new(),
execution_cursor: None,
parent_run_id: None,
parent_pr_number: None,
stack_position: 0,
stack_group_id: wave.id().to_string(),
stack_status: RunStackStatus::Active,
lineage_inferred: false,
target_branch: "main".to_string(),
repair_of: None,
pr: None,
}
}
#[tokio::test]
async fn boot_with_no_wave_row_creates_it_and_registers() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let repo = tmp.path().join("repo");
let goal_dir = repo.join("wave/ship");
std::fs::create_dir_all(&goal_dir).expect("wave dir");
std::fs::write(
goal_dir.join("GOAL.md"),
"---\ngoal: keep shipping\n---\nShip.\n",
)
.expect("GOAL.md");
let wave = ensure_wave_row(&store, &repo, "ship")
.await
.expect("row created");
let stored = store
.get_wave_by_name("ship")
.await
.expect("lookup")
.expect("row exists");
assert_eq!(stored.id, wave.id);
assert_eq!(stored.goal, "keep shipping", "goal from GOAL.md");
assert_eq!(stored.repo(), repo.display().to_string());
let config = registry_config(store.clone(), wave.clone(), false);
let RegisterOutcome::Registered(_reg) =
register(&config, "127.0.0.1:4").await.expect("register")
else {
panic!("first boot registers");
};
let again = ensure_wave_row(&store, &repo, "ship")
.await
.expect("idempotent");
assert_eq!(again.id, wave.id, "ensure reuses the existing row");
let outcome = register(®istry_config(store, again, false), "127.0.0.1:5")
.await
.expect("attempt");
assert!(
matches!(outcome, RegisterOutcome::Refused { .. }),
"second boot refused"
);
}
#[tokio::test]
async fn ensure_wave_row_without_goal_md_uses_defaults() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = ensure_wave_row(&store, tmp.path(), "ship")
.await
.expect("row created");
assert_eq!(wave.goal, "ship-roadmap");
assert_eq!(wave.primary_flow, "ship-roadmap");
}
#[tokio::test]
async fn register_writes_wave_server_row_and_deregisters_once() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let config = registry_config(store.clone(), wave.clone(), false);
let RegisterOutcome::Registered(registration) =
register(&config, "127.0.0.1:4242").await.expect("register")
else {
panic!("fresh wave must register");
};
let stored = store
.get_control_session(registration.session_id())
.await
.expect("lookup")
.expect("row stored");
assert_eq!(stored.session_use, SessionUse::WaveAgent);
assert_eq!(stored.source, WAVE_SERVER_SOURCE);
assert_eq!(stored.status, SessionStatus::Running);
assert_eq!(
stored.env.get(WAVE_SERVER_ENDPOINT_ENV).map(String::as_str),
Some("127.0.0.1:4242")
);
assert_eq!(
stored.env.get(WAVE_SERVER_PID_ENV).map(String::as_str),
Some(config.pid.to_string().as_str())
);
assert_eq!(stored.cwd, "/tmp/repo.ship");
let live = store
.live_wave_agent_session(wave.id())
.await
.expect("live lookup")
.expect("brain live");
assert_eq!(live.id, *registration.session_id());
registration.deregister().await;
registration.deregister().await; let stored = store
.get_control_session(registration.session_id())
.await
.expect("lookup")
.expect("row stored");
assert_eq!(stored.status, SessionStatus::Succeeded);
}
#[tokio::test]
async fn second_brain_is_refused_naming_the_live_session() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let live = server_session(&wave, std::process::id());
store
.register_session(&live)
.await
.expect("seed live brain");
let config = registry_config(store.clone(), wave, false);
let outcome = register(&config, "127.0.0.1:1").await.expect("attempt");
let RegisterOutcome::Refused { message } = outcome else {
panic!("second brain must be refused");
};
assert!(
message.contains(live.id.as_str()),
"refusal names the live session: {message}"
);
}
#[tokio::test]
async fn dead_pid_row_is_reconciled_and_takeover_needs_no_force() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let crashed = server_session(&wave, 4_000_000);
store
.register_session(&crashed)
.await
.expect("seed crashed brain");
let config = registry_config(store.clone(), wave.clone(), false);
let RegisterOutcome::Registered(registration) =
register(&config, "127.0.0.1:2").await.expect("register")
else {
panic!("dead brain must not block a new server");
};
let old = store
.get_control_session(&crashed.id)
.await
.expect("lookup")
.expect("row kept");
assert_eq!(old.status, SessionStatus::Failed);
let live = store
.live_wave_agent_session(wave.id())
.await
.expect("live lookup")
.expect("new brain live");
assert_eq!(live.id, *registration.session_id());
}
#[tokio::test]
async fn force_takes_over_a_live_brain_by_killing_its_pid() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let mut child = std::process::Command::new("sleep")
.arg("60")
.spawn()
.expect("spawn sleep");
let live = server_session(&wave, child.id());
store
.register_session(&live)
.await
.expect("seed live brain");
let config = registry_config(store.clone(), wave.clone(), true);
let RegisterOutcome::Registered(registration) =
register(&config, "127.0.0.1:3").await.expect("register")
else {
panic!("--force must take over");
};
let old = store
.get_control_session(&live.id)
.await
.expect("lookup")
.expect("row kept");
assert_eq!(old.status, SessionStatus::Canceled);
let now_live = store
.live_wave_agent_session(wave.id())
.await
.expect("live lookup")
.expect("new brain live");
assert_eq!(now_live.id, *registration.session_id());
let status = child.wait().expect("child reaped");
assert!(!status.success(), "sleep was killed, not completed");
}
#[tokio::test]
async fn observer_journals_worker_facts_exactly_once_across_polls() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let runtime =
WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open runtime");
let observer = StoreObserver::new(runtime.clone(), store.clone(), wave.id().clone());
observer.set_liveness_probe(|_| true);
let run_1 = make_run(&wave, "implement", "wire the tail");
store.create_run(&run_1).await.expect("run-1");
let mut sess_1 = worker_session(&wave, &run_1.id);
store.register_session(&sess_1).await.expect("sess-1");
observer.poll_once().await;
observer.poll_once().await; assert!(runtime.worker_known(run_1.id.as_str()));
assert_eq!(runtime.in_flight_workers().len(), 1);
let run_2 = make_run(&wave, "design", "sketch the next item");
store.create_run(&run_2).await.expect("run-2");
store
.register_session(&worker_session(&wave, &run_2.id))
.await
.expect("sess-2");
let mut finished_run = run_1.clone();
finished_run.pr = Some(PullRequest {
url: "https://github.com/x/y/pull/7".to_string(),
number: Some(7),
state: Some("open".to_string()),
title: None,
branch: None,
});
store.update_run(&finished_run).await.expect("run-1 pr");
assert!(sess_1.complete(0));
store
.update_control_session(&sess_1)
.await
.expect("sess-1 terminal");
observer.poll_once().await;
observer.poll_once().await;
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let dispatched: Vec<String> = events
.iter()
.filter_map(|event| match &event.kind {
EventKind::RunObserved {
run_id, flow, task, ..
} => Some(format!("{run_id}/{flow}/{task}")),
_ => None,
})
.collect();
let mut sorted = dispatched.clone();
sorted.sort();
let mut expected = vec![
format!("{}/implement/wire the tail", run_1.id),
format!("{}/design/sketch the next item", run_2.id),
];
expected.sort();
assert_eq!(sorted, expected, "each dispatch journaled exactly once");
let finished: Vec<(String, String)> = events
.iter()
.filter_map(|event| match &event.kind {
EventKind::RunCompleted {
run_id, summary, ..
} => Some((run_id.clone(), summary.clone())),
_ => None,
})
.collect();
assert_eq!(finished.len(), 1, "one finish, once");
assert_eq!(finished[0].0, run_1.id.to_string());
assert!(
finished[0].1.contains("succeeded")
&& finished[0].1.contains("https://github.com/x/y/pull/7"),
"summary carries what the registry knows: {}",
finished[0].1
);
let runtime_2 =
WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("reopen runtime");
let observer_2 = StoreObserver::new(runtime_2, store, wave.id().clone());
observer_2.set_liveness_probe(|_| true);
observer_2.poll_once().await;
let (_, events_after) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
assert_eq!(
events.len(),
events_after.len(),
"restart journals nothing new"
);
}
#[tokio::test]
async fn observer_closes_dead_worker_rows_and_journals_process_gone() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let runtime =
WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open runtime");
let observer = StoreObserver::new(runtime.clone(), store.clone(), wave.id().clone());
observer.set_liveness_probe(|_| false);
let run = make_run(&wave, "implement", "wire the tail");
store.create_run(&run).await.expect("run");
let session = worker_session(&wave, &run.id);
store.register_session(&session).await.expect("session");
observer.poll_once().await;
observer.poll_once().await;
let closed = store
.get_control_session(&session.id)
.await
.expect("lookup")
.expect("row kept");
assert_eq!(
closed.status,
SessionStatus::Failed,
"a dead worker's Running row is closed by the probe"
);
assert!(
runtime.in_flight_workers().is_empty(),
"in_flight frees up when the dead worker is closed"
);
let (_, events) = Journal::open(&journal_path(tmp.path(), "ship")).expect("reopen");
let finished: Vec<(String, WorkerOutcome, String)> = events
.iter()
.filter_map(|event| match &event.kind {
EventKind::RunCompleted {
run_id,
outcome,
summary,
} => Some((run_id.clone(), *outcome, summary.clone())),
_ => None,
})
.collect();
assert_eq!(finished.len(), 1, "one finish, once");
assert_eq!(finished[0].0, run.id.to_string());
assert_eq!(finished[0].1, WorkerOutcome::Failed);
assert_eq!(finished[0].2, "process gone");
}
#[tokio::test]
async fn observer_probe_leaves_pending_workers_alone() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let wave = make_wave("ship");
store.create_wave(&wave).await.expect("seed wave");
let runtime =
WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open runtime");
let observer = StoreObserver::new(runtime, store.clone(), wave.id().clone());
observer.set_liveness_probe(|_| false);
let run = make_run(&wave, "implement", "not yet launched");
store.create_run(&run).await.expect("run");
let mut session = worker_session(&wave, &run.id);
session.status = SessionStatus::Pending;
store.register_session(&session).await.expect("session");
observer.poll_once().await;
let stored = store
.get_control_session(&session.id)
.await
.expect("lookup")
.expect("row kept");
assert_eq!(
stored.status,
SessionStatus::Pending,
"a Pending row is the dispatcher's to launch, not the probe's to close"
);
}
}