Skip to content

Observable queries

A task board should update when someone adds a task, without the browser polling. An observable query serves the current snapshot on an ordinary GET and streams every change to subscribers, from the same route and the same pipeline.

The Tasks sample declares its live list as an ordinary query that returns a source:

@query()
static observeAllTasks(tasks: Tasks): BehaviorSubject<TaskItem[]> { return tasks.observeAll(); }

The sample’s Tasks service keeps a new BehaviorSubject<TaskItem[]>([]) from RxJS and calls next(...) whenever a task is registered.

Arc has to know that a query is observable before it runs: snapshots, server-sent events, WebSocket admission, introspection, and generated clients all depend on it. The sample installs generated artifact metadata, which reads the declared return type, so bare @query() is enough. Without generated metadata, say so explicitly and list the parameters:

@query({ observable: true }, service(Tasks))
static observeAllTasks(tasks: Tasks): BehaviorSubject<TaskItem[]> { return tasks.observeAll(); }
SourceSnapshot behavior
RxJS BehaviorSubject<T>Has a current value: GET answers 200 immediately, and new subscribers receive it first
RxJS Subject<T> or Observable<T>No current value: GET answers 202 with isReady: false; waitForFirstResult=true subscribes until the first value
RxJS ReplaySubject<T>No readable current value: GET answers 202 even after an emission; a waiting GET receives the buffered value
AsyncIterable<T>No current value until the first item
CurrentValueSubject<T> (deprecated)Legacy current/pending source; use RxJS BehaviorSubject or Subject instead

A BehaviorSubject exposes its current value, including undefined; Subject and ReplaySubject do not. The core accepts structural subscribables and async iterables without loading RxJS at runtime, so RxJS is an optional peer dependency for consumers using only async iterables.

A live lookup, such as “the task with this ID”, may have nothing to show yet. Two different situations look alike from the outside:

The sourceA snapshot GET answersA subscriber receives
Has no current value yet (Subject, Observable, an empty async iterable)202 with isReady: falseNothing until the first emission
Emits undefined or null (a BehaviorSubject<TaskItem | undefined> with nothing stored)200, isReady: true, no data propertyA ready result without data

In both cases the subscription stays open. When the document appears and the source emits it, the subscriber receives it like any other update, and when it disappears again the source can emit undefined. Emit undefined for “we looked and nothing is there”, and leave a source silent only when you genuinely do not know yet; a client can tell the two apart by isReady.

An error is neither. A source that errors ends the subscription with a failed result, and a source that completes before its first value returns an error to a waiting GET. Do not turn a storage failure into an undefined emission.

Each subscription runs your query method once, after authorization and validation, and then follows the source it returned:

  • The subscription owns its own service scope for its whole life. Scoped services are not shared between subscribers, and they are disposed when the subscription ends.
  • The method can read the caller with currentContext() from @cratis/arc.core, which returns the subscription’s execution context, and its signal aborts when the subscription ends. Pass it to anything that must stop with the subscriber.
  • The subscription ends when the client disconnects or unsubscribes, when the source completes or errors, or when an emission guard denies an emission. Arc then cancels the source and disposes the scope.

Read the snapshot and subscribe from the terminal

Section titled “Read the snapshot and subscribe from the terminal”

With the Tasks sample running:

Terminal window
curl http://127.0.0.1:3000/api/tasks/listing/observe-all-tasks
curl -N -H 'Accept: text/event-stream' http://127.0.0.1:3000/api/tasks/listing/observe-all-tasks

The first answers 200 with the current tasks in data. The second keeps the connection open and prints a data: <query result JSON> frame now, and another whenever you register a task. Using observable queries with curl covers waiting for a first result and the error codes.

The paging and sorting parameters work on an observable query exactly as on a snapshot query, and they apply to every emission:

Terminal window
curl -N -H 'Accept: text/event-stream' \
'http://127.0.0.1:3000/api/tasks/listing/observe-all-tasks?pageSize=1&page=1&sortBy=title&sortDirection=desc'

With two tasks registered, each frame holds the second task in descending title order, and paging reports {"page":1,"size":1,"totalItems":2,"totalPages":2}. When a third task arrives, the next frame is sorted and cut again, and totalItems follows the whole list. A generated client’s useWithPaging(pageSize) hook sends the same parameters.

Arc pages the arrays your source emits, in memory. Model-bound observable queries cannot return a queryPage; generated metadata rejects that declaration at build, and generated clients do not support observable-page results. Drizzle model-bound queries use observe() for complete small arrays, as MongoDB does. For larger collections, narrow the source with query arguments. Low-level defineObservableQuery without generated metadata or proxies can emit a QueryPage (including Drizzle’s observePage); Core validates and streams each page, but the generated client cannot consume this shape. The paging rules for invalid sizes and sort fields are the same as for snapshots.

An observable query takes the same @roles, @authorize, and @allowAnonymous declarations as any query, and Arc checks them, with validation, when a subscription opens. A denied snapshot GET answers 401 or 403 with isAuthorized: false, following the status code rules; a denied hub subscription receives an Unauthorized frame. The query method never runs for a denied caller.

That check happens once. Arc keeps a copy of the caller’s identity for the life of the subscription and does not re-run authorization on each emission, so a role revoked or a session that expires after the subscription opened does not close it. When access must be re-checked while the stream runs, add an emission guard, which sees every result before delivery and can end the subscription.

A role decides who may subscribe, not which rows they see. For a query like “my tasks”, filter inside the source by the caller’s identity:

MyTasks.ts
import { field } from '@cratis/fundamentals';
import { authorize, currentContext, query, readModel } from '@cratis/arc.core';
import { BehaviorSubject, map, type Observable } from 'rxjs';
const tasks = new BehaviorSubject<OwnedTask[]>([]);
@readModel()
@authorize()
export class OwnedTask {
@field(String) id!: string;
@field(String) title!: string;
@field(String) owner!: string;
@query({ observable: true })
static myTasks(): Observable<OwnedTask[]> {
const caller = currentContext()?.principal?.id;
return tasks.pipe(map(all => all.filter(task => task.owner === caller)));
}
}

@authorize() turns anonymous callers away before myTasks runs, so caller is always an authenticated ID. Each subscriber gets their own filtered stream: when a task owned by ada is added, only Ada’s subscription emits a new list with it. The filter runs in the producer, so rows the caller must not see never enter the result at all. Never leave that filtering to the client.

A browser can subscribe over direct server-sent events or a direct WebSocket, and the @cratis/arc client subscribes through generated proxies over a multiplexed hub. Subscribe to an observable query shows each one.

defineObservableQuery takes an observe callback instead of a decorated method. This complete low-level example needs no build step on Node.js 22.19 or later, which strips types by default; run it with node observable.ts inside the workspace:

observable.ts
import express from 'express';
import { z } from 'zod';
import { ArcServer, defineObservableQuery } from '@cratis/arc.core';
import { BehaviorSubject } from 'rxjs';
import { cratisArc } from '@cratis/arc.express';
const numbers = new BehaviorSubject<number[]>([1]);
const query = defineObservableQuery({
name: 'Numbers',
schema: z.object({}),
observe: () => numbers
});
const server = new ArcServer({ observableQueries: [query] });
const app = express();
const adapter = cratisArc(server);
app.use(adapter);
const listener = app.listen(3000, '127.0.0.1');
adapter.injectWebSocket(listener);
let nextNumber = 2;
const timer = setInterval(() => numbers.next([nextNumber++]), 1000);
process.once('SIGINT', () => {
clearInterval(timer);
void (async () => {
try {
await adapter.close(listener); // Drain WebSockets and SSE, then close the listener.
} finally { await server.dispose(); }
})();
});

curl http://127.0.0.1:3000/api/numbers answers 200 with data: [1] or a later number. observe may resolve services through currentServices().

One route serves a snapshot and a stream. Return a source that has a current value when you have one, emit undefined for “nothing there”, and let errors be errors. Paging and sorting apply to every emission; authorization applies once, when the subscription opens, and emission guards cover what changes after that. Filter rows by the caller inside the source.

Subscribe to an observable query connects a browser or the @cratis/arc client to the route you just declared.