yog 0.0.5

yog: a balls-oriented session manager for lernie loops (egui frontend)
Documentation
//! **The window's second asker lane** (REMOTE §3, §10; DESIGN §7.2; bl-73e7):
//! one connection and one thread, held open on the focused conversation's live
//! tail while the [`asker`](super::asker) keeps making its serial pass over
//! everything else.
//!
//! **Why a lane and not a question on the standing set.** The asker's pass is
//! serial — ask, wait, publish, next — so a read that deliberately never
//! finishes would stall every other surface for its whole duration. REMOTE §9.7
//! priced exactly that when it declined the graduation, and it is the price the
//!2026-08-22 ruling accepted: a lane, so the two cadences never touch. One
//! connection each, one thread each, and no shared state between them but the
//! frame that reads both.
//!
//! **Re-ask is the whole reconnect ladder.** A stream ends for three reasons and
//! the lane treats them alike: the step committed (the engine terminated it),
//! the subject moved (the lane hung up), or the dial failed. Each is followed by
//! a fresh ask. There is no backoff schedule and no liveness protocol, because
//! the fallback is not "nothing" — while the lane is down the seat paints the
//! tail the **pull** `Query::Transcript` folds, at
//! [`ASK_PERIOD`](super::asker::ASK_PERIOD). That is the migration's own
//! behaviour, kept deliberately, so the lane is an improvement that can fail
//! rather than a mechanism that can break the chat.
//!
//! **Two channels, and deliberately no lock** — [`link`](super::link)'s reason
//! exactly: nothing here is shared mutable state, it is a hand-off in each
//! direction. What crosses is the **whole** fold per frame, never a delta, so a
//! frame the lane misses costs nothing and a seat never reassembles anything.
//!
//! **The question is its own key**, again as `link`'s is: a frame carries the
//! envelope text it answers, so a subject that moved cannot land the previous
//! conversation's tail on the new one.

use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, Sender, TryRecvError, channel};
use std::thread::JoinHandle;

use serde_json::Value;

use crate::boundary::reply::Reply;
use crate::git_tree::Stream;
use crate::watch::Repaint;

use super::asker::ASK_PERIOD;
use super::client::Seat;

/// The frame's end of the lane: the conversation it wants followed, and the
/// newest fold that has arrived for it.
pub struct Tail {
    subject: Sender<Option<Value>>,
    frames: Receiver<(String, Option<Stream>)>,
    /// Declared during this frame's render, read at the next settle.
    standing: Option<Value>,
    /// What the lane is on.
    asked: Option<Value>,
    landed: Option<Stream>,
}

/// The lane's end: what to follow, and where a frame goes.
pub struct TailEnd {
    subject: Receiver<Option<Value>>,
    frames: Sender<(String, Option<Stream>)>,
    standing: Option<Value>,
    /// Whether the frame end has gone away. A lane that kept re-asking for a
    /// window nobody is painting would be the one leak a detached thread can
    /// make, so the disconnect is latched rather than re-derived.
    hung_up: bool,
}

/// A fresh pair, minted together for [`link::pair`](super::link::pair)'s reason:
/// neither end is useful alone.
pub fn pair() -> (Tail, TailEnd) {
    let (s_tx, s_rx) = channel();
    let (f_tx, f_rx) = channel();
    (
        Tail {
            subject: s_tx,
            frames: f_rx,
            standing: None,
            asked: None,
            landed: None,
        },
        TailEnd {
            subject: s_rx,
            frames: f_tx,
            standing: None,
            hung_up: false,
        },
    )
}

/// **A lane nobody answers.** The model holds one from the moment it boots, so
/// a seat's read of the tail is the same call whether or not this box got a
/// lane up — and a window with none simply paints the pull fold.
impl Default for Tail {
    fn default() -> Self {
        pair().0
    }
}

impl Tail {
    /// Declare `question` the followed subject and read whatever fold has
    /// landed for it. Called during render, so it does one clone and nothing
    /// else — never blocks and never dials.
    pub fn ask(&mut self, question: &Value) -> Option<Stream> {
        self.standing = Some(question.clone());
        self.landed.clone()
    }

    /// One frame's whole duty: take what landed, tell the lane what is followed
    /// if that changed, and start the next frame's declaration empty. A subject
    /// nobody declared this frame is a subject nobody is watching, which is the
    /// whole of stopping — there is no unfollow to forget.
    pub fn settle(&mut self) {
        let key = self.asked.as_ref().map(Value::to_string);
        for (answered, frame) in self.frames.try_iter() {
            if Some(&answered) == key.as_ref() {
                self.landed = frame;
            }
        }
        if self.standing != self.asked {
            let _ = self.subject.send(self.standing.clone());
            self.asked = self.standing.clone();
            self.landed = None;
        }
        self.standing = None;
    }
}

impl TailEnd {
    /// The subject as of now — the newest declaration the frame sent, or the one
    /// before it when the frame has declared nothing new. `None` once the frame
    /// end is gone, which is what ends a detached lane.
    pub fn standing(&mut self) -> Option<Value> {
        loop {
            match self.subject.try_recv() {
                Ok(newest) => self.standing = newest,
                Err(TryRecvError::Empty) => return self.standing.clone(),
                Err(TryRecvError::Disconnected) => {
                    self.hung_up = true;
                    return None;
                }
            }
        }
    }

    /// Publish one frame against the question it answers. `None` is *this
    /// stream is over* — the seat drops the tail and falls back to the pull
    /// fold, which is also what makes the step boundary a swap rather than a
    /// duplication.
    pub fn publish(&self, question: &str, frame: Option<Stream>) {
        let _ = self.frames.send((question.to_owned(), frame));
    }
}

/// One window's follow lane: its own seat on the wire, its end of the frame's
/// hand-off, and the face to wake when a frame lands.
pub struct Lane {
    seat: Seat,
    end: TailEnd,
    repaint: Arc<dyn Repaint>,
}

impl Lane {
    /// Assemble the lane. Built by [`Engine::lane`](crate::engine::Engine::lane)
    /// so a test can drive [`turn`](Self::turn) by hand — the
    /// [`Asker`](super::asker::Asker) precedent exactly.
    pub fn new(seat: Seat, end: TailEnd, repaint: Arc<dyn Repaint>) -> Self {
        Self { seat, end, repaint }
    }

    /// One turn: hold the line on whatever the frame is following, until the
    /// engine ends the stream, the subject moves, or the dial fails. Answers
    /// whether any frame landed — which is what tells the caller whether the
    /// far end is answering at all, and therefore whether to re-ask at once or
    /// to wait a period first.
    pub fn turn(&mut self) -> bool {
        let Self { seat, end, repaint } = self;
        let Some(question) = end.standing() else {
            return false;
        };
        let key = question.to_string();
        let mut landed = false;
        let _ = seat.followed(&question, &mut |frame| {
            let Ok(Reply::Follow(stream)) = frame else {
                return false;
            };
            end.publish(&key, Some(stream));
            repaint.request();
            landed = true;
            // Stay only while the frame still wants this conversation.
            end.standing().as_ref() == Some(&question)
        });
        // However it ended, the seat is told the stream is over: a tail that
        // stopped growing must not stand painted over a transcript that has
        // since committed it.
        end.publish(&key, None);
        repaint.request();
        landed
    }

    /// Run [`turn`](Self::turn) forever. A turn that landed frames re-asks at
    /// once — the engine holds the next stream open itself, so there is nothing
    /// to pace — and one that landed none waits a period, which is the whole of
    /// the backoff.
    pub fn start(mut self) -> LaneThread {
        let stop = Arc::new(AtomicBool::new(false));
        let flag = Arc::clone(&stop);
        let handle = std::thread::spawn(move || {
            while !flag.load(Ordering::Relaxed) {
                if !self.turn() {
                    std::thread::park_timeout(ASK_PERIOD);
                }
            }
        });
        LaneThread {
            stop,
            handle: Some(handle),
        }
    }
}

/// The lane thread's handle; [`Drop`] signals stop and unparks.
///
/// **It does not join, and that is the one place this differs from every other
/// thread yog owns** (the worker, the asker, the consumer, the listener: stop,
/// unpark, join). Those are parked on a *local* period, so a join costs at most
/// one tick. This one is parked on a socket read whose bound is the **engine's**
/// hold — thirty seconds of a quiet conversation — so joining it would make
/// closing a window wait on a remote timer. It owns nothing but its own
/// connection and its end of two channels; the frame end dropping is what makes
/// [`TailEnd::standing`] answer `None`, so a lane whose window is gone stops
/// asking on its own turn even if the flag never reached it.
pub struct LaneThread {
    stop: Arc<AtomicBool>,
    handle: Option<JoinHandle<()>>,
}

impl Drop for LaneThread {
    fn drop(&mut self) {
        self.stop.store(true, Ordering::Relaxed);
        if let Some(handle) = self.handle.take() {
            handle.thread().unpark();
        }
    }
}

#[cfg(test)]
mod tests;