@repo/datadata/server
Running a datadata server: the engine that validates, authorizes, stores and broadcasts document changes, and the storage adapter interface it runs on.
Classes
AsyncDatadataServer
export declare class AsyncDatadataServer<Registry extends ValidatorRegistry, EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>The Promise-returning twin of DatadataServer, for storage adapters whose methods return promises (AsyncDatadataStorageAdapter — e.g. Postgres). A synchronous adapter satisfies the async interface as-is, so this server also runs over the in-memory and DO-SQLite adapters.
SAME ENGINE, different shell: both servers drive the identical sans-IO generator modules — this facade awaits each storage hop (runAsync) where the sync facade runs on the caller's stack (runSync), and adds ONE piece of machinery the sync server gets for free: the operation queue below. Operation semantics (validation, authorization, broadcast shaping, error categories) are shared by construction; only the error CHANNEL differs — where a sync method throws, the returned promise rejects with the same error.
Method-level contracts (parameters, results, edge cases) are documented once, on DatadataServer — each method here is that operation behind a promise.
[Symbol.asyncDispose](): Promise<void>;Symbol.asyncDispose alias for AsyncDatadataServer.close, so a host can scope the server with await using server = createAsyncServer(...). No timeout bound: close()'s drain covers only already-queued operations, which is finite by construction.
constructor(storageAdapter: AsyncDatadataStorageAdapter<Registry, StorageContext>, schemas: Registry, eventBus: DatadataEventBus, logger: Logger, objectStorageAdapter?: DatadataObjectStorageAdapter<StorageContext>, commands?: Commands);Constructs a new instance of the AsyncDatadataServer class
authorizeBlobRead(context: StorageContext, ref: {
blobId: string;
}): Promise<BlobReadGrant | null>;close(): Promise<void>;Terminal drain: reject every operation submitted from now on, wait for the queued ones to complete, then release the storage adapter (its optional AsyncDatadataStorageAdapter.close). After close() resolves the server has no work in flight — for a storage layer shared across servers (one Postgres database, many folders) the host closes the database only after every folder's server has closed; closing PGlite with a statement still in flight spins its WASM in a microtask loop that starves all timers. Idempotent: every call returns the same promise.
commandDocument(context: EventBusContext & StorageContext & LoggerContext, params: CommandInvocation<Commands>, options?: ObservabilityOptions): Promise<void>;See DatadataServer.commandDocument. The mutator itself stays synchronous.
createBlobUpload(context: StorageContext & LoggerContext, params: BlobStageParams): Promise<{
blobId: string;
}>;createDocument<Type extends keyof Registry>(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
type: Type;
data: TypedDocumentInput<Registry, Type>;
yjsDocs?: Record<string, Uint8Array>;
}, options?: ObservabilityOptions): Promise<void>;createDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
type: string;
data: unknown;
yjsDocs?: Record<string, Uint8Array>;
name?: string | null;
}, options?: ObservabilityOptions): Promise<void>;See DatadataServer.createDocument; the untyped form, for a stored-schema type.
deleteDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
baseSequence?: number;
guard?: DeleteGuardMode;
}, options?: ObservabilityOptions): Promise<void>;drain(): Promise<void>;Settle once every operation enqueued so far has completed — including operations that those operations themselves enqueue (an onDocumentChange callback calling a public method). Purely observational: the server stays open and keeps accepting operations, so a caller racing new work against drain() gets no guarantee about THAT work — use close() for a terminal drain. Never rejects (a failed operation rejects to its own caller only).
exportDocument(context: StorageContext & SystemContext, ref: {
docId: string;
}): Promise<DocumentExport | null>;finalizeBlobUpload(context: StorageContext, params: {
blobId: string;
size: number;
sha256: string;
}): Promise<void>;findInvalidDocuments(context: StorageContext & SystemContext, opts?: {
type?: string;
}): Promise<InvalidDocument[]>;getBlobMetadata(context: StorageContext, ref: {
blobId: string;
}): Promise<BlobMetadata | null>;getDocument<Type extends keyof Registry & string>(context: StorageContext, params: {
docId: string;
type: Type;
}): Promise<{
document: DatadataDocument<TypedDocumentData<Registry, Type>> | null;
invalid?: DocumentInvalid;
}>;getDocument(context: StorageContext, params: {
docId: string;
type?: string;
}): Promise<{
document: DatadataDocument | null;
invalid?: DocumentInvalid;
}>;See DatadataServer.getDocument; the untyped form.
getDocumentInitEvent(context: StorageContext, ref: {
docId: string;
}): Promise<DocumentInitEvent | DocumentNotFoundEvent | DocumentDeletedEvent | DocumentErrorEvent>;getYjsDoc(context: StorageContext, ref: {
docId: string;
yjsId: string;
}): Promise<Y.Doc | null>;handleDisconnect(context: EventBusContext & LoggerContext & ConnectionIdentityContext): Promise<void>;See DatadataServer.handleDisconnect, with one shell divergence: the presence clear is QUEUED here, so a host that must uphold cleared-before-detach ordering awaits the returned promise before detaching the connection from the event bus. A host that detaches immediately (the in-process transport) accepts a bounded window where the departed connection's cells are still visible — the clear's broadcasts route to the remaining subscribers either way.
handleEvent(context: EventBusContext & StorageContext & LoggerContext & ConnectionIdentityContext & {
authority?: never;
connectionId?: string;
}, event: ClientSentEvent, options?: ObservabilityOptions): Promise<void>;Dispatch one client-sent event — see DatadataServer.handleEvent. The returned promise settles when the operation has fully completed (including its broadcasts); a transport that must preserve per-connection ordering beyond this server's own queue (e.g. a Durable Object) should await it. Write failures surface as doc:error to the calling client rather than rejecting; protocol violations (malformed event, oversized event id) still reject.
importDocument(context: EventBusContext & StorageContext & SystemContext, documentExport: DocumentExport, options?: {
validation?: "strict" | "relaxed";
}): Promise<{
indexSequence: number;
}>;onBlobSweepCandidate(callback: BlobSweepCandidateCallback): () => void;Register a hook on "the blob catalog may hold a new GC candidate" — a blob staged (by AsyncDatadataServer.createBlobUpload or an import), a document's blob edges replaced (a validated write, or a lazy read or AsyncDatadataServer.validateDocuments refreshing them after a schema change), or a purge. Hosts arm their sweep schedule from it; it over-approximates, so arming must be idempotent. Invoked synchronously after the triggering storage write, never awaited. Returns the unregister function. See BlobSweepCandidateCallback.
onDocumentChange(callback: DocumentChangeCallback<EventBusContext, StorageContext>): () => void;See DatadataServer.onDocumentChange. The callback is invoked synchronously inside the queued operation and never awaited: it cannot stall the operation, and an async callback's rejection is logged rather than left unhandled. A callback may call public server methods — they enqueue behind the current operation — but their results only exist after this operation completes.
onDocumentRestored(callback: DocumentRestoredCallback<EventBusContext, StorageContext>): () => void;See DatadataServer.onDocumentRestored. Invoked like AsyncDatadataServer.onDocumentChange: synchronously inside the queued restore, never awaited.
onDocumentsPurged(callback: DocumentsPurgedCallback<EventBusContext, StorageContext>): () => void;See DatadataServer.onDocumentsPurged. Invoked like AsyncDatadataServer.onDocumentChange: synchronously inside the queued purge, never awaited.
previewSchemaChange(context: StorageContext & SystemContext, params: {
type: string;
schema: DocumentSchema;
limit?: number;
afterDocId?: string;
}): Promise<SchemaChangePreview>;purgeDeletedDocuments(context: EventBusContext & StorageContext & SystemContext, params: {
deletedBefore: number;
type?: string;
limit?: number;
}, options?: ObservabilityOptions): Promise<{
purged: DeletedDocumentEntry[];
}>;purgeDocuments(context: EventBusContext & StorageContext & SystemContext, params: {
docIds: readonly string[];
}, options?: ObservabilityOptions): Promise<{
purgedAt: number;
}>;reauthorizeConnections(context: EventBusContext & StorageContext): Promise<void>;rebuildDocument(context: StorageContext & SystemContext, ref: {
docId: string;
}): Promise<VerifyResult | null>;renameDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
name: string | null;
}, options?: ObservabilityOptions): Promise<void>;restoreDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
}, options?: ObservabilityOptions): Promise<void>;resyncPresence(context: EventBusContext): Promise<void>;sweepBlobs(context: StorageContext & SystemContext, params: {
retentionMs: number;
limit?: number;
}): Promise<BlobSweepResult>;updateDocument<Type extends keyof Registry & string>(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
type?: Type;
callback: (doc: DatadataDocument<TypedDocumentMutable<Registry, Type>>, yjsDocs: YjsDocsAccessor, ops: DocumentUpdateOps) => void;
expectedSequence?: number;
}, options?: ObservabilityOptions): Promise<void>;See DatadataServer.updateDocument. The callback itself stays synchronous.
validateDocuments(context: StorageContext & SystemContext, params?: {
mode?: "incremental" | "full";
type?: string;
limit?: number;
afterDocId?: string;
}): Promise<ValidateDocumentsResult>;verifyDocument(context: StorageContext & SystemContext, ref: {
docId: string;
}): Promise<VerifyResult | null>;BlobBodyLengthMismatchError
export declare class BlobBodyLengthMismatchError extends ErrorThrown by an object storage adapter's putObject when the body contradicts its declared contentLength — the one putObject failure that is the CLIENT's fault (⇒ 400), where anything else is a store-side failure the upload handler answers 500 for (retryable; the client's request was fine). Part of this transport contract rather than the core adapter seam: only hosts serving untrusted uploads need the distinction.
constructor(message: string);Constructs a new instance of the BlobBodyLengthMismatchError class
DatadataServer
export declare class DatadataServer<Registry extends ValidatorRegistry, EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>The synchronous datadata server engine for one folder: it handles the wire events of connected clients (DatadataServer.handleEvent), applies writes through the schema and authorization gates, broadcasts changes through the event bus, and offers the same operations as direct calls for server-side code, plus maintenance (verify, rebuild, export, import, purge) and the blob lane. Create one with createServer; a network-backed store uses AsyncDatadataServer instead.
constructor(storageAdapter: DatadataStorageAdapter<Registry, StorageContext>, schemas: Registry, eventBus: DatadataEventBus, logger: Logger, objectStorageAdapter?: DatadataObjectStorageAdapter<StorageContext>, commands?: Commands);Constructs a new instance of the DatadataServer class
authorizeBlobRead(context: StorageContext, ref: {
blobId: string;
}): BlobReadGrant | null;The blob download gate: the object-store key (plus response facts) iff the caller may read the blob, else null — and the host answers 404, so a hidden blob is indistinguishable from a nonexistent one (the document read gate's contract, extended to bytes). A live blob is readable iff the principal can read at least one referencing document; a staged blob only by its uploader. See ops/blobs.ts for the full rules.
commandDocument(context: EventBusContext & StorageContext & LoggerContext, params: CommandInvocation<Commands>, options?: ObservabilityOptions): void;Run a domain command directly from the server (the programmatic counterpart of the wire doc:command, as DatadataServer.updateDocument is of doc:update). The command runs as the context's principal, through the same gates a client's command meets: read, then the command's own authorization kind (access.commands.<name>), then the registry the server was created with, the args shape, the mutator (with now stamped here) and the schema. A server agent can so be granted only the commands it runs, with no update. With the server's commands a defineCommands map, params is typed as a session's command is (CommandInvocation): type names the docType, name one of its commands and args that command's shape. It runs against the document as stored, with no guard, whatever the command's contract, and is not recorded for replay. A mutator that changes nothing commits nothing. Throws on failure: a refusal (a RefusedCommandError carrying the mutator's code, or unknownCommand / invalidArgs / mutatorError), a denial (unauthorized), a missing or unreadable document (notFound), a type the document does not have (invalidRequest), or a result that fails the schema (schemaValidation).
createBlobUpload(context: StorageContext & LoggerContext, params: BlobStageParams): {
blobId: string;
};Stage a blob upload: authorize as the intended write (kind + docId + docType, see BlobUploadIntent — the cheapest gate, and the same verdict the handle-committing write will get), refuse an upload no blobRef field of the type could accept (accept/maxBytes), mint the opaque blob id, and record the staged catalog row. The host then streams the request body into the object store under the returned blobId and calls DatadataServer.finalizeBlobUpload — the two-phase upload whose ordering the integrity gate enforces (a handle can only be committed after finalize, so a synced document never points at missing bytes). Requires an object storage adapter.
createDocument<Type extends keyof Registry>(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
type: Type;
data: TypedDocumentInput<Registry, Type>;
yjsDocs?: Record<string, Uint8Array>;
}, options?: ObservabilityOptions): void;Create a document directly from the server (the programmatic counterpart of the wire doc:create), typed by the registry's Type. The data is validated against the type's schema and the create is broadcast like a client's. Throws on failure: a taken docId (alreadyExists), a soft-deleted one (deleted), data that fails validation (schemaValidation), an invalid docId (invalidRequest), or a type with no schema (unknownType).
createDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
type: string;
data: unknown;
yjsDocs?: Record<string, Uint8Array>;
name?: string | null;
}, options?: ObservabilityOptions): void;Untyped form of DatadataServer.createDocument, for a type defined only by its stored sys:schema:<type> document. name sets the document's index-level name.
deleteDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
baseSequence?: number;
guard?: DeleteGuardMode;
}, options?: ObservabilityOptions): void;Soft-delete a document directly from the server (the programmatic counterpart of the wire doc:delete). Optional opt-in CAS via guard: "sequence" + baseSequence, mirroring updateDocument's expectedSequence. Throws on failure (preconditionFailed guard miss, never-existed, sys: docId); deleting an already-deleted document is benign (acks, no new folder event).
exportDocument(context: StorageContext & SystemContext, ref: {
docId: string;
}): Promise<DocumentExport | null>;Export a document as a portable DocumentExport — its type plus full event log — for replica/DR or moving it to another store. The log is self-describing for reconstruction (PerDocumentReplay), so the export needs nothing else. Returns null if the document does not exist.
Deliberately EXEMPT from the per-operation read authorization, mirroring importDocument's write exemption: the bulk replica/DR path is a maintenance surface — it requires system authority (SystemContext) instead.
finalizeBlobUpload(context: StorageContext, params: {
blobId: string;
size: number;
sha256: string;
}): Promise<void>;Record the streamed bytes' facts (size, sha256 — as returned by the object adapter's putObject) against the staged row, completing the upload. Host-internal: called by the same request handler that staged and streamed; enforces maxBlobSizeBytes. Async even on this sync facade: when the staged row was swept while the host streamed (finalize answers notFound), the just-stored bytes are deleted before the failure surfaces — otherwise they would sit in the object store with no catalog row, invisible to the catalog-driven sweep forever.
findInvalidDocuments(context: StorageContext & SystemContext, opts?: {
type?: string;
}): InvalidDocument[];Enumerate documents discovered (lazily, on read) to be invalid against the current schema — the foundation for an admin "resolve" tool. Cheap: backed by the invalid_schema_sequence column. A document is resolved by a normal updateDocument whose callback produces conforming data (the write path stays strict and clears the flag).
The listing is the union of two sources, each entry attributed by kind: the stored schema-shape flags (discovered lazily on read — a lower bound until a DatadataServer.validateDocuments incremental sweep finishes with remaining: false), and the DERIVED reference entries — documents whose indexed reference edges dangle against the live document set, computed fresh per call so target deletes and restores reflect instantly with nothing to heal. The reference half is complete with respect to the reference index itself; a document last written before the index existed contributes no edges until a validated write or a validation sweep backfills them.
A maintenance surface like verifyDocument / purgeDeletedDocuments: deliberately EXEMPT from per-operation read authorization (the listing exposes docIds and types across the whole folder), so it requires system authority (SystemContext) instead.
getBlobMetadata(context: StorageContext, ref: {
blobId: string;
}): BlobMetadata | null;The blob's catalog facts (content type, size, sha256, state), gated exactly like DatadataServer.authorizeBlobRead — null for a blob the caller may not read, indistinguishable from absence. Serves HEAD /blob/:id.
getDocument<Type extends keyof Registry & string>(context: StorageContext, params: {
docId: string;
type: Type;
}): {
document: DatadataDocument<TypedDocumentData<Registry, Type>> | null;
invalid?: DocumentInvalid;
};Read a document. Returns a { document, invalid? } WRAPPER — not the document directly: reach for result.document (it is null when the document does not exist), never result.data.
The wrapper exists to carry validity alongside the document. Reads are relaxed: a document that no longer conforms to its current schema is returned with its best-effort migrated data rather than throwing, alongside an invalid marker describing the failure (absent when the document is valid). The invalidity is also recorded for findInvalidDocuments. Callers decide whether to use the data, warn, or gate editing until an updateDocument resolves it (the write path stays strict).
Note: when invalid is set, document.data typed via the registry overload is best-effort — it failed validation and may not fully conform to the type. This mirrors the client, which delivers the document and surfaces validity separately (via sys:client-docs-status); branch on invalid before trusting typed fields.
getDocument(context: StorageContext, params: {
docId: string;
type?: string;
}): {
document: DatadataDocument | null;
invalid?: DocumentInvalid;
};Untyped form of DatadataServer.getDocument: type, when given, is checked against the stored document's type (a mismatch throws); the data is typed unknown.
getDocumentInitEvent(context: StorageContext, ref: {
docId: string;
}): DocumentInitEvent | DocumentNotFoundEvent | DocumentDeletedEvent | DocumentErrorEvent;Get a document as a DocumentInitEvent or DocumentNotFoundEvent. This is useful for HTTP endpoints that need to return the same data structure as WebSocket subscriptions.
getYjsDoc(context: StorageContext, ref: {
docId: string;
yjsId: string;
}): Y.Doc | null;Get a Y.Doc from storage by document ID and Yjs ID. This is useful for agents that need to read Yjs content without subscribing to updates.
Returns
The Y.Doc instance, or null if not found
handleDisconnect(context: EventBusContext & LoggerContext & ConnectionIdentityContext): void;Clear every presence entry a closing connection holds, broadcasting the departure to each document's remaining subscribers. The transport host calls this when a connection closes (e.g. webSocketClose/webSocketError, or an in-process client's dispose), BEFORE detaching the connection from the event bus — the broadcasts route to the other subscribers either way, but the ordering keeps "presence cleared" inside the connection's lifetime. A no-op for connections that never published presence.
handleEvent(context: EventBusContext & StorageContext & LoggerContext & ConnectionIdentityContext & {
authority?: never;
connectionId?: string;
}, event: ClientSentEvent, options?: ObservabilityOptions): void;Dispatch one client-sent event. Fully synchronous — every write (including its authorization) completes in-handler. Write failures surface as doc:error to the calling client rather than throwing; protocol violations (malformed event, oversized event id) still throw.
importDocument(context: EventBusContext & StorageContext & SystemContext, documentExport: DocumentExport, options?: {
validation?: "strict" | "relaxed";
}): Promise<{
indexSequence: number;
}>;Import a DocumentExport into this (fresh) store: replay its log to derive the snapshot, validate that snapshot against the schema at its stamped conformed sequence (see #validateImportedSnapshot), then write the document with its full history via the adapter. The reconstructed document is logically equal to the source and carries its history, so it can itself be re-exported or serve reads. Returns the folder's new index sequence; throws alreadyExists if the document is already present.
A faithful same-lineage export always passes the validation — a stored snapshot conforms to its conformed cursor by construction — so a strict (default) rejection means the events were produced under a schema this folder never held. validation: "relaxed" accepts such an import anyway (logged); the first read then bring-forwards and flags the document invalid like any other non-conforming document. Relaxed also accepts a conformed sequence that OUTRUNS the stored schema (the stamps are foreign coordinates anyway): the snapshot is then validated against the newest schema the store holds and the stored cursor clamped to it, so the stored state keeps SchemaSequenceWithinRange. Only the no-schema-at-all check is unconditional — a document whose type resolves no validator could never be read.
Deliberately EXEMPT from the per-operation write authorization (authorizeWrite): imported events keep their ORIGINAL per-event attribution (subject/actor), so the importing principal is not the author and evaluating it against per-write policy would be meaningless. The bulk -operator path requires system authority (SystemContext) instead.
NOT exempt from the presence existence-transition: an import brings the document into existence like create/restore do, so retained presence subscriptions (taken while the docId was absent, hence never read-gated) are re-inited or evicted the same way — otherwise the first presence publish after an import of a read-hidden document would broadcast to a non-reader. This is why the context is event-bus-capable, unlike export/verify/rebuild's.
onBlobSweepCandidate(callback: BlobSweepCandidateCallback): () => void;Register a hook on "the blob catalog may hold a new GC candidate" — a blob staged (by DatadataServer.createBlobUpload or an import), a document's blob edges replaced (a validated write, or a lazy read or DatadataServer.validateDocuments refreshing them after a schema change), or a purge. Hosts arm their sweep schedule from it; it over-approximates, so arming must be idempotent. Invoked synchronously after the triggering storage write, never awaited. Returns the unregister function. See BlobSweepCandidateCallback.
onDocumentChange(callback: DocumentChangeCallback<EventBusContext, StorageContext>): () => void;Register a callback to receive every stored document's changes (doc:init, doc:patch, and doc:deleted) — creates, updates and deletes; a restore is reported by DatadataServer.onDocumentRestored instead. This is useful for server-side agents that need to observe all changes without explicit subscriptions. The lazy schema bring-forward a read performs persists migrated data but is not an accepted document write, and does not fire it (the sys:schema:<type> change does). The folder listings (sys:index / sys:trash) never fire it: watch them with an in-process client instead. See DocumentChangeCallback.
Returns
Unsubscribe function to remove the callback
onDocumentRestored(callback: DocumentRestoredCallback<EventBusContext, StorageContext>): () => void;Register a callback invoked once per committed restore — from DatadataServer.restoreDocument or a wire doc:restore — with the restored document's trash entry. A restore does not fire DatadataServer.onDocumentChange. See DocumentRestoredCallback.
Returns
Unsubscribe function to remove the callback
onDocumentsPurged(callback: DocumentsPurgedCallback<EventBusContext, StorageContext>): () => void;Register a callback invoked once per purge batch — from DatadataServer.purgeDocuments, DatadataServer.purgeDeletedDocuments or a wire doc:purge — with the purged tombstones' entries. See DocumentsPurgedCallback.
Returns
Unsubscribe function to remove the callback
previewSchemaChange(context: StorageContext & SystemContext, params: {
type: string;
schema: DocumentSchema;
limit?: number;
afterDocId?: string;
}): SchemaChangePreview;Dry-run a schema change: what committing schema as the type's sys:schema:<type> document would do to the type's live documents, without writing anything.
The candidate goes through the schema write's own gates first (the schema document validator, migration stamping at the sequence the commit would produce, the blobRef gate), so a change the write would refuse throws the same error here. Then every live document of the type gets the read path's verdict twice, under the current schema and under the candidate: its bring-forward from its own cursor, then shape and reference validation. When its shape passes under the candidate, the blobRef handles in the validated value (defaults included) are also looked up in the blob catalog, since the commit cannot register a handle naming no finalized blob. Documents are read from the stored snapshot, not through the gated read, so nothing is persisted and no invalid flag moves.
Advisory: the schema write itself does not consult it. A host that wants a confirm step before a breaking schema change builds it on this result.
limit, a positive integer when given, bounds one invocation (loop while nextAfterDocId is non-null, passing it back as afterDocId) — the full validation sweep's paging, and the same budget advice for a Durable Object. Requires system authority, like its maintenance siblings. Throws for a sys: type and for a type with no stored schema document.
purgeDeletedDocuments(context: EventBusContext & StorageContext & SystemContext, params: {
deletedBefore: number;
type?: string;
limit?: number;
}, options?: ObservabilityOptions): {
purged: DeletedDocumentEntry[];
};Purge every soft-deleted document whose tombstone is older than deletedBefore (epoch ms) — the retention sweep behind "trash empties after N days", or with deletedBefore: Date.now() an "empty trash". The library owns only this mechanism: the retention window and WHEN the sweep runs (a host alarm, a cron trigger, opportunistically on wake) are application policy. Oldest tombstones purge first; limit bounds one invocation's work (a caller loops while purged.length === limit), and each invocation is one batch — one trash-sequence bump, one removal patch, one storage TRANSACTION. That atomicity cuts both ways: omitting limit purges every eligible tombstone in one synchronous transaction and one broadcast fan-out, which on a large backlog can exceed a Durable Object's per-invocation budget and then fails whole (nothing purges, no incremental progress). Unlimited sweeps are only safe for small trash volumes — a retention job should always pass a conservative limit and loop. Returns the purged entries. Requires system authority, like purgeDocuments.
purgeDocuments(context: EventBusContext & StorageContext & SystemContext, params: {
docIds: readonly string[];
}, options?: ObservabilityOptions): {
purgedAt: number;
};Hard-delete (purge) a batch of soft-deleted documents: destroy their event logs, Y.Doc states and snapshots, and RELEASE their docIds — afterwards each id behaves exactly like one that never existed, so a create under it starts a fresh generation (what makes purge safe for deterministic ids). Only soft-deleted documents can be purged: delete first, even for erasure — purge is strictly how trash stops being retained. Its guarantee is that THIS folder destroyed its copy; exports and replicas elsewhere are outside it, and importing an old export afterwards is simply a legal create of a new generation.
The whole batch is one trash mutation: the trash counter advances ONCE and subscribers receive ONE removal patch, however many documents purge. The documents' own subscribers are NOT notified: they already hold deleted status from the delete broadcast, restore just became impossible (the sys:trash entry vanishing is the observable signal), and a resubscribe now answers doc:notfound.
This direct call is deliberately EXEMPT from the per-operation write authorization: destroying content irreversibly from host code is an operator/retention decision, mirroring importDocument/exportDocument, so it runs under system authority (SystemContext).
All-or-nothing: throws if ANY document is live (storageError — delete it first) or not found, and a throwing batch purges nothing; never-existed and already-purged are deliberately the same case, because a purge releases the id.
reauthorizeConnections(context: EventBusContext & StorageContext): void;Re-run the read-revocation sweep for host-driven authorization-fact changes the engine cannot observe — e.g. a connection's grants updated in the transport's connection metadata having moved. Commits to sys:access / sys:schema:* trigger the sweep automatically; everything else is the host's call. Also refreshes every sys:principal subscriber with its connection's CURRENT principal — the one leg the automatic sweeps skip (only the host can change a connection's metadata principal, so only this entry point needs it) — which is what makes a mid-connection principal change (grant widening, demotion) reach the client's prediction and capability layer.
rebuildDocument(context: StorageContext & SystemContext, ref: {
docId: string;
}): VerifyResult | null;Rebuild a document's stored snapshot (data + Y.Docs) by replaying its event log — the source of truth — and overwriting the derived snapshot in place. The recovery / replica-construction primitive: it makes the log win when the snapshot is lost or corrupt. The log and sequence are left untouched (replay is computed from them), so this never loses history.
Off the hot path (admin/maintenance/recovery). Requires system authority (SystemContext). Returns the post-rebuild verification result (ok unless the store mutated concurrently), or null if the document does not exist.
renameDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
name: string | null;
}, options?: ObservabilityOptions): void;Rename a document directly from the server (the programmatic counterpart of the wire doc:rename). A name is index-level metadata, not content, so this patches sys:index only — it never touches the document's data, sequence or event log, and the document's own subscribers are not notified. null clears the name. Throws on a sys: docId or a document that is not live.
restoreDocument(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
}, options?: ObservabilityOptions): void;Restore a soft-deleted document directly from the server (the programmatic counterpart of the wire doc:restore). Subscribers receive a fresh doc:init and the document re-enters sys:index. Restoring a live document is benign (acked with its current state); throws if the document never existed.
resyncPresence(context: EventBusContext): void;Send doc:resync to every live presence document's subscribers. The transport host calls this when it constructs a fresh server instance while connections are still open — a Durable Object waking from hibernation (or restarting) with live WebSockets: the previous instance's in-memory cell registry died with it, so each still-subscribed presence document is nudged and its subscribers republish their cells, rebuilding the registry in one round-trip. doc:resync is used INSTEAD of an empty doc:init so a client does not flicker its peers' cursors empty while republishes arrive (it arms its staleness ledger instead). This is how presence survives hibernation WITHOUT touching storage. Harmless against a populated registry (republishes overwrite identical cells), so a host that can't distinguish cold from first start may call it whenever restored connections exist.
sweepBlobs(context: StorageContext & SystemContext, params: {
retentionMs: number;
limit?: number;
}): Promise<BlobSweepResult>;The blob GC sweep: tombstone-then-delete every staged blob created more than retentionMs ago (abandoned uploads) and every blob that lost its last reference more than retentionMs ago (grace elapsed), retrying tombstones earlier sweeps left behind. The retention (see BlobSweepPolicy, checked by assertBlobSweepPolicy) and WHEN the sweep runs (a DO alarm, an interval, an explicit call in tests) are host policy, like DatadataServer.purgeDeletedDocuments. limit bounds one invocation's reclaimed blobs; the result's nextSweepAt says when the next sweep is due under the same retention — now, a later moment, or not until the catalog changes — the host's scheduling hint, see BlobSweepResult. Requires system authority and an object storage adapter. Async even on this sync facade: byte deletion is an object-store round trip.
updateDocument<Type extends keyof Registry & string>(context: EventBusContext & StorageContext & LoggerContext, params: {
docId: string;
type?: Type;
callback: (doc: DatadataDocument<TypedDocumentMutable<Registry, Type>>, yjsDocs: YjsDocsAccessor, ops: DocumentUpdateOps) => void;
expectedSequence?: number;
}, options?: ObservabilityOptions): void;Update a document directly from the server. This bypasses the client subscription model and is useful for server-side agents that need to modify documents.
validateDocuments(context: StorageContext & SystemContext, params?: {
mode?: "incremental" | "full";
type?: string;
limit?: number;
afterDocId?: string;
}): ValidateDocumentsResult;The completeness sweep over lazy validation: force discovery over documents no read has touched, so findInvalidDocuments stops being a lower bound. Validates by READING — each candidate runs the ordinary gated read (bring forward, validate, persist a conforming change and advance the cursor, or flag invalid and leave it stale), so a sweep is exactly "as if someone had read every candidate".
Two modes: - incremental (default, cheap): only documents whose conformed cursor trails their type's current schema sequence and that are not already flagged at it. When it finishes (remaining: false), findInvalidDocuments is COMPLETE for schema-shape invalidity — "is the migration done?" is answerable. It deliberately does NOT revisit documents flagged at the current sequence, so reference invalidity that appeared or healed on a cursor-current document (a target deleted or restored after the flagging read) is not re-checked here. - full: every live document of a resolvable app type, in stable docId order behind a nextAfterDocId cursor — the honest answer for reference invalidity too, complete as of the pass, and the healing pass for stale flags.
limit bounds one invocation's work (loop while remaining, threading nextAfterDocId in full mode) — the same budget discipline as purgeDeletedDocuments; scheduling is host policy. Omitting limit sweeps everything in ONE synchronous invocation: each document's persist commits independently (a mid-sweep abort loses no data, unlike purge), but the invocation still runs as one uninterrupted task ahead of live client operations and can exceed a Durable Object's budget on a large backlog — a production host should always pass a conservative limit and loop. Validation failures of individual documents surface in errors and keep remaining true rather than aborting the batch (each processed document commits independently, so partial progress is durable). Requires system authority, like its maintenance siblings.
verifyDocument(context: StorageContext & SystemContext, ref: {
docId: string;
}): VerifyResult | null;Verify that a document's stored snapshot (data + Y.Docs) is exactly the replay of its own event log, at logical equality. The executable form of the reconstruction guarantee (host.allium ReplayableLog / core.allium PerDocumentReplay): a passing result means the document can be rebuilt from its log alone with no loss. Compares the *raw* stored snapshot against the *stored* log (no migration-on-read), so it checks storage consistency, not the read path.
A recovery/maintenance guard — keep it off the hot read path. Returns null if the document does not exist, a VerifyResult for logical drift, and *throws* if the log is structurally corrupt and cannot be folded (an unparseable patch, a create/update event missing its blob, an unsupported Yjs type).
Deliberately EXEMPT from the per-operation read authorization (like exportDocument / rebuildDocument / importDocument): a maintenance surface — it requires system authority (SystemContext) instead.
DocumentExportParseError
export declare class DocumentExportParseError extends ErrorA serialized export whose text is malformed — not JSON, or not this shape.
constructor(message: string);Constructs a new instance of the DocumentExportParseError class
DocumentNotFoundError
export declare class DocumentNotFoundError extends DocumentRetrievalErrorThrown when a write targets a docId that has neither a live document nor a tombstone. Callers that want to fall back to creating the document (upsert flows) can catch this with instanceof instead of matching the message string. The message is caller-supplied so existing wording is preserved.
Extends DocumentRetrievalError with the category pinned to notFound (the StorageConflictError pattern), so the wire scaffold answers a classifiable doc:error(notFound) instead of the fallback unknown — which used to mark a merely-missing document's status as an engine error on the client.
constructor(docId: string, message?: string);Constructs a new instance of the DocumentNotFoundError class
readonly docId: string;The id of the document that was not found.
DocumentRetrievalError
export declare class DocumentRetrievalError extends ErrorError thrown during document retrieval with an explicit category. This allows error handling code to use the category directly without string matching.
constructor(
category: DocumentErrorCategory, message: string, options?: {
docId?: string;
});Constructs a new instance of the DocumentRetrievalError class
readonly category: DocumentErrorCategory;What kind of failure this is; the category the wire doc:error carries.
readonly docId?: string;The document the failure is about, when the throw site can name one. Batch operations (doc:purge) rely on this to answer with the OFFENDING docId rather than the batch's first; single-target paths may leave it unset (their emit sites already know the docId).
FolderRegistry
export declare class FolderRegistry<TFolder extends FolderRegistryEntry>Lazily-initialized folder entries keyed by folderId, with coordinated shutdown. Hosts supply the per-folder init (construct adapter + server over the shared database, seed schemas/access); the registry guarantees:
- ONE init per folder, shared by concurrent first requests (the init runs behind a cached promise), and a FAILED init is evicted so the next request retries instead of caching the failure forever. - evict (the host's idle/eviction policy hook — the registry imposes no policy of its own) closes the folder's server and, crucially, ORDERS a re-open behind that close: a get racing an evict initializes only after the evicted server has fully drained and released its adapter, so a single-writer-per-folder adapter never sees two live instances. - close is the coordinated shutdown: every folder's server closes — draining its operation FIFO — before close() resolves, so the host can then close the shared database with nothing in flight.
[Symbol.asyncDispose](): Promise<void>;Symbol.asyncDispose alias for FolderRegistry.close, which enables await using registry = ....
constructor(initFolder: (folderId: string) => Promise<TFolder> | TFolder);Constructs a new instance of the FolderRegistry class
close(): Promise<void>;Coordinated shutdown: refuse new folders, then close every folder's server — each close drains that folder's operation FIFO — and wait for any in-flight evictions. When this resolves no folder has work in flight, so the host may close the shared database. Idempotent.
Safe to race with FolderRegistry.evict: new gets are already refused when this runs, so an eviction landing while settled() is awaited closes the SAME server this method closes (server close is idempotent — one shutdown, two awaiters), and one landing after the folder map is cleared finds nothing and no-ops.
evict(folderId: string): Promise<void>;Remove a folder and close its server (draining its queue and releasing its adapter). The registry has NO eviction policy of its own — when to call this (idle timers, LRU, never) is the host's decision. Safe to race with FolderRegistry.get: a concurrent re-open waits for this close to finish. Unknown folderIds are a no-op.
A REJECTED close (a storage adapter whose release throws) rejects this eviction and any re-open racing it — deliberately: the folder's teardown did not finish cleanly, so the racing open fails loudly with that error instead of initializing over half-released state. Only the racing attempt is poisoned; a failed init self-evicts (see get), so a LATER get() retries a fresh init.
get(folderId: string): Promise<TFolder>;The folder's entry, initializing it on first request. Concurrent calls share one init; a rejected init propagates to every waiter and is evicted so a later call retries. Rejects after FolderRegistry.close.
settled(): Promise<TFolder[]>;Every folder whose init has completed successfully, waiting out in-flight inits (failed ones are simply absent). The shutdown-adjacent enumeration: a host terminates each folder's transport connections from this before calling FolderRegistry.close.
StorageConflictError
export declare class StorageConflictError extends DocumentRetrievalErrorThrown by a storage adapter when a write's optimistic-concurrency assumption is violated: the CAS sequence guard found the document advanced past the sequence the engine read (a concurrent write won), or the write's target vanished between the engine's read and the write (a concurrent delete won). Both mean "the world moved — re-running the whole operation from a fresh read may succeed", which is exactly what the engine's withRetry combinator does: it retries THIS class and nothing else, so an adapter must throw it (not a generic error) for a lost race to be retried.
Extends DocumentRetrievalError with the category pinned to storageError, so every existing category-based path — the wire doc:error, the client's transient-storageError re-send — treats a conflict exactly as before; the subclass only adds the retryability signal.
constructor(message: string, options?: {
docId?: string;
});Constructs a new instance of the StorageConflictError class
Functions
assertBlobSweepPolicy
export declare function assertBlobSweepPolicy(policy: {
retentionMs: number;
limit?: number | undefined;
}, source: string): void;Refuse a sweep policy that is not one, naming source (where the policy came from) and the field. The engine checks every sweep call with it, and a host checks its configured policy with it BEFORE deriving anything else from it (an alarm or a timer), since those uses never reach the engine:
- retentionMs must be an integer from 0 to MAX_BLOB_RETENTION_MS. Negative sweeps uploads still in flight; NaN or an infinity makes every cutoff comparison false (nothing is ever reclaimed) and every derived deadline invalid; a fraction binds to an integer column as a fraction. Zero is legitimate: everything eligible as soon as it exists. - limit, when given, must be a positive integer: a limit of 0 sweeps nothing yet leaves the same candidates due, so a host draining while the sweep says more is due would ask forever. An unbounded sweep omits it.
blobCacheControl
export declare function blobCacheControl(policy: BlobCachePolicy, outcome: "granted" | "denied"): string;The Cache-Control value a policy prescribes for a granted response or a denial. A malformed policy throws on EVERY call, denials included, so a misconfigured host fails on its first blob response of any kind rather than on its first authorized read. Exported for app-owned routes that derive from a blob (the playground's image variants) and so restate the blob lane's header rules — one mapping, not one per route.
blobRequestFailure
export declare function blobRequestFailure(error: unknown): BlobRequestFailure;Map an expected blob-operation failure into the outcome shape BlobRouteServer carries: a DocumentRetrievalError becomes { ok: false, category, message }; anything else stays a throw — a genuine 500 the transport surfaces as such. Shared by every host that fronts the engine's throwing blob facade with the outcome-shaped seam (the DO base class's RPC helpers, the Node harness's in-process wiring).
connectionIdOf
export declare function connectionIdOf(context: {
connectionId?: ConnectionIdOrServerInitiated;
}): string | null;The connection that made this call — null for a server-initiated one, and for a context that carries no identity at all. Everything that needs a socket (a caller ack, a subscription, a presence cell's owner) reads the context through this, never context.connectionId directly.
createAsyncServer
export declare function createAsyncServer<Registry extends ValidatorRegistry, EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>(params: {
storageAdapter: AsyncDatadataStorageAdapter<Registry, StorageContext>;
dataEventBus: DatadataEventBus<EventBusContext>;
logger: Logger;
objectStorageAdapter?: DatadataObjectStorageAdapter<StorageContext>;
commands?: Commands;
}): AsyncDatadataServer<Registry, EventBusContext, StorageContext, Commands>;The async twin of createServer — see its doc comments for the options.
createDocumentSizeGuards
export declare function createDocumentSizeGuards(limits: DatadataServerLimits): DocumentSizeGuards;Build the DocumentSizeGuards for a set of resolved limits.
createServer
export declare function createServer<Registry extends ValidatorRegistry, EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry>(params: {
storageAdapter: DatadataStorageAdapter<Registry, StorageContext>;
dataEventBus: DatadataEventBus<EventBusContext>;
logger: Logger;
objectStorageAdapter?: DatadataObjectStorageAdapter<StorageContext>;
commands?: Commands;
}): DatadataServer<Registry, EventBusContext, StorageContext, Commands>;Create a DatadataServer over a synchronous storage adapter and an event bus. Only the built-in sys:* types are validated statically; every app type is validated from its stored sys:schema:<type> document, so Registry types the API without registering validators.
handleBlobDownload
export declare function handleBlobDownload<Context>(request: Request, options: {
blobId: string;
server: Pick<BlobRouteServer, "authorizeBlobRead">;
objectStorage: DatadataObjectStorageAdapter<Context>;
context: Context;
cache: BlobCachePolicy;
}): Promise<Response>;GET …/blob/:blobId — the read gate then the bytes. A denied or missing blob answers the same 404 ("a hidden blob is indistinguishable from a nonexistent one"). Serves strong ETags (the sha256), If-None-Match, and single-range Range requests, under the pinned security headers and the host's BlobCachePolicy.
handleBlobHead
export declare function handleBlobHead(_request: Request, options: {
blobId: string;
server: Pick<BlobRouteServer, "getBlobMetadata">;
cache: BlobCachePolicy;
}): Promise<Response>;HEAD …/blob/:blobId — the catalog facts as headers, gated exactly like the download (404 for anything the caller may not read), under the same cache policy as the GET for that URL.
handleBlobUpload
export declare function handleBlobUpload<Context>(request: Request, options: {
server: Pick<BlobRouteServer, "createBlobUpload" | "finalizeBlobUpload">;
objectStorage: DatadataObjectStorageAdapter<Context>;
context: Context;
}): Promise<Response>;POST …?kind=<create|update>&docId=<docId>&docType=<docType>[&path=<json>] with the file bytes as the body and its Content-Type + Content-Length headers (a body without a declared length is refused with 411 — see the guard below) — the two-phase upload as one request: stage (authorized as the intended write the query names, see BlobUploadIntent), stream the body into the object store, finalize. Answers 200 { blobId, size, contentType, sha256 }; the client then commits the handle through an ordinary document write.
parseDocumentExport
export declare function parseDocumentExport(json: string): DocumentExport;Parse a serializeDocumentExport string back into a live DocumentExport. Throws DocumentExportParseError on text that is not an export.
Structural validation is deliberately SHALLOW — the envelope and the blob spelling, the parts this module owns. The deep questions (does the log replay? does the snapshot conform to its schema?) belong to importDocument, which answers them for every export whatever its source.
serializeDocumentExport
export declare function serializeDocumentExport(documentExport: DocumentExport): string;The plain-text form of documentExport: JSON with each Yjs blob spelled as base64. The inverse of parseDocumentExport.
serverInitiated
export declare function serverInitiated(name: string): ServerInitiated;The ServerInitiated marker for a call the server itself starts, to put where a context's connectionId goes. name says what started it, for logs and observability.
Interfaces
BlobMetadata
export interface BlobMetadataThe catalog facts a read-gated caller may see (no uploader identity).
blobId: string;The blob's id.
contentType: string;The MIME type declared at upload start.
createdAt: number;When the upload was staged, in epoch milliseconds.
sha256: string | null;Lowercase-hex SHA-256 of the bytes, recorded at finalize; null while in flight.
size: number | null;Byte length, recorded at finalize; null while the upload is in flight.
state: "staged" | "live";staged until a committed document write references the blob, then live. See StoredBlobMetadata.state.
BlobReadGrant
export interface BlobReadGrantA granted blob download: the blob id — which is also the object-store key the host streams from (a store that prefixes keys per folder does so inside its adapter) — plus the facts a host needs to answer the HTTP response (Content-Type, ETag via sha256, Content-Length).
blobId: string;The blob's id and object-store key.
contentType: string;The MIME type to answer as Content-Type.
sha256: string | null;Lowercase-hex SHA-256 of the bytes, usable as an ETag; null while in flight.
size: number | null;Byte length, for Content-Length; null while the upload is in flight.
BlobRequestFailure
export interface BlobRequestFailureAn expected blob request failure, carried as data instead of a thrown error.
category: DocumentErrorCategory;The failure's category, which the upload handler maps to an HTTP status.
message: string;A human-readable description of the failure.
ok: false;Always false: distinguishes the failure from a successful outcome.
BlobRouteServer
export interface BlobRouteServerWhat the handlers need from the engine, already bound to the requesting principal. In the Worker→DO deployment the DO base class's four blob RPC helpers satisfy these member for member, so the Worker wires { createBlobUpload: (params) => stub.createBlobUpload(principal, params), … }; an in-process host (the Node harness) wires them straight to its server facade, converting expected throws with blobRequestFailure.
authorizeBlobRead(ref: {
blobId: string;
}): Awaitable<BlobReadGrant | null>;The download gate: the grant when the principal may read the blob, else null (the handler answers 404). See DatadataServer.authorizeBlobRead.
createBlobUpload(params: BlobStageParams): Awaitable<BlobStageOutcome>;Authorize the upload as its intended write and stage a catalog row, answering the minted blob id the host then streams the body under. See DatadataServer.createBlobUpload.
finalizeBlobUpload(params: {
blobId: string;
size: number;
sha256: string;
}): Awaitable<BlobFinalizeOutcome>;Record the stored bytes' size and SHA-256, completing the upload. See DatadataServer.finalizeBlobUpload.
getBlobMetadata(ref: {
blobId: string;
}): Awaitable<BlobMetadata | null>;The blob's catalog facts, gated like BlobRouteServer.authorizeBlobRead; serves HEAD requests. See DatadataServer.getBlobMetadata.
BlobStageParams
export interface BlobStageParams extends BlobUploadIntentWhat staging an upload takes: the bytes' declared facts plus the BlobUploadIntent the upload is authorized as.
contentLength?: number;The declared body length, when the transport knows it (the HTTP route requires it).
contentType: string;The MIME type of the bytes, stored on the catalog row and served on download.
BlobSweepHorizon
export interface BlobSweepHorizonWhat the blob catalog holds for a sweep, independent of any cutoff — as answered by getBlobSweepHorizon. Every row the sweep condition can ever match is covered: a tombstone matches under any cutoffs, a staged row once the staged cutoff passes its createdAt, and a live row with no reference edge once the unreferenced cutoff passes its unreferencedAt. So the earliest moment anything becomes sweepable is derivable from these three facts plus the retention policy, which stays the caller's.
hasTombstones: boolean;Whether any tombstone awaits its byte deletion (sweepable now).
oldestStagedAt: number | null;The oldest staged row's createdAt, or null when nothing is staged.
oldestUnreferencedAt: number | null;The oldest unreferencedAt among live rows with no reference edge, or null when none. A live row whose unreferencedAt is set but that still has an edge is referenced, and never counts.
BlobSweepPolicy
export interface BlobSweepPolicyHow a host runs the blob GC sweep: retentionMs is how long a staged upload may stay unfinalized-or-unreferenced and how long an unreferenced blob is kept (its grace window), limit how many blobs one sweep call reclaims at most. Retention is the host's policy, like purgeDeletedDocuments' deletedBefore; the engine only checks it.
limit: number;The most blobs one sweep call reclaims; a positive integer.
retentionMs: number;The retention and grace window, in milliseconds (0 to MAX_BLOB_RETENTION_MS).
BlobSweepResult
export interface BlobSweepResultWhat one sweep did, and when the next one is due — the second half is how a host schedules without a catalog query of its own.
nextSweepAt: number | null;When the catalog next holds a sweepable blob under the same retention, read from the catalog after the sweep: - at or before Date.now() — candidates are due already (this call's limit was hit, or rows came due while it ran); sweep again now. - later — the moment the oldest staged upload or unreferenced blob passes its retention; schedule a sweep for then. - null — nothing pending at all: no sweep is needed until the catalog changes (see onBlobSweepCandidate), so an idle folder need not wake.
sweptCount: number;How many blobs this sweep reclaimed (catalog row and bytes both removed).
ConnectionIdentityContext
export interface ConnectionIdentityContextThe connection identity the server reads off a handleEvent context. Both fields are optional because the server is generic over its contexts.
connectionId?: ConnectionIdOrServerInitiated;The calling connection's id: minted by the in-process bus, routed by the WebSocket bus. A presence write without one is rejected. A direct call the server itself started carries the ServerInitiated marker instead.
principal?: Principal;The caller's identity, the same one the storage adapter stamps on stored events. Presence stamps its subject onto cells so a client cannot claim another user's identity. Absent means anonymous.
DatadataEventBus
export interface DatadataEventBus<EventBusContext extends LoggerContext = LoggerContext>The transport seam between the server and its connections: it owns which connection is subscribed to which document, and delivers server-sent events to the caller, to one connection, or to a document's subscribers. The context of each call identifies the calling connection.
addSubscription(context: EventBusContext, ref: {
docId: string;
}): unknown;Subscribe the context's connection to the document, so broadcasts for it reach it.
connections(context: EventBusContext): Iterable<EventBusRecipient & {
subscribedDocIds: readonly string[];
}>;Every live connection with its principal and current subscriptions — the read-revocation sweep's enumeration and the per-subscriber view refresh. Callers materialize the iterable before mutating subscriptions.
hasSubscription(context: EventBusContext, ref: {
docId: string;
}): boolean;Whether the context's connection currently holds a subscription to the document. The server consults this where an operation is only meaningful for subscribers (presence): the bus owns the subscription state, so only it can answer.
removeSubscription(context: EventBusContext, ref: {
docId: string;
}): unknown;Unsubscribe the context's connection from the document.
sendEvents(context: EventBusContext, events: ServerSentEvent[], options?: {
respondToCaller?: true;
toConnection?: string;
recipientFilter?: (recipient: EventBusRecipient) => boolean;
}): void;Deliver events. By default each event is broadcast to the subscribers of its document (an event with no docId reaches every connection); respondToCaller sends them to the context's own connection instead.
subscribedDocIds(context: EventBusContext): Iterable<string>;Every distinct docId that currently has at least one subscriber, across all connections. Bus-wide, not per-connection (the context is for logging / consistency only). The server uses this on a cold start to find the live presence documents whose subscribers must be nudged to republish — see DatadataServer.resyncPresence. Order is unspecified; callers treat it as a set.
DatadataObjectStorageAdapter
export interface DatadataObjectStorageAdapter<Context = unknown>Byte transport for immutable blobs, keyed by the blob id. Keys are opaque server-minted ids ("an id scheme must not embed data that would itself need erasing"); multi-folder deployments prefix keys per folder at construction time (the adapter instance is folder-scoped, like the storage adapter).
Blobs are immutable: putObject on an existing key must not occur (the catalog's finalize-once rule prevents it upstream) and no overwrite operation exists.
deleteObject(context: Context, key: string): Promise<void>;Remove the bytes. Idempotent: deleting a missing key succeeds.
getObject(context: Context, key: string, options?: {
range?: {
offset: number;
length?: number;
};
}): Promise<{
body: ReadableStream<Uint8Array>;
size: number;
contentType: string;
} | null>;The stored bytes with their content type, or null for a missing key. range serves HTTP Range requests (offset + optional length).
headObject(context: Context, key: string): Promise<{
size: number;
contentType: string;
} | null>;The stored object's facts without the body, or null for a missing key.
putObject(context: Context, key: string, body: ReadableStream<Uint8Array> | Uint8Array, options: {
contentType: string;
contentLength?: number;
}): Promise<{
size: number;
sha256: string;
}>;Store the body under key, computing its byte length and lowercase-hex SHA-256 while streaming — the facts finalizeBlob records. When the caller knows the body length up front it passes contentLength so the store can reject an oversized body before buffering it.
DatadataServerLimits
export interface DatadataServerLimitsPer-value byte caps. Folder size has no cap, but it does have a design target: folders of up to a few thousand documents. The folder-enumeration paths — the sys:index / sys:trash projections built on subscribe, plus the access and validity sweeps — grow linearly with the folder. Beyond that target, revisit the design (a folder-level change feed: advisor plan 027) rather than tuning these paths.
maxBlobSizeBytes: number;Maximum blob size in bytes (the immutable-content lane). Enforced by the blob catalog at createBlob (when the transport declares a length) and at finalizeBlob (always) — before/as soon as bytes could be stored, never after a handle is committed.
maxDocumentSizeBytes: number;Maximum document size in bytes (based on UTF-8 encoded byte length of JSON.stringify output)
maxEventSizeBytes: number;Maximum event size in bytes (based on UTF-8 encoded byte length of JSON.stringify output of the patch)
maxYjsDocSizeBytes: number;Maximum Yjs document state size in bytes
DatadataStorageAdapter
export interface DatadataStorageAdapter<Registry extends ValidatorRegistry = ValidatorRegistry, Context = unknown>The synchronous storage contract the server engine drives: documents, both event-log lanes, the folder index, the trash, conformance markers, the reference index and the blob catalog, all in one transactional store. Each method is one transaction.
Context is the per-call context the host passes to the server (for example a logger context with the principal); the adapter may read it for attribution. Network-backed stores implement AsyncDatadataStorageAdapter instead.
createBlob(context: Context, options: {
blobId: string;
contentType: string;
subject: string | null;
actor: string;
contentLength?: number;
}): void;Stage a new blob row at upload start. Throws alreadyExists for a taken blobId, and sizeLimitExceeded when contentLength (the declared body size, when the transport knows it) exceeds maxBlobSizeBytes — the cheapest rejection point, before any bytes are stored.
createDocument<Type extends keyof Registry & string>(context: Context, options: {
docId: string;
type: Type;
data: TypedDocumentData<Registry, Type>;
patch: Operation[];
yjsUpdates: YjsUpdateRaw[];
yjsDocStates?: Map<string, Uint8Array>;
generation: string;
conformedSchemaSequence: number;
schemaSequence: number;
name?: string | null;
references?: readonly ReferenceEdge[];
blobReferences?: readonly string[];
clientEventId?: string | null;
}): {
indexSequence: number;
};Insert a new document at sequence 1 with its genesis event, Y.Doc states, folder membership (a 'create' folder event), reference edges, blob reference edges and idempotency record, all in one transaction. Throws alreadyExists for a live docId and deleted for a soft-deleted one. Returns the folder's new index sequence.
deleteBlobRecord(context: Context, ref: {
blobId: string;
}): void;Remove a tombstoned blob's catalog row — the sweep's last step, after the object-store bytes are gone. Throws if the row exists and is NOT a tombstone (the sweep must CAS first); no-op for a missing row.
deleteDocument(context: Context, ref: {
docId: string;
clientEventId?: string | null;
}): {
indexSequence: number;
trashSequence: number;
deletedAt: number;
name: string | null;
};Soft-delete a document: tombstone it (deleted_at) so it leaves every read path (getDocument, getDocumentIndex, getDocumentEvents) while its row, event log and Y.Doc states are retained for restore. Appends a 'delete' folder event — a lifecycle operation lives in the folder log, never the document's own event log, so the document's sequence is untouched. Throws storageError if the document does not exist live (the server resolves already-deleted / never-existed before calling). Returns the folder's new index sequence, the trash projection's new sequence (delete/restore/purge are the only ops that advance it — see getDeletedDocuments), the tombstone timestamp, and the document's name (index metadata the stored document doesn't carry — the server puts it on the sys:trash entry without enumerating the folder).
deleteYjsDoc(context: Context, ref: {
docId: string;
yjsId: string;
}): void;Delete a specific Y.Doc from storage. Called when a field with a YjsRef is removed from the document.
Parameters
contextStorage context
refReference to the Y.Doc to delete (docId and yjsId)
finalizeBlob(context: Context, ref: {
blobId: string;
size: number;
sha256: string;
}): void;Record the stored bytes' facts after the object-store put completes. Only a finalized blob may be referenced by a document write. Throws notFound for a missing/tombstoned row, alreadyExists if already finalized (blobs are immutable — there is no overwrite), and sizeLimitExceeded when size exceeds maxBlobSizeBytes.
findProcessedWrite(context: Context, ref: {
docId: string;
clientEventId: string;
}): {
sequence: number;
} | null;Idempotency for the JSON-lane write path (WritesAreIdempotentByEventId). A client that replays an unacked doc:update on reconnect re-sends it under its original event id; the server recognizes the duplicate here and re-acks at the recorded sequence instead of applying the patch a second time (which would double-apply a non-idempotent op like an array append).
Returns the result sequence recorded for this (docId, clientEventId), or null if it has not been applied within the adapter's retention window. The key is scoped by docId because one storage backend serves a whole folder. Each recording write — updateDocument, createDocument, deleteDocument, restoreDocument, renameDocument — records it in the SAME transaction as its own durable effect (via its clientEventId option), so a crash can't leave one durable without the other; a write with no durable effect is recorded by recordNoOpWrite. The recorded number is the write's result sequence (document sequence for updates/creates, index sequence for the lifecycle ops); only the update and rename re-acks carry it back to the client — for the others it is an existence marker. The adapter prunes entries older than its retention window so the store stays bounded (dedup is only needed for the brief replay window, never permanently).
findSweepableBlobs(context: Context, options: {
stagedBefore: number;
unreferencedBefore: number;
limit: number;
}): SweepableBlob[];Enumerate GC candidates: staged rows created before stagedBefore (abandoned uploads), live rows unreferenced since before unreferencedBefore, and tombstones left by interrupted sweeps. Cutoffs are caller-computed absolute timestamps — retention windows are application policy, like purgeDocuments' deletedBefore. Ordered by blobId; limit bounds one call.
getBlobMetadata(context: Context, ref: {
blobId: string;
}): StoredBlobMetadata | null;The catalog row, or null for a missing or tombstoned blob.
getBlobReferencingDocIds(context: Context, ref: {
blobId: string;
}): string[];Every non-purged document (trash included) holding a reference edge to the blob — the read-authorization walk's input: a blob is readable iff the principal can read at least one of these.
getBlobSweepHorizon(context: Context): BlobSweepHorizon;The catalog's sweep horizon (see BlobSweepHorizon): the facts the sweep needs to say when a sweep is next due, including rows that are not sweepable yet. Takes no cutoffs, but keeps the sweep condition's state and reference-edge predicates exactly.
getDanglingReferenceDocuments(context: Context, opts?: {
type?: string;
}): {
docId: string;
type: string;
}[];The DERIVED half of invalid-document enumeration: live documents holding at least one reference edge whose target is missing, tombstoned, or of the wrong type — a pure function of (reference index, live documents), so it reacts instantly to deletes and restores of targets with nothing to heal. Deduped by source docId, optionally filtered by the SOURCE's type, ordered by docId. Sources are live documents only (a tombstoned referrer is out of every read path, its dangling refs included). Completeness follows the index: a document last written before the index existed has no edges until a validated write or a validation sweep backfills them, and two rarer gaps share the same sweep-closes-it story — a schema change that adds reference rules without changing a document's data indexes that document only via a sweep once its cursor is current, and a schema change that REMOVES rules leaves stale edges (possible false dangling entries) behind until the sweep's unconditional refresh clears them.
getDeletedDocument(context: Context, ref: {
docId: string;
}): DeletedDocumentEntry | null;Point lookup of a soft-deleted document (null when the doc is live or never existed). The cheap "is this id tombstoned?" check the server uses to make deletion observable — answering subscribes with doc:deleted instead of doc:notfound, and rejecting create/update with the deleted category.
getDeletedDocuments(context: Context, opts?: {
type?: string;
}): {
documents: DeletedDocumentEntry[];
trashSequence: number;
};Enumerate the folder's soft-deleted documents (the trash listing), optionally filtered by docType (mirroring getInvalidDocuments), with the trash projection's own sequence. That counter advances ONLY on delete/restore/purge — unlike the index sequence, which every folder op bumps — so sys:trash subscribers see contiguous patch sequences no matter how many creates/renames interleave (a shared counter would trip the client's per-doc gap guard into a needless resync). sys:trash always lists everything; the type filter is for direct adapter consumers (e.g. an admin trash tool).
getDocument(context: Context, ref: {
docId: string;
}): StoredDocument | null;The live document, or null when it is missing or soft-deleted.
getDocumentBlobIds(context: Context, ref: {
docId: string;
}): string[];The document's stored blob reference edges (export/import and GC support).
getDocumentEvents(context: Context, ref: {
docId: string;
afterSequence?: number;
limit?: number;
}): StoredDocumentEvent[] | null;The JSON lane's log in sequence order, or null when the document does not exist live. Paged by the optional range: only events with sequence > afterSequence, at most limit of them (doc:get-events reads one row past its page to learn whether more remain). Absent = the whole log, which replay, verify and export read.
getDocumentIndex(context: Context): {
documents: Array<{
docId: string;
type: string;
name: string | null;
}>;
indexSequence: number;
};The folder's live documents (id, type and name) and its current index sequence: the content the synthesized sys:index document carries.
getDocuments(context: Context, refs: {
docIds: readonly string[];
}): (StoredDocument | null)[];Batch form of DatadataStorageAdapter.getDocument: answer refs.docIds POSITIONALLY — result[i] is the document for docIds[i], null for a missing/tombstoned id, duplicates answered per position. MUST be observably equivalent to mapping getDocument over the ids (same visibility, same rows); what an implementation may change is the number of storage round trips — the whole point: a SQL adapter answers the batch in one query where N getDocument hops would each pay a round trip. The engine calls this where one operation fans out over independent ids (the cross-reference existence checks of a validated write).
getDocumentsBehindSchema(context: Context, ref: {
type: string;
schemaSequence: number;
limit: number;
}): string[];DocIds of live documents of type whose data may not conform to the type's current schema (the governing schema doc's sequence, resolved by the server and passed in): the conformed cursor is behind AND the doc is not already flagged at that same sequence. A doc flagged AT the current sequence is excluded on purpose — its shape verdict is already recorded, so revalidating it can learn nothing new (reference healing is the full sweep's job) — and the exclusion is what makes the incremental sweep's fetch-process loop terminate: processing a returned doc either advances its cursor to the current sequence or flags it there, dropping it from this query either way. Ordered by docId; limit bounds one call.
getInvalidDocuments(context: Context, opts?: {
type?: string;
}): InvalidDocument[];Enumerate documents flagged invalid (invalid_schema_sequence IS NOT NULL), optionally filtered by docType. Cheap; surfaces invalidity discovered lazily on read. A lower bound, not a complete answer: a cold document that became invalid but was never read since carries no flag and is not listed (see DatadataServer.findInvalidDocuments for the full caveat).
getReferencingDocuments(context: Context, ref: {
docId: string;
}): {
docId: string;
type: string;
}[];Live documents holding a reference edge to docId — the inbound direction of the reference index for ONE target, whether the edge currently resolves or dangles. The server reads it when the target's existence changes (create, delete, restore) to refresh the referrers' subscribers, whose derived reference validity just flipped.
getYjsEvents(context: Context, ref: {
docId: string;
afterYjsSequence?: number;
limit?: number;
}): StoredYjsEvent[] | null;The Yjs lane's log in yjsSequence order (delete tombstones included), or null when the document does not exist live. The JSON-lane twin is getDocumentEvents; the two logs share nothing but the document. Paged the same way, by yjsSequence > afterYjsSequence and limit.
hasPurgedGeneration(context: Context, ref: {
docId: string;
generation: string;
}): boolean;Whether generation is a PURGED generation of docId — one a purge released the id from. The create gate's read: a create naming a purged generation (the replay of the create that minted it, whose processed write the purge cleared) is refused, so a generation is never live twice. Answered from the record purge keeps per generation.
importDocument(context: Context, options: {
docId: string;
type: string;
data: Record<string, unknown>;
yjsDocStates: Map<string, Uint8Array>;
events: StoredDocumentEvent[];
yjsEvents: ExportedYjsEvent[];
generation: string;
conformedSchemaSequence: number;
blobReferences?: readonly string[];
}): {
indexSequence: number;
};Insert a document with its FULL history into a fresh store: the documents row (at the latest event's sequence and the latest yjs event's yjsSequence), every JSON-lane document_events row, every Yjs-lane yjs_events row, the derived yjs_docs snapshot, and a folder-membership record (so the synthesized sys:index includes it). The replay-derived snapshot (data + yjsDocStates) and conformedSchemaSequence are computed by the caller (importDocument); the adapter only writes. Throws alreadyExists if the document is already present — import targets a store that does not yet have it. Returns the folder's new index sequence.
overwriteDocumentSnapshot(context: Context, ref: {
docId: string;
}, snapshot: {
data: Record<string, unknown>;
yjsDocStates: Map<string, Uint8Array>;
}): void;Overwrite a document's derived snapshot (data + full Y.Doc states) in place, WITHOUT appending to the event log or changing the sequence. Used by rebuildDocument to restore the snapshot from a replay of the log (the source of truth). Replaces the document's Y.Doc set wholesale with exactly the supplied states. Throws if the document does not exist.
purgeDocuments(context: Context, ref: {
docIds: readonly string[];
}): {
trashSequence: number;
purgedAt: number;
};Hard-delete (purge) a batch of soft-deleted documents' content and RELEASE their ids. For each document: destroys both lanes' event logs and the Y.Doc states and blanks the snapshot; the per-generation record the adapter keeps (as a purge audit anchor, and for DatadataStorageAdapter.hasPurgedGeneration) must be unreachable by any docId read — after a purge the id behaves exactly like one that never existed, so a create under it starts a fresh generation. Deliberately appends NO folder event: folder-event sequences are the sys:index patch stream and a purged document is already out of the index, so a purge advances ONLY the trash counter (purge is the third trash-projection mutation alongside delete/restore) — and it advances ONCE for the whole batch, so N removals reach every trash subscriber as one patch at one sequence. All-or-nothing: throws storageError if ANY document in the batch is not currently soft-deleted (the server resolves live / never-existed — which includes already-purged — before calling), and a throwing batch must destroy nothing.
recordNoOpWrite(context: Context, ref: {
docId: string;
clientEventId: string;
sequence: number;
}): void;Remember a client write that turned out to be a no-op — a doc:update whose patch left the document as it was, so nothing was stored for it — so that findProcessedWrite knows it and a replay is re-acked instead of applied to whatever the document has become since (WritesAreIdempotentByEventId). sequence is the document sequence the write was answered at. An id already recorded keeps its first record. Pruned with the rest. Called only on that path, where it is the one storage write: an update that does store something records its id inside updateDocument.
renameDocument(context: Context, ref: {
docId: string;
name: string | null;
clientEventId?: string | null;
}): {
indexSequence: number;
};Rename a document: update the materialised name on the document row and append a 'rename' folder event carrying the new name. A name is index-level metadata, not content, so this never touches the document's data, sequence or event log. null clears the name. Throws storageError if the document does not exist live (the server resolves never-existed / deleted before calling). Returns the folder's new index sequence.
restoreDocument(context: Context, ref: {
docId: string;
clientEventId?: string | null;
}): {
indexSequence: number;
trashSequence: number;
};Undo a soft delete: clear the tombstone so the document re-enters the read path at the sequence it was deleted at, and append a 'restore' folder event. Throws storageError if the document is not currently soft-deleted. Returns the folder's new index sequence and the trash projection's new sequence.
setConformedSchemaSequence(context: Context, ref: {
docId: string;
}, conformedSchemaSequence: number): void;Advance a document's stored conformed_schema_sequence — the conformed schema-doc sequence. Used by the lazy read path after bringing a document forward, including when bring-forward left the data unchanged, so the cursor still moves forward and the doc drops out of the "needs migration" query.
setDocumentBlobReferences(context: Context, ref: {
docId: string;
}, blobIds: readonly string[]): void;Replace a document's stored blob reference-edge set — the blob twin of DatadataStorageAdapter.setDocumentReferences, derived from the validated data's schema membership (collectBlobRefs). The WRITE paths never call this: they carry the edge set into createDocument/updateDocument/importDocument's blobReferences so content and membership commit in one transaction. This standalone form is the validation sweep's refresh/backfill lever. Referencing a missing, unfinalized, or tombstoned blob throws preconditionFailed — the transactional half of "a synced document never points at missing bytes" (a sweep that tombstoned the blob first wins, and this write loses benignly). Newly referenced blobs flip staged → live; a blob losing its last edge records unreferencedAt and starts the GC grace clock.
setDocumentReferences(context: Context, ref: {
docId: string;
}, references: readonly ReferenceEdge[]): void;Replace a live document's stored outbound reference-edge set outside a content write — the validation sweep's refresh/backfill lever (organic reads never touch edges; writes carry them transactionally via createDocument/updateDocument's references). No-op for a missing or tombstoned doc.
setInvalidSchemaSequence(context: Context, ref: {
docId: string;
}, invalidSchemaSequence: number | null): void;Set (or clear, with null) a document's invalid_schema_sequence — the lazily discovered validity marker the read path writes when a doc fails / passes validation against the current schema.
tombstoneBlob(context: Context, ref: {
blobId: string;
}, cutoffs: {
stagedBefore: number;
unreferencedBefore: number;
}): boolean;The sweep's compare-and-set: re-verify the sweep condition against the SAME cutoffs and, only if it still holds, transition the row to its terminal tombstone state (invisible to getBlobMetadata and to setDocumentBlobReferences, which now rejects the handle). Returns false — and destroys nothing — when the world moved: a write committed a reference, or a finalize refreshed a staged row past the cutoff. This CAS, not the grace period, is what makes the sweep race-free against concurrent writes; the sweep deletes bytes only after it returns true. Idempotently true for an existing tombstone (byte deletion retries).
updateDocument(context: Context, options: {
event: DocumentPatchEvent;
data: Record<string, unknown>;
yjsDocStates?: Map<string, Uint8Array>;
clientEventId?: string;
schemaSequence: number;
conformance?: DocumentConformance;
references?: readonly ReferenceEdge[];
blobReferences?: readonly string[];
}): void;Apply one committed write: store the new data (and Y.Doc states, when given) and append event to the JSON lane's log, together with the optional idempotency record, conformance markers and reference edges, all in one transaction.
DocumentConformance
export interface DocumentConformanceA write's verdict on the data it is storing — the pair of conformance markers, always set together (see updateDocument's conformance).
The two cursors answer different questions, which is why a document can carry both at the same sequence:
- conformedSchemaSequence — how far the STORED data has been brought forward. It is the bring-forward's filter, so it must move whenever migrated data is persisted, valid or not: a remap is not idempotent, and replaying one over data that already carries it transforms it twice. - invalidSchemaSequence — whether that data CONFORMS at the sequence it was carried to, and null when it does.
So { S, null } is the ordinary validated write (and the conforming bring-forward), while { S, S } is the relaxed read's persist: carried to S, does not conform there.
conformedSchemaSequence: number;The schema sequence the stored data has been brought forward to.
invalidSchemaSequence: number | null;The schema sequence the data fails to conform at, or null when it conforms.
DocumentExport
export interface DocumentExportA portable, store-agnostic snapshot of a single document: its type and both lanes' full logs — the JSON-patch events (in sequence order) and the Yjs events (in log order). Importing it into a fresh store replays the logs to reproduce the document — logically equal, with its history intact. The reconstruction unit for replica/DR and portability (core.allium PerDocumentReplay).
A LIVE value, not a wire format: Yjs blobs are Uint8Array like everywhere else, so it crosses structured-clone boundaries (DO RPC, postMessage) as-is but must NOT be fed to JSON.stringify (a Uint8Array silently mangles into an index-keyed object). The blessed plain-text form for files and HTTP is serializeDocumentExport / parseDocumentExport, which spell each blob as base64.
The Yjs events carry no yjsSequence: array order IS the order, and order is the only thing the log's numbering ever meant (replay folds by order, and only across a delete/recreate boundary does even that matter). Nor does the export carry the source document's yjsSequence: that number versions the lane WITHIN one store and is only ever compared against later values of itself, so an importing store renumbers from its own log. Nothing in the stream is a number the importer must honour, which is what makes it portable.
Portable, NOT canonical: a store whose write-behind coalesces a burst logs one merged Yjs event where a store that does not logs several, so the same document exported from two adapters can differ in row count and in blob bytes. Both replay to the same Y.Doc (CRDT merge is associative), and verifyReplay compares logical state rather than the log. Do not diff two exports for equality; import them and compare the documents.
blobs: ExportedBlobContent[];The blobs the document's current data references (schema-derived membership), carried with their bytes so the export is self-contained and an import can re-upload them into the target's object store. Bytes are Uint8Array in the live value (same rule as the Yjs blobs above; base64 only in the serialized text form). Empty for a document with no blob references — and always empty from a server without an object storage adapter. Blobs above MAX_INLINE_EXPORT_BLOB_BYTES make the export fail rather than silently dropping content; a streaming export form is future work.
docId: string;The exported document's id.
events: StoredDocumentEvent[];The JSON lane's full event log, in sequence order.
invalidSchemaSequence?: number;The schema sequence at which this document's data was last found NOT to conform, or absent when it conformed (or was never checked).
Carried because a document can legitimately be exported while invalid, and the log alone cannot say so: the relaxed read persists its bring-forward, so the last event's schemaSequence stamp is a sequence at which the snapshot does not validate. Without this field an import could not tell that from an export produced under a schema this folder never held (the cross-folder overlap validateImportedSnapshot exists to catch), and would have to reject both. Nothing is given up by trusting it: non-conforming data validates at NO sequence, so there is no verdict the importer could have reached on its own.
type: string;The exported document's type.
yjsEvents: ExportedYjsEvent[];The Yjs lane's full event log, in log order.
DocumentSizeGuards
export interface DocumentSizeGuardsByte-length pre-checks shared by every storage adapter. Adapters call these BEFORE writing so an oversized payload is rejected with a well-formed DocumentRetrievalError (sizeLimitExceeded) rather than surfacing as an opaque backend error — or, on backends without a size backstop, silently succeeding. The resolved DatadataServerLimits are exposed on limits for adapter-specific guards (e.g. Yjs event deltas) that reuse the same caps.
assertBlobSizeWithinLimits(blobId: string, byteLength: number): void;Throw sizeLimitExceeded when a blob's byte length exceeds maxBlobSizeBytes.
assertDocumentSizeWithinLimits(docId: string, jsonString: string): void;Throw sizeLimitExceeded when the document's JSON exceeds maxDocumentSizeBytes (UTF-8).
assertEventSizeWithinLimits(docId: string, eventJson: string): void;Throw sizeLimitExceeded when the serialized event exceeds maxEventSizeBytes (UTF-8).
assertYjsDocSizeWithinLimits(docId: string, yjsId: string, state: Uint8Array): void;Throw sizeLimitExceeded when the encoded Y.Doc state exceeds maxYjsDocSizeBytes.
readonly limits: DatadataServerLimits;The resolved limits the guards check against.
EventBusRecipient
export interface EventBusRecipientOne live connection as the bus knows it — the unit of per-recipient delivery.
connectionId: string;The connection's id, as the transport assigned it.
principal: Principal;The connection's identity, from the transport's connection state (WebSocket attachment metadata / in-process registration) — the same principal the connection's handleEvent contexts carry.
ExportedBlobContent
export interface ExportedBlobContentOne exported blob: its catalog facts plus the full bytes.
blobId: string;The blob's id.
bytes: Uint8Array;The blob's full content.
contentType: string;The blob's MIME type.
sha256: string;Lowercase-hex SHA-256 of the bytes.
size: number;The byte length.
FolderRegistryEntry
export interface FolderRegistryEntryWhat the registry stores per folder: whatever the host's init produced — typically the server plus its transport-side companions (event bus, loggers) — as long as the server is reachable for coordinated shutdown.
server: FolderRegistryServer;The folder's server, closed by eviction and by the registry's shutdown.
FolderRegistryServer
export interface FolderRegistryServerThe slice of AsyncDatadataServer the registry manages. Structural, so a host's folder entry can carry any concrete server instantiation (the registry never dispatches operations — it only ends lives).
close(): Promise<void>;Drain the server's queued operations and release its storage adapter; idempotent.
InvalidDocument
export interface InvalidDocumentA document invalid against the current schema, with the cause attributed: schema — its shape failed validation (a lazily-persisted flag, recorded at invalidSchemaSequence); reference — a cross-document reference dangles (DERIVED from the reference index against the live document set at enumeration time — never stored, so it cannot go stale). A document broken both ways appears once per kind.
docId: string;The invalid document's id.
invalidSchemaSequence?: number;Present for kind: "schema": the schema sequence the failure was recorded at.
kind: "schema" | "reference";The cause: a failed shape validation (schema) or a dangling reference (reference).
type: string;The invalid document's type.
Logger
export interface LoggerThe structured logger datadata writes to, in the pino calling convention: each method takes either a message string, or an object of fields followed by an optional message. A pino logger satisfies it as is; see createBrowserConsoleLogger for a console-backed one.
child: (bindings: object) => Logger;Return a logger that adds bindings to every entry it writes.
debug: (obj: object | string, msg?: string) => void;Log at debug level.
error: (obj: object | string, msg?: string) => void;Log at error level.
info: (obj: object | string, msg?: string) => void;Log at info level.
warn: (obj: object | string, msg?: string) => void;Log at warn level.
LoggerContext
export interface LoggerContextA context that carries the logger for the call it belongs to.
logger: Logger;The logger for this call.
ObservabilityOptions
export interface ObservabilityOptionsOptions for observability on server operations
onObservability?: ObservabilityCallback;Receives the operation's observability event (success or failure); absent = none emitted.
ReferenceEdge
export interface ReferenceEdgeOne outbound cross-document reference edge, extracted from a document's VALIDATED data (collectReferenceEdges): this document's data names targetDocId, optionally constrained to a requiredType (the rule's toType). The rows of the reference index — refreshed as a set on every JSON-validated write and by the validation sweep, and read in the inbound direction ("whose edges dangle?") by getDanglingReferenceDocuments.
requiredType?: string;The type the target must have; absent when the reference rule accepts any type.
targetDocId: string;The id of the document the data names.
SchemaChangePreview
export interface SchemaChangePreviewWhat committing a candidate schema would do to its type's live documents — see DatadataServer.previewSchemaChange. Each document checked lands in at most one of stranded, stillInvalid and healed, by comparing its shape verdict under the current schema with its verdict under the candidate; danglingReferences lists documents whose shape passes under the candidate but whose references do not.
candidateSequence: number;The schema sequence the commit would produce: the stored schema document's sequence + 1.
checked: number;Live documents of the type this invocation checked.
danglingReferences: SchemaChangePreviewEntry[];Shape-valid under the candidate, but a reference dangles. Reference invalidity is derived, never flagged (see findInvalidDocuments), and mostly independent of the schema change, so it is reported apart. Not a delta: a reference that already dangles under the current schema is listed too, so this is not "what the change would break".
healed: string[];Invalid under the current schema, valid under the candidate.
missingBlobs: SchemaChangePreviewMissingBlobs[];Documents whose data, brought forward under the candidate, holds blobRef handles naming no finalized blob — say a string field retyped to blobRef over values that were never uploaded. The commit would not flag these: registering their blob references fails, so a read or a sweep errors on them instead. Handles are collected from the validated value, so a schema default that backfills a handle counts. Only documents whose shape passes under the candidate are checked: a shape-invalid one is flagged without its blob references being replaced, so it never lands here.
nextAfterDocId: string | null;Resume cursor: pass back as afterDocId to continue; null once every live document of the type has been checked.
stillInvalid: SchemaChangePreviewEntry[];Invalid under the current schema and under the candidate.
stranded: SchemaChangePreviewEntry[];Valid under the current schema, invalid under the candidate: what the change would strand.
SchemaChangePreviewEntry
export interface SchemaChangePreviewEntryOne document's verdict under the candidate schema — see SchemaChangePreview.
docId: string;The document's id.
invalid: DocumentInvalid;Why it fails under the candidate: the marker a read would deliver after the commit.
SchemaChangePreviewMissingBlobs
export interface SchemaChangePreviewMissingBlobsA document whose blobRef handles would not resolve — see SchemaChangePreview.
blobIds: string[];The handles naming no finalized blob, sorted.
docId: string;The document's id.
ServerInitiated
export interface ServerInitiatedWhat stands where a connectionId would on the context of a call the server itself started — a seed at init, a sweep, a host's direct write. It names what started the call ("do-init", "blob-sweep"), for logs and observability, and being an object rather than a string it cannot be mistaken for a socket: nothing is acked to it, subscribed for it or routed by it.
A transport's context type requires a connectionId, and a host used to fill it with a placeholder string for such calls. The engine then took the placeholder for a caller — a sys:access write acked it, and the bus logged an error for the socket it could not find.
serverInitiated: string;What started the call, for example "do-init" or "blob-sweep".
StoredBlobMetadata
export interface StoredBlobMetadataOne blob catalog row: the platform facts about an immutable blob whose bytes live in the object store (DatadataObjectStorageAdapter) under the blob id. Display metadata (alt text, dimensions, crop) is the app's business and lives in the referencing document's own JSON; the catalog owns only what the engine needs for lifecycle, integrity, and audit.
Lifecycle: a row is created staged at upload start (createBlob), completed by finalizeBlob (size/sha256 recorded — only a finalized blob may be referenced by a write), flips to live when a committed write references it, and is destroyed by the GC sweep (tombstone → bytes → record; see tombstoneBlob). Tombstoned rows are invisible here — getBlobMetadata answers null — matching "a deleted blob is indistinguishable from one that never existed".
actor: string;The uploader's actor label (audit attribution, like stored events).
blobId: string;The blob's id, which is also its key in the object store.
contentType: string;The MIME type declared at upload start.
createdAt: number;When the upload was staged, in epoch milliseconds.
finalizedAt: number | null;When the upload was finalized, in epoch milliseconds; null while it is in flight.
sha256: string | null;Lowercase-hex SHA-256 of the bytes, recorded at finalize; null while in flight.
size: number | null;Byte length, recorded at finalize; null while the upload is in flight.
state: "staged" | "live";staged — uploaded (possibly not yet finalized), no committed reference yet; readable only by its uploader. live — referenced by ≥1 non-purged document (documents in trash count: delete→restore must round-trip).
subject: string | null;The uploader's identity (Principal.subject), for staged-blob read authorization and audit. null = anonymous — never matches any reader, so an anonymous upload is readable only once a document references it.
unreferencedAt: number | null;When the blob's last reference edge was removed (a live blob with no remaining referencing documents); null while referenced or never referenced. The GC grace period counts from here.
StoredDocument
export interface StoredDocument<T = unknown> extends DatadataDocument<T>A document as held in storage: the public DatadataDocument plus server-internal tracking.
conformedSchemaSequence is the schema sequence (the governing sys:schema:<type> document's sequence) this document's data has been carried to: the bring-forward cursor the lazy read path advances, which every stored document has and a client's copy of a non-sys: one is told. invalidSchemaSequence stays server-side; client-bound events carry the invalid marker instead.
conformedSchemaSequence: number;Every stored document has one — see DatadataDocument.conformedSchemaSequence.
generation: string;Every stored document has one — see DatadataDocument.generation.
invalidSchemaSequence: number | null;The schema sequence at which a read last found the data non-conforming, or null when it was last seen valid or never checked.
yjsSequence: number;The Yjs lane's version, owned by STORAGE and never surfaced by any read API — not the wire, not getDocument. It is the length of this document's Yjs log, equivalently its latest event's position.
Clients need no Yjs version: the lane is a CRDT, and where ordering does matter — a sub-document's deletion — sequence supplies it, because a Y.Doc lives exactly as long as a yjsRef names it and the orphan reap rides the JSON-changing write. Shipping a Yjs version forced the ENGINE to allocate it at broadcast, before write-behind had persisted anything, which is precisely what drove the lane's version apart from its log's positions. Storage owning it collapses them back into one number.
SweepableBlob
export interface SweepableBlobOne GC-sweep candidate, as enumerated by findSweepableBlobs. The reason names which sweep condition matched — the sweep re-verifies it under tombstoneBlob's compare-and-set before destroying anything, so a candidate is a hint, never a verdict.
blobId: string;The candidate blob's id.
reason: "staleStaged" | "unreferenced" | "tombstone";staleStaged — staged past the caller's stagedBefore cutoff (an abandoned upload); unreferenced — live with zero reference edges since before unreferencedBefore; tombstone — already tombstoned by an earlier sweep whose byte deletion did not complete (retry).
ValidateDocumentsResult
export interface ValidateDocumentsResultOne validation sweep invocation's outcome — see DatadataServer.validateDocuments.
checked: number;Documents this invocation validated (successfully read).
errors: {
docId: string;
message: string;
}[];Documents whose validation THREW (storage/corruption errors — not invalidity, which lands in flagged). They stay candidates and keep remaining true on every run until fixed: the sweep never silently writes a failing document off, because "migration done" must not be claimable past a document nobody could check.
flagged: InvalidDocument[];Documents invalid after this pass (newly discovered or still failing).
migrated: number;Of those, how many were persisted/advanced by lazy bring-forward.
nextAfterDocId: string | null;Full mode's iteration cursor: pass back as afterDocId to continue the pass; null means the pass covered every live document. Always null in incremental mode (its candidate set shrinks instead of paginating).
remaining: boolean;Whether candidates remain beyond this invocation — loop while true.
VerifyResult
export interface VerifyResultThe outcome of DatadataServer.verifyDocument: whether a document's event log reproduces its stored snapshot.
dataMatches: boolean;Whether replaying the JSON lane reproduces the stored data (deep equality).
ok: boolean;true when the log reproduces the snapshot exactly: dataMatches and no Yjs drift.
yjsDrift: YjsDrift[];Each Yjs sub-document whose replayed state differs from the snapshot; empty when none do.
YjsDrift
export interface YjsDriftA single Yjs sub-document whose replayed state diverges from the snapshot.
kind: "missing" | "extra" | "mismatch";Drift of the snapshot relative to the log (the source of truth): "missing" — produced by the log but absent from the snapshot; "extra" — present in the snapshot but not produced by the log; "mismatch" — present in both but not logically equal.
yjsId: string;The id of the Yjs sub-document that drifted.
YjsUpdateRaw
export interface YjsUpdateRawOne Yjs update as the client sent it, passed to DatadataStorageAdapter.createDocument for the event log and audit trail.
action: "create" | "update" | "delete";Whether the update creates the Y.Doc, updates it, or deletes it.
update?: Uint8Array;The encoded Yjs update; absent for a delete.
yjsId: string;The id of the Y.Doc (the yjsRef field's id) the update applies to.
Types
AsyncDatadataServerFor
export type AsyncDatadataServerFor<S extends SchemaRegistry, EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry> = AsyncDatadataServer<RegistryFor<S>, EventBusContext, StorageContext, Commands>;The async-server type for a schema map — the schema-facing name for AsyncDatadataServer<RegistryFor<S>, …>, mirroring DatadataServerFor.
AsyncDatadataStorageAdapter
export type AsyncDatadataStorageAdapter<Registry extends ValidatorRegistry = ValidatorRegistry, Context = unknown> = {
[M in keyof DatadataStorageAdapter<Registry, Context>]: DatadataStorageAdapter<Registry, Context>[M] extends (...args: infer P) => infer R ? (...args: P) => Awaitable<R> : never;
} & {
close?(): Awaitable<void>;
};The async-adapter surface, derived from DatadataStorageAdapter: the same methods, each result optionally promised. A Postgres (or any network-backed) adapter implements this and is driven by AsyncDatadataServer (the engine awaits every hop via runAsync); a synchronous adapter satisfies it as-is, so the async server runs over the in-memory and DO-SQLite adapters unchanged. Derived by mapped type, so the two surfaces cannot drift. (The derivation instantiates createDocument's per-call Type parameter at its constraint — irrelevant to implementors, and the engine dispatches effects untyped.)
Awaitable
export type Awaitable<T> = T | PromiseLike<T>;A value the async shell may await: the result, or a promise of it.
BlobCachePolicy
export type BlobCachePolicy = "revalidate" | {
mode: "immutable";
maxAge: number;
};How a granted blob response may be cached — the host's statement of how it wired the read seam, because the two are coupled and only the host knows:
- "revalidate": authorization is evaluated per REQUEST (cookies on the same URL, the authorizeBlobRead gate on every hit). Even the private browser cache must ask again before reusing a response — a revoked reader would otherwise keep the bytes — so success is private, no-cache. The blob is immutable with a strong sha256 ETag, so revalidation is a byte-free 304; only the authz check re-runs. - { mode: "immutable", maxAge }: the URL carries its own authorization (a signed URL) or none (a secret or public URL whose bare random blob id is the capability). Nothing about the response depends on who asked, so success is public, max-age=<maxAge>, immutable — cacheable by browsers and shared caches alike for maxAge seconds.
Denials are private, no-store under either policy: a cached 404 must not outlive a later granted read (a staged blob goes live, a document is restored). The security headers are pinned regardless of policy.
BlobFinalizeOutcome
export type BlobFinalizeOutcome = {
ok: true;
} | BlobRequestFailure;Finalize outcome, same contract as BlobStageOutcome.
BlobStageOutcome
export type BlobStageOutcome = {
ok: true;
blobId: string;
} | BlobRequestFailure;Staging outcome carried over a host's RPC boundary (e.g. Worker → Durable Object). Expected failures ride as data (ok: false + the error's category) because a thrown DocumentRetrievalError loses its category over RPC serialization, and the upload handler maps categories to HTTP statuses.
BlobSweepCandidateCallback
export type BlobSweepCandidateCallback = () => void;A host's hook on "the blob catalog may hold a new GC candidate" — the moment a retention clock can start: - a blob row was staged, by an upload or by an import re-uploading an export's blobs (a staged row ages out through the TTL sweep whether or not its upload finishes); - a document's blob edges were replaced, which may have dropped a blob's last reference: a validated write of a blobRef-declaring type, a lazy read or a validation sweep refreshing the edges after a schema change; - documents were purged, taking their edges with them.
Deliberately over-approximate — a replacement that kept every handle fires too — because hosts only arm a schedule from it, and arming is idempotent. Nothing about the blob is passed, since the schedule only needs to know that something may need sweeping.
ConnectionIdOrServerInitiated
export type ConnectionIdOrServerInitiated = string | ServerInitiated;A real connection's id, or the marker of a call no connection made.
CreateObservabilityEvent
export type CreateObservabilityEvent = ObservabilityEventBase<"create", {
docId: string;
type: string;
patchOperationCount: number;
yjsUpdateCount: number;
}>;Create event - fired when a document is created
DatadataServerFor
export type DatadataServerFor<S extends SchemaRegistry, EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown, Commands extends CommandRegistry = CommandRegistry> = DatadataServer<RegistryFor<S>, EventBusContext, StorageContext, Commands>;The server type for a schema map — the schema-facing name for DatadataServer<RegistryFor<S>, …>, so callers parameterise it by their schemas rather than by a validator registry.
DeleteObservabilityEvent
export type DeleteObservabilityEvent = ObservabilityEventBase<"delete", {
docId: string;
type: string;
}>;Delete event - fired when a document is soft-deleted
DocumentChangeCallback
export type DocumentChangeCallback<EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown> = (context: DocumentChangeCallbackContext<EventBusContext, StorageContext>, docId: string, document: ReadonlyDatadataDocument | null, previousDocument: ReadonlyDatadataDocument | null, event: DocumentInitEvent | DocumentPatchEvent | DocumentDeletedEvent) => void | Promise<void>;document/previousDocument are not guaranteed to be fresh clones from a storage read — a write path may hand the callback the very object it just built (or a document it is still using for its own observability/broadcast bookkeeping) to avoid a redundant migrating read. A mutation reaching either value (e.g. via a cast past the read-only type) would corrupt the write's own in-flight state and be visible to every OTHER callback in the same fan-out (callbacks run synchronously, in registration order, over the same objects) — copy first if a transformation is needed.
On a delete (event.type === "doc:deleted") document is null and previousDocument is the document AS LAST STORED — the delete path reads the raw stored state rather than the migrating read, so it is NOT brought forward to the current schema (a bring-forward would write to a document that is about to be tombstoned). An idempotent re-delete of an already tombstoned document changes nothing and does not fire.
Callbacks report accepted document writes — creates, updates and deletes — not the lazy schema bring-forward. A read that migrates a document persists the brought-forward data at a new sequence, but fires no document-change callback: the change is schema-driven and surfaces on a read (any getDocument, subscribe or validation pass, or the sweep a schema commit runs over the type's live-subscribed documents, which broadcasts the migration patch to their subscribers — see sweepSchemaMigrations), and the schema write that caused it already fired for sys:schema:<type>. A consumer that mirrors document data takes that schema change as its cue to re-read the type's documents (getDocument always answers the current shape); the next real write to a migrated document reports the migrated state as its previousDocument.
Callbacks fire for stored documents only — never for the synthesized folder listings (sys:index / sys:trash). A rename changes listing metadata only, so it fires nothing; a purge acts on a tombstone whose delete already fired, and is reported through DocumentsPurgedCallback. A restore is reported through DocumentRestoredCallback, which carries no document and so fires before the restore's read-back can fail. A host that needs the listings themselves subscribes an in-process client to them, which keeps them current from patches without re-reading the folder.
Callbacks are invoked synchronously during the operation and NOT awaited: a returned promise is observed only to LOG its rejection, so an async callback can never stall or fail the operation that triggered it.
Callbacks report exactly what was stored. A write's outcome is its storage commit: once it commits, the operation succeeds and its callbacks fire, even if a later step — the orphan Y.Doc reclaim, the subscriber broadcast, the listing broadcasts — fails; that failure is logged and treated as a lost frame. A write that fails before its commit rejects and fires nothing. Every document handed over is one the write path already holds (built from the write's input, or read before the commit), so no post-commit failure can leave a callback without its payload.
DocumentChangeCallbackContext
export type DocumentChangeCallbackContext<EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown> = EventBusContext & StorageContext & LoggerContext;The context a DocumentChangeCallback receives: the context of the operation that made the change, so the callback sees the caller's logger and identity.
DocumentInvalid
export type DocumentInvalid = InferValue<typeof DocumentInvalidValue>;The invalid marker a relaxed read attaches to a non-conforming document — on the wire (doc:init) and on the direct server reads (getDocument). The single source shared by every server signature that carries the marker.
DocumentRestoredCallback
export type DocumentRestoredCallback<EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown> = (context: DocumentChangeCallbackContext<EventBusContext, StorageContext>, entry: Readonly<DeletedDocumentEntry>) => void | Promise<void>;A host's hook on "this tombstoned document was restored" — one call per committed restore, with the entry as it stood in the trash. It fires straight after the restore commits, BEFORE the engine reads the document back (bringing it forward to the current schema) and broadcasts it: the callback needs nothing that read produces, so a restore whose read-back or broadcast then fails is still reported. That is why a restore has its own hook rather than firing DocumentChangeCallback — that callback hands over the document, which only the read-back has. A host that needs the restored content reads it (getDocument answers it brought forward). An idempotent restore of a live document changes nothing and does not fire. Invoked and error-contained like DocumentChangeCallback.
DocumentsPurgedCallback
export type DocumentsPurgedCallback<EventBusContext extends LoggerContext = LoggerContext, StorageContext = unknown> = (context: DocumentChangeCallbackContext<EventBusContext, StorageContext>, entries: readonly Readonly<DeletedDocumentEntry>[]) => void | Promise<void>;A host's hook on "these tombstoned documents were purged" — one call per purge batch, after the sys:trash removal is broadcast, with the entries as they stood in the trash. The document-change callbacks already reported each document's delete, and a purge changes no stored document they could hand over, so this is the only server-side signal that a purge happened. The entries are read-only: the same objects go on to later callbacks, the purge observability events and purgeDeletedDocuments' result. Invoked and error-contained like DocumentChangeCallback, and like it fires for every committed purge, even one whose broadcast then fails.
ErrorObservabilityEvent
export type ErrorObservabilityEvent = ObservabilityEventBase<"error", {
docId: string;
operation: "subscribe" | "create" | "update" | "command" | "delete" | "restore" | "rename" | "purge";
category: DocumentErrorCategory;
message: string;
wasInvalid?: boolean;
}>;Error event - fired when an error occurs during any operation
ObservabilityCallback
export type ObservabilityCallback = (event: ObservabilityEvent) => void;Callback function for observability events. Called synchronously but should not throw (errors are logged and swallowed). Intended to be fire-and-forget.
ObservabilityEvent
export type ObservabilityEvent = SubscribeObservabilityEvent | UnsubscribeObservabilityEvent | CreateObservabilityEvent | UpdateObservabilityEvent | DeleteObservabilityEvent | RestoreObservabilityEvent | PurgeObservabilityEvent | RenameObservabilityEvent | ErrorObservabilityEvent;Union of all observability events
PurgeObservabilityEvent
export type PurgeObservabilityEvent = ObservabilityEventBase<"purge", {
docId: string;
type: string;
}>;Purge event - fired when a soft-deleted document's content is destroyed and its id released (the hard-delete maintenance surface).
RenameObservabilityEvent
export type RenameObservabilityEvent = ObservabilityEventBase<"rename", {
docId: string;
type: string;
}>;Rename event - fired when a document's index-level name is changed.
The new name is deliberately NOT carried here: a document name is free-form user content that may be confidential, and observability events are routinely logged, exported and retained. docId/type identify the operation without leaking content. (Consumers needing the name should read sys:index, not tap the observability stream.)
RestoreObservabilityEvent
export type RestoreObservabilityEvent = ObservabilityEventBase<"restore", {
docId: string;
type: string;
}>;Restore event - fired when a soft-deleted document is restored
SubscribeObservabilityEvent
export type SubscribeObservabilityEvent = ObservabilityEventBase<"subscribe", {
docId: string;
clientSequence: number | null;
serverSequence: number | null;
sequenceGap: number | null;
outcome: "init" | "resume" | "notfound" | "deleted" | "error";
yjsUpdateCount: number;
migrated: boolean;
invalid: boolean;
}>;Subscribe event - fired when a client subscribes to a document
UnsubscribeObservabilityEvent
export type UnsubscribeObservabilityEvent = ObservabilityEventBase<"unsubscribe", {
docId: string;
}>;Unsubscribe event - fired when a client unsubscribes from a document
UpdateObservabilityEvent
export type UpdateObservabilityEvent = ObservabilityEventBase<"update", {
docId: string;
type: string;
previousSequence: number;
newSequence: number;
patchOperationCount: number;
yjsUpdateCount: number;
source: "client" | "server";
clientEventId?: string;
command?: string;
migrated: boolean;
yjsSizeBytes: number;
wasInvalid: boolean;
}>;Update event - fired when a document is updated
Variables
DEFAULT_BLOB_SWEEP_POLICY
DEFAULT_BLOB_SWEEP_POLICY: BlobSweepPolicyThe default BlobSweepPolicy: a day's retention, 100 blobs per sweep.
MAX_BLOB_RETENTION_MS
MAX_BLOB_RETENTION_MS: numberThe longest retention a policy may name: ten years. Bounded so every deadline a host derives from it — a catalog timestamp plus the retention plus one — stays a safe integer, and so stays exact through a storage driver's integer column and a platform timer.
MAX_INLINE_EXPORT_BLOB_BYTES
MAX_INLINE_EXPORT_BLOB_BYTES = 16777216Blobs above this size cannot ride an inline DocumentExport — the export would balloon memory on a cold maintenance path (base64 in the text form makes it worse). Exporting a document referencing a larger blob fails loudly rather than silently dropping content; a streaming export form is future work.
SERVER_YJS_AUTHORING_CLIENT_ID
SERVER_YJS_AUTHORING_CLIENT_ID = 3665484410The Yjs clientID every datadata server authors under (0xDA7ADA7A — "data data"). Hard-coded, not random: a document has exactly ONE authoritative server, and successive instances of it (a Durable Object waking from hibernation is a NEW instance) write sequentially, never concurrently — temp Y.Docs continue the constant's clock from stored state. One shared constant therefore keeps the server's state-vector footprint at ONE entry for a document's whole life, and makes server-authored structs recognisable when inspecting a Y.Doc. Predictability costs nothing: clients submit raw encoded updates and could author under any clientID regardless. Clients and StagingSessions must NOT author under this value (their random ids make an accidental collision negligible).
Known narrow window, accepted: if a crash drops an unflushed write-behind burst and the revived server authors into the same yjsId BEFORE a client holding the lost broadcasts reconnects, it reuses clocks; yjs dedups structs by id, so the state-vector self-heal cannot restore those structs and that client stays diverged until it refetches from scratch.