Skip to content

[FEATURE]: A session's event journal and late attach #779

Description

@justintime4tea

Summary

On nightly (PR #738) a run has one channel: RunContext holds the sender, and AgentRun::take_agent_events() hands its mpsc::Receiver<AgentEvent> to whoever observes the run — once. PR #719 moves that receiver onto the Agent itself (begin_run creates the channel through RunContext::channel_for_agent), which changes nothing here: there is still exactly one receiver, and it is still the journal that takes it. That is the shape of a request-bound agent: one consumer, present from the first byte, and if it stops polling the bounded channel fills and the producer waits. #229 wedged tool execution at the 33rd unread event that way, and #230 pinned a process-wide lock behind one stuck send. Attaching to a run already underway means seeing what it has done, and a channel holds nothing.

Give each session a journal: an append-only sequence of SessionEvents across the session's runs, stored behind a session-store capability whose in-memory implementation is today's behavior. The live run's task drains the run's receiver into it, so the run never waits on a reader. Subscribing replays from a cursor and then follows live, and a subscriber that falls behind is told so rather than holding anything up. Because the stream is the session's, an agent's whole history is one load, and a run boundary is a position in it.

Goals

  • The live run's task reads the receiver nightly provides, stamps each AgentEvent into a SessionEvent with the session id, run id, next sequence number, and timestamp, and appends it. Lifecycle events from the runtime ([FEATURE]: A runtime that owns sessions and their runs #780) go through the same append, so the sequence is one sequence.
  • Subscribe from a cursor: the start of the session, a sequence number, or live-only. Replay then live, as one stream, with no seam visible to the consumer. "Attach to this run" is a cursor at the run's Started.
  • A slow subscriber lags and is told so — a gap marker carrying the sequence it resumed at — and neither the run nor its other subscribers wait. The fell-behind contract bus_bridge gives its relays, applied to every subscriber. This is the structural end of [BUG]: tool execution hangs when custom events are off #229 and [BUG]: one stuck progress send freezes progress for all requests #230's class.
  • Bounded by entry count per session. Oldest entries drop past the bound; a cursor older than the bound replays from the bound and says so. No coalescing: sequence numbers stay dense and every entry is what was emitted.
  • Content is in the journal. Nightly's tee_content (PR Deliver a run's events through its handle and retire the request registries #738) copies text, reasoning, and the final Completed onto the run's channel; its doc notes no consumer reads the copy yet. The journal is that consumer, so a late observer sees what the agent said and not only what it did.
  • The journal holds nothing the SSE stream did not already expose. A late observer sees what the original client saw, and no more.
  • The runtime's own state is not rebuilt from the journal. It is for observers. The facts the runtime keeps ([FEATURE]: A runtime that owns sessions and their runs #780) are nonetheless each a fold over the stream, so a durable backend can rebuild them later without the runtime ever doing so.
  • One writer per session at a time. append is serialized: it takes the journal's lock, mints seq, hands the event to the store, broadcasts it, and releases. The live run's task and the runtime both call it; nothing else does, and there is no side channel between them. In one process that is the lock plus the rule that a session has at most one live run ([FEATURE]: A runtime that owns sessions and their runs #780). Across instances it is the session claim ([FEATURE]: Atomic Fence for VFS Claims #581): the instance holding the claim is the one appending. The store therefore sees each session's sequence strictly increasing, from one writer at a time.
  • Storage is behind a trait from day one. JournalStore has an in-memory implementation that is today's behavior and is registered on the session store beside ApprovalStore, TaskStore, SkillInvocationStore, and EventBus — the pattern e54d2bc6 just used for skills: one trait, three backends, a versioned record probed before parse, TTL where the backend has it. A durable JournalStore is a second impl of the same trait; the journal, the runtime, and every seam are unchanged by it.
  • Subscribe is written so an async store changes nothing. A subscriber joins the live broadcast first, then reads the store up to the first live sequence number and drops any live event it already replayed. Sequence numbers make the dedupe exact. No lock is held across the read, so a store that awaits is the same algorithm as one that does not.
  • Every event is stored. No elision of text deltas, no coalescing. A Gap means one thing — the subscriber missed events — whether it came from a slow subscriber lagging the broadcast or a cursor older than what the store still holds.

Data structures

Proposed. The storage trait lives in aura::session_store beside ApprovalStore, SkillInvocationStore, and EventBus; the journal itself in the aura crate beside run_context. The shape is chosen so that a durable JournalStore is a second implementation and nothing else moves — the reasoning is in #743, The journal, shaped for durability now.

// aura::session_store — a capability, like ApprovalStore and EventBus.
// The in-memory impl is today's behavior; a Redis/Postgres impl is later work.

/// Append-only log of one session's events. One writer at a time, ever.
#[async_trait]
pub trait JournalStore: Send + Sync {
    /// Stores `event`, whose `seq` the caller minted. Returns when the event is
    /// durable to this store's standard: at once in memory, on commit elsewhere.
    async fn append(&self, event: &SessionEvent) -> Result<(), StoreError>;

    /// Events of `session` with `seq` in `from..=to` that the store still holds,
    /// in order. `Truncated` names the earliest sequence still available when
    /// `from` is older than that.
    async fn read(&self, session: &SessionId, from: SequenceNumber, to: SequenceNumber)
        -> Result<Vec<SessionEvent>, ReadError>;

    async fn latest(&self, session: &SessionId) -> Result<Option<SequenceNumber>, StoreError>;
}
// No `retire`: a store bounds itself. The memory impl caps entries per session
// and sessions held (least-recently-touched out, as `InMemorySkillInvocationStore`
// does); a durable impl expires by TTL. The runtime's retention window (#780)
// governs the live record, not the store.

pub enum ReadError { Truncated { earliest: SequenceNumber }, Store(StoreError) }

/// The in-memory impl: a bounded ring per session, a bounded number of sessions.
pub struct MemoryJournalStore {
    sessions: RwLock<HashMap<SessionId, VecDeque<Arc<SessionEvent>>>>,
    max_entries_per_session: usize,
    max_sessions: usize,                        // least-recently-touched evicted
}
// aura::journal — the in-process half: sequencing and fan-out. One per live session.

pub struct SessionJournal {
    session_id: SessionId,
    store: Arc<dyn JournalStore>,
    append: tokio::sync::Mutex<()>,             // serializes mint → store → broadcast
    latest: AtomicU64,                          // last seq the store confirmed; `append` mints latest + 1
    live: broadcast::Sender<Arc<SessionEvent>>, // fan-out; a lagging receiver becomes a Gap
}

pub enum Cursor {
    Start,                   // replay everything the store holds, then live
    After(SequenceNumber),   // replay events after this one, then live
    Live,                    // no replay
}

/// What a subscriber reads.
pub enum JournalItem {
    Event(Arc<SessionEvent>),
    /// The subscriber missed events; `resumed_at` is the first one it sees
    /// again. Emitted for a lagged broadcast receiver and for a cursor older
    /// than the store holds. One meaning, two causes.
    Gap { resumed_at: SequenceNumber },
    /// A stored entry this binary cannot deserialize — a newer variant after a
    /// rollback. Replay continues past it; `raw` is kept for a reader that can.
    Unknown { seq: SequenceNumber, raw: Bytes },
}

impl SessionJournal {
    /// Opens the session's journal, reading `latest` from the store so a session
    /// that outlives a process continues its sequence rather than restarting it.
    pub(crate) async fn open(session_id: SessionId, store: Arc<dyn JournalStore>) -> Result<Self, StoreError>;

    /// The one append path, and the only place `seq` is minted. Takes the lock,
    /// mints `latest + 1`, awaits the store, broadcasts, releases. Called by the
    /// live run's task for agent events and by the runtime for lifecycle events.
    pub(crate) async fn append(&self, run_id: Option<RunId>, payload: SessionEventPayload)
        -> Result<Arc<SessionEvent>, StoreError>;

    /// Joins the broadcast first, then reads `from..=first_live_seq - 1` from
    /// the store and drops any live event with `seq` at or below the last one
    /// replayed. Holds no lock across the read.
    pub fn subscribe(&self, from: Cursor) -> impl Stream<Item = JournalItem>;

    pub fn latest(&self) -> Option<SequenceNumber>;
}

The live run's task (#780) appends Started before the agent produces anything, drains the run's mpsc::Receiver<AgentEvent> and appends each, and appends the terminal event after the agent's channel closes. The runtime appends observer and liveness events directly through the same append, whether or not a run is live. There is no appender task, no lifecycle side channel, and no second sequence.

Additional Context

Per-instance and in-memory. A durable journal, and replay for a session another instance holds, belong with #210 / #325. The cursor and gap semantics here are what a store-backed journal would honour, so nothing at a seam changes when one arrives.

The receiver nightly provides is bounded, and RunContext::emit returns false when it cannot deliver. With the run's task the only reader, that path is reachable only if the task itself dies; the journal treats it as the run failing, not as a lost event.

A slow store slows append, a slow append fills the run's channel, and a full channel makes the agent wait on emit. That is the correct back-pressure — a run should not outpace what its journal can record — and it is invisible in memory. It is the one place durability changes behavior rather than shape, and it is a store tuning question, not a refactor.

Keying by session is what lets the stream double as the session's history: Started { prompt } and the Completed content of each run are the user and assistant turns in order. Whether that is how the server keeps history, or whether #210 stores messages beside the stream, is #210's call; this issue only makes the derivation possible.

Searched Issues

  • No similar issues found

Code of Conduct

  • I agree to follow this project's Code of Conduct

Activity

  1. justintime4tea commented on Oct 9, 2026

    @justintime4tea
    CollaboratorAuthor

    From Jake's (@jakedipity) review of #794 (#794 (comment)): SequenceNumber::follows returns false for a gap, a repeat, and an out-of-order arrival alike. The doc now says so (53e19a1e), but nothing tells them apart.

    This journal's subscribe loop is where the distinction is actually made: event.seq < expected drops a repeat, and advance() turns a jump into a Gap. Worth lifting that comparison onto SequenceNumber as a three-way answer (next / gap, with what was missed / already seen), so the loop and any consumer that attaches through #785 use one definition rather than each reimplementing expected. Noting it here so the review point lands with the code that needs it.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions