Skip to content
//! 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);
    }
}