Skip to content

[FEATURE]: Observers declare whether they collect, claim, or attend #781

Description

@justintime4tea

Summary

A run has one observer and it is implicit: whoever took the receiver from AgentRun::take_agent_events() — or, after PR #719, from the Agent that owns it. Its arrival starts the run and its departure ends it, and those two facts stand in for everything the runtime might want to know about who is watching. The epic distinguishes what an observer can be, and a run behaves differently depending on which are present.

Make the subscription carry it. An observer attaches to a session, collects or claims, and either may attend. The vocabulary — Observer, ObserverKind, presence — is already in aura_events::run (#778); this gives it behavior.

Goals

  • Two kinds. A collecting observer reads and has no effect on a run's lifetime — an exporter, an inspecting CLI, a dashboard. A claiming observer holds the session's right to keep its run going: a chat client, or the A2A task that owns an execution. When the last claimant detaches while a run is live, the runtime cancels the run. That is today's behavior, made explicit; [FEATURE]: Liveness policy when a run's last claim detaches #784 makes it one policy among three.
  • Presence as a property of any observer, not a third kind. A human at the helm may be attached passively: the epic's "zoom in and have presence passively" is a collecting observer with presence; "fully attach/claim" is a claiming one with it. HITL ([FEATURE]: HITL routing consults presence #786) reads presence; liveness reads claims.
  • Observers attach to a session, not a run. The subscription is the session's stream ([FEATURE]: A session's event journal and late attach #779) from a cursor — a run's Started for "attach to this run", the start for "everything this agent has done", live for "from now". A claim or presence held on a session applies to whatever run is live in it, including one that starts after the attach, so a chat client that attaches to its session and then starts a run is the claimant from the run's first event.
  • At most one claimant at a time. A second claim is refused, or takes over with the first told why, by policy — never two holders of one session. Collecting observers are unlimited.
  • Attaching and detaching are lifecycle events ([FEATURE]: A session's event envelope carries run lifecycle alongside agent activity #778) carrying the Observer, and detaching carries why — released, expired, displaced — so a session's own stream records who watched it and how each left.
  • Dropping a subscription detaches. In one process a subscription is a guard, the way AgentRun::into_events is a guard for the run; its drop is the Released detach. A lease — an attachment with a TTL that must be renewed — is the right shape only where a process boundary makes drop invisible: the cross-instance relay in [FEATURE]: Endpoints to list, inspect, attach to, and control sessions and runs #785 and the session claim in [FEATURE]: Atomic Fence for VFS Claims #581. Those are the two places Expired can come from. In-process observers are never leases.
  • Counts by kind and presence held on the session's record ([FEATURE]: A runtime that owns sessions and their runs #780) and derivable from the stream by any projection. Two of them — any present observer, any claimant — are mirrored onto the live run's RunContext as atomics on every attach and detach and when a run starts, 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 worker's run sees its session's presence and claim. Orchestration workers and the coordinator run on children of the orchestration's RunContext (RunContext::child, PR refactor(agent): separate a prepared agent from its runs #719): the same id and observer, cancelled with it, with tool state of their own. Presence and claim belong to the session, not to an agent within it, so a child shares its parent's two atomics rather than starting with its own. The runtime mirrors onto the live run once, and every worker's HITL gate ([FEATURE]: HITL routing consults presence #786) and liveness check reads the same value. A worker's gate reads presence from the worker's own run once begin_run_within binds it, so without sharing, a worker in adaptive mode would see no one present while a human is attached.
  • attach is the only way to read a session. There is no raw journal reader that bypasses the observer set, so every subscriber is one the session knows about and counts.

Data structures

Observer and ObserverKind are as implemented in PR #730; DetachCause is proposed there (#778). The rest is proposed:

// aura_events::run — as implemented (PR #730)
pub struct Observer { pub id: ObserverId, pub kind: ObserverKind, pub presence: bool }
pub enum ObserverKind { Collecting, Claiming }
// proposed on the rebase (#778)
pub enum DetachCause { Released, Expired, Displaced }

// proposed — the runtime side
pub struct Attach {
    pub kind: ObserverKind,
    pub presence: bool,
    pub from: Cursor,                 // #779 — Start, After(seq), or Live
    pub claim: ClaimPolicy,           // what to do if the session already has a claimant
}

pub enum ClaimPolicy {
    Refuse,                           // Err(AlreadyClaimed { by: ObserverId })
    TakeOver,                         // the previous claimant is detached with Displaced
}

/// Held by the seam for as long as it is attached; dropping it detaches with
/// `Released`. The only channel end a seam ever holds is the stream inside it.
pub struct Subscription {
    session: SessionId,
    observer: Observer,
    events: Pin<Box<dyn Stream<Item = JournalItem> + Send>>,
    runtime: Weak<Runtime>,           // for detach-on-drop
}

impl Stream for Subscription { type Item = JournalItem; /* ... */ }

pub struct ObserverSet {
    collecting: HashMap<ObserverId, Observer>,
    claimant: Option<Observer>,       // at most one, design question 3
}

impl ObserverSet {
    pub fn present(&self) -> bool;    // any observer with presence
    pub fn claimed(&self) -> bool;
    pub fn counts(&self) -> ObserverCounts;
}

/// What a summary (#785) reports; instance-local, never persisted.
#[derive(Serialize)]
pub struct ObserverCounts { pub collecting: u32, pub claiming: u32, pub present: u32 }

// On the live run's RunContext (#780), kept current by attach, detach, and run start.
// Shared by its children (RunContext::child), so workers read the session's values:
//   present: Arc<AtomicBool>   — ObserverSet::present(), for the gate
//   claimed: Arc<AtomicBool>   — ObserverSet::claimed(), for liveness

impl Runtime {
    /// The one way to read a session. Registers the observer, mirrors the counts
    /// onto the live run's context, appends `ObserverAttached`, and returns the
    /// subscription — a guard whose drop does the reverse with `Released`.
    pub fn attach(&self, session: &SessionId, how: Attach) -> Result<Subscription, AttachError>;
}

pub enum AttachError { NoSuchSession, AlreadyClaimed { by: ObserverId } }

Additional Context

The one-claimant rule is a V1 choice matching the epic's "at most once per each agent". It is what makes "who is driving this agent" a question with an answer. Relaxing it later is a policy change, not a schema change.

ObserverSet is never persisted. An observer is a live connection; a restarted process has none, and a stored "present" observer would route an approval to nobody (#786). Observer serializes only because it rides in ObserverAttached / ObserverDetached events as history. The exception is deliberate and lives elsewhere: an attachment that crosses a process boundary is a lease in the session store's LeaseStore (#581 defines it for the session claim; #785's relay reuses it), with a TTL, because the holder can die unseen. In-process observers get the exact signal a guard gives; cross-process ones get the approximate signal a lease gives, and nothing pays for a boundary it does not cross.

Presence is self-declared. An observer that attaches saying a human is at it is believed; the runtime has no way to check. That is the trust mode = "conversational" already places in a client, made per-observer rather than per-deployment. ClaimPolicy::TakeOver is a runtime capability; #785's endpoint does not expose it in V1 and refuses a second claim.

Attaching to a session with no live run is allowed and ordinary: a collecting observer waiting for the agent's next run, or a chat client claiming its session before starting one. AttachError therefore has no Ended; a session that has retired is NoSuchSession.

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 7, 2026

    @justintime4tea
    CollaboratorAuthor

    Design questions need to be answered before implementing. Questions are under the runtime design discussion.

    one claimant (index question 3);
    presence as a flag (question 4);
    guards in-process, leases at process boundaries (question 9)

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