loopflow 0.12.9

Run steps and flows with coding agents
Documentation
//! The resident wire: what crosses between the listener (the wave server,
//! vendor-free — hear / check / fold / tell) and the resident (the loop, its
//! own `lf` process owning the vendor harness).
//!
//! Two directions, two transports:
//!
//! - **Resident → listener**: ordered turn deltas over `POST /resident/deltas`
//!   (the resident sends serially and awaits each response, so per-turn order
//!   is the connection's order). The old in-process `TurnSink` vocabulary —
//!   Opened / Text / Item / Finished — IS this wire, promoted to full
//!   DTO discipline, plus the consumption markers (`TurnOpened.answers`,
//!   `TurnSteered.answers` — the RESIDENT decides what a turn answers; the
//!   listener validates against its queue fold and journals), the resident's
//!   reported loop state, and the body's provider session id. The single writer stays
//!   with the listener: the resident never touches journal files.
//! - **Listener → resident**: the resident consumes its own wave's `/events`
//!   subscription with `?inbox=true` — `inbox` SSE frames ([`InboxFrame`])
//!   carry queued messages, typed Task observations, steer, and interrupt ops
//!   (replayed pending queue on connect, then live).
//!
//! DTO discipline: every field is required or explicitly `Option` — no serde
//! defaults anywhere. The round-trip fixture lives at
//! `tests/fixtures/dto/resident_deltas.json`. Carve-out: unlike the other DTO
//! fixtures this wire is Rust↔Rust only (listener and resident are the same
//! binary); Swift and Python do not consume it, so only the Rust fixture test
//! pins it.
//!
//! # Auth
//! The resident door is gated by a per-boot bearer token
//! ([`RESIDENT_TOKEN_HEADER`]): the listener generates it at bind, passes it
//! to a spawned resident via [`RESIDENT_TOKEN_ENV`], and writes it to
//! `wave/<name>/.wave-resident-token` beside the endpoint pointer for
//! the internal resident — the same filesystem-trust domain as the discovery
//! file. This is a stopgap: when a remote Loop can hold the resident seat, the
//! token becomes a credential the gatekeeper issues, not a file the repo trusts.

use serde::{Deserialize, Serialize};

use crate::chat::types::{ConversationItem, Lifecycle};
use crate::project::ProjectObservation;
use crate::task::TaskObservation;
use crate::wave::journal::{DiscordMessageSource, MessageOp};
use crate::wave::playhead::{BodyProvenance, PlayheadView, StepOutcome};

/// Header carrying the resident token on every `/resident/*` request.
pub const RESIDENT_TOKEN_HEADER: &str = "x-lf-resident-token";

/// Env var the listener sets on a spawned resident.
pub const RESIDENT_TOKEN_ENV: &str = "LF_WAVE_RESIDENT_TOKEN";

/// Basename of the token file beside `.wave-endpoint` (attached residents).
pub const RESIDENT_TOKEN_FILE: &str = ".wave-resident-token";

/// One ordered increment from the resident's harness stream, applied by the
/// listener's fold ([`crate::wave::runtime::WaveRuntime::apply_resident_delta`]).
///
/// Per-turn order is the transport's order; turn ids never ride the wire —
/// the listener mints them from its journal seq exactly as before.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum ResidentDelta {
    /// A vendor turn opened. `answers` is the consumption declaration: the
    /// queued message ids this turn's input consumed (the resident drained
    /// its queue). The listener validates each id against its pending fold —
    /// unknown or already-consumed ids are dropped with a warning — and
    /// journals the valid set in `TurnStarted.answers`.
    TurnOpened { answers: Vec<String> },
    /// A prose fragment of the open turn.
    TurnText { text: String },
    /// A non-prose item (tool / command / file / thought) of the open turn.
    TurnItem { item: ConversationItem },
    /// The open turn finalized. The listener already holds the turn's content
    /// grown from the Text/Item deltas.
    TurnFinished {
        status: Lifecycle,
        reason: Option<String>,
    },
    /// Mid-turn consumption: the harness accepted these queued messages as
    /// steering input, so the CURRENT turn answers them (journaled as
    /// `TurnSteered.answers`). Validated like `TurnOpened.answers`.
    TurnSteered { answers: Vec<String> },
    /// The undo of a consumption claim: these message ids were declared
    /// consumed (`TurnSteered`) but the vendor never received the input —
    /// the harness send failed AFTER the claim was journaled. The listener
    /// returns them to its pending fold so the next resident's replay
    /// re-delivers them. The claim rides first, the undo is explicit:
    /// at-most-once to the vendor, never a silent redelivery.
    MessagesRequeued { ids: Vec<String> },
    /// A fresh body took the current logical playhead step.
    BodyStarted { body: BodyProvenance },
    /// The harness announced its provider session after the body opened.
    BodySessionUpdated { body_id: String, session_id: String },
    /// A body ended. Only completed/skipped outcomes advance the playhead;
    /// failure/interruption leave the logical step selected for retry.
    BodyFinished {
        body_id: String,
        outcome: StepOutcome,
        reason: String,
    },
    /// The resident's reported loop state ([`ResidentStateTo`]). `Turning`
    /// and the boundary `Idle` are DERIVED by the listener from
    /// `TurnOpened`/`TurnFinished`; only the transitions the turn deltas
    /// can't express ride here.
    LoopState { to: ResidentStateTo, reason: String },
}

/// Destination of a reported [`ResidentDelta::LoopState`] transition. The
/// listener supplies the turn id for `Interrupting` (the current open turn —
/// the resident never learns journal-minted ids).
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ResidentStateTo {
    /// Cancel fired for the open turn (cooperative interrupt in flight).
    Interrupting,
    /// The loop itself died (harness terminal error, failure cap). The
    /// resident reports this and exits; the listener's supervisor owns the
    /// respawn ladder from there.
    Failed,
}

/// `POST /resident/deltas` request. Deltas apply in order; the resident sends
/// batches serially (one in flight at a time), so order is total.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PostDeltasRequest {
    pub deltas: Vec<ResidentDelta>,
}

/// `POST /resident/deltas` response.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PostDeltasResponse {
    /// Number of deltas applied (dropped-by-guard deltas still count as
    /// accepted — the guard is the listener's business, not a retry signal).
    pub accepted: u64,
}

/// `POST /resident/attach` request: the resident's first call. Registers the
/// resident's pid for liveness (the listener probes attached residents; a
/// spawned child is watched by process exit) and revives a `failed` loop
/// state — a fresh resident IS the revival.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AttachRequest {
    pub pid: u32,
}

/// `POST /resident/attach` response: the bootstrap the resident needs before
/// its first turn. The queued messages arrive via the `/events?inbox=true`
/// subscription replay, not here — one door per direction.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AttachResponse {
    /// The wave this listener serves — the resident refuses a mismatch.
    pub wave: String,
}

/// `GET /resident/context` response: the pre-turn snapshot the resident folds
/// into its prompts. Serving it freshens the listener's child observations.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ContextResponse {
    pub playhead: PlayheadView,
    pub provider_session: Option<ProviderSessionRef>,
}

/// A provider thread is only meaningful to the harness that minted it.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ProviderSessionRef {
    pub harness: String,
    pub session_id: String,
}

/// One `inbox` SSE frame on `/events?inbox=true` — a resident-directed op.
/// A `Message` carries its journaled id (`"msg-<seq>"`). `Interrupt` and `Skip`
/// carry none: nothing is journaled, because the op is control, not content.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum InboxFrame {
    Message {
        id: String,
        op: MessageOp,
        text: String,
        source: Option<DiscordMessageSource>,
    },
    Task {
        observation: TaskObservation,
    },
    Project {
        observation: ProjectObservation,
    },
    Promotion {
        parent_wave_id: crate::id::WaveId,
        parent: String,
    },
    Interrupt,
    Skip,
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn resident_delta_round_trips_every_variant() {
        let deltas = vec![
            ResidentDelta::TurnOpened {
                answers: vec!["msg-1".into(), "msg-2".into()],
            },
            ResidentDelta::TurnText {
                text: "thinking".into(),
            },
            ResidentDelta::TurnItem {
                item: ConversationItem::Tool {
                    id: "t-1".into(),
                    name: "Bash".into(),
                    status: Lifecycle::Completed,
                    input: None,
                    output: Some("ok".into()),
                },
            },
            ResidentDelta::TurnFinished {
                status: Lifecycle::Completed,
                reason: None,
            },
            ResidentDelta::TurnSteered {
                answers: vec!["msg-3".into()],
            },
            ResidentDelta::MessagesRequeued {
                ids: vec!["msg-3".into()],
            },
            ResidentDelta::LoopState {
                to: ResidentStateTo::Failed,
                reason: "harness disconnected".into(),
            },
        ];
        for delta in deltas {
            let value = serde_json::to_value(&delta).expect("serialize");
            let decoded: ResidentDelta = serde_json::from_value(value).expect("deserialize");
            assert_eq!(decoded, delta);
        }
    }

    /// No serde defaults: an absent REQUIRED field is a parse error, never a
    /// silent fill-in. Absent `Option` fields decode as `None` — explicitly
    /// Optional is the one sanctioned absence.
    #[test]
    fn absent_required_fields_are_parse_errors() {
        for bad in [
            serde_json::json!({ "kind": "turn_opened" }),
            serde_json::json!({ "kind": "turn_finished" }),
            serde_json::json!({ "kind": "turn_text" }),
            serde_json::json!({ "kind": "loop_state", "to": "failed" }),
            serde_json::json!({ "kind": "messages_requeued" }),
        ] {
            assert!(
                serde_json::from_value::<ResidentDelta>(bad.clone()).is_err(),
                "must reject {bad}"
            );
        }
        assert!(serde_json::from_value::<InboxFrame>(serde_json::json!({
            "kind": "message",
            "text": "hi"
        }))
        .is_err());

        let frame: InboxFrame = serde_json::from_value(serde_json::json!({
            "kind": "message",
            "id": "msg-1",
            "op": "message",
            "text": "hi"
        }))
        .expect("message frame parses");
        assert!(matches!(frame, InboxFrame::Message { .. }));
    }
}