Skip to content

[FEATURE]: A runtime that owns sessions and their runs #780

Description

@justintime4tea

Summary

A run is a request. The handler prepares an agent, begins a run, calls stream, and holds the AgentRun until the client goes away; A2A does the same inside an execution. Nightly has retired the request-keyed registries (PR #738), so what is left tracking runs is ActiveRequestTracker — a count and a list of join handles so shutdown can drain and abort — A2A's TaskCancelState, and the park-mode PARK_CELLS (#732). Nothing answers "what is running for this session", and nothing can, because the run and the request are the same object.

Introduce the runtime: one owner of every live session on the instance and of the run each is executing. It prepares the session's agent, begins the run, drives it on a task of its own, writes the session's journal, and hands out subscriptions. Callers describe a run and attach to its session; they do not prepare, hold, or drive anything.

Goals

  • Start a run from a description — agent name, the request's forwarded headers and client tools, the prompt and history, RunOptions — not from an agent the caller built. The runtime finds or creates the session, prepares its agent (RigBuilder::prepare_agent) if the session has none, begins the run (PreparedAgent::begin_run, which creates the RunContext through channel_for_agent and holds its event receiver), streams it (Agent::stream), and takes the receiver for the run's task — all on a task the runtime owns. The RunId is minted here, once, and is what begin_run receives as the run's id. Preparing in the runtime is what lets [FEATURE]: A prepared agent stays warm across a session's runs #787 keep the prepared agent across the session's runs without the seams knowing.
  • start returns the RunId as soon as the record exists, in Preparing. Started follows on the journal once the agent is ready and begin_run succeeds; a caller that needs it attaches and waits for it. On a session whose agent is already prepared, Preparing lasts as long as begin_run; on a cold session it includes MCP discovery, which a headless start therefore does not block on. Shutdown counts the run from the first moment ([BUG]: graceful shutdown ignores requests still in per-request MCP init #736).
  • One live run per session. A second start while the session has a live run is refused with the running RunId (StartError::RunInProgress), which the caller can attach to. This is what PR refactor(agent): separate a prepared agent from its runs #719's RunLease enforces per prepared agent, lifted to the session so it holds before [FEATURE]: A prepared agent stays warm across a session's runs #787 caches anything. Today two requests on one chat_session_id both run, each with its own agent; under the runtime a session serves one run at a time. Whether a second start waits instead of being refused is open (discussion question 10).
  • The run's task writes the agent's events; the runtime writes the lifecycle; both through one append. The task appends Started, then Agent(McpStatus) from the prepared agent's server status, then drains the run's channel and appends each event. When the channel closes it appends the terminal lifecycle event and writes the run's final facts in one RunStore::finish call, so a store that can make the two atomic does. Observer and liveness events go through the same serialized SessionJournal::append ([FEATURE]: A session's event journal and late attach #779) from the runtime directly — no side channel into the task.
  • Storage is behind the session store's traits. The runtime takes Arc<dyn RunStore> and Arc<dyn JournalStore> from SessionStore::runs() and SessionStore::journals(), beside approvals(), tasks(), skills(), and bus(). The in-memory implementations are today's behavior; list reads through RunStore, so a durable backend is a configuration change.
  • Read a run's facts, a session's state, and the list of runs by session, from anywhere that has the runtime. The live record — the session's agent, journal, observer set, and the run's context and task — stays inside the runtime; what crosses to a caller is RunFacts, an AgentState, and observer counts, with usage snapshotted from the live counters while the run is running.
  • Attaching is [FEATURE]: Observers declare whether they collect, claim, or attend #781's, and it attaches to a session. The runtime exposes no reader that bypasses the observer set, so every subscriber is an observer the session knows about.
  • Cancel by run id, through RunContext::cancel_token, with a reason that becomes the lifecycle event.
  • Per-session state the runtime holds as plain values, not derived from the journal: the agent slot, the live run, the last run, the observer set ([FEATURE]: Observers declare whether they collect, claim, or attend #781). Per-run: status, usage, pending approval decisions, the park reference. These are what a summary ([FEATURE]: Endpoints to list, inspect, attach to, and control sessions and runs #785) reads, and each is also a fold over the stream, so a durable backend could rebuild them; the runtime never does. Presence and claim are mirrored onto the live run's RunContext as atomics, so work inside the run — the HITL gate ([FEATURE]: HITL routing consults presence #786), the liveness timer ([FEATURE]: Liveness policy when a run's last claim detaches #784) — reads them without a reference back to the runtime.
  • A finished, failed, cancelled, or parked run stays live for a retention window — its record on the session, the session's journal still subscribable — so a subscriber arriving just after the end sees the ending rather than a miss. Then the run record is dropped; the session record stays while its agent is cached ([FEATURE]: A prepared agent stays warm across a session's runs #787) or an observer is attached, and is dropped after. What the stores keep afterwards is their own policy: the memory stores bound by a cap, as SkillInvocationStore does; a durable store by TTL.
  • A start that asks for a liveness policy the agent cannot honour — Park on an agent with no checkpoint support — is refused, not silently downgraded.
  • Shutdown drains through the runtime with the two-phase behavior ActiveRequestTracker provides today: refuse new, cancel every live run through its token with RunCancelReason::Shutdown so each writes a terminal event, wait the grace period, then abort what has not exited. The tracker itself retires in [FEATURE]: Chat completions run on the runtime #782, once requests no longer bypass the runtime.
  • Nothing calls it yet, except tests. Chat completions moves over in [FEATURE]: Chat completions run on the runtime #782, A2A in [FEATURE]: A2A tasks run on the runtime #783 — the alongside-then-over discipline [FEATURE]: Agent event schema and broker adapter #618 used.

Data structures

The per-run context is implemented (nightly, extended by PR #719); #781 adds two atomics to it. The runtime and its stores are proposed.

// crates/aura/src/run_context.rs — as implemented on PR #719's head, abbreviated.
// "One run — what its own work needs to correlate, where its events go, and
// the state an agent's tools keep for it."
pub struct RunContext {
    id: Arc<str>,                               // becomes RunId (#778) in PR #730's rebase
    tool_calls: Mutex<VecDeque<ToolCallId>>,
    events: mpsc::Sender<AgentEvent>,
    cancel: CancellationToken,
    scratchpad_budget: Option<ContextBudget>,
    turn_nudge: Option<Arc<TurnNudgeState>>,
    skill_recorder: Option<Arc<SkillInvocationRecorder>>,   // e54d2bc6
    // proposed by #781, written by the runtime on attach, detach, and run start:
    present: Arc<AtomicBool>,                   // some observer is present — read by the HITL gate (#786); shared with child runs
    claimed: Arc<AtomicBool>,                   // a claimant exists — read by liveness (#784); shared with child runs
}
impl RunContext {
    /// A run carrying the state a prepared agent's tools keep for it, and the
    /// receiver its observer reads. What `PreparedAgent::begin_run` calls.
    pub fn channel_for_agent(id, cancel, scratchpad_budget, turn_nudge, skill_recorder)
        -> (Arc<Self>, mpsc::Receiver<AgentEvent>);
    /// A run within `parent` for one agent of it — an orchestration worker or
    /// coordinator: same id, observer, and cancellation; tool state of its own.
    pub fn child(parent: &Arc<Self>, scratchpad_budget, turn_nudge, skill_recorder) -> Arc<Self>;
}

/// A run in progress on a prepared agent. The `Agent` holds one; the
/// `PreparedAgent` holds a `Weak` to it and refuses `begin_run` while it lives.
pub struct RunLease { run: Arc<RunContext> }

/// The slot tools resolve the current run through. One per prepared agent.
pub struct BoundRun(Mutex<Option<Arc<RunContext>>>);

The runtime does not carry a session on RunContext; the session is the record the run hangs off, and the journal stamps it onto each SessionEvent.

// proposed — aura::runtime

// struct name? See discussion question 1 @ https://github.com/mezmo/aura/discussions/743#discussioncomment-18800261
pub struct Runtime {                            // name: discussion 743 question 1
    configs: Arc<Vec<aura_config::Config>>,     // what prepare_agent builds from
    sessions: RwLock<HashMap<SessionId, Arc<SessionRecord>>>,  // sessions this process holds
    runs: RwLock<HashMap<RunId, SessionId>>,    // index: which session a run belongs to
    run_store: Arc<dyn RunStore>,               // facts and the session index; memory today
    journals: Arc<dyn JournalStore>,            // event storage; memory today
    retention: Duration,                        // how long an ended run stays live; the stores bound themselves
}

// aura::session_store — beside ApprovalStore, TaskStore, SkillInvocationStore, EventBus
#[async_trait]
pub trait RunStore: Send + Sync {
    async fn put(&self, facts: &RunFacts) -> Result<(), StoreError>;
    async fn get(&self, run: &RunId) -> Result<Option<RunFacts>, StoreError>;
    async fn list(&self, filter: &RunFilter) -> Result<Vec<RunFacts>, StoreError>;
    /// The terminal write: the run's last journal event and its final facts,
    /// atomic where the backend can make them so.
    async fn finish(&self, terminal: &SessionEvent, facts: &RunFacts) -> Result<(), StoreError>;
}
// No `retire`: the memory impl bounds itself by a cap, a durable impl by TTL.

pub struct RunFilter { pub session: Option<SessionId>, pub status: Option<RunStatus>, pub pending_approval: Option<bool>, pub limit: usize, pub after: Option<RunId> }

/// Live process state for one session. Crate-private and never `Serialize`.
pub(crate) struct SessionRecord {
    id: SessionId,
    agent_id: String,                           // AgentInfo::id
    agent: Mutex<AgentSlot>,                    // #787 keeps this Ready across runs
    journal: Arc<SessionJournal>,               // #779
    observers: RwLock<ObserverSet>,             // #781 — instance-local, never persisted
    run: RwLock<Option<Arc<RunRecord>>>,        // the live run; at most one
    last_run: RwLock<Option<RunId>>,
}
pub(crate) enum AgentSlot { Absent, Preparing, Ready(Arc<PreparedAgent>) }

/// Live process state for one run. Crate-private and never `Serialize`: it
/// holds a task, locks, atomics, and the run's context. What crosses to a
/// caller is `facts`.
pub(crate) struct RunRecord {
    id: RunId,
    context: Arc<RunContext>,
    task: JoinHandle<()>,                       // drives the AgentRun to completion and appends its events
    usage: UsageState,                          // live counters from the AgentRun
    facts: RwLock<RunFacts>,
}

/// The durable half of a run: everything about it that outlives the process
/// and is safe to serialize. `RunSummary` (#785) is this plus observer counts.
#[derive(Serialize, Deserialize)]
pub struct RunFacts {
    pub run_id: RunId,
    pub session_id: SessionId,
    pub agent: String,                          // AgentInfo::id
    pub status: RunStatus,
    pub liveness: Liveness,                     // aura_events::run, #784
    pub started_at: Timestamp,
    pub ended_at: Option<Timestamp>,
    pub usage: TokenUsage,                      // the live counters while running; final at the terminal write
    pub first_seq: Option<SequenceNumber>,      // the run's `Started` in the session stream; the attach cursor
    pub latest_seq: Option<SequenceNumber>,
    pub pending_approvals: Vec<DecisionId>,     // ids only; the decisions live in ApprovalStore
    pub checkpoint: Option<CheckpointRef>,      // a store key, never a local path
    pub continues: Option<RunId>,               // the run this one resumed or followed
}

#[derive(Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RunStatus { Preparing, Running, Finished, Failed, Cancelled, Parked }

/// The session's outward state — constructed, never stored, never an event.
/// What #785's session summary reports.
#[derive(Serialize)]
#[serde(tag = "state", rename_all = "snake_case")]
pub enum AgentState {
    Preparing,                                  // agent slot is Preparing: MCP discovery, later the session claim
    Running,                                    // a live run with nothing pending
    Blocked { on: BlockReason },                // a live run waiting on something outside the model
    Parked,                                     // last run parked; nothing in memory; resumable
    Done,                                       // no live run; the agent comes back with history on the next start
}
#[derive(Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum BlockReason {
    Approval { decisions: Vec<DecisionId> },    // pending_approvals non-empty, not yet parked
    ProviderRetry { attempt: u32 },             // the run's last Retrying event (#778) with no turn since
}

/// What a seam hands the runtime. No agent in it: the runtime prepares one.
pub struct StartRun {
    pub agent: Option<String>,                  // AgentInfo::id; the default agent when absent
    pub req_headers: Option<HashMap<String, String>>, // forwarded per `headers_from_request`
    pub client_tools: Option<Vec<ClientTool>>,
    pub additional_tools: Vec<Box<dyn ToolDyn>>,
    pub session: Option<SessionId>,             // a new session when absent
    pub continues: Option<RunId>,               // the run this one follows, if any
    pub prompt: String,
    pub history: Vec<Message>,
    pub options: RunOptions,                    // PR #710: timeout, parent token
    pub liveness: Liveness,                     // #784; the seam supplies its config default
}

pub enum StartError {
    NoSuchAgent(String),
    RunInProgress(RunId),                       // the session's live run; attach to it
    PolicyUnsupported { policy: LivenessPolicy, agent: String },   // e.g. Park without checkpoint support
    Store(StoreError),
}

impl Runtime {
    /// Mints the RunId, writes initial facts (`Preparing`), spawns the run's
    /// task, and returns at once. `Started` follows on the session's journal
    /// once the agent is ready and the run begun.
    pub async fn start(&self, req: StartRun) -> Result<RunId, StartError>;
    /// Facts with usage snapshotted from the live counters while running.
    pub fn facts(&self, run: &RunId) -> Option<RunFacts>;
    /// The session's outward state, constructed from the agent slot, the live
    /// run's facts, and its pending approvals and retry signal.
    pub fn session_state(&self, session: &SessionId) -> Option<AgentState>;
    pub fn observer_counts(&self, session: &SessionId) -> Option<ObserverCounts>;   // #781
    pub async fn list(&self, filter: &RunFilter) -> Result<Vec<RunFacts>, StoreError>;
    pub async fn cancel(&self, run: &RunId, reason: RunCancelReason, message: Option<String>) -> bool;
    /// Shutdown: cancel every live run with `Shutdown`, wait `grace`, abort the rest.
    pub async fn drain(&self, grace: Duration);
}

How session_state reads the records, so inspection endpoints agree:

condition, first match wins state
agent slot is Preparing Preparing
live run and its pending_approvals is non-empty Blocked { Approval }
live run and its last agent event is Retrying Blocked { ProviderRetry }
live run Running
last run's status is Parked Parked
otherwise — no live run, agent ready or absent, last run finished, failed, cancelled, or none Done

Additional Context

Lives in the aura crate, beside session_store and run_context, so the standalone CLI can hold one without a web server. AppState holds an Arc of it. The runtime is the one in-process facade: seams call it and nothing below it; its public form is the HTTP API in #785. Nothing below the runtime — an agent, a RunContext, an MCP client, a channel end — is reachable from a seam.

Naming. AgentRuntimeConfig is the TOML, per #628, and HitlRuntime is per-request HITL state, so AgentRuntime from the epic is taken twice over. This draft says Runtime; the type wants a name before it lands. Separately, PR #719's Agent type is "one run of a prepared agent", while the domain's agent is the config entry (design question 12).

A prepared agent serves one run at a time — PreparedAgent holds a Weak<RunLease> and begin_run returns RunInProgress while it is alive (PR #719). The runtime enforces the same rule one level up, per session, and #787 is where the prepared agent is kept across a session's runs.

RunId is a UUID (v7 preferred), minted here. HITL parses the orchestration run id as a UUID today and the park owner key is run:{run_id}; unifying the ids means this one has to satisfy both. Making RunContext::id the RunId newtype touches run_context, the MCP binding, the HITL gate, the park owner key, and begin_run's signature; it belongs to PR #730's rebase, not to this issue. RunFacts is what a RunStore beside ApprovalStore, TaskStore, and SkillInvocationStore holds — in memory now, durable when #210 / #325 provide a backend; e54d2bc6 is the worked example of adding such a capability, versioned record and all. The persistence rules the type family follows are in 578-design.md.

AgentState is Mike's outward typestate list — Preparing, Running, Blocked, Parked, Done — with Idle folded into Done (a prepared agent that has run nothing yet is Done with no last run). None of them is an event; each is read off the records above, and a different implementation could construct them differently without the inspection API noticing.

The session claim (#581) is consulted at start when it exists, and it is what makes "one writer per session" true across instances: the instance holding the claim is the one whose runtime appends to that session's journal. The claim is expected to carry the write claim for the session's memory directory when that exists. Nothing here anticipates its shape.

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)): SessionId is documented as one agent's identity over time, and this issue prepares one agent per session. The web server's chat_session_id doesn't obey that today: a client can keep it and switch models, which is why handlers.rs:293 keys skill logs by (session, agent). When the runtime finds-or-creates the session for a start, it has to decide what that id maps to. Either (a) the runtime's session is (chat_session_id, agent), so switching models starts a second session that shares the chat id, matching the skill-log precedent and keeping one agent per stream; or (b) a session may change agent between runs, and the stream is one agent's only per run, via Started.agent. (a) seems right to me; recording it here so the start path settles it rather than inheriting the ambiguity.

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