Skip to content

Observe SQL tables in process

Experimental, in-process only. For committed changes from other processes on PostgreSQL, see PostgreSQL LISTEN/NOTIFY observation. This mode sees only writes announced with notifyChanged in this application process. Other processes, database clients, triggers, and unannounced writes are invisible. Host-owned transactions around Arc commands and transactions owned by an outer command runner are not covered by command completion: announce those writes after commit. Read from a connection that sees freshly committed data, not a long-lived repeatable-read snapshot or a lagging replica. Pages are eventually consistent, not atomic snapshots of count and rows.

Install rxjs alongside @cratis/arc.drizzle: it is a required peer dependency, loaded even when observation is disabled. Enable DrizzleObservation.InProcess explicitly. Without it, starting an observation fails; ordinary read APIs and validated notifyChanged calls still work.

This single file keeps the command, read model, query, table, and setup together. It uses SQLite through sql.js; your application owns the database and migrations. The command writes and announces the same registered table object.

Tasks.ts
import initSqlJs from 'sql.js';
import { drizzle, type SQLJsDatabase } from 'drizzle-orm/sql-js';
import { sqliteTable, text } from 'drizzle-orm/sqlite-core';
import { field, Guid } from '@cratis/fundamentals';
import { ArcApplication, command, inject, key, query, readModel, service } from '@cratis/arc.core';
import { DrizzleDialect, DrizzleObservation, drizzleDatabase, drizzleReadModel,
guidCodec, sqliteColumn, type DrizzleHandle, type DrizzleReadModels } from '@cratis/arc.drizzle';
export const tasks = sqliteTable('tasks', {
id: sqliteColumn(guidCodec(DrizzleDialect.SQLite))('id').primaryKey(),
title: text('title').notNull()
});
export class TaskRecord {
@field(Guid) @key() id!: Guid;
@field(String) title!: string;
}
@command()
export class AddTask {
@field(Guid) @key() id!: Guid;
@field(String) title!: string;
@inject(drizzleDatabase<SQLJsDatabase>())
handle(database: DrizzleHandle<SQLJsDatabase>): void {
database.native.insert(tasks).values({ id: this.id, title: this.title }).run();
database.notifyChanged(tasks);
}
}
@readModel()
export class TaskListing {
@query({ observable: true }, service(drizzleReadModel(TaskRecord)))
static all(records: DrizzleReadModels<TaskRecord>) {
return records.observe();
}
}
const Sql = await initSqlJs();
const native = new Sql.Database();
native.run('create table tasks (id text primary key, title text not null)');
const database = drizzle(native);
const builder = ArcApplication.createBuilder({ tenancy: { resolve: () => 'default' } });
builder.add(AddTask, TaskListing).withDrizzle({
dialect: DrizzleDialect.SQLite, database,
readModels: [{ type: TaskRecord, table: tasks }],
observation: DrizzleObservation.InProcess
});
const app = await builder.build();
await app.run();

Once running, request a snapshot or open a stream at the query route:

Terminal window
curl 'http://127.0.0.1:3000/api/task-listing/all?page=0&pageSize=10'
curl -N -H 'Accept: text/event-stream' \
'http://127.0.0.1:3000/api/task-listing/all?page=0&pageSize=10'

Check the route generated by your application; namespaces or route configuration can change the URL. Plain GET returns 200 from current(). Arc pages the emitted array in memory for each request; SSE starts with the same page and emits again after an announced change. Unsubscribing or closing the request releases the listener. If you call current() in a long-lived background scope without subscribing, dispose the scope or call close() on the observable to release its listener.

notifyChanged(table) accepts only the registered table object; notifyChanged(TaskRecord) resolves to that table. A different table object with the same name is not registered and throws. The handle takes its tenant from the Arc execution scope; tenants cannot announce each other’s changes, even if databaseFactory points both at one physical database. The tenant key is lowercased just as it is for the factory.

Inside an Arc command, announcements from nested commands for the same tenant are combined and published once when that tenant’s outer command execution finishes, whether it succeeds or fails. An outer command runner means one registered earlier than withDrizzle, including options.commandExecutionRunner; it may own a transaction that outlives the Drizzle runner. That follows any transaction awaited inside that command, but does not imply a transaction owned outside Arc has committed. Outside a command (including background jobs), notifyChanged publishes immediately: call it only after commit. A failed command can leave a nontransactional write behind, so failure also flushes announced changes.

observePage(filter, options) is available only through low-level defineObservableQuery, without generated artifact metadata or client proxies; model-bound observable queries must use observe() because generated metadata rejects observable pages. See observable queries. It counts, sorts, limits, and offsets in SQL for each emission. Count and row selection use separate statements: their lengths are checked, retried up to three times on disagreement, then the stream fails. Equal lengths do not prove an atomic snapshot. An announced write racing a read triggers another read, so pages converge on committed data. For small complete collections, observe(filter?) reads at most maxPageSize + 1 rows in one statement and fails with QueryPagingRequired on overflow; it never emits a truncated list. observeById(key) emits null for a missing or deleted row and can emit that row again if recreated.

A listener is registered before the initial read. Bursts coalesce to one trailing read while a read is in flight; there is no polling or time-based debounce. Read failures and paging failures end the subscription; reconnect by subscribing again to get a fresh snapshot. Core’s outbound queue limit (256 pending emissions by default) also fails a slow subscription explicitly. There is no cross-process delivery in this mode; MongoDB change streams have a different change source.