@repo/datadata/in-memory
The whole engine in one process: an in-memory client and storage for demos, tests and prototypes that need no server.
Classes
LatencyLink
export declare class LatencyLinkDelivers a harness's queued events once they have waited latencyMs, and models unplugging the network cable. Call LatencyLink.pump on every tick (e.g. from requestAnimationFrame) with the current time.
constructor(harness: LinkHarness, latencyMs?: number);Wrap harness with a link of latencyMs milliseconds (default 0), initially plugged.
latencyMs: number;How long, in milliseconds, an event waits before delivery. Can be changed at any time.
plug(): void;Reconnect after LatencyLink.unplug. A no-op when already plugged.
get plugged(): boolean;Whether the link is connected; starts true.
pump(now: number): void;Deliver every event, in each direction, that has waited out the latency. now is a timestamp in milliseconds on any monotonic clock, used consistently across calls. Does nothing while unplugged.
unplug(): void;Drop the connection, discarding events in flight. A no-op when already unplugged.
Functions
connectInMemoryClient
export declare function connectInMemoryClient<S extends SchemaRegistry, C extends CommandRegistry = CommandRegistry>(connection: InMemoryServerConnection<ValidatorsFromSchemas<S> & SystemValidatorRegistry>, schemas: S, options?: ConnectInMemoryClientOptions<C>): {
client: import("../client/types.js").DatadataClientFor<S, C>;
bridge: ClientSessionBridge;
connectionId: string;
flushClientEvents: () => ClientSentEvent[];
flushServerEvents: () => ServerSentEvent[];
flushOneClientEvent: () => ClientSentEvent | null;
flushOneServerEvent: () => ServerSentEvent | null;
deliverServerEvent: (event: ServerSentEvent) => void;
events: (ClientSentEvent | ServerSentEvent)[];
pendingClientEvents: ClientSentEvent[];
pendingServerEvents: ServerSentEvent[];
dropConnection: () => void;
reestablishConnection: () => void;
reportConnected: () => void;
setConnectionPrincipal: (next: Principal) => void;
};Connect a new client to an in-memory server and return a harness around it: the client, its event queues (pendingClientEvents, pendingServerEvents), flush methods that deliver queued events one at a time or all at once, and dropConnection / reestablishConnection / setConnectionPrincipal to model the socket and the host. Events queue until flushed unless ConnectInMemoryClientOptions autoProcess is set.
createFlushAll
export declare function createFlushAll(clients: Array<{
flushClientEvents: () => ClientSentEvent[];
flushServerEvents: () => ServerSentEvent[];
}>): () => void;Build a function that delivers every pending event of clients in both directions, repeating until no queue has anything left, so all clients and the server reach a settled state. Takes the harnesses returned by connectInMemoryClient (any object with the two flush methods).
createFuzzedFlushAll
export declare function createFuzzedFlushAll(clients: Array<{
pendingClientEvents: ClientSentEvent[];
pendingServerEvents: ServerSentEvent[];
flushOneClientEvent: () => ClientSentEvent | null;
flushOneServerEvent: () => ServerSentEvent | null;
}>, schedule: number[] | {
seed: number;
}, options?: FuzzedFlushOptions): () => void;The fuzzed-order counterpart to createFlushAll: the returned function drains every queue to the same settled state, but at each step schedule picks which ready (client, direction) delivers its next single event. Another client's event can therefore land at the server between two events of one client's batch, while each connection's own order is preserved, as on a WebSocket.
schedule is either an array of non-negative integers (cycled; pass fast-check-generated values so a failing order shrinks to a minimal schedule) or { seed } for a reproducible pseudo-random order. The pick cursor persists across calls, so successive flushes keep consuming the schedule rather than replaying its prefix. Throws if delivery has not settled after 100,000 steps.
createInMemoryClient
export declare function createInMemoryClient<S extends SchemaRegistry, C extends CommandRegistry = CommandRegistry>(options: InMemoryDatadataClientOptions<S, C>): InMemoryDatadataClient<RegistryFor<S>, C>;Creates a datadata client backed by an embedded in-memory server.
This provides a drop-in replacement for the real datadata client that runs entirely in memory. Perfect for unit tests, demos, and development.
Example
import { createInMemoryClient } from "@repo/datadata/in-memory";
const { client, flush } = createInMemoryClient({
schemas: { myDoc: MyDocSchema },
autoFlush: true
});
// Use just like the real client: document I/O goes through the session
// surface (client.live is the ambient session).
client.live.prepareDocument("doc1");
client.live.createDocument({
docId: "doc1",
type: "myDoc",
data: { name: "My Document" }
});
// With autoFlush: true, changes are immediately available
const doc = client.live.getDocument({ docId: "doc1", type: "myDoc" });
console.log(doc?.data.name); // "My Document"
createInMemoryObjectStorageAdapter
export declare function createInMemoryObjectStorageAdapter<Context = unknown>(): DatadataObjectStorageAdapter<Context>;An object storage adapter that keeps blob bytes in a Map, for tests and demos that use blobs without an object store. Objects are immutable: writing an existing key throws. Bytes live only as long as the adapter.
createInMemoryServer
export declare function createInMemoryServer<S extends SchemaRegistry, EventBusContext extends InMemoryEventBusContext = InMemoryEventBusContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>(args: {
schemas: S;
initialStorage?: InMemoryStorage;
logger: Logger;
skipValidation?: boolean;
limits?: Partial<DatadataServerLimits>;
objectStorageAdapter?: DatadataObjectStorageAdapter<StorageContext>;
commands?: Commands;
}): InMemoryServerConnection<ValidatorsFromSchemas<S> & SystemValidatorRegistry, EventBusContext, StorageContext, Commands>;Create a real datadata server over in-memory storage, wired to an in-process event bus. schemas are seeded as sys:schema:<docType> documents (unless initialStorage already holds them), and initialStorage is checked for consistency unless skipValidation is set. Throws when a schema declares blobRef fields and no objectStorageAdapter is given.
This overload infers the registry from schemas.
export declare function createInMemoryServer<Registry extends ValidatorRegistry = ValidatorRegistry, EventBusContext extends InMemoryEventBusContext = InMemoryEventBusContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>(args: {
schemas?: Record<string, DocumentSchema>;
initialStorage?: InMemoryStorage;
logger: Logger;
skipValidation?: boolean;
limits?: Partial<DatadataServerLimits>;
objectStorageAdapter?: DatadataObjectStorageAdapter<StorageContext>;
commands?: Commands;
}): InMemoryServerConnection<Registry, EventBusContext, StorageContext, Commands>;Create a real datadata server over in-memory storage; see the schema-inferring overload for the behavior.
This overload takes the registry as an explicit type argument, or none for storage-only setups with no schema map to infer from.
createInMemoryStorage
export declare function createInMemoryStorage<S extends SchemaRegistry>(schemas: S): InMemoryStorageBuilder<ValidatorsFromSchemas<S>>;Start an InMemoryStorage fixture whose documents are validated against schemas (and the built-in sys:* validators). Add documents with InMemoryStorageBuilder.doc and finish with InMemoryStorageBuilder.build. Schema documents for schemas are not seeded here; the in-memory server seeds any that are missing.
Interfaces
FuzzedFlushOptions
export interface FuzzedFlushOptionsOptions for createFuzzedFlushAll.
actions?: Array<() => void>;Work to AUTHOR mid-flight: each thunk competes with the ready deliveries for the schedule's next pick, so an op can be authored while earlier ops' acks are still in transit — the windows a settle-between-ops drive can never produce. Thunks run at most once, in array order (only the queue head is eligible), and the flush does not reach its fixpoint until every thunk has run and the queues have drained.
onStep?: (step: FuzzedFlushStep) => void;Invoked after EVERY single delivered event. This is the observation point for transient-state invariants: a self-healing inconsistency (e.g. a listing hole between an ack and its index patch) is visible only between two frames, never at the drained fixpoint — assert it here, not after the flush returns. A throw propagates out of the flush (failing the enclosing property, which then shrinks the schedule).
FuzzedFlushStep
export interface FuzzedFlushStepOne delivered event during a fuzzed flush — the per-frame observation point.
clientIndex: number;Index into the clients array the event was delivered for.
direction: "client" | "server";"client": a client→server event landed at the server; "server": a server→client frame was processed by the client.
event: ClientSentEvent | ServerSentEvent;The event that was delivered.
InMemoryBlobRow
export interface InMemoryBlobRowOne row of the in-memory blob catalog (InMemoryStorage.blobs). Times are milliseconds since the epoch.
actor: string;The actor of the principal that created the blob.
contentType: string;The MIME type declared when the blob was created.
createdAt: number;When the blob was created.
finalizedAt: number | null;When the upload was finalized, or null while it is still staged. Blobs finalize once.
sha256: string | null;Hex SHA-256 of the bytes, or null until the upload is finalized.
size: number | null;The byte length, or null until the upload is finalized.
state: "staged" | "live" | "tombstone";"staged" from creation until a document first references it, then "live". "tombstone" is the garbage-collection sweep's terminal state: the blob is invisible to reads while its bytes and record are deleted.
subject: string | null;The subject of the principal that created the blob (null = anonymous).
unreferencedAt: number | null;When a live blob lost its last referencing document, starting the garbage-collection grace period; null while it is referenced.
InMemoryDatadataClient
export interface InMemoryDatadataClient<Registry extends ValidatorRegistry, Commands extends CommandRegistry = CommandRegistry>What createInMemoryClient returns: a client wired to an embedded in-memory server, plus controls over the simulated network between them.
bridge: ClientSessionBridge;The WHITEBOX door: the session bridge the client hands its sessions (with this harness's auto-flush already composed in). For tests whose subject is a bridge internal — see the policy in test/whitebox-client.ts. App-modelling tests never need it.
client: DatadataClientInterface<Registry, Commands>;The client, exposed through its public interface: document reads and writes go through the session surface (client.live, createLiveSession, openStagingSession), exactly as against a real server. Tests that need to drive the client's internal surface directly use InMemoryDatadataClient.bridge (or the mock system's connectInMemoryClient, which exposes the same door).
disconnect: () => void;Simulate network disconnection - frames in flight and subscriptions are lost, and what the client sends is queued until reconnect, which sends it behind the resubscribes (like the browser host's socket)
flush: () => void;Flush all pending events between client and server to sync changes immediately
getEvents: () => (ClientSentEvent | ServerSentEvent)[];Get all events that have been processed (for testing/debugging)
getPendingClientEvents: () => ClientSentEvent[];Get pending client events that haven't been flushed yet
getPendingServerEvents: () => ServerSentEvent[];Get pending server events that haven't been flushed yet
hasPendingEvents: () => boolean;Check if there are any pending events
reconnect: () => void;Simulate network reconnection - client will resubscribe to documents
server: DatadataServer<Registry, InMemoryEventBusContext>;Access the embedded server directly (useful for exercising server-side operations)
InMemoryDatadataClientOptions
export interface InMemoryDatadataClientOptions<S extends SchemaRegistry, C extends CommandRegistry = CommandRegistry>Options for createInMemoryClient.
autoFlush?: boolean;Auto-flush events immediately after each operation (default: true)
commands?: C;The app's domain commands (defineCommands), given to the embedded server and the client, whose session command calls they type.
config?: ClientConfig;Client config forwarded to createClient — e.g. { dynamicSchemas: true }.
initialDocuments?: Array<InitialDocument<RegistryFor<S>>>;Pre-populate the server with initial documents. type/data admit the sys:* docTypes alongside the app schemas', so a caller can seed its own sys:schema:<type> document — mirroring the storage builder's .doc(), including its create-input data shape: defaulted fields may be omitted and are backfilled, as createDocument would.
limits?: Partial<DatadataServerLimits>;Override default size limits for the in-memory storage
logger?: Logger;Logger used by the in-memory server (default: a browser console logger). Must be vitest-free when the factory runs outside tests — apps use it at runtime (e.g. the playground's in-memory pages on Cloudflare Workers), where vitest's fake-timers break module init.
principal?: Principal;The identity the client connection's writes run as (absent = anonymous user).
schemas: S;App docTypes' schemas, keyed by docType. The client's runtime validators are built from these, and they are seeded as sys:schema:<docType> documents so the server can resolve them (it has no static app-schema fallback).
InMemoryEventAttribution
export interface InMemoryEventAttributionWho wrote a stored event and when, fixed at write time. Every field is optional in stored fixtures: an absent field reads as an unattributed event (subject and actor null) written at the epoch (createdAt 0).
actor?: string | null;The writing principal's actor; null = unattributed.
createdAt?: number;When the event was written, in milliseconds since the epoch.
subject?: string | null;The writing principal's subject; null = anonymous or unattributed.
InMemoryEventBusContext
export interface InMemoryEventBusContext extends LoggerContextThe per-call context the in-memory server's event bus hands to server handlers.
connectionId: ConnectionIdOrServerInitiated;The connection the event arrived on, or the server-initiated marker for server-side writes.
principal?: Principal;The identity this connection's writes run as (absent = anonymous user).
InMemoryServerConnection
export interface InMemoryServerConnection<Registry extends ValidatorRegistry, EventBusContext extends InMemoryEventBusContext = InMemoryEventBusContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>An in-memory server and the hooks to attach clients to it, returned by createInMemoryServer. Pass it to connectInMemoryClient once per client.
bus: InProcessEventBus<EventBusContext>;The in-process bus the server is wired to — attach real clients via connectInProcessClient.
clearSubscriptions: (connectionId: string) => void;Drop every document subscription of a connection, as a closed socket does.
getNextEventId: () => string;Mint a client event id unique within this server (event1, event2, …).
logger: Logger;The logger the server was created with.
registerClient: (eventHandler: (event: ServerSentEvent) => void, principal: Principal) => string;Register a connection on the bus that delivers server events to eventHandler and runs as principal; returns its connectionId.
server: DatadataServer<Registry, EventBusContext, StorageContext, Commands>;The server, over in-memory storage.
InMemoryStorage
export interface InMemoryStorageThe plain-object state behind the in-memory storage adapter. Build one with createInMemoryStorage, or write a fixture by hand: only documents is required. The server shares it by reference, so a server re-created over the same storage sees everything the previous one wrote.
blobEdges?: Map<string, Set<string>>;The blobIds each document references, keyed by the referencing docId. Optional; created on first use.
blobs?: Map<string, InMemoryBlobRow>;The blob catalog, keyed by blobId. Held on the storage rather than the adapter so an adapter rebuilt around the same storage keeps documents and their blobs together. Optional; created on first use.
documents: {
[docId: string]: {
type: string;
data: Record<string, unknown>;
generation?: string;
conformedSchemaSequence?: number;
invalidSchemaSequence?: number | null;
yjsDocs: {
[yjsId: string]: Uint8Array;
};
events: ({
patch: Operation[];
schemaSequence: number;
} & InMemoryEventAttribution)[];
yjsEvents?: (YjsUpdateRaw & InMemoryEventAttribution)[];
referenceEdges?: ReferenceEdge[];
deletedAt?: number | null;
name?: string | null;
};
};Every stored document, keyed by docId, with its JSON and Yjs event logs.
purgedGenerations?: Map<string, Set<string>>;The generations a purge released each docId from, keyed by docId. Optional; created on first purge.
InMemoryStorageBuilder
export interface InMemoryStorageBuilder<Registry extends ValidatorRegistry>Fluent builder for an InMemoryStorage fixture, returned by createInMemoryStorage.
build(): InMemoryStorage;Return the storage built so far, ready to pass as a server's initialStorage.
doc<Type extends keyof (Registry & SystemValidatorRegistry) & string>(params: {
docId: string;
type: Type;
name?: string | null;
} & ({
data: TypedDocumentInput<Registry & SystemValidatorRegistry, Type>;
skipValidation?: false;
} | {
data: unknown;
skipValidation: true;
})): InMemoryStorageBuilder<Registry>;Add a document as if it had been created earlier, and return the builder. data is validated against the docType's schema, with declared defaults backfilled as a create would, unless skipValidation: true stores it as given (to model a document that no longer conforms). type admits the sys:* docTypes as well as the app's, so a caller can install its own sys:schema document. An existing document with the same docId is replaced.
LinkHarness
export interface LinkHarnessThe slice of an in-memory client harness the link drives.
dropConnection: () => void;Close the connection: in-flight events are lost and the client goes offline.
flushOneClientEvent: () => unknown;Deliver the oldest pending client→server event to the server.
flushOneServerEvent: () => unknown;Deliver the oldest pending server→client event to the client.
pendingClientEvents: readonly object[];Client→server events not yet delivered, oldest first.
pendingServerEvents: readonly object[];Server→client events not yet delivered, oldest first.
reestablishConnection: () => void;Open the connection again; the client resubscribes and replays its offline writes.
Types
ConnectInMemoryClientOptions
export type ConnectInMemoryClientOptions<C extends CommandRegistry = CommandRegistry> = Omit<ClientConfig, "principal"> & {
autoProcess?: boolean;
commands?: C;
principal?: Principal | null;
dropServerEvent?: (event: ServerSentEvent) => boolean;
startDisconnected?: boolean;
};Options for connectInMemoryClient: the client config plus harness options. principal here is the connection's identity on the server, not the client config's principal used for client-side prediction.
InitialDocument
export type InitialDocument<Registry extends ValidatorRegistry> = {
[Type in keyof Registry & string]: {
docId: string;
type: Type;
data: TypedDocumentInput<Registry, Type>;
};
}[keyof Registry & string];One document to seed an in-memory server with. data is typed by the seed's own type (a union discriminated on type), so a seed carrying another docType's data fails to compile. data takes the create-input shape: fields with a declared default may be omitted.