Skip to main content

Processors

A processor receives operations after the write phase and performs side effects: analytics rollups, relational indexing, webhooks, dispatching follow-up actions. Its onOperations method runs post-ready — after the reactor emits JOB_READ_READY for a chain, when the document is already readable and consistent.

This is the contrast with read models. Read models index pre-ready: their indexOperations runs before JOB_READ_READY fires, so they gate read-after-write consistency. A reader holding a consistency token is unblocked only once the pre-ready models finish. Processors are not on that critical path; they observe operations only after readers can already see them. For the full pre-ready/post-ready ordering see Working with the Reactor, and for read models see Advanced Reactor Usage.

This page is the reactor-side reference: the types you implement, the manager that drives them, and how catch-up works. For the editor-side authoring flow (scaffolding a processor in a package) see the processor tutorial.

:::warning Not stable yet Processors run only on the in-process reactor path. On the SharedWorker reactor path, processor registration is not bridged — Connect logs a warning and skips every factory (apps/connect/src/store/reactor.ts). Worker-hosted reactors do not run processors today. :::

Registering a processor​

Register a factory with the processor manager. The manager interface is IProcessorManager:

import type { IProcessorManager, ProcessorFactory } from "@powerhousedao/reactor";

interface IProcessorManager {
registerFactory(identifier: string, factory: ProcessorFactory): Promise<void>;
unregisterFactory(identifier: string): Promise<void>;
get(processorId: string): TrackedProcessor | undefined;
getAll(): TrackedProcessor[];
}

identifier is the registration key and the factoryId prefix of every processor the factory produces. Real callers pass the package name. Registering the same identifier twice unregisters the prior factory first, so re-registration is idempotent.

The manager is not a method on IReactor. It lives on the in-process reactor module as processorManager. In Connect you reach it through the browser-path reactor module:

const reactorModule =
reactorClientModule.kind === "browser"
? reactorClientModule.reactorModule
: undefined;

await reactorModule?.processorManager.registerFactory(id, factory);

What the factory receives and returns​

A ProcessorFactory is called once per drive. It receives the drive header and returns the processors for that drive:

A drive must exist before any processor runs. The manager calls the factories when it sees a CREATE_DOCUMENT operation whose document type is a drive container (powerhouse/document-drive or powerhouse/reactor-drive by default, see withDriveContainerTypes). Until a drive exists, a registered factory is never invoked and no operations are routed, even for plain documents that match a filter. In a standalone reactor, create a drive before you expect onOperations to fire.

import type { PHDocumentHeader } from "document-model";

type ProcessorFactory = (
driveHeader: PHDocumentHeader,
processorApp?: ProcessorApp,
) => Promise<ProcessorRecord[]> | ProcessorRecord[];

Return [] to create no processors for a drive. The factory may be sync or async. The processorApp? parameter is part of the type but the manager never supplies it — it calls factory(driveHeader) with one argument. Read the running app from module.processorApp instead (see the processor host).

If the factory throws, the manager catches it, logs Factory '<id>' failed for drive '<driveId>', and skips that drive. Registration does not crash.

Each element of the returned array is a ProcessorRecord — the exact shape the factory must produce:

type ProcessorRecord = {
processor: IProcessor;
filter: ProcessorFilter;
startFrom?: "beginning" | "current";
};

processor and filter are required. startFrom is covered under catch-up and ordering.

The IProcessor interface​

interface IProcessor {
onOperations(operations: OperationWithContext[]): Promise<void>;
onDisconnect(): Promise<void>;
}

onOperations is called post-ready with the operations that matched this processor's filter, during both live routing and backfill. Each OperationWithContext is { operation, context }; context carries documentId, documentType, scope, branch, the global ordinal, and optionally resultingState.

onDisconnect runs when the factory is unregistered or its drive is deleted. The manager wraps this call: a throw is caught and logged, never propagated.

Most processors subclass RelationalDbProcessor rather than implementing IProcessor by hand. Its constructor takes (namespace, filter, relationalDb), and you implement onOperations, initAndUpgrade, and onDisconnect. See Storage and scaling for the relational DB surface.

A complete factory builder​

Packages export a ProcessorFactoryBuilder, not a bare factory. The builder receives the host module and returns a ProcessorFactory:

type ProcessorFactoryBuilder = (
module: IProcessorHostModule,
) => Promise<ProcessorFactory> | ProcessorFactory;

This is the shape codegen emits for a relational processor:

import type {
IProcessorHostModule,
ProcessorApp,
ProcessorFactoryBuilder,
ProcessorFilter,
ProcessorRecord,
} from "@powerhousedao/reactor-browser";
import type { PHDocumentHeader } from "document-model";
import { TodoIndexer } from "./processor.js";

export const todoIndexerFactoryBuilder: ProcessorFactoryBuilder =
(module: IProcessorHostModule) =>
async (driveHeader: PHDocumentHeader, processorApp?: ProcessorApp) => {
const namespace = TodoIndexer.getNamespace(driveHeader.id);
const store = await module.relationalDb.createNamespace<TodoIndexer>(namespace);

const filter: ProcessorFilter = {
branch: ["main"],
documentId: ["*"],
documentType: ["powerhouse/todo"],
scope: ["global"],
};

const processor = new TodoIndexer(namespace, filter, store);
return [{ processor, filter }];
};

Wiring at startup is two stages: build the host module, call the builder to get a factory, then register it.

const factory = await todoIndexerFactoryBuilder(processorHostModule);
await reactorModule.processorManager.registerFactory("@my-org/my-package", factory);

Processor types are re-exported from @powerhousedao/reactor-browser as types only. RelationalDbProcessor is re-exported as a value. IProcessorManager, ProcessorStatus, TrackedProcessor, and the concrete ProcessorManager class are exported from @powerhousedao/reactor, not from reactor-browser.

ProcessorFilter​

The filter decides which operations reach a processor:

type ProcessorFilter = {
documentType?: string[];
scope?: string[];
branch?: string[];
documentId?: string[];
};

Every field is an optional array. Matching is per field: a field that is undefined or [] matches everything; a field with values matches an operation whose corresponding context value is in the array. All present fields must match.

// Global-scope operations on any todo document, main branch only.
const filter: ProcessorFilter = {
documentType: ["powerhouse/todo"],
scope: ["global"],
branch: ["main"],
documentId: ["*"],
};

The "*" wildcard is honored only in documentId. There it means "any document". In documentType, scope, and branch there is no wildcard special-case: a literal "*" matches only an operation whose value is the string "*". To match all scopes, omit scope or set it to [] — do not write scope: ["*"].

:::warning Latent inconsistency The codegen analytics template emits scope: ["*"]. Under the matcher this matches a scope literally named "*", which matches nothing real. Set scope to the concrete scopes you want, or omit it. :::

startFrom on the ProcessorRecord controls the catch-up starting point, not which operations match — see catch-up and ordering.

The processor host​

The host module is the context a processor gets through its factory builder. The host-agnostic core is IProcessorHostModuleBase:

interface IProcessorHostModuleBase {
analyticsStore: IAnalyticsStore;
relationalDb: IRelationalDb;
processorApp: ProcessorApp;
dispatch: IProcessorDispatch;
getReadModel<T>(name: string): T;
config?: Map<string, unknown>;
}
  • analyticsStore — the analytics store for time-series rollups.
  • relationalDb — an IRelationalDb (a Kysely instance plus createNamespace / queryNamespace) for relational indexing. See Storage and scaling.
  • processorApp — "connect" | "switchboard". How a processor learns which app hosts it. Read this rather than the factory's processorApp? argument.
  • dispatch — writes back to the reactor via dispatch.execute(docId, branch, actions, signal?, meta?), returning { id, status }. In Connect this is wired to the reactor client's async execute. Do not await client.execute() from inside onOperations. The promise resolves only at READ_READY, which the current batch cannot reach while your callback is still running. Use dispatch and handle the result in a later onOperations call (see Processor best practices).
  • getReadModel(name) — looks up a registered read model by its name. Reactor-registered names are typed: getReadModel("document-view") returns IDocumentView and "document-indexer" returns IDocumentIndexer; other names take an explicit type argument. Connect's implementation throws Read model "<name>" not found when there is no match.
  • config? — optional Map<string, unknown> of host config.
  • client — the IReactorClient for reading documents and drives. See IReactorClient.
  • attachments — the IAttachmentClient for uploading and downloading attachments. See the Attachment service.

These six core fields are IProcessorHostModuleBase in @powerhousedao/shared. @powerhousedao/reactor adds client and the typed getReadModel as IReactorProcessorHostModuleBase; @powerhousedao/reactor-browser and @powerhousedao/reactor-api add attachments as IProcessorHostModule, since neither shared nor reactor can depend on reactor-attachments. Connect sets processorApp: "connect".

IProcessorDispatch and ProcessorDispatchResult are defined in shared but are not part of the public reactor export surface; treat the dispatch handle on the module as the supported entry point.

Catch-up and ordering​

Each processor has a cursor, lastOrdinal: the highest operation ordinal it has taken, pulled back below any batch it failed to take. The manager tracks this per processor as TrackedProcessor:

type TrackedProcessor = {
processorId: string; // `${factoryId}:${driveId}:${index}`
factoryId: string;
driveId: string;
processorIndex: number;
record: ProcessorRecord;
lastOrdinal: number;
status: ProcessorStatus; // "active" | "errored"
lastError: string | undefined;
lastErrorTimestamp: Date | undefined;
retry: () => Promise<void>;
};

Ordering. Batches reach the manager one document at a time, and different documents are projected in parallel, so a processor is not called in global ordinal order: a call may carry an ordinal lower than one an earlier call carried. Within one document's scope and branch, operations arrive in ordinal order; across documents there is no ordering guarantee. A processor receives one onOperations call at a time: each processor has its own delivery queue, and the next call begins after the previous one resolves. Batches that queue up behind a call are delivered together in the next one, in the order they arrived. A slow processor delays only its own queue. On each batch the manager applies the filter and queues the matches, minus any a backfill has already delivered. On success it raises lastOrdinal to the batch's maximum ordinal (including unmatched operations) when that is higher. Delivery is at-least-once: a processor may see an operation again after a restart or a retry, or when a backfill reads an operation whose live batch is still on its way to the manager. A crash between two concurrently projected documents, after the higher ordinal's cursor was persisted, can leave the lower ordinals unreplayed.

startFrom. When a processor is first created and no cursor row exists yet:

  • "beginning" (the default) starts the cursor at ordinal 0, so the processor backfills the drive's full history.
  • "current" starts the cursor just below the drive's creation when the processor is created from the drive's creation batch, and at the manager's current ordinal when a late registerFactory creates it, so the processor sees only its drive's history or only what comes next.

startFrom applies only on first creation. Once a cursor row exists, it is ignored — the persisted cursor wins. Restarting the reactor never re-runs startFrom.

Backfill and replay. When a processor's lastOrdinal is behind the manager, the manager pages through history with operationIndex.getSinceOrdinal(lastOrdinal), filters each page, calls onOperations, advances the cursor to the page's max ordinal, persists it, and follows the continuation until exhausted. This runs when a factory registers against a drive that already has history, and after a restart to replay anything missed while the reactor was down. Backfill runs on the processor's queue, so other documents keep indexing and other processors keep receiving; live batches for a processor that is backfilling wait behind it and skip whatever the backfill already delivered. registerFactory resolves once the factory has run and its processors are bound, before their backfills finish; to wait for a backfill, watch the processor's lastOrdinal through get() or getAll().

Persistence. Cursors are stored in the ProcessorCursor table, keyed by processorId, and survive restarts. On startup the manager rehydrates every cursor, then catch-up replays the gap.

Failure isolation. Processors run in parallel, each on its own queue. If onOperations throws, or a backfill page cannot be read, that processor goes to status: "errored", its lastError and lastErrorTimestamp are recorded, and its cursor is set below the lowest ordinal of the batch that failed. Other processors are unaffected. There is no automatic retry: an errored processor stays errored, and later batches that match it are skipped and pull its cursor below them likewise, until you call tracked.retry() (which re-activates it and replays from the cursor) or re-register its factory. See Error handling.

To inspect state, use the manager: get(processorId) for one processor or getAll() for every tracked processor across all drives.

For how processors relate to sync and remote drives, see Synchronization. For the document types a filter can target, see the Document model registry.