This commit is contained in:
482
packages/engine/src/replay-transport.ts
Normal file
482
packages/engine/src/replay-transport.ts
Normal file
@@ -0,0 +1,482 @@
|
||||
import type { DefinedNetworkedGame } from "./define-networked-game.js";
|
||||
import type { NetworkedServerStepResult } from "./networked-server.js";
|
||||
import type { SnapshotBatch } from "./networked-types.js";
|
||||
import type { InputDecision } from "./server.js";
|
||||
import {
|
||||
type ReplayRecording,
|
||||
type ReplayStateHash,
|
||||
type TimeTravelAuthoritativeEngine,
|
||||
type TimeTravelNetworkedGame,
|
||||
type TimeTravelServerOptions,
|
||||
} from "./time-travel.js";
|
||||
import type {
|
||||
ClientStateReport,
|
||||
InputPacket,
|
||||
PlayerId,
|
||||
StateSnapshot,
|
||||
ValidationResult,
|
||||
} from "./types.js";
|
||||
|
||||
export interface ProjectedReplayFrame<ClientState, PerceptionEvent> {
|
||||
tick: number;
|
||||
state: ClientState;
|
||||
events: PerceptionEvent[];
|
||||
}
|
||||
|
||||
export interface ReplayTicket<ClientState, PerceptionEvent> {
|
||||
ticketId: number;
|
||||
requesterId: PlayerId;
|
||||
perspectiveId: PlayerId;
|
||||
issuedAtTick: number;
|
||||
fromTick: number;
|
||||
toTick: number;
|
||||
playbackRate: number;
|
||||
frames: Array<ProjectedReplayFrame<ClientState, PerceptionEvent>>;
|
||||
}
|
||||
|
||||
export interface ReplayTicketPlan {
|
||||
requesterId: PlayerId;
|
||||
perspectiveId: PlayerId;
|
||||
fromTick: number;
|
||||
toTick: number;
|
||||
playbackRate?: number;
|
||||
}
|
||||
|
||||
export interface ReplayTriggerContext<AuthorityState> {
|
||||
tick: number;
|
||||
authorityState: Readonly<AuthorityState>;
|
||||
connectedPlayerIds: readonly PlayerId[];
|
||||
}
|
||||
|
||||
export interface ReplayAuthorizationContext<AuthorityState>
|
||||
extends ReplayTicketPlan {
|
||||
currentTick: number;
|
||||
authorityState: Readonly<AuthorityState>;
|
||||
connectedPlayerIds: readonly PlayerId[];
|
||||
}
|
||||
|
||||
export interface ReplayTransportDefinition<AuthorityState, AuthorityEvent> {
|
||||
/** Maximum age of the projected frame ring and private deterministic log. */
|
||||
historySeconds: number;
|
||||
/** Baseline frame rate. Event ticks are captured even between baseline frames. */
|
||||
captureRateHz?: number;
|
||||
/** Perspectives worth retaining. Defaults to connected network players. */
|
||||
listPerspectives?(
|
||||
authorityState: Readonly<AuthorityState>,
|
||||
connectedPlayerIds: readonly PlayerId[],
|
||||
): Iterable<PlayerId>;
|
||||
/** Trusted server hook that can create automatic tickets from authority events. */
|
||||
createTickets?(
|
||||
event: Readonly<AuthorityEvent>,
|
||||
context: ReplayTriggerContext<AuthorityState>,
|
||||
): ReplayTicketPlan | readonly ReplayTicketPlan[] | null;
|
||||
/** Every automatic or manual ticket passes through this policy. */
|
||||
authorizeReplay?(
|
||||
context: ReplayAuthorizationContext<AuthorityState>,
|
||||
): boolean;
|
||||
maxPendingTicketsPerPlayer?: number;
|
||||
}
|
||||
|
||||
interface NormalizedReplayTransportDefinition<AuthorityState, AuthorityEvent>
|
||||
extends ReplayTransportDefinition<AuthorityState, AuthorityEvent> {
|
||||
historyTicks: number;
|
||||
captureEveryTicks: number;
|
||||
maxPendingTicketsPerPlayer: number;
|
||||
}
|
||||
|
||||
export type ReplayTransportNetworkedGame<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash extends ReplayStateHash,
|
||||
> = Omit<
|
||||
TimeTravelNetworkedGame<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
>,
|
||||
"createServer"
|
||||
> & {
|
||||
createServer(
|
||||
options?: TimeTravelServerOptions<Seed>,
|
||||
): ReplayTransportAuthoritativeEngine<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
>;
|
||||
};
|
||||
|
||||
/**
|
||||
* Keeps client-safe historical projections around a private authoritative
|
||||
* engine. Raw authority checkpoints never cross this boundary.
|
||||
*/
|
||||
export class ReplayTransportAuthoritativeEngine<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash extends ReplayStateHash,
|
||||
> {
|
||||
readonly game: DefinedNetworkedGame<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent
|
||||
>;
|
||||
private readonly engine: TimeTravelAuthoritativeEngine<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
>;
|
||||
private readonly replay: NormalizedReplayTransportDefinition<
|
||||
AuthorityState,
|
||||
AuthorityEvent
|
||||
>;
|
||||
private readonly history = new Map<
|
||||
PlayerId,
|
||||
Array<ProjectedReplayFrame<ClientState, PerceptionEvent>>
|
||||
>();
|
||||
private readonly pendingTickets = new Map<
|
||||
PlayerId,
|
||||
Array<ReplayTicket<ClientState, PerceptionEvent>>
|
||||
>();
|
||||
private nextTicketId = 1;
|
||||
|
||||
constructor(
|
||||
game: DefinedNetworkedGame<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent
|
||||
>,
|
||||
engine: TimeTravelAuthoritativeEngine<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
>,
|
||||
replay: NormalizedReplayTransportDefinition<
|
||||
AuthorityState,
|
||||
AuthorityEvent
|
||||
>,
|
||||
) {
|
||||
this.game = game;
|
||||
this.engine = engine;
|
||||
this.replay = replay;
|
||||
}
|
||||
|
||||
get tick(): number {
|
||||
return this.engine.tick;
|
||||
}
|
||||
|
||||
get currentState(): Readonly<AuthorityState> {
|
||||
return this.engine.currentState;
|
||||
}
|
||||
|
||||
get playerIds(): PlayerId[] {
|
||||
return this.engine.playerIds;
|
||||
}
|
||||
|
||||
addPlayer(playerId: PlayerId): void {
|
||||
this.engine.addPlayer(playerId);
|
||||
this.capture([], true);
|
||||
}
|
||||
|
||||
removePlayer(playerId: PlayerId): void {
|
||||
this.engine.removePlayer(playerId);
|
||||
this.pendingTickets.delete(playerId);
|
||||
this.capture([], true);
|
||||
}
|
||||
|
||||
submitInput(playerId: PlayerId, packet: InputPacket<Input>): InputDecision {
|
||||
return this.engine.submitInput(playerId, packet);
|
||||
}
|
||||
|
||||
submitStateReport(
|
||||
playerId: PlayerId,
|
||||
report: ClientStateReport<ClientState>,
|
||||
): ValidationResult | null {
|
||||
return this.engine.submitStateReport(playerId, report);
|
||||
}
|
||||
|
||||
step(): NetworkedServerStepResult<AuthorityEvent> {
|
||||
const result = this.engine.step();
|
||||
const shouldCapture =
|
||||
result.tick % this.replay.captureEveryTicks === 0 ||
|
||||
result.events.length > 0;
|
||||
if (shouldCapture) this.capture(result.events, true);
|
||||
this.pruneHistory();
|
||||
|
||||
if (this.replay.createTickets) {
|
||||
const context: ReplayTriggerContext<AuthorityState> = {
|
||||
tick: result.tick,
|
||||
authorityState: this.currentState,
|
||||
connectedPlayerIds: this.playerIds,
|
||||
};
|
||||
for (const event of result.events) {
|
||||
const plans = this.replay.createTickets(event, context);
|
||||
if (!plans) continue;
|
||||
for (const plan of Array.isArray(plans) ? plans : [plans]) {
|
||||
this.issueReplay(plan);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
createSnapshot(
|
||||
playerId: PlayerId,
|
||||
serverTime: number,
|
||||
): StateSnapshot<ClientState> {
|
||||
return this.engine.createSnapshot(playerId, serverTime);
|
||||
}
|
||||
|
||||
createSnapshotBatches(serverTime: number): SnapshotBatch<ClientState>[] {
|
||||
return this.engine.createSnapshotBatches(serverTime);
|
||||
}
|
||||
|
||||
createPerceptions(
|
||||
playerId: PlayerId,
|
||||
events: readonly AuthorityEvent[],
|
||||
): PerceptionEvent[] {
|
||||
return this.engine.createPerceptions(playerId, events);
|
||||
}
|
||||
|
||||
exportRecording(): ReplayRecording<
|
||||
AuthorityState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
> {
|
||||
return this.engine.exportRecording();
|
||||
}
|
||||
|
||||
/** Trusted-server API. The configured authorization policy still applies. */
|
||||
issueReplay(
|
||||
plan: ReplayTicketPlan,
|
||||
): ReplayTicket<ClientState, PerceptionEvent> | null {
|
||||
if (!this.playerIds.includes(plan.requesterId)) return null;
|
||||
if (!validTick(plan.fromTick) || !validTick(plan.toTick)) return null;
|
||||
if (plan.toTick < plan.fromTick || plan.toTick > this.tick) return null;
|
||||
|
||||
const authorization: ReplayAuthorizationContext<AuthorityState> = {
|
||||
...plan,
|
||||
currentTick: this.tick,
|
||||
authorityState: this.currentState,
|
||||
connectedPlayerIds: this.playerIds,
|
||||
};
|
||||
const authorized = this.replay.authorizeReplay
|
||||
? this.replay.authorizeReplay(authorization)
|
||||
: plan.requesterId === plan.perspectiveId;
|
||||
if (!authorized) return null;
|
||||
|
||||
const perspectiveHistory = this.history.get(plan.perspectiveId);
|
||||
if (!perspectiveHistory) return null;
|
||||
const selected = perspectiveHistory.filter(
|
||||
({ tick }) => tick >= plan.fromTick && tick <= plan.toTick,
|
||||
);
|
||||
if (selected.length === 0) return null;
|
||||
|
||||
const first = selected[0]!;
|
||||
const last = selected[selected.length - 1]!;
|
||||
const ticket: ReplayTicket<ClientState, PerceptionEvent> = {
|
||||
ticketId: this.allocateTicketId(),
|
||||
requesterId: plan.requesterId,
|
||||
perspectiveId: plan.perspectiveId,
|
||||
issuedAtTick: this.tick,
|
||||
fromTick: first.tick,
|
||||
toTick: last.tick,
|
||||
playbackRate: clamp(plan.playbackRate ?? 1, 0.1, 4),
|
||||
frames: selected.map((frame) => this.cloneFrame(frame)),
|
||||
};
|
||||
|
||||
const queue = this.pendingTickets.get(plan.requesterId) ?? [];
|
||||
queue.push(ticket);
|
||||
while (queue.length > this.replay.maxPendingTicketsPerPlayer) queue.shift();
|
||||
this.pendingTickets.set(plan.requesterId, queue);
|
||||
return ticket;
|
||||
}
|
||||
|
||||
drainReplayTickets(
|
||||
playerId: PlayerId,
|
||||
): Array<ReplayTicket<ClientState, PerceptionEvent>> {
|
||||
const tickets = this.pendingTickets.get(playerId) ?? [];
|
||||
this.pendingTickets.delete(playerId);
|
||||
return tickets;
|
||||
}
|
||||
|
||||
private capture(events: readonly AuthorityEvent[], replaceSameTick: boolean): void {
|
||||
const connected = this.playerIds;
|
||||
const perspectives = this.replay.listPerspectives
|
||||
? this.replay.listPerspectives(this.currentState, connected)
|
||||
: connected;
|
||||
const uniquePerspectives = new Set(perspectives);
|
||||
|
||||
for (const perspectiveId of uniquePerspectives) {
|
||||
if (!validPlayerId(perspectiveId)) continue;
|
||||
const snapshot = this.engine.createSnapshot(perspectiveId, 0);
|
||||
const perceptions = this.engine.createPerceptions(perspectiveId, events);
|
||||
const frame: ProjectedReplayFrame<ClientState, PerceptionEvent> = {
|
||||
tick: this.tick,
|
||||
state: this.cloneState(snapshot.state),
|
||||
events: perceptions.map((event) => this.cloneEvent(event)),
|
||||
};
|
||||
const frames = this.history.get(perspectiveId) ?? [];
|
||||
if (replaceSameTick && frames.at(-1)?.tick === this.tick) {
|
||||
frames[frames.length - 1] = frame;
|
||||
} else {
|
||||
frames.push(frame);
|
||||
}
|
||||
this.history.set(perspectiveId, frames);
|
||||
}
|
||||
}
|
||||
|
||||
private pruneHistory(): void {
|
||||
const earliestTick = Math.max(0, this.tick - this.replay.historyTicks);
|
||||
for (const [perspectiveId, frames] of this.history) {
|
||||
const firstRetained = frames.findIndex(({ tick }) => tick >= earliestTick);
|
||||
if (firstRetained === -1) {
|
||||
this.history.delete(perspectiveId);
|
||||
continue;
|
||||
}
|
||||
if (firstRetained > 0) frames.splice(0, firstRetained);
|
||||
}
|
||||
}
|
||||
|
||||
private cloneFrame(
|
||||
frame: ProjectedReplayFrame<ClientState, PerceptionEvent>,
|
||||
): ProjectedReplayFrame<ClientState, PerceptionEvent> {
|
||||
return {
|
||||
tick: frame.tick,
|
||||
state: this.cloneState(frame.state),
|
||||
events: frame.events.map((event) => this.cloneEvent(event)),
|
||||
};
|
||||
}
|
||||
|
||||
private cloneState(state: ClientState): ClientState {
|
||||
return this.game.codecs.state.decode(this.game.codecs.state.encode(state));
|
||||
}
|
||||
|
||||
private cloneEvent(event: PerceptionEvent): PerceptionEvent {
|
||||
const codec = this.game.codecs.event;
|
||||
if (!codec) {
|
||||
throw new Error("Replay perceptions require an event codec");
|
||||
}
|
||||
return codec.decode(codec.encode(event));
|
||||
}
|
||||
|
||||
private allocateTicketId(): number {
|
||||
const ticketId = this.nextTicketId;
|
||||
this.nextTicketId = this.nextTicketId === 0xffff_ffff
|
||||
? 1
|
||||
: this.nextTicketId + 1;
|
||||
return ticketId;
|
||||
}
|
||||
}
|
||||
|
||||
export function withReplayTransport<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash extends ReplayStateHash,
|
||||
>(
|
||||
game: TimeTravelNetworkedGame<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
>,
|
||||
definition: ReplayTransportDefinition<AuthorityState, AuthorityEvent>,
|
||||
): ReplayTransportNetworkedGame<
|
||||
AuthorityState,
|
||||
ClientState,
|
||||
Input,
|
||||
AuthorityEvent,
|
||||
PerceptionEvent,
|
||||
Seed,
|
||||
StateHash
|
||||
> {
|
||||
const replay = normalizeReplayDefinition(game, definition);
|
||||
|
||||
return Object.freeze({
|
||||
...game,
|
||||
createServer(options: TimeTravelServerOptions<Seed> = {}) {
|
||||
const engine = game.createServer({
|
||||
...options,
|
||||
recordingHistoryTicks:
|
||||
options.recordingHistoryTicks ?? replay.historyTicks,
|
||||
});
|
||||
return new ReplayTransportAuthoritativeEngine(engine.game, engine, replay);
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
function normalizeReplayDefinition<AuthorityState, AuthorityEvent>(
|
||||
game: { tickRateHz: number; snapshotRateHz: number },
|
||||
definition: ReplayTransportDefinition<AuthorityState, AuthorityEvent>,
|
||||
): NormalizedReplayTransportDefinition<AuthorityState, AuthorityEvent> {
|
||||
if (!Number.isFinite(definition.historySeconds) || definition.historySeconds <= 0) {
|
||||
throw new RangeError("historySeconds must be positive");
|
||||
}
|
||||
const captureRateHz = definition.captureRateHz ?? game.snapshotRateHz;
|
||||
if (
|
||||
!Number.isInteger(captureRateHz) ||
|
||||
captureRateHz <= 0 ||
|
||||
game.tickRateHz % captureRateHz !== 0
|
||||
) {
|
||||
throw new RangeError("captureRateHz must divide tickRateHz evenly");
|
||||
}
|
||||
const maxPendingTicketsPerPlayer =
|
||||
definition.maxPendingTicketsPerPlayer ?? 2;
|
||||
if (!Number.isInteger(maxPendingTicketsPerPlayer) || maxPendingTicketsPerPlayer <= 0) {
|
||||
throw new RangeError("maxPendingTicketsPerPlayer must be a positive integer");
|
||||
}
|
||||
return {
|
||||
...definition,
|
||||
historyTicks: Math.ceil(definition.historySeconds * game.tickRateHz),
|
||||
captureEveryTicks: game.tickRateHz / captureRateHz,
|
||||
maxPendingTicketsPerPlayer,
|
||||
};
|
||||
}
|
||||
|
||||
function validPlayerId(value: number): boolean {
|
||||
return Number.isInteger(value) && value >= 0 && value <= 0xffff_ffff;
|
||||
}
|
||||
|
||||
function validTick(value: number): boolean {
|
||||
return Number.isInteger(value) && value >= 0 && value <= 0xffff_ffff;
|
||||
}
|
||||
|
||||
function clamp(value: number, minimum: number, maximum: number): number {
|
||||
return Math.max(minimum, Math.min(maximum, value));
|
||||
}
|
||||
Reference in New Issue
Block a user