//! The job registry: accepts a [`JobKind`], runs `once`, streams its output.
use std::collections::{HashMap, VecDeque};
use std::path::PathBuf;
use std::process::Stdio;
use std::sync::Arc;
use std::time::Duration;
use jiff::Timestamp;
use tokio::io::{AsyncBufReadExt as _, BufReader};
use tokio::process::Command;
use tokio::sync::{Mutex, broadcast, mpsc};
use twice_core::{JobId, JobKind, JobOutcome, JobView, LogLine, LogStream};
use crate::ansi;
use crate::error::ApiError;
use crate::once_cli;
/// Jobs retained in memory before the oldest finished one is dropped.
const MAX_RETAINED_JOBS: usize = 50;
/// Log lines retained per job.
const MAX_LOG_LINES: usize = 2_000;
/// Hard ceiling on a single `once` invocation.
const JOB_TIMEOUT: Duration = Duration::from_secs(30 * 60);
/// Buffered events per subscriber before it is considered lagging.
const EVENT_BUFFER: usize = 256;
/// A live update about a job.
#[derive(Clone, Debug)]
pub(crate) enum JobEvent {
/// The child produced a line of output.
Line(LogLine),
/// The child finished; no further events will be sent.
Finished(JobOutcome),
}
struct Slot {
view: JobView,
events: broadcast::Sender<JobEvent>,
}
#[derive(Default)]
struct Registry {
slots: HashMap<JobId, Slot>,
order: VecDeque<JobId>,
}
/// Owns every job this server has run since start-up.
#[derive(Debug)]
pub(crate) struct JobRegistry {
once_binary: PathBuf,
namespace: String,
inner: Mutex<Registry>,
}
impl std::fmt::Debug for Registry {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Registry")
.field("slots", &self.slots.len())
.field("order", &self.order.len())
.finish()
}
}
impl JobRegistry {
/// Creates an empty registry bound to a specific `once` binary and namespace.
pub(crate) fn new(once_binary: PathBuf, namespace: String) -> Arc<Self> {
Arc::new(Self {
once_binary,
namespace,
inner: Mutex::new(Registry::default()),
})
}
/// Accepts a job, rejecting it when the same application already has one
/// in flight, and starts the child process in the background.
pub(crate) async fn submit(self: &Arc<Self>, kind: JobKind) -> Result<JobView, ApiError> {
let mut guard = self.inner.lock().await;
if let Some(running) = guard
.slots
.values()
.find(|slot| !slot.view.outcome.is_terminal() && slot.view.kind.host() == kind.host())
{
return Err(ApiError::JobInFlight(running.view.kind.label()));
}
let view = JobView {
id: JobId::new(),
kind: kind.clone(),
outcome: JobOutcome::Running,
started_at: Timestamp::now(),
finished_at: None,
log: Vec::new(),
};
let (events, _) = broadcast::channel(EVENT_BUFFER);
let id = view.id;
guard.order.push_back(id);
guard.slots.insert(
id,
Slot {
view: view.clone(),
events,
},
);
prune(&mut guard);
drop(guard);
let registry = Arc::clone(self);
tokio::spawn(async move { registry.run(id, kind).await });
Ok(view)
}
/// Returns a snapshot of one job.
pub(crate) async fn get(&self, id: JobId) -> Option<JobView> {
let guard = self.inner.lock().await;
let view = guard.slots.get(&id).map(|slot| slot.view.clone());
drop(guard);
view
}
/// Returns every retained job, newest first.
pub(crate) async fn list(&self) -> Vec<JobView> {
let guard = self.inner.lock().await;
let mut jobs: Vec<JobView> = guard.slots.values().map(|slot| slot.view.clone()).collect();
drop(guard);
jobs.sort_by_key(|job| std::cmp::Reverse(job.started_at));
jobs
}
/// Returns the job as it stands plus a stream of everything that happens next.
pub(crate) async fn subscribe(
&self,
id: JobId,
) -> Option<(JobView, broadcast::Receiver<JobEvent>)> {
let guard = self.inner.lock().await;
let subscription = guard
.slots
.get(&id)
.map(|slot| (slot.view.clone(), slot.events.subscribe()));
drop(guard);
subscription
}
async fn run(self: Arc<Self>, id: JobId, kind: JobKind) {
let outcome = match self.spawn_and_pump(id, &kind).await {
Ok(outcome) => outcome,
Err(err) => JobOutcome::Failed {
exit_code: None,
message: err.to_string(),
},
};
let mut guard = self.inner.lock().await;
if let Some(slot) = guard.slots.get_mut(&id) {
slot.view.outcome = outcome.clone();
slot.view.finished_at = Some(Timestamp::now());
let _ = slot.events.send(JobEvent::Finished(outcome));
}
}
async fn spawn_and_pump(
self: &Arc<Self>,
id: JobId,
kind: &JobKind,
) -> Result<JobOutcome, std::io::Error> {
let mut child = Command::new(&self.once_binary)
.arg("-n")
.arg(&self.namespace)
.args(once_cli::argv(kind))
// `once` opens a bubbletea TUI when it has a terminal; a null stdin
// and a dumb terminal keep it in non-interactive mode.
.env("NO_COLOR", "1")
.env("TERM", "dumb")
.env("CI", "1")
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true)
.spawn()?;
let (tx, mut rx) = mpsc::channel::<(LogStream, String)>(EVENT_BUFFER);
if let Some(stdout) = child.stdout.take() {
spawn_reader(stdout, LogStream::Stdout, tx.clone());
}
if let Some(stderr) = child.stderr.take() {
spawn_reader(stderr, LogStream::Stderr, tx.clone());
}
drop(tx);
let mut seq = 0_u64;
while let Some((stream, text)) = rx.recv().await {
seq = seq.saturating_add(1);
let line = LogLine { seq, stream, text };
self.append(id, line).await;
}
let Ok(status) = tokio::time::timeout(JOB_TIMEOUT, child.wait()).await else {
let _ = child.kill().await;
return Ok(JobOutcome::Failed {
exit_code: None,
message: format!("once exceeded the {}s job timeout", JOB_TIMEOUT.as_secs()),
});
};
let status = status?;
Ok(if status.success() {
JobOutcome::Succeeded
} else {
JobOutcome::Failed {
exit_code: status.code(),
message: format!("once exited with {status}"),
}
})
}
async fn append(&self, id: JobId, line: LogLine) {
let mut guard = self.inner.lock().await;
if let Some(slot) = guard.slots.get_mut(&id) {
if slot.view.log.len() >= MAX_LOG_LINES {
slot.view.log.remove(0);
}
slot.view.log.push(line.clone());
let _ = slot.events.send(JobEvent::Line(line));
}
}
}
fn spawn_reader<R>(pipe: R, stream: LogStream, tx: mpsc::Sender<(LogStream, String)>)
where
R: tokio::io::AsyncRead + Unpin + Send + 'static,
{
tokio::spawn(async move {
let mut lines = BufReader::new(pipe).lines();
while let Ok(Some(raw)) = lines.next_line().await {
let text = ansi::strip(&raw);
if text.is_empty() {
continue;
}
if tx.send((stream, text)).await.is_err() {
break;
}
}
});
}
fn prune(registry: &mut Registry) {
while registry.order.len() > MAX_RETAINED_JOBS {
let Some(oldest) = registry.order.front().copied() else {
break;
};
let finished = registry
.slots
.get(&oldest)
.is_none_or(|slot| slot.view.outcome.is_terminal());
if !finished {
break;
}
registry.order.pop_front();
registry.slots.remove(&oldest);
}
}