@repo/datadata/in-process
Connecting a real client to a server in the same process without a socket, for example inside a Durable Object.
Classes
CompositeEventBus
export declare class CompositeEventBus<Ctx extends InProcessEventBusContext = InProcessEventBusContext> implements DatadataEventBus<Ctx>, InProcessConnectionRegistrarPair a primary DatadataEventBus (e.g. a WebSocketEventBus inside a Durable Object) with an InProcessEventBus, so one server can serve both transports at once. A server given this bus broadcasts every document change to BOTH WebSocket subscribers and in-process subscribers, and routes each connection's subscriptions and awaited-write acks to whichever transport owns it. This is what lets a server-side in-process client (the AI staging edits) and a remote WebSocket client (the human) see each other's committed changes over a single server.
Attach in-process clients by passing this bus to connectInProcessClient (it is an InProcessConnectionRegistrar, delegating to the wrapped InProcessEventBus).
constructor(primary: DatadataEventBus<Ctx>, inProcess: InProcessEventBus<Ctx>);Pair primary (e.g. a WebSocket bus) with inProcess, which owns in-process connections.
addSubscription(context: Ctx, options: {
docId: string;
}): void;Subscribe the calling connection to docId, on whichever transport owns the connection.
connections(context: Ctx): Iterable<EventBusRecipient & {
subscribedDocIds: readonly string[];
}>;Both transports' connections: the primary's first, then the in-process ones.
hasSubscription(context: Ctx, options: {
docId: string;
}): boolean;Whether the calling connection is subscribed to docId, on the transport that owns it.
registerConnection(handler: ServerSentEventHandler, principal: Principal): string;Register a connection on the wrapped InProcessEventBus; returns its connectionId.
removeSubscription(context: Ctx, options: {
docId: string;
}): void;Unsubscribe the calling connection from docId, on the transport that owns it.
sendEvents(context: Ctx, events: ServerSentEvent[], options?: {
respondToCaller?: true;
toConnection?: string;
recipientFilter?: (recipient: EventBusRecipient) => boolean;
}): void;Deliver events: a reply to the caller or a targeted send goes only to the transport that owns the connection; a broadcast goes to both, each delivering to its own subscribers and applying recipientFilter.
subscribedDocIds(context: Ctx): Iterable<string>;The union of both transports' subscribed docIds.
unregisterConnection(connectionId: string): void;Detach an in-process connection from the wrapped InProcessEventBus.
InProcessEventBus
export declare class InProcessEventBus<Ctx extends InProcessEventBusContext = InProcessEventBusContext> implements DatadataEventBus<Ctx>, InProcessConnectionRegistrarA DatadataEventBus whose subscribers are in-process callbacks. Each registered connection has a handler and its own subscription set; sendEvents delivers a broadcast to every connection subscribed to the event's document, and a respondToCaller event only to the originating connection — the same routing a WebSocket bus performs, against JS callbacks instead of sockets.
Delivery is synchronous, and deliberately so — see sendEvents for the contract and how it differs from a socket transport.
addSubscription(context: Ctx, options: {
docId: string;
}): void;Subscribe the calling connection to docId. A no-op for a server-initiated context.
clearSubscriptions(connectionId: string): void;Drop a connection's subscriptions without detaching it (e.g. a simulated reconnect).
connections(_context: Ctx): Iterable<EventBusRecipient & {
subscribedDocIds: readonly string[];
}>;Every registered connection, with its principal and the docIds it is subscribed to.
hasConnection(connectionId: string): boolean;Whether connectionId belongs to this bus — used by CompositeEventBus to route.
hasSubscription(context: Ctx, options: {
docId: string;
}): boolean;Whether the calling connection is subscribed to docId; false when server-initiated.
registerConnection(handler: ServerSentEventHandler, principal: Principal): string;Register an in-process connection. The returned connectionId is what the caller threads into the context it passes to server.handleEvent, so the server can address this connection (subscriptions, awaited-write acks). The principal is what the bus reports for per-recipient delivery (read filtering).
removeSubscription(context: Ctx, options: {
docId: string;
}): void;Unsubscribe the calling connection from docId. A no-op for a server-initiated context.
sendEvents(context: Ctx, events: ServerSentEvent[], options?: {
respondToCaller?: true;
toConnection?: string;
recipientFilter?: (recipient: EventBusRecipient) => boolean;
}): void;Delivery is SYNCHRONOUS: the handler runs inline, so a subscriber has fully processed the event by the time this returns (unlike WebSocketEventBus, whose ws.send lands on a later task). That is a real property of this transport, not a shortcut — there is no wire. Note the asymmetry with the OUTBOUND direction, which connectInProcessClient defers onto a microtask: an originating client's echo therefore arrives after its own write call has returned and its optimistic state has settled, just synchronously with respect to the server's handleEvent.
One subscriber's handler must not be able to abort delivery to the rest, so each invocation is isolated — mirroring trySend on the WebSocket side, where a dead socket's failure never starves the remaining subscribers. This matters most under a CompositeEventBus, where the "rest" includes the other transport.
subscribedDocIds(_context: Ctx): Iterable<string>;Every docId at least one connection on this bus is subscribed to.
unregisterConnection(connectionId: string): void;Detach a connection: it stops receiving events and its subscriptions are dropped.
updateConnectionPrincipal(connectionId: string, principal: Principal): void;Replace a live connection's principal — the transport half of a mid-connection identity change (grant widening, demotion). The caller must keep the contexts it threads into server.handleEvent in step (one identity source), and then run server.reauthorizeConnections so read access is re-evaluated and sys:principal subscribers learn the new identity. No-op for an unknown connectionId.
Functions
connectInProcessClient
export declare function connectInProcessClient<S extends SchemaRegistry, Ctx extends InProcessEventBusContext, StorageContext = unknown, C extends CommandRegistry = CommandRegistry>(params: {
server: {
handleEvent(context: Ctx & StorageContext & PrincipalContext & ConnectionContext, event: ClientSentEvent): void | Promise<void>;
handleDisconnect(context: Ctx & StorageContext & PrincipalContext & ConnectionContext): void | Promise<void>;
};
bus: InProcessConnectionRegistrar;
schemas: S;
logger: Logger;
createEventId: () => string;
principal: Principal;
createContext: (connectionId: string) => Ctx & Omit<StorageContext, keyof PrincipalContext>;
config?: ClientConfig;
commands?: C;
dropServerEvent?: (event: ServerSentEvent) => boolean;
}): {
client: DatadataClientFor<S, C>;
connectionId: string;
dispose: () => void;
};Wire a real DatadataClientImpl to a DatadataServer in the same process via an InProcessEventBus. The client's outbound events are delivered to the server (on a microtask, to avoid re-entrancy — the client emits some events before it has finished updating its own local state), and the server's broadcasts route back into the client. Returns the client plus a dispose that detaches its connection.
params.principal is the identity this connection runs as: it is registered on the bus (so per-recipient delivery and read filtering see it) and merged into every server.handleEvent context (so scope attenuation and the declarative access rules apply), one source for both. To connect anonymously, pass ANONYMOUS_PRINCIPAL explicitly; the default policy denies anonymous writes, so the type makes the choice deliberate.
params.createContext builds the rest of the context server.handleEvent runs under, given the connectionId the bus assigned. It must carry whatever the server's storage and event-bus layers need (the connectionId, a logger, and any storage context), the same context the server's WebSocket path constructs per connection. The principal is merged in on top.
Interfaces
InProcessConnectionRegistrar
export interface InProcessConnectionRegistrarThe subset of an in-process bus that attaches a client: register a callback connection (minting the connectionId the client's sends are addressed from) and detach it. Both InProcessEventBus and CompositeEventBus implement it, so connectInProcessClient works whether the server runs a pure in-process bus or one composed with another transport (e.g. WebSockets).
registerConnection(handler: ServerSentEventHandler, principal: Principal): string;Register a callback connection running as principal — the identity the bus reports for per-recipient delivery (read filtering); it must be the same principal the connection's handleEvent contexts carry.
unregisterConnection(connectionId: string): void;Detach a connection: it stops receiving events and its subscriptions are dropped.
InProcessEventBusContext
export interface InProcessEventBusContext extends LoggerContextThe event-bus context an in-process connection is addressed by: a LoggerContext carrying the connectionId the bus keys its routing on (mirroring WebSocketEventBusContext) — or, on a call the server itself started, the ServerInitiated marker, which no connection answers to.
connectionId: ConnectionIdOrServerInitiated;The calling connection, or the server-initiated marker for a call the server started.
Types
ConnectionContext
export type ConnectionContext = {
connectionId: string;
};What the server's wire entry (handleEvent, handleDisconnect) asks of a context on top of the bus's own: a real connectionId, never the server-initiated marker an InProcessEventBusContext also admits.
ServerSentEventHandler
export type ServerSentEventHandler = (event: ServerSentEvent) => void;A sink for the server-sent events routed to one in-process connection.