/** * The §13.6 SESSION RAILS: the framed protocol + bounded credit window over the two epoch-pinned * eps subjects (cast-in / watch-out). SPLIT out of endpoint-session.ts (P2 item 6) so it carries * ZERO node-only dependencies (no node:crypto) and can be BUNDLED into the browser console session * client — the SAME core flow-control code runs in-browser. The grant mint/verify/redeem + the KV * ledger (node:crypto, KV) stay in endpoint-session.ts, which RE-EXPORTS this module so no * consumer's import path changes. Pure mechanical split: zero signature/behavior change. */ import type { NatsConnection } from "@nats-io/transport-node"; import type { SessionGrant } from "./endpoint-session.js"; /** The bounded flow window (§13.6: declared in the grant; overflow is `resource-exhausted`). */ export declare const SESSION_WINDOW_DEFAULT = 64; export declare const SESSION_WINDOW_MAX = 1024; /** Validate a flow window (1..MAX). Exported so the grant mint/verify in endpoint-session.ts share * the single definition (it lives here, the browser-safe module). */ export declare function assertWindow(v: unknown): number; /** The composite's own tiny framed protocol; `data` is OPAQUE (any JSON value — binary rides * the application's own encoding). `credit`/`close` are CONTROL frames, EXEMPT from the data * window (a full data window must never block the credits that reopen it, else instant * deadlock). `ack` is an ABSOLUTE cumulative watermark (the sender's contiguous-received count * on the OTHER rail): a data frame PIGGYBACKS it, so a lost dedicated credit self-heals on the * next reverse data frame, and any deeper loss recovers on the keepalive re-emit; absolute * (not delta) so any single credit re-advertises the whole position and a duplicate is * harmless. The in-band `close` is advisory (§13.6): revocation authority is the ledger, * never this frame. */ export type SessionFrame = { t: "f"; seq: number; data: unknown; ack?: number; } | { t: "credit"; ack: number; } | { t: "close"; }; export declare function encodeSessionFrame(frame: SessionFrame): Uint8Array; /** Fail-loud frame parse (closed schema): a garbled frame is a PROTOCOL error the rail * surfaces via `onProtocolError` — never silently skipped, never a crash. */ export declare function parseSessionFrame(bytes: Uint8Array): SessionFrame; /** Which rail each role sends on (§13.6: `in` = caller → endpoint; `out` = the reverse). */ export type SessionRole = "caller" | "serving"; export interface SessionRailOpts { nc: NatsConnection; grant: Pick & { serving: { epoch: number; }; }; role: SessionRole; /** Delivered in-order for CONTIGUOUS frames, and the application accepts FIRST: the handler * may be async — it is AWAITED, and the watermark advances and credit emits only after it * RESOLVES, so credit means the receiver's buffer actually freed (back-pressure) and a * rejection refuses the frame exactly like a synchronous throw: the rail breaks (`handler`) * and the refused frame is neither counted delivered nor credited. Acceptance is SERIALIZED * in seq order (one handler in flight; NATS does not serialize callback promises); frames * arriving while a handler is pending queue up to the grant WINDOW — past it the rail breaks * (`flood`), never an unbounded backlog. A handler wedged forever stalls credit, so the * SENDER's window fills and its stall watchdog surfaces the fault. A gap surfaces via * onProtocolError("gap"). */ onData(data: unknown, seq: number): void | Promise; /** The peer's advisory close frame arrived (authoritative close is the ledger's). The local * subscription and timer are torn down before this fires. */ onClose?(): void; /** The session is broken — close and re-establish. `reason` is one of `garbled-frame` | * `gap` | `credit-overrun` | `flood` | `subscription` | `stall` | `handler` | `publish` | * `seq-exhausted`. The rail's subscription and timer are torn down before this fires (a * broken rail holds no resources). */ onProtocolError?(reason: string, detail?: unknown): void; /** Broker payload ceiling for the SEND preflight (like assertFactFits). Default 1 MiB. */ maxPayloadBytes?: number; /** Keepalive credit re-emit interval (ms): while this side has delivered ANY data and the * peer has gone quiet, re-advertise the absolute watermark every tick — including * watermarks already advertised, because this side cannot observe whether an emitted * credit ARRIVED (gating on "newer than last emitted" turns loss of the advertisement * itself into a permanent stall). Absolute acks are idempotent, so the honest recovery is * repetition; the cost is one control frame per quiet tick. 0 disables. Default 1000. */ idleCreditMs?: number; /** Sender stall watchdog (ms): if the data window stays full this long with NO ack advance * (sustained credit loss or a dead peer), the rail breaks with a DETECTABLE `stall` fault — * TIMER-driven, so a sender that stops calling send() still learns its peer is gone; the * send path double-checks as a belt. 0 disables. Default 30000. */ stallTimeoutMs?: number; /** Injectable clock (testability); default Date.now. */ now?: () => number; /** Injectable interval timer (testability); defaults to Node setInterval/clearInterval. */ setIntervalFn?: (fn: () => void, ms: number) => { unref?: () => void; }; clearIntervalFn?: (h: unknown) => void; } export interface SessionRail { /** Send one opaque data frame (piggybacking this side's absolute reverse-rail watermark). * Throws `resource-exhausted` when the window is full (no buffering, §13.6), `contract-invalid` * when the encoded frame exceeds the payload ceiling, and `failed-precondition` once the rail * is closed/broken (including a detected stall or a failed publish). Returns the frame's seq. */ send(data: unknown): number; /** Send the advisory close frame and stop the rail locally. Idempotent. */ close(): void; /** In-memory window state (observability + smoke assertions). */ stats(): { sent: number; ackedThrough: number; delivered: number; inFlight: number; }; } /** * Open one side of an established session over its two core rails. The credentials the * redemption released confine each side to exactly its pub/sub pair; this helper only speaks * the framed protocol and enforces the bounded window — it grants nothing. * * FLOW CONTROL (panel-locked): the data window is bounded and per-direction; control frames * (`credit`, `close`) are EXEMPT (a full window never blocks the credits that reopen it). * RECEIVE-side acceptance is serialized and the (possibly async) handler AWAITED — credit * emits only for frames the application actually accepted — and the pending-frame queue is * bounded by the same window (`flood` past it), so neither side ever buffers unboundedly. * Credits carry an ABSOLUTE cumulative watermark, PIGGYBACKED on reverse data frames, so a lost * dedicated credit self-heals on the next reverse traffic; ANY deeper loss (including loss of * already-emitted threshold credits) recovers on the KEEPALIVE re-emit; sustained loss or a * dead peer surfaces the TIMER-driven `stall` fault (never a silent hang, even for a sender * that stopped calling send). A dropped DATA frame is unrecoverable at this transport (EPS is * at-most-once, core-only) and shows as a seq gap the app reacts to — reliability layers * inside `data` or uses the journal/checkpoint composites. */ export declare function openSessionRail(opts: SessionRailOpts): SessionRail; //# sourceMappingURL=endpoint-session-rail.d.ts.map