import { RootStore, Transceiver } from "atom.io/internal"; import { Json, stringified } from "atom.io/foundations/json"; import { UList } from "atom.io/transceivers/u-list"; import { Subject } from "atom.io/foundations/subject"; import * as AtomIO from "atom.io"; import { JoinToken, Loadable, MutableAtomFamilyToken, MutableAtomToken, ReadableFamilyToken, RegularAtomFamilyToken, WritableToken } from "atom.io"; import { AnyMosaicTransceiver, Clock, GuardedSocket, MosaicAcceptedOperationEnvelope, MosaicAtomAddress, MosaicOperation, MosaicOperationEnvelope, MosaicSnapshot, MosaicSnapshotEnvelope, MosaicTransceiverConstructor, RoomKey, RoomSocketInterface, Socket, SocketKey, StandardSchemaV1, UserKey } from "atom.io/realtime"; import { ChildProcessWithoutNullStreams } from "node:child_process"; import { Hierarchy } from "atom.io/experiments/realms"; import { IncomingHttpHeaders } from "node:http"; import { Server } from "socket.io"; import { Canonical } from "atom.io/foundations/canonical"; import { Readable, Writable } from "node:stream"; import { ParsedUrlQuery } from "node:querystring"; //#region src/realtime-server/ipc-sockets/custom-socket.d.ts type Events = Json.Object; type EventPayload = [string, ...receiveRelay[K]]; declare function isEventPayload(value: unknown): value is [string, ...Json.Array]; interface EventBuffer extends Buffer { toString(): stringified>; } declare abstract class CustomSocket implements Socket { protected listeners: Map void>>; protected globalListeners: Set<(event: string, ...args: Json.Array) => void>; protected globalListenersOutgoing: Set<(event: string, ...args: Json.Array) => void>; protected handleEvent(...args: EventPayload): void; id: string; emit: (event: Event, ...args: I[Event]) => CustomSocket; constructor(emit: (event: Event, ...args: I[Event]) => CustomSocket); on(event: Event, listener: (...args: O[Event]) => void): this; onAny(listener: (event: string, ...args: Json.Array) => void): this; onAnyOutgoing(listener: (event: string, ...args: Json.Array) => void): this; off(event: Event, listener?: (...args: O[Event]) => void): this; offAny(listener?: (event: string, ...args: Json.Array) => void): this; } //#endregion //#region src/realtime-server/ipc-sockets/child-socket.d.ts type ChildProcess = { pid?: number | undefined; stdin: Writable; stdout: Readable; stderr: Readable; }; type StderrLog = [`e` | `i` | `w`, ...Json.Array]; declare class ChildSocket extends CustomSocket { #private; id: string; readonly ready: Promise; proc: P; key: string; logger: Pick; protected handleLog(log: StderrLog): void; constructor(proc: P, key: string, logger?: Pick); dispose(): void; } //#endregion //#region src/realtime-server/ipc-sockets/delimited-json-codec.d.ts declare const IPC_FRAME_DELIMITER = ""; type DelimitedJsonCodecOptions = { delimiter?: string; onMalformed: (frame: string, error: unknown) => void; onValue: (value: T) => void; }; /** Incrementally decodes delimiter-framed JSON without sharing stream state. */ declare class DelimitedJsonCodec { #private; constructor({ delimiter, onMalformed, onValue }: DelimitedJsonCodecOptions); write(chunk: Buffer | string): void; end(chunk?: Buffer | string): void; get pending(): string; } declare function encodeJsonFrame(value: unknown): string; //#endregion //#region src/realtime-server/ipc-sockets/parent-socket.d.ts declare const PROOF_OF_LIFE_SIGNAL = "ALIVE"; declare class SubjectSocket extends CustomSocket { in: Subject>; out: Subject>; id: string; disposalEffects: (() => void)[]; constructor(id: string); dispose(): void; } type ParentProcess = { pid?: number | undefined; stdin: Readable; stdout: Writable; stderr: Writable; exit: (code?: number) => void; }; declare class ParentSocket extends CustomSocket { protected relays: Map>; protected initRelay: (socket: SubjectSocket, userKey: UserKey) => (() => void) | void; proc: P; id: string; protected log(...args: StderrLog): void; logger: { info: (...args: Json.Array) => void; warn: (...args: Json.Array) => void; error: (...args: Json.Array) => void; }; constructor(proc: P); receiveRelay(attachServices: (socket: SubjectSocket, userKey: UserKey) => (() => void) | void): void; } //#endregion //#region src/realtime-server/mosaic/storage.d.ts type MosaicStoredReceipt = { /** The operation as it was originally accepted, including server authorship. */ readonly accepted: MosaicAcceptedOperationEnvelope; /** A canonical fingerprint of the schema-normalized, authenticated proposal. */ readonly fingerprint: string; }; /** A durable checkpoint is independent of any requesting client session. */ type MosaicStoredCheckpoint = Omit; type MosaicStorageRecovery = { /** The latest durable checkpoint, if one has been created. */ readonly checkpoint: MosaicStoredCheckpoint | null; /** Every durable operation after `checkpoint`, in revision order. */ readonly tail: readonly MosaicAcceptedOperationEnvelope[]; /** The current durable stream revision. */ readonly headRevision: number; /** All accepted ids, including compacted receipts. */ readonly receiptIds: readonly string[]; /** Fences concurrent checkpoint/compaction attempts. */ readonly retentionEpoch: number; }; type MosaicStorageAppendRequest = { readonly accepted: MosaicAcceptedOperationEnvelope; readonly expectedRevision: number; readonly fingerprint: string; }; type MosaicStorageAppendResult = { readonly accepted: MosaicAcceptedOperationEnvelope; readonly status: `accepted`; } | { readonly accepted: MosaicAcceptedOperationEnvelope; readonly status: `duplicate`; } | { readonly existing: MosaicStoredReceipt; readonly status: `collision`; } | { readonly actualRevision: number; readonly status: `stale`; }; type MosaicStorageCheckpointRequest = { readonly checkpoint: MosaicStoredCheckpoint; readonly expectedRevision: number; readonly expectedRetentionEpoch: number; }; type MosaicStorageCheckpointResult = { readonly compactedThrough: number; readonly retentionEpoch: number; readonly status: `stored`; } | { readonly actualRevision: number; readonly retentionEpoch: number; readonly status: `stale`; }; type MosaicHeadHint = { readonly atom: MosaicAtomAddress; readonly revision: number; }; type MosaicStorageResult = Promise | Value; /** * Linearizable persistence required by a Mosaic server. * * `append` must atomically enforce both the expected stream revision and the * `(atom address, operation id)` uniqueness constraint. A watch notification is * only a hint: consumers always recover the authoritative contiguous tail. */ interface MosaicStorageAdapter { append(request: MosaicStorageAppendRequest): MosaicStorageResult; checkpoint(request: MosaicStorageCheckpointRequest): MosaicStorageResult; clearSession(atom: MosaicAtomAddress, session: string): MosaicStorageResult; recover(atom: MosaicAtomAddress): MosaicStorageResult; receipt(atom: MosaicAtomAddress, operationId: string): MosaicStorageResult; setSessionWatermark(atom: MosaicAtomAddress, session: string, revision: number): MosaicStorageResult; watchHead?(atom: MosaicAtomAddress, listener: (hint: MosaicHeadHint) => void): (() => void) | Promise<() => void>; } /** A restart-safe fixture when the adapter instance itself is retained. */ declare class InMemoryMosaicStorage implements MosaicStorageAdapter { #private; append(request: MosaicStorageAppendRequest): MosaicStorageAppendResult; checkpoint(request: MosaicStorageCheckpointRequest): MosaicStorageCheckpointResult; clearSession(atom: MosaicAtomAddress, session: string): void; recover(atom: MosaicAtomAddress): MosaicStorageRecovery; receipt(atom: MosaicAtomAddress, operationId: string): MosaicStoredReceipt | null; setSessionWatermark(atom: MosaicAtomAddress, session: string, revision: number): void; watchHead(atom: MosaicAtomAddress, listener: (hint: MosaicHeadHint) => void): () => void; } //#endregion //#region src/realtime-server/mosaic/server.d.ts type MaybePromise = Promise | Value; type MosaicAuthorizationAction = `presence` | `propose` | `read`; type MosaicAuthorizationContext = { readonly action: MosaicAuthorizationAction; readonly actor: string; readonly operation?: Operation; readonly atom: MosaicAtomAddress; readonly session: string; }; type MosaicPresenceContext = { readonly actor: string; readonly atom: MosaicAtomAddress; readonly session: string; readonly view: View; }; /** A mutable transceiver class whose static metadata identifies its wire model. */ type MosaicTransceiverClass = MosaicTransceiverConstructor; type MosaicServerTarget = MutableAtomToken | MutableAtomFamilyToken; type MosaicAtomRegistrationBase = { /** Checkpoint automatically after this many accepted operations. */ readonly checkpointEvery?: number; /** The transceiver class registered for this ordinary mutable atom target. */ readonly class: MosaicTransceiverClass; /** Validate and normalize untrusted model payloads before authorization. */ readonly operationSchema: StandardSchemaV1>; /** Presence is disabled unless a schema is supplied. */ readonly presenceSchema?: StandardSchemaV1; /** Perform model-aware checks such as validating relative anchors. */ readonly validatePresence?: (presence: Presence, context: MosaicPresenceContext) => MaybePromise; }; type MosaicAtomRegistration = MosaicAtomRegistrationBase & ({ /** Validate every dynamic key before opening a family-member stream. */ readonly keySchema: StandardSchemaV1; readonly target: MutableAtomFamilyToken; } | { readonly keySchema?: never; /** Register one standalone atom or concrete family-member token. */ readonly target: MutableAtomToken; }); type MosaicServerConnection = { readonly actor: string; readonly session: string; readonly socket: Socket; }; /** Internal erasure that preserves heterogeneous registration inference. */ type ErasedRegistration = { readonly checkpointEvery?: number; readonly class: MosaicTransceiverClass; readonly keySchema?: StandardSchemaV1; readonly operationSchema: StandardSchemaV1; readonly presenceSchema?: StandardSchemaV1; readonly target: MosaicServerTarget; readonly validatePresence?: (presence: any, context: MosaicPresenceContext) => MaybePromise; }; type MosaicServerOptions = { readonly authorize?: (context: MosaicAuthorizationContext) => MaybePromise; readonly registrations: readonly ErasedRegistration[]; readonly storage?: MosaicStorageAdapter; }; type MosaicServerAtomStatus = { readonly initialized: boolean; readonly revision: number; }; type MosaicServer = { checkpoint(atom: MosaicAtomAddress): Promise; connect(connection: MosaicServerConnection): () => Promise; dispose(): Promise; atomStatus(atom: MosaicAtomAddress): Promise; }; /** Fingerprint the schema-normalized proposal, including actor and session. */ declare const fingerprintMosaicOperation: (operation: MosaicOperationEnvelope) => string; /** * Create a durable, server-authoritative Mosaic operation service. * * One instance may cache projections, but the storage adapter owns ordering. * Head watches are only wake-up hints; every wake-up drains a checked, * contiguous recovery result from the linearizable store. */ declare function createMosaicServer(options: MosaicServerOptions): MosaicServer; /** Preserve transceiver, schema and presence inference for one atom target. */ declare function defineMosaicAtomRegistration(registration: MosaicAtomRegistration): MosaicAtomRegistration; /** Extract the transceiver's JSON-safe durable checkpoint type. */ type MosaicServerSnapshot = MosaicSnapshot; //#endregion //#region src/realtime-server/provide-rooms.d.ts type RoomMap = Map>; declare global { var ATOM_IO_REALTIME_SERVER_ROOMS: RoomMap; } declare const ROOMS: RoomMap; declare const roomMeta: { count: number; }; type RoomTimeouts = { startupMs?: number | undefined; idleMs?: number | undefined; maximumMs?: number | undefined; shutdownMs?: number | undefined; }; declare const DEFAULT_ROOM_TIMEOUTS: { readonly startupMs: 5_000; readonly shutdownMs: 1_000; }; type RoomProcessFactory = (command: string, args: readonly string[], roomKey: RoomKey) => ChildProcessWithoutNullStreams; type SpawnRoomConfig = { clock?: Clock; processFactory?: RoomProcessFactory; store: RootStore; socket: Socket; userKey: UserKey; resolveRoomScript: (roomName: RoomNames) => [string, string[]]; timeouts?: RoomTimeouts; }; declare function spawnRoom({ clock, processFactory, store, socket, userKey, resolveRoomScript, timeouts }: SpawnRoomConfig): (roomName: RoomNames) => Promise>; type ProvideEnterAndExitConfig = { store: RootStore; socket: Socket; roomSocket: GuardedSocket>; userKey: UserKey; }; declare function provideEnterAndExit({ store, socket, roomSocket, userKey }: ProvideEnterAndExitConfig): ((roomKey: RoomKey) => void) & { dispose: () => void; }; type DestroyRoomConfig = { store: RootStore; socket: Socket; userKey: UserKey; }; declare function destroyRoom({ store, socket, userKey }: DestroyRoomConfig): (roomKey: RoomKey) => void; type ProvideRoomsConfig = { resolveRoomScript: (path: RoomNames) => [string, string[]]; roomAdminsToken: ReadableFamilyToken; roomNames: RoomNames[]; roomTimeLimit?: number; roomIdleTimeLimit?: number; roomShutdownTimeLimit?: number; roomStartupTimeLimit?: number; userKey: UserKey; store: RootStore; socket: Socket; }; declare function provideRooms({ resolveRoomScript, roomAdminsToken, roomNames, socket, store, userKey, roomTimeLimit, roomIdleTimeLimit, roomShutdownTimeLimit, roomStartupTimeLimit }: ProvideRoomsConfig): () => void; //#endregion //#region src/realtime-server/realtime-family-provider.d.ts type FamilyProvider = ReturnType; declare function realtimeAtomFamilyProvider({ socket, consumer, store }: ServerConfig): (family: AtomIO.RegularAtomFamilyToken, index: AtomIO.ReadableToken> | null> | Iterable>) => () => void; //#endregion //#region src/realtime-server/realtime-mutable-family-provider.d.ts type MutableFamilyProvider = ReturnType; declare function realtimeMutableFamilyProvider({ socket, consumer, store }: ServerConfig): , K extends Canonical>(family: AtomIO.MutableAtomFamilyToken, index: AtomIO.ReadableToken> | null> | Iterable>) => () => void; //#endregion //#region src/realtime-server/realtime-mutable-provider.d.ts type MutableProvider = ReturnType; declare function realtimeMutableProvider({ socket, consumer, store }: ServerConfig): >(token: AtomIO.MutableAtomToken) => () => void; //#endregion //#region src/realtime-server/realtime-state-provider.d.ts type StateProvider = (clientToken: AtomIO.WritableToken, serverData?: AtomIO.ReadableToken | S) => () => void; declare function realtimeStateProvider({ socket, consumer, store }: ServerConfig): StateProvider; //#endregion //#region src/realtime-server/realtime-state-receiver.d.ts type RealtimeLeaseClock = { clearTimeout: (timer: unknown) => void; now: () => number; setTimeout: (callback: () => void, milliseconds: number) => unknown; }; type RealtimeStateReceiverOptions = { /** How long ownership remains valid without a renewal. Defaults to 10s. */ leaseDurationMs?: number; /** Injectable clock used by deterministic tests and alternate runtimes. */ leaseClock?: RealtimeLeaseClock; /** Suggested client renewal interval. Defaults to one third of the lease. */ renewAfterMs?: number; }; type StateReceiver = (schema: StandardSchemaV1, clientToken: WritableToken, serverToken?: WritableToken, options?: RealtimeStateReceiverOptions) => () => void; declare function realtimeStateReceiver({ socket, consumer, store }: ServerConfig): StateReceiver; //#endregion //#region src/realtime-server/server-config.d.ts type ServerConfig = { socket: Socket; consumer: RoomKey | UserKey; store?: RootStore; }; type UserServerConfig = { socket: Socket; consumer: UserKey; store?: RootStore; }; /** Socket Handshake details--taken from socket.io */ type Handshake = { /** The headers sent as part of the handshake */ headers: IncomingHttpHeaders; /** The date of creation (as string) */ time: string; /** The ip of the client */ address: string; /** Whether the connection is cross-domain */ xdomain: boolean; /** Whether the connection is secure */ secure: boolean; /** The date of creation (as unix timestamp) */ issued: number; /** The request URL string */ url: string; /** The query object */ query: ParsedUrlQuery; /** The auth object */ auth: { [key: string]: any; }; }; declare function realtime(server: Server, auth: (handshake: Handshake) => Loadable, onConnect: (config: UserServerConfig) => Loadable<() => Loadable>, store?: RootStore): () => Promise; //#endregion //#region src/realtime-server/server-socket-state.d.ts type SocketSystemHierarchy = Hierarchy<[{ above: `root`; below: [UserKey, SocketKey, RoomKey]; }]>; declare const socketAtoms: RegularAtomFamilyToken; declare const socketKeysAtom: MutableAtomToken>; declare const onlineUsersAtom: MutableAtomToken>; declare const usersOfSockets: JoinToken<`user`, UserKey, `socket`, SocketKey, `1:n`>; //#endregion export { ChildProcess, ChildSocket, CustomSocket, DEFAULT_ROOM_TIMEOUTS, DelimitedJsonCodec, DelimitedJsonCodecOptions, DestroyRoomConfig, EventBuffer, EventPayload, Events, FamilyProvider, Handshake, IPC_FRAME_DELIMITER, InMemoryMosaicStorage, MosaicAtomRegistration, MosaicAuthorizationAction, MosaicAuthorizationContext, MosaicHeadHint, MosaicPresenceContext, MosaicServer, MosaicServerAtomStatus, MosaicServerConnection, MosaicServerOptions, MosaicServerSnapshot, MosaicServerTarget, MosaicStorageAdapter, MosaicStorageAppendRequest, MosaicStorageAppendResult, MosaicStorageCheckpointRequest, MosaicStorageCheckpointResult, MosaicStorageRecovery, MosaicStorageResult, MosaicStoredCheckpoint, MosaicStoredReceipt, MosaicTransceiverClass, MutableFamilyProvider, MutableProvider, PROOF_OF_LIFE_SIGNAL, ParentProcess, ParentSocket, ProvideEnterAndExitConfig, ProvideRoomsConfig, ROOMS, RealtimeLeaseClock, RealtimeStateReceiverOptions, RoomMap, RoomProcessFactory, RoomTimeouts, ServerConfig, SocketSystemHierarchy, SpawnRoomConfig, StateProvider, StateReceiver, StderrLog, SubjectSocket, UserServerConfig, createMosaicServer, defineMosaicAtomRegistration, destroyRoom, encodeJsonFrame, fingerprintMosaicOperation, isEventPayload, onlineUsersAtom, provideEnterAndExit, provideRooms, realtime, realtimeAtomFamilyProvider, realtimeMutableFamilyProvider, realtimeMutableProvider, realtimeStateProvider, realtimeStateReceiver, roomMeta, socketAtoms, socketKeysAtom, spawnRoom, usersOfSockets }; //# sourceMappingURL=index.d.ts.map