agentplane 0.45.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
//! Streaming A2A updates, read from the journal.
//!
//! # The stream is a view of the journal, not an event bus
//!
//! The obvious way to stream progress is an in-process broadcast channel: a step
//! finishes, it publishes, subscribers receive. It is also wrong here, in three
//! ways that only show up in production.
//!
//! A channel's events live in memory, so a subscriber that reconnects has
//! **missed** whatever happened while it was away, and nothing can tell it what.
//! A channel is per process, so a subscriber attached to the instance that is
//! *not* running the work receives nothing — and which instance that is changes
//! after every failover. And a channel is a second record of what happened,
//! which can disagree with the first.
//!
//! So updates are read from the journal instead. That makes the stream exactly
//! as durable as the run: a client that drops and re-subscribes picks up the
//! current state and continues, any instance can serve it, and the events cannot
//! disagree with history because they *are* history.
//!
//! The cost is polling the journal for each open subscriber. It is a real cost
//! and is stated rather than hidden: one indexed read per subscriber per
//! interval, against a store that is already answering worse queries. What it
//! buys is a stream that survives the things streams are asked to survive.
//!
//! # When the stream ends
//!
//! The spec requires closing on a terminal state. `INPUT_REQUIRED` and
//! `AUTH_REQUIRED` are interrupted rather than terminal, so they remain open:
//! an out-of-band answer may resume the task without another client request.
//! Intermediaries may reap a very idle connection; reconnecting is safe because
//! the stream is rebuilt from the journal rather than resumed from memory.
//!
//! It also ends when **this instance** is stopping, for the same reason and with
//! the same remedy. A subscription outlives any one process by reconnecting, and
//! a stream that did not end would hold the server's graceful shutdown open for
//! as long as the run it watches — which for a run sleeping five working days is
//! until the supervisor gives up and kills it.

use std::sync::Arc;
use std::time::Duration;

use axum::response::sse::{Event, Sse};
use futures_util::stream::{Stream, StreamExt as _};
use serde_json::{Value, json};

use crate::core::{RunId, Seq};
use crate::journal::RecordKind;
use crate::runtime::Runtime;

use super::a2a::{
    A2aArtifact, A2aTask, ArtifactCache, StreamSlot, TaskState, cached_artifacts, sealed_state,
    task_artifacts,
};

/// How often a subscriber re-reads the journal.
///
/// Short enough that progress feels live, long enough that a hundred
/// subscribers are not a hundred reads per millisecond.
const POLL: Duration = Duration::from_millis(200);

/// One `StreamResponse`, as the wire carries it.
///
/// A oneof: exactly one field is present. Built here rather than by the caller
/// so the "exactly one" part cannot be got wrong in four places.
fn stream_response(id: &Value, payload: &Value) -> Event {
    Event::default().data(json!({"jsonrpc": "2.0", "id": id, "result": payload}).to_string())
}

/// `TaskStatusUpdateEvent`.
///
/// `contextId` is required by the schema, and a run without a case has no case
/// id to put there. It carries the **run's own id** rather than an empty string:
/// a standalone run genuinely is its whole context, so this is a true statement
/// about grouping rather than a placeholder a client has to special-case.
pub(super) fn status_update(
    run: RunId,
    case: Option<&str>,
    state: TaskState,
    detail: &str,
) -> Value {
    json!({
        "statusUpdate": {
            "taskId": run.to_string(),
            "contextId": case.unwrap_or(&run.to_string()),
            "status": {
                "state": state,
                "message": {
                    "messageId": format!("{run}-{detail}"),
                    "role": "ROLE_AGENT",
                    "parts": [{"text": detail}],
                    "taskId": run.to_string(),
                },
            },
        }
    })
}

pub(super) fn artifact_update(run: RunId, case: Option<&str>, artifact: &A2aArtifact) -> Value {
    json!({
        "artifactUpdate": {
            "taskId": run.to_string(),
            "contextId": case.unwrap_or(&run.to_string()),
            "artifact": artifact,
            "append": false,
            "lastChunk": true,
        }
    })
}

/// What a journal record says about progress, if anything a caller can use.
///
/// Deliberately not every record: a subscriber wants to know *what is
/// happening*, and a stream that narrates internal bookkeeping is one people
/// stop reading. Records with no caller-visible meaning produce no event.
pub(super) fn progress_of(kind: &RecordKind) -> Option<(TaskState, String)> {
    match kind {
        RecordKind::StepStarted { skill } => Some((TaskState::Working, format!("started {skill}"))),
        RecordKind::StepFinished { outcome } => {
            Some((TaskState::Working, format!("finished: {outcome}")))
        }
        RecordKind::RunSuspended { reason } => {
            Some((TaskState::InputRequired, format!("waiting: {reason}")))
        }
        RecordKind::RunConcluded { outcome, .. } => Some((sealed_state(outcome), outcome.clone())),
        _ => None,
    }
}

/// `StreamResponse` payloads represented by one durable record.
pub(super) async fn payloads_for_record(
    runtime: &Runtime,
    record: &crate::journal::Record,
    case: Option<&str>,
) -> Result<Vec<Value>, crate::core::RuntimeError> {
    let run = record.body.run;
    if let RecordKind::RunConcluded { outcome, .. } = record.kind() {
        let state = sealed_state(outcome);
        let mut payloads = Vec::new();
        if state == TaskState::Completed
            && let Some(artifacts) = task_artifacts(runtime, run, state).await?
        {
            payloads.extend(
                artifacts
                    .iter()
                    .map(|artifact| artifact_update(run, case, artifact)),
            );
        }
        payloads.push(status_update(run, case, state, outcome));
        return Ok(payloads);
    }
    Ok(progress_of(record.kind())
        .map(|(state, detail)| vec![status_update(run, case, state, &detail)])
        .unwrap_or_default())
}

/// Whether a state ends the stream.
///
/// Only terminal states end a subscription. `INPUT_REQUIRED` can receive a
/// later message or out-of-band authorization and therefore remains live under
/// A2A 1.0.
pub(super) const fn closes(state: TaskState) -> bool {
    matches!(
        state,
        TaskState::Completed | TaskState::Failed | TaskState::Canceled | TaskState::Rejected
    )
}

/// Stream a run's progress, starting from `from`.
///
/// The first event is always the `Task` itself, which the spec requires: a
/// subscriber must be able to learn the current state without having been
/// present for the events that produced it.
///
/// `slot` is held for as long as the stream lives, so the caller's count of
/// open streams falls when it ends or when the client goes away.
#[allow(clippy::too_many_arguments)]
pub(super) fn tail(
    runtime: Arc<Runtime>,
    cache: Arc<std::sync::Mutex<ArtifactCache>>,
    slot: StreamSlot,
    run: RunId,
    case: Option<String>,
    id: Value,
    first: A2aTask,
    from: Seq,
) -> Sse<impl Stream<Item = Result<Event, std::convert::Infallible>>> {
    Sse::new(
        frames(runtime, cache, run, case, id, first, from).map(move |event| {
            let _held = &slot;
            event
        }),
    )
}

/// The event sequence behind [`tail`], before axum wraps it.
///
/// Split out because `Sse` hands back no way to read what it will send. The
/// request-level tests reach a terminal task only through `SubscribeToTask`,
/// which refuses one before it gets here, so the `already_over` return below
/// is pinned by a test that drives this function directly.
fn frames(
    runtime: Arc<Runtime>,
    cache: Arc<std::sync::Mutex<ArtifactCache>>,
    run: RunId,
    case: Option<String>,
    id: Value,
    first: A2aTask,
    from: Seq,
) -> impl Stream<Item = Result<Event, std::convert::Infallible>> {
    // Whether the run is *already* over when the stream opens. Checked here and
    // not only in the loop below: the record that ended it was consumed before
    // this subscriber existed, so the loop would never see it and would poll a
    // finished run forever. Subscribing to a completed task is the ordinary
    // case — a client reconnecting after a drop does exactly that.
    let already_over = closes(first.status.state);
    let stream = async_stream::stream! {
        yield Ok(stream_response(&id, &json!({ "task": first })));
        if already_over {
            return;
        }

        let mut next = from;
        loop {
            // A read that fails is not a run that failed. The stream ends rather
            // than reporting a terminal state the run never reached — a client
            // that reconnects gets the truth from the journal.
            let Ok(records) = runtime.journal().read(run, next).await else {
                return;
            };

            let mut done = false;
            for record in &records {
                next = record.body.seq + 1;
                if let RecordKind::RunConcluded { outcome, .. } = record.kind() {
                    let state = sealed_state(outcome);
                    if state == TaskState::Completed {
                        match cached_artifacts(&runtime, &cache, run, state).await {
                            Ok(Some(artifacts)) => {
                                for artifact in artifacts {
                                    yield Ok(stream_response(
                                        &id,
                                        &artifact_update(run, case.as_deref(), &artifact),
                                    ));
                                }
                            }
                            Ok(None) => {}
                            Err(_) => return,
                        }
                    }
                    yield Ok(stream_response(
                        &id,
                        &status_update(run, case.as_deref(), state, outcome),
                    ));
                    done = true;
                    continue;
                }
                if let Some((state, detail)) = progress_of(record.kind()) {
                    yield Ok(stream_response(
                        &id,
                        &status_update(run, case.as_deref(), state, &detail),
                    ));
                    done |= closes(state);
                }
            }
            if done {
                return;
            }
            // A subscription to a run that sleeps for five working days is a
            // connection that never closes, and a server draining gracefully
            // waits for every connection — so without this an ordinary deploy
            // would hang until the supervisor's `SIGKILL`, which is the one stop
            // nothing can drain. Ending here is the failure this module is
            // already shaped for: the stream is a view of history, so the client
            // reconnects, is told the current state from the journal, and is
            // served by whichever instance is up.
            if runtime.is_draining() {
                return;
            }
            tokio::time::sleep(POLL).await;
        }
    };

    // Deliberately **no** keep-alive. It was tried, and with it the response
    // body did not end when the stream did: the connection outlived the task,
    // which is precisely the failure this whole module is shaped to avoid — a
    // client holding a socket open for a run that already finished.
    //
    // The trade is that a very idle stream may be reaped by an intermediary. That
    // is the better failure: the client reconnects and is told the current state
    // from the journal, because the stream is a view of history rather than a
    // subscription to memory. A connection that never ends cannot be recovered
    // from by anybody.
    stream
}

#[cfg(all(test, feature = "redb"))]
mod tests {
    use super::{Event, RunId, TaskState, frames};
    use futures_util::StreamExt as _;
    use serde_json::json;
    use std::sync::Arc;

    /// **A stream opened on an already-finished task ends.**
    ///
    /// It yields the task snapshot the spec requires and then stops, rather than
    /// polling a run whose closing record was written before this subscriber
    /// existed — a loop that would never see it and never end.
    ///
    /// `SubscribeToTask` is refused on a terminal task before it reaches the
    /// stream, so a test of that refusal says nothing about this path; the
    /// method that does reach it is `SendStreamingMessage`, which has no such
    /// pre-check.
    #[tokio::test]
    async fn a_stream_on_an_already_finished_task_ends() {
        let store = Arc::new(crate::store::RedbStore::open_in_memory().expect("store"));
        let runtime = crate::runtime::Runtime::builder(
            Arc::clone(&store) as Arc<dyn crate::journal::JournalStore>
        )
        .build();
        let run = RunId::generate();

        // Terminal before anybody subscribed, which is the whole condition.
        let finished = crate::api::a2a::task_of(run, TaskState::Completed, "succeeded", None);

        // Bounded on purpose. Without the guard this stream never ends, so a
        // plain `collect()` would **hang** rather than fail — and a hanging test
        // is a CI timeout with no message, which is a worse report than a red
        // assertion. The deadline turns "never ends" into a sentence.
        let collected: Vec<Result<Event, std::convert::Infallible>> = tokio::time::timeout(
            std::time::Duration::from_secs(5),
            frames(runtime, Arc::default(), run, None, json!(1), finished, 1).collect(),
        )
        .await
        .expect(
            "a stream on an already-finished task never ended: its closing record was \
             written before this subscriber existed, so the loop will never see it",
        );

        assert_eq!(
            collected.len(),
            1,
            "a stream on a finished task must yield the snapshot and stop; it \
             produced {} events",
            collected.len()
        );
    }

    /// **A stream ends when this instance is stopping.**
    ///
    /// A subscription to a run that has not concluded polls indefinitely by
    /// design, and a server draining gracefully waits for every open connection
    /// — so without this an ordinary deploy hangs until the supervisor's
    /// `SIGKILL`, which is the one stop nothing can drain. The bound below turns
    /// "never ends" into a sentence rather than a CI timeout, exactly as the
    /// test above does.
    ///
    /// What this does **not** cover is the client's side: reconnecting is what
    /// makes ending here safe, and that is the module's standing contract rather
    /// than something a drain introduced.
    #[tokio::test]
    async fn a_stream_ends_when_the_instance_is_draining() {
        let store = Arc::new(crate::store::RedbStore::open_in_memory().expect("store"));
        let runtime = crate::runtime::Runtime::builder(
            Arc::clone(&store) as Arc<dyn crate::journal::JournalStore>
        )
        .build();
        let run = RunId::generate();

        // Working, so the loop would otherwise poll a run that never concludes.
        let working = crate::api::a2a::task_of(run, TaskState::Working, "accepted", None);
        runtime.drain(std::time::Duration::ZERO).await;

        let collected: Vec<Result<Event, std::convert::Infallible>> = tokio::time::timeout(
            std::time::Duration::from_secs(5),
            frames(runtime, Arc::default(), run, None, json!(1), working, 1).collect(),
        )
        .await
        .expect(
            "a stream did not end while the instance was draining: it would hold the \
             server's graceful shutdown open for as long as the run it watches",
        );

        assert_eq!(
            collected.len(),
            1,
            "a draining instance must yield the snapshot and stop; it produced {} events",
            collected.len()
        );
    }
}