/** * T6 · Mux session — the multiplexer over ONE WebSocket tunnel. Assembles the T2 codec, * T3 stream state machine, T4 credit flow-control, and T5 heartbeat into a demux/dispatch loop. * * Security (INV2/INV5/INV11): DATA payloads are copied OPAQUE and never inspected; buffers are * per-stream, per-session (no global pool → no cross-tenant buffer bleed) and freed on CLOSE; * this module imports NO terminal/ANSI parser. An inbound DATA for an unknown/closed stream, an * illegal transition, or a `payloadLen > maxFrameBytes` frame RSTs that ONE stream — the tunnel * stays up (INV13). Ordering is preserved and never reordered/deduped (so P4's `seq` holds). */ import { concatBytes, GOAWAY_CODE_TO_REASON, type MuxOpen } from 'relay-contracts' import { MUX_HEADER_BYTES, encodeMuxFrame, decodeHeader, encodeOpen, decodeOpen, encodeWindowUpdate, decodeWindowUpdate, encodeGoaway, type MuxFrameHeader, type MuxFrameType, } from './frame-codec.js' import { assertWithinFrameCeiling } from './frame-guards.js' import { initialStreamState, nextStreamState, isTerminal, type StreamState, } from './stream.js' import { createFlowController, type FlowController } from './flow-control.js' import { createHeartbeat, type Heartbeat, type ScheduleHandle, } from './heartbeat.js' export interface MuxStreamHandle { readonly streamId: number writeData(payload: Uint8Array): boolean // false ⇒ backpressured (buffered until WINDOW_UPDATE) close(rst?: boolean): void onData(cb: (payload: Uint8Array) => void): void // inbound opaque bytes for THIS stream (INV2) onClose(cb: (rst: boolean) => void): void // remote/RST close of THIS stream } export type MuxRole = 'relay' | 'agent' export interface MuxSessionDeps { role: MuxRole sendWire(frame: Uint8Array): void onOpen(open: MuxOpen, stream: MuxStreamHandle): void onData(streamId: number, payload: Uint8Array): void // fallback if no per-stream subscriber onClose(streamId: number, rst: boolean): void onDead(): void maxFrameBytes: number initialWindowBytes: number // P1-owned testability seam (NOT a frozen contract): inject the heartbeat timer. schedule?: (fn: () => void, ms: number) => ScheduleHandle heartbeatIntervalMs?: number } export interface MuxSession { openStream(open: MuxOpen): MuxStreamHandle onWire(buf: Uint8Array): void drain(lastStreamId: number, reason: number): void close(): void } interface StreamCtx { state: StreamState readonly handle: MuxStreamHandle sendBuffer: Uint8Array[] dataSub: ((payload: Uint8Array) => void) | null closeSub: ((rst: boolean) => void) | null } const EMPTY = new Uint8Array(0) export function createMuxSession(deps: MuxSessionDeps): MuxSession { const flow: FlowController = createFlowController() const streams = new Map() let pending: Uint8Array = EMPTY let nextStreamId = 1 // relay allocates monotonic, never reused let draining = false let closed = false const heartbeat: Heartbeat = createHeartbeat({ sendPing: (token) => sendFrame('ping', 0, token, false, false), onDead: () => deps.onDead(), ...(deps.schedule ? { schedule: deps.schedule } : {}), ...(deps.heartbeatIntervalMs !== undefined ? { intervalMs: deps.heartbeatIntervalMs } : {}), }) function sendFrame( type: MuxFrameType, streamId: number, payload: Uint8Array, fin: boolean, rst: boolean, ): void { if (closed) return const header: MuxFrameHeader = { version: 1, type, fin, rst, streamId, payloadLen: payload.length, } deps.sendWire(encodeMuxFrame(header, payload)) } function makeHandle(streamId: number): MuxStreamHandle { return { streamId, writeData: (payload) => writeData(streamId, payload), close: (rst = false) => localClose(streamId, rst), onData: (cb) => { const ctx = streams.get(streamId) if (ctx) ctx.dataSub = cb }, onClose: (cb) => { const ctx = streams.get(streamId) if (ctx) ctx.closeSub = cb }, } } function registerStream(streamId: number, state: StreamState): StreamCtx { const handle = makeHandle(streamId) const ctx: StreamCtx = { state, handle, sendBuffer: [], dataSub: null, closeSub: null } streams.set(streamId, ctx) flow.registerStream(streamId, deps.initialWindowBytes) return ctx } function writeData(streamId: number, payload: Uint8Array): boolean { const ctx = streams.get(streamId) if (ctx === undefined || isTerminal(ctx.state)) return false if (payload.length > deps.maxFrameBytes) return false // never emit an over-ceiling frame if (!flow.canSend(streamId, payload.length)) { ctx.sendBuffer.push(payload) // backpressure: buffer until WINDOW_UPDATE return false } flow.consumeSendCredit(streamId, payload.length) sendFrame('data', streamId, payload, false, false) return true } function flushBuffer(streamId: number): void { const ctx = streams.get(streamId) if (ctx === undefined) return while (ctx.sendBuffer.length > 0) { const next = ctx.sendBuffer[0]! if (!flow.canSend(streamId, next.length)) break ctx.sendBuffer.shift() flow.consumeSendCredit(streamId, next.length) sendFrame('data', streamId, next, false, false) } } function localClose(streamId: number, rst: boolean): void { const ctx = streams.get(streamId) if (ctx === undefined) return sendFrame('close', streamId, EMPTY, !rst, rst) cleanupStream(streamId) } function cleanupStream(streamId: number): void { flow.releaseStream(streamId) streams.delete(streamId) } /** RST exactly ONE stream (illegal transition / unknown stream / ceiling) — tunnel stays up. */ function rstStream(streamId: number): void { const ctx = streams.get(streamId) sendFrame('close', streamId, EMPTY, false, true) if (ctx) { if (ctx.closeSub) ctx.closeSub(true) else deps.onClose(streamId, true) } cleanupStream(streamId) } function deliverData(ctx: StreamCtx, streamId: number, payload: Uint8Array): void { if (ctx.dataSub) ctx.dataSub(payload) else deps.onData(streamId, payload) // Receiver-side replenishment: grant credit back once the threshold is crossed. const replenish = flow.onDelivered(streamId, payload.length) if (replenish > 0) sendFrame('windowUpdate', streamId, encodeWindowUpdate(replenish), false, false) } function dispatch(header: MuxFrameHeader, payload: Uint8Array): void { const { type, streamId, fin, rst } = header // Connection-level control (streamId 0). if (streamId === 0) { if (type === 'ping') sendFrame('pong', 0, payload, false, false) else if (type === 'pong') heartbeat.onPong(payload) else if (type === 'goaway') draining = true return } if (type === 'open') { if (deps.role !== 'agent') return // relay never receives OPEN if (streams.get(streamId) !== undefined) { rstStream(streamId) // re-OPEN of a live stream is illegal return } const open = decodeOpen(payload) const ctx = registerStream(streamId, 'open') deps.onOpen(open, ctx.handle) return } const ctx = streams.get(streamId) if (ctx === undefined) { // DATA/CLOSE/WU for an unknown/closed stream → RST that stream only. sendFrame('close', streamId, EMPTY, false, true) return } const transition = nextStreamState(ctx.state, type, fin, 'inbound') if ('illegal' in transition) { rstStream(streamId) return } ctx.state = transition.next if (type === 'data') { deliverData(ctx, streamId, payload) } else if (type === 'windowUpdate') { flow.grantCredit(streamId, decodeWindowUpdate(payload)) flushBuffer(streamId) } else if (type === 'close') { if (ctx.closeSub) ctx.closeSub(rst) else deps.onClose(streamId, rst) cleanupStream(streamId) } } return { openStream(open) { if (draining || closed) throw new Error('session draining/closed: cannot open new stream') const streamId = nextStreamId++ const ctx = registerStream(streamId, 'open') const full: MuxOpen = { ...open, streamId } sendFrame('open', streamId, encodeOpen(full), false, false) heartbeat.start() return ctx.handle }, onWire(buf) { if (closed) return heartbeat.start() pending = pending.length === 0 ? buf : concatBytes([pending, buf]) for (;;) { if (pending.length < MUX_HEADER_BYTES) break let header: MuxFrameHeader try { header = decodeHeader(pending) } catch { // Unrecoverable framing error — drop the tunnel's buffer (do NOT parse further). pending = EMPTY break } try { assertWithinFrameCeiling(header.payloadLen, deps.maxFrameBytes) } catch { if (header.streamId > 0) rstStream(header.streamId) else deps.onDead() pending = EMPTY break } const total = MUX_HEADER_BYTES + header.payloadLen if (pending.length < total) break // wait for the rest of the payload const payload = pending.slice(MUX_HEADER_BYTES, total) pending = pending.slice(total) dispatch(header, payload) } }, drain(lastStreamId, reason) { draining = true const label = GOAWAY_CODE_TO_REASON[reason] if (label !== undefined) sendFrame('goaway', 0, encodeGoaway(lastStreamId, label), false, false) }, close() { closed = true heartbeat.stop() streams.clear() pending = EMPTY }, } }