401 lines
13 KiB
TypeScript
401 lines
13 KiB
TypeScript
/**
|
|
* Transports: the byte pipe under a lockstep session.
|
|
*
|
|
* A transport is deliberately primitive — send bytes, receive bytes, learn when
|
|
* the pipe dies. Framing, ordering, reliability and delivery *guarantees* are the
|
|
* transport's job, not the session's, because WebSocket and RTCDataChannel are
|
|
* both message-oriented and both reliable: a lost tick is a stall that recovers,
|
|
* never a silently corrupted world.
|
|
*
|
|
* Two implementations ship here. {@link memoryTransportPair} is an in-process
|
|
* pair used by the tests, and {@link socketTransport} wraps anything that looks
|
|
* like a `WebSocket` — which includes the browser's and Node's global, so the
|
|
* same code path is exercised headlessly against a real socket.
|
|
*/
|
|
|
|
/** A reliable, message-oriented byte pipe to one peer. */
|
|
export interface Transport {
|
|
/**
|
|
* Send one message.
|
|
*
|
|
* @param data - the encoded message.
|
|
*/
|
|
send: (data: Uint8Array) => void
|
|
/**
|
|
* Register the message handler. Only one handler is supported.
|
|
*
|
|
* @param handler - called for each received message.
|
|
*/
|
|
onMessage: (handler: (data: Uint8Array) => void) => void
|
|
/**
|
|
* Register the close handler. Only one handler is supported.
|
|
*
|
|
* @param handler - called once when the pipe closes.
|
|
*/
|
|
onClose: (handler: () => void) => void
|
|
/** Close the pipe. Idempotent. */
|
|
close: () => void
|
|
/** Whether the pipe is currently usable. */
|
|
readonly open: boolean
|
|
}
|
|
|
|
/** Options shared by the transports in this file. */
|
|
export interface TransportOptions {
|
|
/** Maximum accepted message size, in bytes. Defaults to 4096. */
|
|
readonly maxMessageBytes?: number
|
|
}
|
|
|
|
/** Largest message accepted by default. */
|
|
const DEFAULT_MAX_MESSAGE = 4096
|
|
|
|
/**
|
|
* Build a connected pair of in-process transports.
|
|
*
|
|
* Delivery is deferred to {@link MemoryTransport.flush} rather than immediate, so
|
|
* tests can hold messages back to model a slow or stalled link — the case that
|
|
* lockstep has to survive.
|
|
*
|
|
* @param options - size limits.
|
|
* @returns the two ends, indexed by peer.
|
|
*/
|
|
export function memoryTransportPair(options: TransportOptions = {}): [MemoryTransport, MemoryTransport] {
|
|
const a = new MemoryTransport('a', options)
|
|
const b = new MemoryTransport('b', options)
|
|
a.link = b
|
|
b.link = a
|
|
return [a, b]
|
|
}
|
|
|
|
/**
|
|
* One end of an in-process transport pair.
|
|
*/
|
|
export class MemoryTransport implements Transport {
|
|
/** The other end. */
|
|
link: MemoryTransport | null = null
|
|
/** Messages sent but not yet delivered. */
|
|
private readonly pending: Uint8Array[] = []
|
|
/** A label, used in error messages. */
|
|
readonly name: string
|
|
private readonly maxMessageBytes: number
|
|
private messageHandler: ((data: Uint8Array) => void) | null = null
|
|
private closeHandler: (() => void) | null = null
|
|
private isOpen = true
|
|
/** How many messages this end has discarded because the link was closed. */
|
|
dropped = 0
|
|
/**
|
|
* When true, this end's queue is not delivered.
|
|
*
|
|
* Models a link that is slower than the game's tick rate — the case lockstep
|
|
* exists to survive, and the case that must stall the world rather than corrupt
|
|
* it.
|
|
*/
|
|
hold = false
|
|
|
|
/**
|
|
* @param name - a label, used in error messages.
|
|
* @param options - size limits.
|
|
*/
|
|
constructor(name: string, options: TransportOptions = {}) {
|
|
this.name = name
|
|
this.maxMessageBytes = options.maxMessageBytes ?? DEFAULT_MAX_MESSAGE
|
|
}
|
|
|
|
/** Whether the pipe is usable. */
|
|
get open(): boolean {
|
|
return this.isOpen && this.link !== null && this.link.isOpen
|
|
}
|
|
|
|
/** Messages waiting to be delivered to the far end. */
|
|
get queued(): number {
|
|
return this.pending.length
|
|
}
|
|
|
|
/**
|
|
* Queue a message for the far end.
|
|
*
|
|
* @param data - the bytes.
|
|
*/
|
|
send(data: Uint8Array): void {
|
|
if (!this.isOpen || this.link === null || !this.link.isOpen) {
|
|
this.dropped += 1
|
|
return
|
|
}
|
|
if (data.byteLength > this.maxMessageBytes) throw new Error(`${this.name}: message of ${String(data.byteLength)} bytes exceeds ${String(this.maxMessageBytes)}`)
|
|
this.pending.push(data.slice())
|
|
}
|
|
|
|
/** @param handler - message handler. */
|
|
onMessage(handler: (data: Uint8Array) => void): void {
|
|
this.messageHandler = handler
|
|
}
|
|
|
|
/** @param handler - close handler. */
|
|
onClose(handler: () => void): void {
|
|
this.closeHandler = handler
|
|
}
|
|
|
|
/** Close this end; the far end observes the close. */
|
|
close(): void {
|
|
if (!this.isOpen) return
|
|
this.isOpen = false
|
|
this.pending.length = 0
|
|
this.closeHandler?.()
|
|
this.link?.close()
|
|
}
|
|
|
|
/**
|
|
* Deliver everything queued, on both ends.
|
|
*
|
|
* Called by the test harness once per simulated tick. Because each delivery can
|
|
* queue more messages, this drains until the pair is quiet or the round cap is
|
|
* hit — a cap rather than a `while (true)` so a handler that replies to itself
|
|
* fails loudly instead of hanging.
|
|
*
|
|
* @param rounds - maximum delivery rounds. Defaults to 8.
|
|
*/
|
|
flush(rounds = 8): void {
|
|
for (let round = 0; round < rounds; round += 1) {
|
|
const a = this.drain(this)
|
|
const b = this.link === null ? 0 : this.drain(this.link)
|
|
if (a + b === 0) return
|
|
}
|
|
throw new Error('memory transport did not settle')
|
|
}
|
|
|
|
/**
|
|
* Deliver this end's queue.
|
|
*
|
|
* @param side - the end whose queue should be delivered.
|
|
* @returns how many messages were delivered.
|
|
*/
|
|
private drain(side: MemoryTransport): number {
|
|
let delivered = 0
|
|
while (!side.hold && side.pending.length > 0) {
|
|
const data = side.pending.shift()!
|
|
const target = side.link
|
|
if (target === null || !target.isOpen) { side.dropped += 1; continue }
|
|
target.messageHandler?.(data)
|
|
delivered += 1
|
|
}
|
|
return delivered
|
|
}
|
|
}
|
|
|
|
/**
|
|
* An in-process hub: the same shape as the relay, without the sockets.
|
|
*
|
|
* Three or more peers cannot be built from pairs, and the interesting failures in
|
|
* a session of four are the ones a two-peer test cannot reproduce at all — a late
|
|
* joiner, one peer that stops answering while the others carry on. This mirrors
|
|
* the relay's one rule (everyone hears everything, except the sender) so those
|
|
* cases can be driven a tick at a time.
|
|
*/
|
|
export class MemoryHub {
|
|
private readonly ends: HubTransport[] = []
|
|
private readonly queue: { readonly from: number; readonly data: Uint8Array }[] = []
|
|
private readonly maxMessageBytes: number
|
|
/**
|
|
* Ends whose sends are withheld, to model a link slower than the game.
|
|
*
|
|
* Withheld, not lost: a WebSocket and a data channel are both reliable, so the
|
|
* honest model of a slow link is a message that arrives late. Dropping it would
|
|
* instead model a transport this game never runs on, and would leave the peers
|
|
* that waited for it stuck forever rather than briefly behind.
|
|
*/
|
|
readonly held = new Set<number>()
|
|
|
|
/**
|
|
* @param options - size limits.
|
|
*/
|
|
constructor(options: TransportOptions = {}) {
|
|
this.maxMessageBytes = options.maxMessageBytes ?? DEFAULT_MAX_MESSAGE
|
|
}
|
|
|
|
/**
|
|
* Attach an end, as a peer joining the game would.
|
|
*
|
|
* @returns the new end; its `index` is its peer index.
|
|
*/
|
|
attach(): HubTransport {
|
|
const end = new HubTransport(this, this.ends.length)
|
|
this.ends.push(end)
|
|
return end
|
|
}
|
|
|
|
/** The attached ends, indexed by peer. */
|
|
get transports(): HubTransport[] {
|
|
return this.ends
|
|
}
|
|
|
|
/**
|
|
* Queue a message from one end.
|
|
*
|
|
* @param from - the sending end's index.
|
|
* @param data - the bytes.
|
|
*/
|
|
enqueue(from: number, data: Uint8Array): void {
|
|
if (data.byteLength > this.maxMessageBytes) throw new Error(`hub message of ${String(data.byteLength)} bytes exceeds ${String(this.maxMessageBytes)}`)
|
|
this.queue.push({ from, data: data.slice() })
|
|
}
|
|
|
|
/**
|
|
* Deliver everything queued.
|
|
*
|
|
* Delivery can itself queue more messages (an acknowledgement answering a
|
|
* hello), so this drains until quiet or the round cap is hit.
|
|
*
|
|
* @param rounds - maximum delivery rounds. Defaults to 8.
|
|
*/
|
|
flush(rounds = 8): void {
|
|
for (let round = 0; round < rounds; round += 1) {
|
|
if (this.queue.length === 0) return
|
|
const batch = this.queue.splice(0, this.queue.length)
|
|
const withheld: { readonly from: number; readonly data: Uint8Array }[] = []
|
|
let delivered = 0
|
|
for (const message of batch) {
|
|
// A held end's messages wait their turn rather than disappearing.
|
|
if (this.held.has(message.from)) { withheld.push(message); continue }
|
|
for (const end of this.ends) {
|
|
if (end.index === message.from) continue
|
|
end.deliver(message.data)
|
|
}
|
|
delivered += 1
|
|
}
|
|
this.queue.push(...withheld)
|
|
if (delivered === 0) return
|
|
}
|
|
throw new Error('memory hub did not settle')
|
|
}
|
|
|
|
/** Close every end. */
|
|
close(): void {
|
|
for (const end of this.ends) end.close()
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Build a hub with a given number of attached ends.
|
|
*
|
|
* @param size - how many peers.
|
|
* @param options - size limits.
|
|
* @returns the ends, indexed by peer.
|
|
*/
|
|
export function memoryHub(size: number, options: TransportOptions = {}): { hub: MemoryHub; ends: HubTransport[] } {
|
|
const hub = new MemoryHub(options)
|
|
for (let index = 0; index < size; index += 1) hub.attach()
|
|
return { hub, ends: hub.transports }
|
|
}
|
|
|
|
/** One end of a {@link MemoryHub}. */
|
|
export class HubTransport implements Transport {
|
|
/** This end's peer index. */
|
|
readonly index: number
|
|
private readonly hub: MemoryHub
|
|
private messageHandler: ((data: Uint8Array) => void) | null = null
|
|
private closeHandler: (() => void) | null = null
|
|
private isAttached = true
|
|
|
|
/**
|
|
* @param hub - the hub it belongs to.
|
|
* @param index - its peer index.
|
|
*/
|
|
constructor(hub: MemoryHub, index: number) {
|
|
this.hub = hub
|
|
this.index = index
|
|
}
|
|
|
|
/** Whether this end can still send and receive. */
|
|
get open(): boolean {
|
|
return this.isAttached
|
|
}
|
|
|
|
/** @param data - the bytes. */
|
|
send(data: Uint8Array): void {
|
|
if (!this.isAttached) return
|
|
this.hub.enqueue(this.index, data)
|
|
}
|
|
|
|
/** @param handler - message handler. */
|
|
onMessage(handler: (data: Uint8Array) => void): void {
|
|
this.messageHandler = handler
|
|
}
|
|
|
|
/** @param handler - close handler. */
|
|
onClose(handler: () => void): void {
|
|
this.closeHandler = handler
|
|
}
|
|
|
|
/** Close this end. */
|
|
close(): void {
|
|
if (!this.isAttached) return
|
|
this.isAttached = false
|
|
this.closeHandler?.()
|
|
}
|
|
|
|
/**
|
|
* Deliver one message from the hub.
|
|
*
|
|
* @param data - the bytes.
|
|
*/
|
|
deliver(data: Uint8Array): void {
|
|
if (!this.isAttached) return
|
|
this.messageHandler?.(data)
|
|
}
|
|
}
|
|
|
|
/** The subset of `WebSocket` this file uses. */
|
|
export interface SocketLike {
|
|
readonly readyState: number
|
|
binaryType?: string
|
|
send: (data: ArrayBufferLike | ArrayBufferView | string) => void
|
|
close: () => void
|
|
addEventListener: (type: string, listener: (event: { data?: unknown }) => void) => void
|
|
}
|
|
|
|
/** Socket ready state meaning "open". */
|
|
const SOCKET_OPEN = 1
|
|
|
|
/**
|
|
* Wrap a socket in a {@link Transport}.
|
|
*
|
|
* Both the browser `WebSocket` and Node's global one satisfy {@link SocketLike},
|
|
* so the network path can be tested without a browser. Only binary messages are
|
|
* accepted: a text frame on the wire means the peer is not speaking this
|
|
* protocol, and treating it as bytes anyway would feed the decoder garbage.
|
|
*
|
|
* @param socket - the socket.
|
|
* @param options - size limits.
|
|
* @returns the transport.
|
|
*/
|
|
export function socketTransport(socket: SocketLike, options: TransportOptions = {}): Transport {
|
|
const maxMessageBytes = options.maxMessageBytes ?? DEFAULT_MAX_MESSAGE
|
|
socket.binaryType = 'arraybuffer'
|
|
let messageHandler: ((data: Uint8Array) => void) | null = null
|
|
let closeHandler: (() => void) | null = null
|
|
let opened = socket.readyState === SOCKET_OPEN
|
|
|
|
socket.addEventListener('open', () => { opened = true })
|
|
socket.addEventListener('message', event => {
|
|
const data = event.data
|
|
const bytes = data instanceof ArrayBuffer
|
|
? new Uint8Array(data)
|
|
: ArrayBuffer.isView(data) ? new Uint8Array(data.buffer, data.byteOffset, data.byteLength) : null
|
|
if (bytes === null) throw new Error('socket transport received a non-binary message')
|
|
if (bytes.byteLength > maxMessageBytes) throw new Error(`socket message of ${String(bytes.byteLength)} bytes exceeds ${String(maxMessageBytes)}`)
|
|
messageHandler?.(bytes)
|
|
})
|
|
socket.addEventListener('close', () => { opened = false; closeHandler?.() })
|
|
socket.addEventListener('error', () => { opened = false; closeHandler?.() })
|
|
|
|
return {
|
|
send: (data: Uint8Array): void => {
|
|
if (!opened) return
|
|
socket.send(data)
|
|
},
|
|
onMessage: (handler): void => { messageHandler = handler },
|
|
onClose: (handler): void => { closeHandler = handler },
|
|
close: (): void => { if (opened) socket.close() },
|
|
get open(): boolean { return opened },
|
|
}
|
|
}
|