yog 0.0.1

yog: a balls-oriented session manager for lernie loops (egui frontend)
Documentation
//! The **streamed-piped** spawn class (DESIGN §8, §8.3): line-buffered stdout
//! read live off a running child and delivered to the invoking surface as whole
//! lines arrive, plus stderr capture and the terminal exit — the shape
//! `bz --login`'s device flow renders through (§5.3 whitelists the live lines as
//! instance-local RAM: a device code is for the human at *this* keyboard).
//!
//! Non-blocking by construction: [`Streamed::poll`] drains only what the pump
//! threads have already buffered and returns at once, so the egui frame loop
//! never stalls waiting on the child. The stream itself is never logged
//! line-by-line — it **converges to one outcome row** at exit (§4.2); appending
//! that row is the caller's job ([`crate::login`]), keeping this crate free of
//! the `opslog` seam (the `actions`/`login` layer bridges the two).

use super::{Chunk, ExitInfo, Stream, StreamPoll};

/// Hard cap on a single rendered line's byte length — the streamed class's
/// *bounded line length* (§8). A physical line reaching this cap stops accreting
/// (further bytes drop) and flushes with a [`TRUNC_MARK`] suffix, so a device
/// flow that never emits a newline cannot grow RAM without bound. Device codes
/// and URLs are far shorter; the cap only bites pathological input.
const MAX_LINE: usize = 4096;

/// Suffix stamped on a line truncated at [`MAX_LINE`] (§8), so the truncation is
/// visible in the pane rather than silent.
const TRUNC_MARK: &str = "…[truncated]";

/// A pure newline splitter with a bounded retained line (§8). Bytes accrete into
/// `partial`; each `\n` flushes it as one line. A line that reaches [`MAX_LINE`]
/// stops accreting and flushes with [`TRUNC_MARK`], bounding RAM. Flush is
/// lossy-UTF8 — a device code is ASCII, and a multibyte char torn at a read
/// boundary is a display glyph, never an error.
#[derive(Debug, Default)]
struct LineBuf {
    partial: Vec<u8>,
    truncated: bool,
}

impl LineBuf {
    /// Feed `bytes`, returning every line completed by a `\n` within them, in
    /// order. A trailing unterminated remainder stays buffered for the next feed
    /// (or [`finish`](Self::finish) at EOF).
    fn push(&mut self, bytes: &[u8]) -> Vec<String> {
        let mut lines = Vec::new();
        for &b in bytes {
            if b == b'\n' {
                lines.push(self.flush());
            } else if self.partial.len() < MAX_LINE {
                self.partial.push(b);
            } else {
                self.truncated = true;
            }
        }
        lines
    }

    /// The buffered remainder as a final line at EOF, or `None` when empty — a
    /// device code printed without a trailing newline still surfaces.
    fn finish(&mut self) -> Option<String> {
        (!self.partial.is_empty()).then(|| self.flush())
    }

    /// Take `partial` as one lossy-UTF8 line, appending [`TRUNC_MARK`] when it was
    /// capped, and reset for the next line.
    fn flush(&mut self) -> String {
        let mut line = String::from_utf8_lossy(&self.partial).into_owned();
        if self.truncated {
            line.push_str(TRUNC_MARK);
        }
        self.partial.clear();
        self.truncated = false;
        line
    }
}

/// One non-blocking read of a [`Streamed`] (§8): new whole lines arrived, nothing
/// is ready yet, or the child exited. The terminal [`Done`](StreamedPoll::Done)
/// carries any final flushed line, the exit code, and the captured stderr the
/// caller folds into the single outcome row (§4.2).
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StreamedPoll {
    Lines(Vec<String>),
    Pending,
    Done(StreamedOutcome),
}

/// The terminal fact of a streamed run (§4.2 outcome row): the shell-convention
/// exit code, the captured stderr, and any final unterminated line flushed at
/// EOF.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StreamedOutcome {
    pub lines: Vec<String>,
    pub exit: i32,
    pub stderr: String,
}

/// A running child consumed as the streamed-piped class (§8): stdout line-buffered
/// live, stderr accumulated, exit captured. Wraps a [`Stream`] whose Drop still
/// SIGTERM/SIGKILLs the child — closing the login surface aborts the device flow,
/// consistent with its instance-local nature (§5.3).
pub struct Streamed {
    stream: Stream,
    buf: LineBuf,
    stderr: Vec<u8>,
}

impl Streamed {
    /// Wrap a freshly-spawned [`Stream`] (stdin already null via
    /// [`Cli::run`](super::Cli::run)) for line-buffered live consumption.
    pub fn new(stream: Stream) -> Self {
        Self {
            stream,
            buf: LineBuf::default(),
            stderr: Vec::new(),
        }
    }

    /// Non-blocking: drain whatever the pumps have buffered right now, then return.
    /// New stdout lines come back as [`StreamedPoll::Lines`]; nothing-yet as
    /// [`Pending`](StreamedPoll::Pending); the child's exit as
    /// [`Done`](StreamedPoll::Done) with any final line + exit + stderr. Idempotent
    /// after `Done` (the underlying [`Stream::try_next`] stays `Pending`), so a
    /// stray extra poll yields `Pending`, never a second `Done`.
    pub fn poll(&mut self) -> StreamedPoll {
        let mut lines = Vec::new();
        loop {
            match self.stream.try_next() {
                StreamPoll::Ready(Chunk::Stdout(b)) => lines.extend(self.buf.push(&b)),
                StreamPoll::Ready(Chunk::Stderr(b)) => self.stderr.extend_from_slice(&b),
                StreamPoll::Ready(Chunk::Exited(e)) => {
                    lines.extend(self.buf.finish());
                    return StreamedPoll::Done(self.outcome(lines, e));
                }
                StreamPoll::Pending => {
                    return if lines.is_empty() {
                        StreamedPoll::Pending
                    } else {
                        StreamedPoll::Lines(lines)
                    };
                }
            }
        }
    }

    /// Build a [`Streamed`] over a bare receiver with no child — the seam the
    /// streamed-class and login unit tests drive [`poll`](Self::poll) through
    /// deterministically (chunks sent by hand, no process to race). `pub(crate)`
    /// so [`crate::login`]'s tests reach it; test-only, so it costs no coverage.
    #[cfg(test)]
    pub(crate) fn from_rx(rx: std::sync::mpsc::Receiver<Chunk>) -> Self {
        Self::new(Stream::from_rx(rx))
    }

    /// Fold the accumulated stderr and the exit into the outcome row's fields.
    fn outcome(&self, lines: Vec<String>, exit: ExitInfo) -> StreamedOutcome {
        StreamedOutcome {
            lines,
            exit: exit.shell_code(),
            stderr: String::from_utf8_lossy(&self.stderr).into_owned(),
        }
    }
}