diff --git a/ContractTests/Client/package.json b/ContractTests/Client/package.json index afc57158..39dee106 100644 --- a/ContractTests/Client/package.json +++ b/ContractTests/Client/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.core-client-contract", - "version": "0.37.0", + "version": "0.38.0", "private": true, "type": "module", "dependencies": { diff --git a/Documentation/index.md b/Documentation/index.md index 237ac5fc..47062ec5 100644 --- a/Documentation/index.md +++ b/Documentation/index.md @@ -8,7 +8,7 @@ Arc for TypeScript is a Node.js server implementation of [Arc](/arc/), the Crati Without it, a Node.js backend for an Arc frontend means writing every route, request parser, validation response, and status code by hand, then keeping all of it in step with the frontend. With it, commands and queries run through one pipeline that owns those concerns, the wire behavior follows Arc on .NET, and the proxy generator writes the typed frontend client from your source. :::caution[Source preview, no full parity] -No package is published to npm; the manifests are at version 0.37.0 for a source preview. Arc for TypeScript does **not** have full parity with Arc on .NET, and package names and APIs can still change. The [capability reference](reference/capabilities.md) is the single place for status and evidence. +No package is published to npm; the manifests are at version 0.38.0 for a source preview. Arc for TypeScript does **not** have full parity with Arc on .NET, and package names and APIs can still change. The [capability reference](reference/capabilities.md) is the single place for status and evidence. ::: ## What it looks like diff --git a/Documentation/queries/observable-queries.md b/Documentation/queries/observable-queries.md index ff3ce02d..04e70e57 100644 --- a/Documentation/queries/observable-queries.md +++ b/Documentation/queries/observable-queries.md @@ -78,7 +78,7 @@ curl -N -H 'Accept: text/event-stream' \ 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. An observable query cannot return a `queryPage`; generated metadata rejects that declaration at build. For a collection too large to emit whole, narrow what the source emits with query arguments, or use a database integration that observes a query, such as [MongoDB change streams](../mongodb/observing-collections.md). The [paging rules](model-bound/paging.md#request-parameters) for invalid sizes and sort fields are the same as for snapshots. +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](model-bound/paging.md#request-parameters) for invalid sizes and sort fields are the same as for snapshots. ## Authorize a live query diff --git a/Documentation/reference/capabilities.md b/Documentation/reference/capabilities.md index 696730e2..f4d752e1 100644 --- a/Documentation/reference/capabilities.md +++ b/Documentation/reference/capabilities.md @@ -116,7 +116,8 @@ Evidence paths are relative to the repository root. Spec folders follow `for_ pack --out ` and install the tarballs with npm. Use `yarn pack`: it rewrites `workspace:^` dependencies to version ranges, and `npm pack` does not. `yarn check:consumers` installs packed packages this way to check NodeNext and Bundler consumers. diff --git a/Documentation/sql/getting-started.md b/Documentation/sql/getting-started.md index 8921f158..bb2237db 100644 --- a/Documentation/sql/getting-started.md +++ b/Documentation/sql/getting-started.md @@ -98,11 +98,12 @@ export class AddTask { @inject(drizzleDatabase()) handle(database: DrizzleHandle): void { database.native.insert(tasks).values({ id: this.id, title: this.title }).run(); + database.notifyChanged(tasks); } } ``` -`drizzleDatabase()` resolves to a `DrizzleHandle` whose `native` property is the tenant's Drizzle database, typed as you declare it. Register `AddTask` with `builder.add(...)` like the query. +`drizzleDatabase()` resolves to a `DrizzleHandle` whose `native` property is the tenant's Drizzle database, typed as you declare it. Register `AddTask` with `builder.add(...)` like the query. `notifyChanged(tasks)` validates the registered table; it publishes changes only when you opt in to [in-process observation](observing-tables.md). Inside the command, publication follows execution (including a transaction awaited inside it), not a transaction owned by a host or outer runner. For those, call it after commit. ## Load a read model in a command @@ -146,6 +147,7 @@ Queries take `drizzleReadModel(Model)`, commands take `drizzleDatabase()`. The r | `findOne(filter)` | The first match in primary-key order, or `undefined` | | `findById(key)` | The row matching the single primary key, or `null` | | `table` | The Drizzle table | +| `observe(filter?)`, `observePage(filter, options)`, `observeById(key)` | Experimental opt-in in-process SQL observations; [limits](observing-tables.md) | It has no write methods and does not expose the writable database. A filter is a Drizzle `SQL` expression such as `eq(tasks.title, 'a')`, built with bound parameters; never interpolate request input into SQL text. diff --git a/Documentation/sql/index.md b/Documentation/sql/index.md index f7faba4c..ebfe9ec4 100644 --- a/Documentation/sql/index.md +++ b/Documentation/sql/index.md @@ -17,6 +17,7 @@ Your read models live in SQL tables. Every query needs the right database for th | Store GUIDs, concepts, dates, times, durations, and JSON per dialect | [Column types](column-types.md) | | Map model fields to columns, generate and apply migrations, and choose database credentials | [Map the schema and own migrations](schema-and-migrations.md) | | Count, sort, and page in SQL | [Paging and sorting](paging.md) | +| Announce writes and observe tenant-scoped SQL results in process (Experimental) | [Observe tables](observing-tables.md) | | Route each tenant to its own database | [Tenancy](tenancy.md) | ## Why Drizzle @@ -28,7 +29,7 @@ Drizzle 0.45 is before 1.0. The peer dependency `^0.45.0` accepts 0.45 releases, ## What it does not do - **No schema management.** `withDrizzle` never creates tables, adds columns, or runs migrations. See [Own the schema](getting-started.md#own-the-schema). -- **No live queries.** There is no `observe()`, and Arc does not refresh an observable query when a table changes. SQLite has no cross-process change notification here, and Drizzle does not announce writes. PostgreSQL `LISTEN`/`NOTIFY` would need managed triggers, a listener connection per tenant, resubscription after reconnects, a race-free first read, and tested shutdown; none of that is included. If your application has a reliable, tenant-scoped change source of its own, an Arc [observable query](../queries/observable-queries.md) can consume it. An in-process event after a command write does not see changes made by other processes. +- **Only announced in-process changes.** Opt-in `DrizzleObservation.InProcess` observes registered tables after explicit `notifyChanged` calls. Other processes and unannounced writes are invisible; host-owned outer transactions are not covered by the command completion boundary. Observed pages are eventually consistent, not atomic count-and-row snapshots. PostgreSQL `LISTEN`/`NOTIFY` is not included. See [Observe tables](observing-tables.md). - **No transactions or change tracking.** There is no unit of work shared with command execution. Use a Drizzle transaction in your command when several writes must succeed together. - **One registration per application.** A second `withDrizzle` fails at build with a duplicate service. Several tenants use one registration with `databaseFactory`. - **No EF-only features.** Arc on .NET's Entity Framework integration also has SQL Server, spatial Point, LineString, and Polygon types, several DbContexts, and automatic concept conversion in `BaseDbContext`. None of those are part of this package. diff --git a/Documentation/sql/observing-tables.md b/Documentation/sql/observing-tables.md new file mode 100644 index 00000000..e4963522 --- /dev/null +++ b/Documentation/sql/observing-tables.md @@ -0,0 +1,85 @@ +--- +title: Observe SQL tables in process +description: Announce committed writes to tenant-scoped Drizzle read models and stream small SQL results. +--- + +**Experimental, in-process only.** Observation 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. + +## A task slice + +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*. + +```typescript title="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()) + handle(database: DrizzleHandle): 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) { + 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: + +```bash +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. + +## Timing, tenants, and limits + +`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](../queries/observable-queries.md#page-and-sort-a-live-list). 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](../mongodb/observing-collections.md) have a different change source. diff --git a/Documentation/sql/toc.yml b/Documentation/sql/toc.yml index dfa251d9..c2d55be6 100644 --- a/Documentation/sql/toc.yml +++ b/Documentation/sql/toc.yml @@ -8,5 +8,7 @@ href: schema-and-migrations.md - name: Paging and sorting href: paging.md +- name: Observe tables in process + href: observing-tables.md - name: Tenancy href: tenancy.md diff --git a/README.md b/README.md index c370a8a6..19766c47 100644 --- a/README.md +++ b/README.md @@ -54,7 +54,7 @@ export class TaskItem { | `@cratis/arc.chronicle` | [`Source/Chronicle`](Source/Chronicle) | **Experimental.** `builder.withChronicle` appends returned events and resolves registered read models by command key; nested command returns join one event-log batch. In-memory command assertions are available under `@cratis/arc.chronicle/testing`. SDK 6.14.0 imports natively and infers read models from projections/reducers; an opt-in kernel suite covers aggregate replay and reactor commands. Full .NET transaction parity remains unverified. | | `@cratis/cratis` | [`Source/Cratis`](Source/Cratis) | **Experimental source preview.** `CratisApplication.createBuilder()` and `builder.addCratis()` compose Arc and a Chronicle client without installing authentication; not yet published to npm. | -Every package manifest is at version 0.37.0. That is the version of this source preview, not an npm release, and the Chronicle package is experimental. The packages ship ES modules only, and schemas use Zod 4. The default core entry, host adapters, MongoDB, and Drizzle packages need Node.js 22 or later. The Fetch entry has a neutral bundle with `node:async_hooks` as its only Node import; its command, query, and SSE paths run in a Next.js App Router route handler on the Node.js runtime, with Bun and Deno smoke checks; Cloudflare Workers and the Next.js Edge runtime are not supported. See [Fetch API runtimes](Documentation/hosts/fetch-runtimes.md). The root workspace needs Node.js 22.19 or later, because it installs the Chronicle SDK; Node.js 24 LTS is recommended. +Every package manifest is at version 0.38.0. That is the version of this source preview, not an npm release, and the Chronicle package is experimental. The packages ship ES modules only, and schemas use Zod 4. The default core entry, host adapters, MongoDB, and Drizzle packages need Node.js 22 or later. The Fetch entry has a neutral bundle with `node:async_hooks` as its only Node import; its command, query, and SSE paths run in a Next.js App Router route handler on the Node.js runtime, with Bun and Deno smoke checks; Cloudflare Workers and the Next.js Edge runtime are not supported. See [Fetch API runtimes](Documentation/hosts/fetch-runtimes.md). The root workspace needs Node.js 22.19 or later, because it installs the Chronicle SDK; Node.js 24 LTS is recommended. ## Try it diff --git a/Source/Chronicle/package.json b/Source/Chronicle/package.json index 67feaebe..d07b1cdf 100644 --- a/Source/Chronicle/package.json +++ b/Source/Chronicle/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.chronicle", - "version": "0.37.0", + "version": "0.38.0", "publishConfig": { "access": "public" }, @@ -34,8 +34,8 @@ "README.md" ], "peerDependencies": { - "@cratis/arc.core": "^0.37.0", - "@cratis/arc.testing": "^0.37.0", + "@cratis/arc.core": "^0.38.0", + "@cratis/arc.testing": "^0.38.0", "@cratis/chronicle": "^6.7.0", "@cratis/fundamentals": "^7.19.6", "rxjs": "^7.8.2", diff --git a/Source/CodeAnalysis/package.json b/Source/CodeAnalysis/package.json index 17d6df50..640f1c6b 100644 --- a/Source/CodeAnalysis/package.json +++ b/Source/CodeAnalysis/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/eslint-plugin-arc-core", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "description": "ESLint diagnostics for Arc for TypeScript server artifacts", diff --git a/Source/Core/package.json b/Source/Core/package.json index 62516894..677ae79b 100644 --- a/Source/Core/package.json +++ b/Source/Core/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.core", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Core/queries/observable/for_ObservableQuerySession/when_streaming_a_low_level_page.ts b/Source/Core/queries/observable/for_ObservableQuerySession/when_streaming_a_low_level_page.ts new file mode 100644 index 00000000..f8e3b81e --- /dev/null +++ b/Source/Core/queries/observable/for_ObservableQuerySession/when_streaming_a_low_level_page.ts @@ -0,0 +1,41 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { should } from 'vitest'; +import { z } from 'zod'; +import { ArcServer } from '../../../ArcServer.js'; +import { queryPage } from '../../QueryPage.js'; +import { CurrentValueSubject } from '../CurrentValueSubject.js'; +import { defineObservableQuery } from '../defineObservableQuery.js'; + +should(); + +describe('when streaming a low-level observable page through SSE', () => { + let results: { data: { id: string }[]; paging: { totalItems: number; size: number } }[]; + + beforeEach(async () => { + const source = new CurrentValueSubject(queryPage([{ id: 'a' }], 1)); + const server = new ArcServer({ observableQueries: [defineObservableQuery({ + name: 'Items', schema: z.object({}), observe: () => source + })] }); + const response = await server.handle(new Request('http://localhost/api/items?page=0&pageSize=1', { + headers: { accept: 'text/event-stream' } + })); + const reader = response!.body!.getReader(); + const readFrame = async () => { + const frame = new TextDecoder().decode((await reader.read()).value); + return JSON.parse(frame.split('\n').find(line => line.startsWith('data:'))!.slice(5)) as typeof results[number]; + }; + try { + const first = await readFrame(); + source.next(queryPage([{ id: 'b' }], 2)); + results = [first, await readFrame()]; + } finally { await reader.cancel(); await server.dispose(); } + }); + + it('should deliver the initial and changed page with database totals', () => { + results.map(result => ({ data: result.data, total: result.paging.totalItems, size: result.paging.size })).should.deep.equal([ + { data: [{ id: 'a' }], total: 1, size: 1 }, + { data: [{ id: 'b' }], total: 2, size: 1 } + ]); + }); +}); diff --git a/Source/Cratis/package.json b/Source/Cratis/package.json index fc4d9523..9fe67377 100644 --- a/Source/Cratis/package.json +++ b/Source/Cratis/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/cratis", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "description": "Arc and experimental Chronicle composition for Node.js", diff --git a/Source/Drizzle/DrizzleChangeNotifications.ts b/Source/Drizzle/DrizzleChangeNotifications.ts new file mode 100644 index 00000000..b57a0d59 --- /dev/null +++ b/Source/Drizzle/DrizzleChangeNotifications.ts @@ -0,0 +1,69 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { AsyncLocalStorage } from 'node:async_hooks'; +import { getTableName } from 'drizzle-orm'; +import type { Table } from 'drizzle-orm'; + +/** Application-local, tenant-isolated change bus. Table keys are registered object identities. */ +export class DrizzleChangeNotifications { + readonly #listeners = new Map void>>>(); + readonly #frames = new AsyncLocalStorage; active: boolean }>>(); + constructor(private readonly tables: ReadonlyMap object, Table>, readonly enabled: boolean) {} + + /** Validate even in disabled mode, then queue within a command or publish immediately. */ + notify(tenant: string, targets: readonly (Table | (new () => object))[]): void { + const resolved = targets.map(target => { + const table = this.tables.get(target as new () => object) ?? target as Table; + if (![...this.tables.values()].includes(table)) { + let name: string; + try { name = typeof target === 'function' ? target.name : getTableName(target as Table); } + catch { name = target?.constructor?.name ?? 'unknown'; } + throw new Error(`Unknown Drizzle read-model table: ${name}`); + } + return table; + }); + if (!this.enabled) return; + const frame = this.#frames.getStore()?.get(tenant); + if (frame?.active) { + for (const table of resolved) frame.pending.add(table); + } else for (const table of new Set(resolved)) this.publish(tenant, table); + } + + /** Nested calls for this tenant join the ambient frame; independent calls never share one. */ + async run(tenant: string, execute: () => Promise): Promise { + const active = this.#frames.getStore(); + if (active?.get(tenant)?.active) return execute(); + const frame = { pending: new Set(), active: true }; + const frames = new Map(active); + frames.set(tenant, frame); + return this.#frames.run(frames, async () => { + try { return await execute(); } + finally { + frame.active = false; + for (const table of frame.pending) this.publish(tenant, table); + frame.pending.clear(); + } + }); + } + + listen(tenant: string, table: Table, listener: () => void): () => void { + let tables = this.#listeners.get(tenant); + if (!tables) { tables = new Map(); this.#listeners.set(tenant, tables); } + let listeners = tables.get(table); + if (!listeners) { listeners = new Set(); tables.set(table, listeners); } + listeners.add(listener); + return () => { + listeners.delete(listener); + if (!listeners.size) tables.delete(table); + if (!tables.size) this.#listeners.delete(tenant); + }; + } + + listenerCount(tenant: string): number { + return [...(this.#listeners.get(tenant)?.values() ?? [])].reduce((sum, listeners) => sum + listeners.size, 0); + } + + private publish(tenant: string, table: Table): void { + for (const listener of [...(this.#listeners.get(tenant)?.get(table) ?? [])]) listener(); + } +} diff --git a/Source/Drizzle/DrizzleHandle.ts b/Source/Drizzle/DrizzleHandle.ts index a55e93fe..0a3c28ed 100644 --- a/Source/Drizzle/DrizzleHandle.ts +++ b/Source/Drizzle/DrizzleHandle.ts @@ -1,8 +1,16 @@ // Copyright (c) Cratis. All rights reserved. // Licensed under the MIT license. See LICENSE file in the project root for full license information. import type { DrizzleDatabase } from './DrizzleDatabase.js'; +import type { DrizzleChangeNotifications } from './DrizzleChangeNotifications.js'; +import type { Table } from 'drizzle-orm'; /** Scoped reference to an application-owned connection or pool; Arc never disposes the native database. */ export class DrizzleHandle { - constructor(readonly native: T) {} + constructor(readonly native: T, private readonly notifications?: DrizzleChangeNotifications, private readonly tenant?: string) {} + + /** Announce writes to registered tables or read-model types. Call after a host-owned transaction commits. */ + notifyChanged(...targets: (Table | (new () => object))[]): void { + if (!this.notifications || !this.tenant) throw new Error('Drizzle change notification requires withDrizzle'); + this.notifications.notify(this.tenant, targets); + } } diff --git a/Source/Drizzle/DrizzleObservable.ts b/Source/Drizzle/DrizzleObservable.ts new file mode 100644 index 00000000..61f6ab58 --- /dev/null +++ b/Source/Drizzle/DrizzleObservable.ts @@ -0,0 +1,51 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { Observable } from 'rxjs'; +import { DrizzleObservationSession } from './DrizzleObservationSession.js'; + +/** Experimental lazy SQL observation; call current() for a snapshot or subscribe for changes. */ +export class DrizzleObservable extends Observable { + #prime?: DrizzleObservationSession; + #closed = false; + readonly #sessions = new Set>(); + + constructor(private readonly open: (onClose: () => void) => DrizzleObservationSession, + private readonly canStart: () => void, private readonly onActive: () => void = () => {}, + private readonly onIdle: () => void = () => {}) { + super(subscriber => { + if (this.#closed) { subscriber.complete(); return; } + try { this.canStart(); } + catch (error) { subscriber.error(error); return; } + const session = this.#prime ?? this.start(); + this.#prime = undefined; + session.adopt(subscriber); + return () => session.close(); + }); + } + + private start(): DrizzleObservationSession { + this.canStart(); + const session = this.open(() => { + this.#sessions.delete(session); + if (!this.#sessions.size) this.onIdle(); + }); + this.#sessions.add(session); + this.onActive(); + return session; + } + + /** Prime one listener before reading; the first subscription adopts that snapshot. */ + current(): Promise<{ hasValue: true; value: T }> { + if (this.#closed) return Promise.reject(new Error('Drizzle observation was closed')); + try { this.#prime ??= this.start(); return this.#prime.current(); } + catch (error) { return Promise.reject(error); } + } + + /** Release an unused prime or active subscriptions. */ + close(): void { + if (this.#closed) return; + this.#closed = true; + for (const session of [...this.#sessions]) session.close(); + this.#prime = undefined; + } +} diff --git a/Source/Drizzle/DrizzleObservation.ts b/Source/Drizzle/DrizzleObservation.ts new file mode 100644 index 00000000..a11705f5 --- /dev/null +++ b/Source/Drizzle/DrizzleObservation.ts @@ -0,0 +1,7 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +/** Available SQL observation modes. Experimental: in-process notifications only. */ +export enum DrizzleObservation { + InProcess = 'in-process' +} diff --git a/Source/Drizzle/DrizzleObservationSession.ts b/Source/Drizzle/DrizzleObservationSession.ts new file mode 100644 index 00000000..c06a9cff --- /dev/null +++ b/Source/Drizzle/DrizzleObservationSession.ts @@ -0,0 +1,116 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { AsyncResource } from 'node:async_hooks'; +import type { Subscriber } from 'rxjs'; + +/** Single-flight observation with an adoptable, retained initial snapshot. */ +export class DrizzleObservationSession { + #release?: () => void; + #subscriber?: Subscriber; + #dirty = false; + #running = false; + #closed = false; + #hasValue = false; + #scheduled = false; + #released = false; + readonly #readContext = new AsyncResource('DrizzleObservation'); + readonly #cancellations = new Set<(reason: Error) => void>(); + readonly #initial: Promise; + + constructor(private readonly read: () => Promise, listen: (changed: () => void) => () => void, + private readonly signal: AbortSignal | undefined, private readonly onClose: () => void) { + if (signal?.aborted) { + this.#closed = true; + this.#initial = Promise.reject(new Error('Drizzle observation was closed')); + void this.#initial.catch(() => {}); + return; + } + this.#release = listen(() => { + if (this.#closed) return; + this.#dirty = true; + if (this.#subscriber) this.schedule(); + }); + signal?.addEventListener('abort', this.abort, { once: true }); + this.#initial = this.read().then(value => { + if (this.#closed) throw new Error('Drizzle observation was closed'); + this.#hasValue = true; + return value; + }, error => { + if (!this.#closed) this.release(); + throw error; + }); + void this.#initial.catch(() => {}); + } + + get closed(): boolean { return this.#closed; } + + async current(): Promise<{ hasValue: true; value: T }> { + if (this.#closed) throw new Error('Drizzle observation was closed'); + let cancel!: (reason: Error) => void; + const canceled = new Promise((_, reject) => { cancel = reject; }); + this.#cancellations.add(cancel); + try { + const value = await Promise.race([this.#initial, canceled]); + if (this.#closed) throw new Error('Drizzle observation was closed'); + return { hasValue: true, value }; + } finally { this.#cancellations.delete(cancel); } + } + + adopt(subscriber: Subscriber): void { + if (this.#closed) { subscriber.complete(); return; } + this.#subscriber = subscriber; + void this.#initial.then(value => { + if (this.#closed || subscriber.closed) return; + subscriber.next(value); + if (this.#dirty) this.schedule(); + }, error => { + if (!this.#closed && !subscriber.closed) subscriber.error(error); + this.close(); + }); + } + + private schedule(): void { + if (this.#scheduled || this.#running || this.#closed) return; + this.#scheduled = true; + setImmediate(() => { + this.#scheduled = false; + if (!this.#closed) this.#readContext.runInAsyncScope(() => { void this.pump(); }); + }); + } + + private async pump(): Promise { + if (this.#running || this.#closed || !this.#hasValue || !this.#subscriber) return; + this.#running = true; + try { + while (this.#dirty && !this.#closed) { + this.#dirty = false; + const value = await this.read(); + if (this.#closed) break; + this.#subscriber?.next(value); + } + } catch (error) { + if (!this.#closed) { this.#subscriber?.error(error); this.close(); } + } finally { this.#running = false; } + } + + private readonly abort = (): void => { this.close(); }; + + private release(): void { + if (this.#released) return; + this.#released = true; + this.#release?.(); + this.#release = undefined; + this.signal?.removeEventListener('abort', this.abort); + this.#readContext.emitDestroy(); + this.onClose(); + } + + close(): void { + if (this.#closed) return; + this.#closed = true; + this.release(); + for (const cancel of this.#cancellations) cancel(new Error('Drizzle observation was closed')); + this.#cancellations.clear(); + this.#subscriber?.complete(); + } +} diff --git a/Source/Drizzle/DrizzleOptions.ts b/Source/Drizzle/DrizzleOptions.ts index fda3fed2..87b356ce 100644 --- a/Source/Drizzle/DrizzleOptions.ts +++ b/Source/Drizzle/DrizzleOptions.ts @@ -4,6 +4,7 @@ import { DrizzleDialect } from './DrizzleDialect.js'; import type { ExecutionContext } from '@cratis/arc.core'; import type { Table } from 'drizzle-orm'; import type { DrizzleDatabase } from './DrizzleDatabase.js'; +import type { DrizzleObservation } from './DrizzleObservation.js'; /** Storage location is selected for each Arc execution, never from a caller-supplied query argument. */ export interface DrizzleOptions { @@ -14,4 +15,6 @@ export interface DrizzleOptions { databaseFactory?: (tenant: string, context: ExecutionContext) => DrizzleDatabase | Promise; readModels?: readonly { type: new () => object; table: Table }[]; maxPageSize?: number; + /** Experimental: receive only changes explicitly announced in this process. */ + observation?: DrizzleObservation; } diff --git a/Source/Drizzle/DrizzleReadModels.ts b/Source/Drizzle/DrizzleReadModels.ts index 8707c219..3b4747f9 100644 --- a/Source/Drizzle/DrizzleReadModels.ts +++ b/Source/Drizzle/DrizzleReadModels.ts @@ -8,6 +8,9 @@ import type { QueryOptions, QueryPage } from '@cratis/arc.core'; import type { DrizzleDatabase } from './DrizzleDatabase.js'; import type { DrizzleFilter } from './DrizzleFilter.js'; import { DrizzleModelCodec } from './DrizzleModelCodec.js'; +import type { DrizzleChangeNotifications } from './DrizzleChangeNotifications.js'; +import { DrizzleObservable } from './DrizzleObservable.js'; +import { DrizzleObservationSession } from './DrizzleObservationSession.js'; // Each Drizzle driver implements these query-builder operations and maps selected columns on execute. type SelectQuery = { @@ -27,14 +30,19 @@ export class DrizzleReadModels { private readonly codec: DrizzleModelCodec; private readonly keys: Column[]; private readonly columns: Record; + readonly #observations = new Set>(); + #disposed = false; constructor(private readonly database: DrizzleDatabase, readonly table: Table, type: new () => T, - private readonly maxPageSize = 100, codec?: DrizzleModelCodec) { + private readonly maxPageSize = 100, codec?: DrizzleModelCodec, + private readonly notifications?: DrizzleChangeNotifications, private readonly tenant?: string, + private readonly signal?: AbortSignal) { if (!Number.isSafeInteger(maxPageSize) || maxPageSize <= 0 || maxPageSize > 10000) throw new RangeError('maxPageSize must be between 1 and 10000'); this.columns = getTableColumns(table); this.keys = Object.values(this.columns).filter(column => column.primary); if (!this.keys.length) throw new Error('A Drizzle read model requires a primary key for stable paging'); this.codec = codec ?? new DrizzleModelCodec(type, this.columns); + signal?.addEventListener('abort', this.abort, { once: true }); } private get db(): ReadDatabase { return this.database as ReadDatabase; } @@ -55,6 +63,70 @@ export class DrizzleReadModels { return rows.map(row => this.codec.deserialize(row)); } + private validatePage(options: QueryOptions): { page: number; pageSize: number } { + const { page, pageSize } = options.paging ?? { page: 0, pageSize: this.maxPageSize }; + if (!Number.isSafeInteger(page) || page < 0 || !Number.isSafeInteger(pageSize) || pageSize < 1 || + !Number.isSafeInteger(page * pageSize)) throw new RangeError('Invalid Drizzle page'); + if (pageSize > this.maxPageSize) throw new QueryPagingRequired(this.maxPageSize); + this.order(options.sorting); + return { page, pageSize }; + } + + private observeWith(read: () => Promise): DrizzleObservable { + const canStart = (): void => { + if (this.#disposed || this.signal?.aborted) throw new Error('Drizzle read models have been disposed'); + if (!this.notifications?.enabled || !this.tenant) throw new Error( + 'Drizzle observation is not enabled; set observation: DrizzleObservation.InProcess in withDrizzle'); + }; + const observable = new DrizzleObservable(onClose => new DrizzleObservationSession(read, + changed => this.notifications!.listen(this.tenant!, this.table, changed), this.signal, onClose), canStart, + () => this.#observations.add(observable), () => this.#observations.delete(observable)); + return observable; + } + + /** Experimental: observe a complete small result; overflow fails instead of emitting a truncated list. */ + observe(filter?: DrizzleFilter): DrizzleObservable { + return this.observeWith(async () => { + const items = await this.select(filter, undefined, this.maxPageSize + 1); + if (items.length > this.maxPageSize) throw new QueryPagingRequired(this.maxPageSize, true, + `The result exceeds the maximum of ${this.maxPageSize} items; narrow the filter, or use observePage through defineObservableQuery for paged results`); + return items; + }); + } + + /** Experimental: observe a fixed SQL-sorted page. Count and rows converge after announced concurrent writes. */ + observePage(filter: DrizzleFilter, options: QueryOptions): DrizzleObservable> { + const { page, pageSize } = this.validatePage(options); + return this.observeWith(async () => { + for (let attempt = 0; attempt < 3; attempt++) { + const total = await this.db.$count(this.table, filter); + if (!Number.isSafeInteger(total) || total < 0) throw new Error('Invalid Drizzle count'); + if (!options.paging && total > this.maxPageSize) throw new QueryPagingRequired(this.maxPageSize, true); + const items = await this.select(filter, options.sorting, pageSize, page * pageSize); + if (items.length === Math.min(pageSize, Math.max(0, total - page * pageSize))) + return queryPage(items, total, options.sorting); + } + throw new Error('Drizzle observation could not read a consistent page'); + }); + } + + /** Experimental: emit null for a missing or deleted single-key row. */ + observeById(key: string): DrizzleObservable { + if (this.keys.length !== 1) throw new Error('Drizzle command read models require a single primary key'); + return this.observeWith(() => this.findById(key)); + } + + private readonly abort = (): void => { void this[Symbol.asyncDispose](); }; + + /** End and release all observations owned by this read-model scope. */ + async [Symbol.asyncDispose](): Promise { + if (this.#disposed) return; + this.#disposed = true; + this.signal?.removeEventListener('abort', this.abort); + for (const observation of this.#observations) observation.close(); + this.#observations.clear(); + } + /** Return at most maxPageSize matches, never an unbounded result. */ async find(filter: DrizzleFilter, sorting?: QueryOptions['sorting']): Promise { this.order(sorting); @@ -80,11 +152,7 @@ export class DrizzleReadModels { /** Push a typed predicate, count, sort and page into SQL; never load the unbounded result before paging. */ async queryPage(filter: DrizzleFilter, options: QueryOptions): Promise> { - const { page, pageSize } = options.paging ?? { page: 0, pageSize: this.maxPageSize }; - if (!Number.isSafeInteger(page) || page < 0 || !Number.isSafeInteger(pageSize) || pageSize < 1 || - !Number.isSafeInteger(page * pageSize)) throw new RangeError('Invalid Drizzle page'); - if (pageSize > this.maxPageSize) throw new QueryPagingRequired(this.maxPageSize); - this.order(options.sorting); + const { page, pageSize } = this.validatePage(options); const total = await this.db.$count(this.table, filter); if (!Number.isSafeInteger(total) || total < 0) throw new Error('Invalid Drizzle count'); if (!options.paging && total > this.maxPageSize) throw new QueryPagingRequired(this.maxPageSize, true); diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/given/a_notification_bus.ts b/Source/Drizzle/for_DrizzleChangeNotifications/given/a_notification_bus.ts new file mode 100644 index 00000000..d9dfefd1 --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/given/a_notification_bus.ts @@ -0,0 +1,13 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { sqliteTable, text } from 'drizzle-orm/sqlite-core'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; + +export class Task {} +export const table = sqliteTable('tasks', { id: text('id').primaryKey() }); +export class a_notification_bus { + bus = new DrizzleChangeNotifications(new Map([[Task, table]]), true); + hits: string[] = []; + constructor() { this.bus.listen('a', table, () => this.hits.push('a')); + this.bus.listen('b', table, () => this.hits.push('b')); } +} diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/when_a_command_fails_after_notifying.ts b/Source/Drizzle/for_DrizzleChangeNotifications/when_a_command_fails_after_notifying.ts new file mode 100644 index 00000000..7aebb630 --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/when_a_command_fails_after_notifying.ts @@ -0,0 +1,13 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../given.js'; +import { Task, a_notification_bus } from './given/a_notification_bus.js'; + +describe('when a command fails after announcing a write', given(a_notification_bus, context => { + beforeEach(async () => { + context.hits.length = 0; + try { await context.bus.run('a', async () => { context.bus.notify('a', [Task]); throw Error('failed'); }); } + catch { /* A nontransactional write may have persisted. */ } + }); + it('should publish the announcement despite the command failure', () => { context.hits.should.deep.equal(['a']); }); +})); diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/when_an_unregistered_table_notifies.ts b/Source/Drizzle/for_DrizzleChangeNotifications/when_an_unregistered_table_notifies.ts new file mode 100644 index 00000000..74dddadf --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/when_an_unregistered_table_notifies.ts @@ -0,0 +1,21 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { sqliteTable, text } from 'drizzle-orm/sqlite-core'; +import { DrizzleHandle } from '../DrizzleHandle.js'; +import { given } from '../given.js'; +import { a_notification_bus, table } from './given/a_notification_bus.js'; + +describe('when a different table object with the same name is announced', given(a_notification_bus, context => { + let notify: () => void; + beforeEach(() => { + const handle = new DrizzleHandle({}, context.bus, 'a'); + const other = sqliteTable('tasks', { id: text('id').primaryKey() }); + notify = () => handle.notifyChanged(other); + }); + it('should reject the unregistered identity', () => { notify.should.throw('Unknown Drizzle read-model table: tasks'); }); +})); + +describe('when a registered table is announced outside a command', given(a_notification_bus, context => { + beforeEach(() => { context.hits.length = 0; new DrizzleHandle({}, context.bus, 'a').notifyChanged(table); }); + it('should publish immediately', () => { context.hits.should.deep.equal(['a']); }); +})); diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/when_independent_commands_overlap.ts b/Source/Drizzle/for_DrizzleChangeNotifications/when_independent_commands_overlap.ts new file mode 100644 index 00000000..43c4d484 --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/when_independent_commands_overlap.ts @@ -0,0 +1,20 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../given.js'; +import { Task, a_notification_bus } from './given/a_notification_bus.js'; + +describe('when two independent commands for the same tenant overlap', given(a_notification_bus, context => { + let whileFirstIsBlocked: string[]; + beforeEach(async () => { + context.hits.length = 0; + let unblock!: () => void; + const blocked = new Promise(resolve => { unblock = resolve; }); + const first = context.bus.run('a', async () => { context.bus.notify('a', [Task]); await blocked; }); + await context.bus.run('a', async () => { context.bus.notify('a', [Task]); }); + whileFirstIsBlocked = [...context.hits]; + unblock(); + await first; + }); + it('should flush the second command independently', () => { whileFirstIsBlocked.should.deep.equal(['a']); }); + it('should flush the first command when it finishes', () => { context.hits.should.deep.equal(['a', 'a']); }); +})); diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/when_nesting_different_tenants.ts b/Source/Drizzle/for_DrizzleChangeNotifications/when_nesting_different_tenants.ts new file mode 100644 index 00000000..03fc851d --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/when_nesting_different_tenants.ts @@ -0,0 +1,21 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../given.js'; +import { Task, a_notification_bus } from './given/a_notification_bus.js'; + +describe('when a command for tenant B nests inside a command for tenant A', given(a_notification_bus, context => { + let beforeACompletes: string[]; + beforeEach(async () => { + context.hits.length = 0; + await context.bus.run('a', async () => { + context.bus.notify('a', [Task]); + await context.bus.run('b', async () => { + context.bus.notify('a', [Task]); + context.bus.notify('b', [Task]); + }); + beforeACompletes = [...context.hits]; + }); + }); + it('should not publish A while B completes', () => { beforeACompletes.should.deep.equal(['b']); }); + it('should publish A once at its own completion', () => { context.hits.should.deep.equal(['b', 'a']); }); +})); diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/when_notifying_from_a_late_timer.ts b/Source/Drizzle/for_DrizzleChangeNotifications/when_notifying_from_a_late_timer.ts new file mode 100644 index 00000000..549e5fcf --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/when_notifying_from_a_late_timer.ts @@ -0,0 +1,24 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../given.js'; +import { Task, a_notification_bus } from './given/a_notification_bus.js'; + +describe('when a timer inherited from a completed command announces a change', given(a_notification_bus, context => { + let afterCommand: string[]; + beforeEach(async () => { + context.hits.length = 0; + let notifyLater!: () => void; + const late = new Promise(resolve => { notifyLater = resolve; }); + const timer = new Promise(resolve => { + void context.bus.run('a', async () => { + setTimeout(() => { void late.then(() => { context.bus.notify('a', [Task]); resolve(); }); }, 0); + }); + }); + await new Promise(resolve => setImmediate(resolve)); + afterCommand = [...context.hits]; + notifyLater(); + await timer; + }); + it('should not publish before the timer fires', () => { afterCommand.should.deep.equal([]); }); + it('should publish immediately after completion', () => { context.hits.should.deep.equal(['a']); }); +})); diff --git a/Source/Drizzle/for_DrizzleChangeNotifications/when_parallel_children_notify.ts b/Source/Drizzle/for_DrizzleChangeNotifications/when_parallel_children_notify.ts new file mode 100644 index 00000000..dee12eff --- /dev/null +++ b/Source/Drizzle/for_DrizzleChangeNotifications/when_parallel_children_notify.ts @@ -0,0 +1,18 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../given.js'; +import { Task, a_notification_bus } from './given/a_notification_bus.js'; + +describe('when parallel nested children announce the same table', given(a_notification_bus, context => { + let beforeParentCompletes: string[]; + beforeEach(async () => { + context.hits.length = 0; + await context.bus.run('a', async () => { + await Promise.all([context.bus.run('a', async () => { await Promise.resolve(); context.bus.notify('a', [Task]); }), + context.bus.run('a', async () => { await Promise.resolve(); context.bus.notify('a', [Task]); })]); + beforeParentCompletes = [...context.hits]; + }); + }); + it('should keep both notifications buffered until the parent completes', () => { beforeParentCompletes.should.deep.equal([]); }); + it('should publish only once', () => { context.hits.should.deep.equal(['a']); }); +})); diff --git a/Source/Drizzle/for_DrizzleHandle/when_notifying_changes/with_nested_server_commands.ts b/Source/Drizzle/for_DrizzleHandle/when_notifying_changes/with_nested_server_commands.ts new file mode 100644 index 00000000..4b6deb51 --- /dev/null +++ b/Source/Drizzle/for_DrizzleHandle/when_notifying_changes/with_nested_server_commands.ts @@ -0,0 +1,65 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { describe, it, should } from 'vitest'; +import { command, inject, ArcApplication, Severity } from '@cratis/arc.core'; +import { DrizzleDialect, DrizzleObservation, drizzleDatabase } from '../../index.js'; +import type { DrizzleHandle } from '../../DrizzleHandle.js'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import { a_sqlite_database } from '../../for_DrizzleReadModels/given/a_sqlite_database.js'; +import { TaskRecord } from '../../for_DrizzleReadModels/given/TaskRecord.js'; + +should(); +let table: a_sqlite_database['table']; +@command() class InnerNotification { + @inject(drizzleDatabase()) + handle(handle: DrizzleHandle): void { handle.notifyChanged(table); } +} +@command() class FailingNotification { + @inject(drizzleDatabase()) + handle(handle: DrizzleHandle): void { handle.notifyChanged(TaskRecord); throw Error('write failed later'); } +} +@command() class OuterNotification { + @inject(drizzleDatabase()) + handle(handle: DrizzleHandle): void { handle.notifyChanged(TaskRecord); } +} + +describe('when notifying from nested server commands', () => { + it('should publish once after the outermost execution finishes', async () => { + const fixture = new a_sqlite_database(); + await fixture.establish(); + table = fixture.table; + const builder = ArcApplication.createBuilder(); + builder.withDrizzle({ dialect: DrizzleDialect.SQLite, database: fixture.database, + readModels: [{ type: TaskRecord, table }], observation: DrizzleObservation.InProcess }); + let deliveries = 0; + builder.addCommandExecutionRunner(async (context, execute) => { + if (context.operationName?.endsWith('OuterNotification')) { + const identity = { tenantId: context.tenantId, correlationId: context.correlationId, + principal: context.principal, allowedSeverity: context.allowedSeverity, signal: context.signal }; + const nested = await application.server.executeCommand('InnerNotification', {}, identity); + nested.isSuccess.should.equal(true); + deliveries.should.equal(0); + } + return execute(); + }); + builder.add(InnerNotification, OuterNotification, FailingNotification); + const application = await builder.build(); + const scope = application.server.services.createScope({ tenantId: 'default', principal: undefined, + allowedSeverity: Severity.Warning, correlationId: crypto.randomUUID(), signal: new AbortController().signal }); + try { + const notifications = await scope.resolve(DrizzleChangeNotifications); + const release = notifications.listen('default', table, () => { deliveries++; }); + const outcome = await application.server.executeCommand('OuterNotification', {}, { + tenantId: 'DeFaUlT', principal: undefined, allowedSeverity: Severity.Warning, + correlationId: crypto.randomUUID(), signal: new AbortController().signal }); + outcome.isSuccess.should.equal(true); + deliveries.should.equal(1); + const failed = await application.server.executeCommand('FailingNotification', {}, { + tenantId: 'default', principal: undefined, allowedSeverity: Severity.Warning, + correlationId: crypto.randomUUID(), signal: new AbortController().signal }); + failed.isSuccess.should.equal(false); + deliveries.should.equal(2); + release(); + } finally { await scope.dispose(); await application.dispose(); fixture.close(); } + }); +}); diff --git a/Source/Drizzle/for_DrizzleHandle/when_notifying_without_withDrizzle.ts b/Source/Drizzle/for_DrizzleHandle/when_notifying_without_withDrizzle.ts new file mode 100644 index 00000000..705bec45 --- /dev/null +++ b/Source/Drizzle/for_DrizzleHandle/when_notifying_without_withDrizzle.ts @@ -0,0 +1,12 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { DrizzleHandle } from '../DrizzleHandle.js'; +import { TaskRecord } from '../for_DrizzleReadModels/given/TaskRecord.js'; + +describe('when notifying from a directly constructed Drizzle handle', () => { + let notify: () => void; + beforeEach(() => { notify = () => new DrizzleHandle({}).notifyChanged(TaskRecord); }); + it('should require withDrizzle instead of silently ignoring the change', () => { + notify.should.throw('Drizzle change notification requires withDrizzle'); + }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/given/AddObservedTask.ts b/Source/Drizzle/for_DrizzleReadModels/given/AddObservedTask.ts new file mode 100644 index 00000000..b71a4143 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/given/AddObservedTask.ts @@ -0,0 +1,21 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { SQLJsDatabase } from 'drizzle-orm/sql-js'; +import { field, Guid } from '@cratis/fundamentals'; +import { command, inject, key } from '@cratis/arc.core'; +import { drizzleDatabase } from '../../drizzleToken.js'; +import type { DrizzleHandle } from '../../DrizzleHandle.js'; +import { taskTable } from './a_sqlite_database.js'; + +/** HTTP fixture uses the same registered table object for writes and reads. */ +@command() +export class AddObservedTask { + @field(Guid) @key() id!: Guid; + @field(String) title!: string; + + @inject(drizzleDatabase()) + handle(database: DrizzleHandle): void { + database.native.insert(taskTable).values({ id: this.id, title: this.title }).run(); + database.notifyChanged(taskTable); + } +} diff --git a/Source/Drizzle/for_DrizzleReadModels/given/ObservedTaskQueries.ts b/Source/Drizzle/for_DrizzleReadModels/given/ObservedTaskQueries.ts new file mode 100644 index 00000000..6e81d137 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/given/ObservedTaskQueries.ts @@ -0,0 +1,15 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { query, queryOptions, readModel, service } from '@cratis/arc.core'; +import type { QueryOptions } from '@cratis/arc.core'; +import { drizzleReadModel } from '../../drizzleToken.js'; +import type { DrizzleReadModels } from '../../DrizzleReadModels.js'; +import { TaskRecord } from './TaskRecord.js'; + +@readModel() +export class ObservedTaskQueries { + @query({ observable: true }, service(drizzleReadModel(TaskRecord)), queryOptions()) + static page(tasks: DrizzleReadModels, options: QueryOptions) { + return tasks.observePage(undefined, options); + } +} diff --git a/Source/Drizzle/for_DrizzleReadModels/given/a_sqlite_database.ts b/Source/Drizzle/for_DrizzleReadModels/given/a_sqlite_database.ts index d9c5e50a..63617777 100644 --- a/Source/Drizzle/for_DrizzleReadModels/given/a_sqlite_database.ts +++ b/Source/Drizzle/for_DrizzleReadModels/given/a_sqlite_database.ts @@ -9,13 +9,15 @@ import { Guid } from '@cratis/fundamentals'; import { guidCodec } from '../../ColumnCodec.js'; import { sqliteColumn } from '../../columns.js'; +export const taskTable = sqliteTable('tasks', { + id: sqliteColumn(guidCodec(DrizzleDialect.SQLite))('id').primaryKey(), + title: text('title').notNull() +}); + /** An isolated in-memory SQLite database running in WebAssembly (no native binding). */ export class a_sqlite_database { native!: Database; - readonly table = sqliteTable('tasks', { - id: sqliteColumn(guidCodec(DrizzleDialect.SQLite))('id').primaryKey(), - title: text('title').notNull() - }); + readonly table = taskTable; database!: ReturnType; async establish() { const SQL = await initSqlJs(); diff --git a/Source/Drizzle/for_DrizzleReadModels/given/an_observation_scope.ts b/Source/Drizzle/for_DrizzleReadModels/given/an_observation_scope.ts new file mode 100644 index 00000000..b3d9a39f --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/given/an_observation_scope.ts @@ -0,0 +1,35 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { ArcApplication, Severity } from '@cratis/arc.core'; +import { DrizzleDialect } from '../../DrizzleDialect.js'; +import { DrizzleObservation } from '../../DrizzleObservation.js'; +import type { DrizzleHandle } from '../../DrizzleHandle.js'; +import type { DrizzleReadModels } from '../../DrizzleReadModels.js'; +import { drizzleDatabase, drizzleReadModel } from '../../drizzleToken.js'; +import '../../index.js'; +import { a_sqlite_database } from './a_sqlite_database.js'; +import { TaskRecord } from './TaskRecord.js'; + +export class an_observation_scope { + fixture = new a_sqlite_database(); + app!: Awaited['build']>>; + scope!: ReturnType; + models!: DrizzleReadModels; + handle!: DrizzleHandle; + async establish(maxPageSize?: number): Promise { + await this.fixture.establish(); + const builder = ArcApplication.createBuilder(); + builder.withDrizzle({ dialect: DrizzleDialect.SQLite, database: this.fixture.database, + readModels: [{ type: TaskRecord, table: this.fixture.table }], observation: DrizzleObservation.InProcess, + maxPageSize }); + this.app = await builder.build(); + this.scope = this.app.server.services.createScope({ tenantId: 'DeFaUlT', principal: undefined, + allowedSeverity: Severity.Warning, signal: new AbortController().signal, correlationId: crypto.randomUUID() }); + this.models = await this.scope.resolve(drizzleReadModel(TaskRecord)); + this.handle = await this.scope.resolve(drizzleDatabase()); + } + async dispose(): Promise { await this.scope?.dispose(); await this.app?.dispose(); this.fixture.close(); } +} +export async function settle(): Promise { + for (let index = 0; index < 8; index++) await new Promise(resolve => setImmediate(resolve)); +} diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/given/a_gated_observation.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/given/a_gated_observation.ts new file mode 100644 index 00000000..50066eba --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/given/a_gated_observation.ts @@ -0,0 +1,30 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { sqliteTable, text } from 'drizzle-orm/sqlite-core'; +import { DrizzleChangeNotifications } from '../../../DrizzleChangeNotifications.js'; +import { DrizzleObservable } from '../../../DrizzleObservable.js'; +import { DrizzleObservationSession } from '../../../DrizzleObservationSession.js'; + +class Item {} +const table = sqliteTable('items', { id: text('id').primaryKey() }); +export function deferred() { + let resolve!: (value: number) => void; + let reject!: (reason: Error) => void; + const promise = new Promise((yes, no) => { resolve = yes; reject = no; }); + return { promise, resolve, reject }; +} +export async function flush(): Promise { + for (let index = 0; index < 6; index++) await new Promise(resolve => setImmediate(resolve)); +} +export class a_gated_observation { + bus = new DrizzleChangeNotifications(new Map([[Item, table]]), true); + reads = 0; + observation?: DrizzleObservable; + open(read: () => Promise): DrizzleObservable { + return this.observation = new DrizzleObservable(closed => new DrizzleObservationSession( + () => { this.reads++; return read(); }, notify => this.bus.listen('tenant', table, notify), undefined, closed), () => {}); + } + notify(): void { this.bus.notify('tenant', [Item]); } + close(): void { this.observation?.close(); } + get listeners(): number { return this.bus.listenerCount('tenant'); } +} diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_a_prime_fails_before_adoption.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_a_prime_fails_before_adoption.ts new file mode 100644 index 00000000..61c6e558 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_a_prime_fails_before_adoption.ts @@ -0,0 +1,28 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { sqliteTable, text } from 'drizzle-orm/sqlite-core'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import { DrizzleObservable } from '../../DrizzleObservable.js'; +import { DrizzleObservationSession } from '../../DrizzleObservationSession.js'; + +class Task {} +const table = sqliteTable('tasks', { id: text('id').primaryKey() }); + +describe('when a prime fails and is never adopted', () => { + let error: string; + let listeners: number; + let idle: number; + beforeEach(async () => { + const bus = new DrizzleChangeNotifications(new Map([[Task, table]]), true); + idle = 0; + const observation = new DrizzleObservable(closed => new DrizzleObservationSession( + () => Promise.reject(new Error('read failed')), changed => bus.listen('tenant', table, changed), undefined, closed), + () => {}, () => {}, () => { idle++; }); + error = await observation.current().then(() => '', (reason: Error) => reason.message); + listeners = bus.listenerCount('tenant'); + // Do not adopt or close: the failed prime must leave owner tracking on its own. + }); + it('should retain the read failure for current', () => { error.should.equal('read failed'); }); + it('should release the listener', () => { listeners.should.equal(0); }); + it('should notify its owner when the failed prime terminates', () => { idle.should.equal(1); }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_a_websocket_hub_disconnects.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_a_websocket_hub_disconnects.ts new file mode 100644 index 00000000..750267e6 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_a_websocket_hub_disconnects.ts @@ -0,0 +1,42 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { z } from 'zod'; +import { ArcServer } from '../../../Core/ArcServer.js'; +import { defineObservableQuery } from '../../../Core/queries/observable/defineObservableQuery.js'; +import { HubConnection } from '../../../Core/queries/observable/HubConnection.js'; +import { HubSubscriptionOutcome } from '../../../Core/queries/observable/HubSubscriptionOutcome.js'; +import { Severity } from '../../../Core/validation/Severity.js'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import { DrizzleReadModels } from '../../DrizzleReadModels.js'; +import { TaskRecord } from '../given/TaskRecord.js'; +import { a_sqlite_database } from '../given/a_sqlite_database.js'; + +describe('when a WebSocket hub disconnects from a SQL observation', () => { + let before: number; + let after: number; + let outcome: HubSubscriptionOutcome; + beforeEach(async () => { + const fixture = new a_sqlite_database(); + await fixture.establish(); + const bus = new DrizzleChangeNotifications(new Map([[TaskRecord, fixture.table]]), true); + const models = new DrizzleReadModels(fixture.database, fixture.table, TaskRecord, 100, undefined, bus, 'default'); + const server = new ArcServer({ observableQueries: [defineObservableQuery({ name: 'Tasks', schema: z.object({}), + observe: () => models.observe() })] }); + const controller = new AbortController(); + const output = { signal: controller.signal, lastActivity: Date.now(), send: async () => {}, + close: () => controller.abort() }; + const connection = new HubConnection(server, 'WebSocket', output, + { tenantId: 'default', signal: controller.signal, allowedSeverity: Severity.Warning, + correlationId: crypto.randomUUID(), principal: undefined }, 0, () => {}, () => {}); + try { + await connection.connect(); + outcome = await connection.subscribe('tasks', 1, { queryName: 'Tasks' }); + before = bus.listenerCount('default'); + await connection.close(); + after = bus.listenerCount('default'); + } finally { await connection.close(); await models[Symbol.asyncDispose](); await server.dispose(); fixture.close(); } + }); + it('should admit the SQL stream', () => { outcome.should.equal(HubSubscriptionOutcome.Accepted); }); + it('should register one listener', () => { before.should.equal(1); }); + it('should release it on disconnect', () => { after.should.equal(0); }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_aborting_an_unused_prime.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_aborting_an_unused_prime.ts new file mode 100644 index 00000000..23e702c8 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_aborting_an_unused_prime.ts @@ -0,0 +1,31 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { sqliteTable, text } from 'drizzle-orm/sqlite-core'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import { DrizzleObservable } from '../../DrizzleObservable.js'; +import { DrizzleObservationSession } from '../../DrizzleObservationSession.js'; + +class Task {} +const table = sqliteTable('tasks', { id: text('id').primaryKey() }); + +describe('when aborting the signal of an unused prime', () => { + let error: string; + let listeners: number; + let idle: number; + beforeEach(async () => { + const bus = new DrizzleChangeNotifications(new Map([[Task, table]]), true); + const controller = new AbortController(); + idle = 0; + const observable = new DrizzleObservable(closed => new DrizzleObservationSession( + () => new Promise(() => {}), changed => bus.listen('tenant', table, changed), controller.signal, closed), + () => {}, () => {}, () => { idle++; }); + const pending = observable.current().then(() => '', (reason: Error) => reason.message); + controller.abort(); + error = await pending; + listeners = bus.listenerCount('tenant'); + observable.close(); + }); + it('should reject the pending current read', () => { error.should.equal('Drizzle observation was closed'); }); + it('should release the listener', () => { listeners.should.equal(0); }); + it('should leave owner tracking', () => { idle.should.equal(1); }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_adopting_a_primed_change.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_adopting_a_primed_change.ts new file mode 100644 index 00000000..87410105 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_adopting_a_primed_change.ts @@ -0,0 +1,24 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../../given.js'; +import { an_observation_scope, settle } from '../given/an_observation_scope.js'; +import type { TaskRecord } from '../given/TaskRecord.js'; + +describe('when subscribing after a change to a primed snapshot', given(an_observation_scope, context => { + let initial: number; + let emissions: TaskRecord[][]; + beforeEach(async () => { + await context.establish(); + const observed = context.models.observe(); + initial = (await observed.current()).value.length; + context.fixture.native.run("insert into tasks values ('21112233-4455-6677-8899-aabbccddeeff', 'new')"); + context.handle.notifyChanged(context.fixture.table); + emissions = []; + const subscription = observed.subscribe(rows => emissions.push(rows)); + await settle(); + subscription.unsubscribe(); + }); + afterEach(async () => { await context.dispose(); }); + it('should have primed the original snapshot', () => { initial.should.equal(2); }); + it('should read the change after adoption', () => { emissions.map(rows => rows.length).should.deep.equal([2, 3]); }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_another_tenant_notifies.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_another_tenant_notifies.ts new file mode 100644 index 00000000..9c7cac1d --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_another_tenant_notifies.ts @@ -0,0 +1,45 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { ArcApplication, Severity } from '@cratis/arc.core'; +import { DrizzleDialect } from '../../DrizzleDialect.js'; +import { DrizzleObservation } from '../../DrizzleObservation.js'; +import { drizzleDatabase, drizzleReadModel } from '../../drizzleToken.js'; +import '../../index.js'; +import { TaskRecord } from '../given/TaskRecord.js'; +import { a_sqlite_database } from '../given/a_sqlite_database.js'; +import { settle } from '../given/an_observation_scope.js'; + +describe('when a second tenant announces a SQL change', () => { + let emissions: number[]; + let afterOtherTenant: number[]; + beforeEach(async () => { + const own = new a_sqlite_database(); + const other = new a_sqlite_database(); + await own.establish(); await other.establish(); + const builder = ArcApplication.createBuilder(); + builder.withDrizzle({ dialect: DrizzleDialect.SQLite, + databaseFactory: tenant => tenant === 'default' ? own.database : other.database, + readModels: [{ type: TaskRecord, table: own.table }], observation: DrizzleObservation.InProcess }); + const app = await builder.build(); + const identity = (tenantId: string) => ({ tenantId, principal: undefined, allowedSeverity: Severity.Warning, + signal: new AbortController().signal, correlationId: crypto.randomUUID() }); + const scope = app.server.services.createScope(identity('DeFaUlT')); + const otherScope = app.server.services.createScope(identity('OTHER')); + try { + const models = await scope.resolve(drizzleReadModel(TaskRecord)); + emissions = []; + const subscription = models.observe().subscribe(rows => emissions.push(rows.length)); + await settle(); + other.native.run("insert into tasks values ('31112233-4455-6677-8899-aabbccddeeff', 'other')"); + (await otherScope.resolve(drizzleDatabase())).notifyChanged(TaskRecord); + await settle(); + afterOtherTenant = [...emissions]; + own.native.run("insert into tasks values ('31112233-4455-6677-8899-aabbccddeeff', 'ours')"); + (await scope.resolve(drizzleDatabase())).notifyChanged(TaskRecord); + await settle(); + subscription.unsubscribe(); + } finally { await scope.dispose(); await otherScope.dispose(); await app.dispose(); own.close(); other.close(); } + }); + it('should ignore the other tenant even with mixed-case IDs', () => { afterOtherTenant.should.deep.equal([2]); }); + it('should re-read the original tenant after its own change', () => { emissions.should.deep.equal([2, 3]); }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_changes_race_reads.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_changes_race_reads.ts new file mode 100644 index 00000000..3708549c --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_changes_race_reads.ts @@ -0,0 +1,107 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../../given.js'; +import { a_gated_observation, deferred, flush } from './given/a_gated_observation.js'; + +describe('when a change arrives during the initial observation read', given(a_gated_observation, context => { + let emissions: number[]; + beforeEach(async () => { + const first = deferred(); + context.reads = 0; + emissions = []; + const observation = context.open(() => context.reads === 1 ? first.promise : Promise.resolve(2)); + const subscription = observation.subscribe(value => emissions.push(value)); + context.notify(); + first.resolve(1); + await flush(); + subscription.unsubscribe(); + }); + afterEach(() => context.close()); + it('should read once more after the notification', () => { emissions.should.deep.equal([1, 2]); }); + it('should not leave a listener after unsubscribe', () => { context.listeners.should.equal(0); }); +})); + +describe('when a burst arrives during an in-flight reread', given(a_gated_observation, context => { + let emissions: number[]; + beforeEach(async () => { + const second = deferred(); + context.reads = 0; + emissions = []; + const observation = context.open(() => context.reads === 2 ? second.promise : Promise.resolve(context.reads)); + const subscription = observation.subscribe(value => emissions.push(value)); + await flush(); + context.notify(); + await flush(); + for (let index = 0; index < 100; index++) context.notify(); + second.resolve(2); + await flush(); + subscription.unsubscribe(); + }); + afterEach(() => context.close()); + it('should coalesce the burst into one trailing read', () => { context.reads.should.equal(3); }); + it('should emit the final state', () => { emissions.should.deep.equal([1, 2, 3]); }); +})); + +describe('when a later observation read fails', given(a_gated_observation, context => { + let reported: unknown; + let readsAfterAnotherNotification: number; + const failure = new Error('read failed'); + beforeEach(async () => { + context.reads = 0; + reported = undefined; + context.open(() => context.reads === 1 ? Promise.resolve(1) : Promise.reject(failure)) + .subscribe({ error: error => { reported = error; } }); + await flush(); + context.notify(); + await flush(); + context.notify(); + readsAfterAnotherNotification = context.reads; + }); + afterEach(() => context.close()); + it('should deliver the original error', () => { (reported === failure).should.equal(true); }); + it('should release the listener', () => { context.listeners.should.equal(0); }); + it('should not reopen after another notification', () => { readsAfterAnotherNotification.should.equal(2); }); +})); + +describe('when the initial read fails before adoption', given(a_gated_observation, context => { + let errors: string[]; + let currentFailed: boolean; + beforeEach(async () => { + const first = deferred(); + const observation = context.open(() => first.promise); + const current = observation.current().then(() => false, () => true); + first.reject(new Error('database failed')); + currentFailed = await current; + errors = []; + observation.subscribe({ error: error => errors.push((error as Error).message) }); + await flush(); + }); + afterEach(() => context.close()); + it('should reject current', () => { currentFailed.should.equal(true); }); + it('should retain the error for its first subscriber', () => { errors.should.deep.equal(['database failed']); }); +})); + +describe('when a primed observation closes before adoption', given(a_gated_observation, context => { + let completed: boolean; + beforeEach(async () => { + const observation = context.open(() => Promise.resolve(1)); + await observation.current(); + observation.close(); + completed = false; + observation.subscribe({ complete: () => { completed = true; } }); + }); + it('should complete the late subscriber', () => { completed.should.equal(true); }); +})); + +describe('when a primed read fails after closure', given(a_gated_observation, context => { + let canceled: boolean; + beforeEach(async () => { + const first = deferred(); + const observation = context.open(() => first.promise); + const current = observation.current().then(() => false, () => true); + observation.close(); + first.reject(new Error('late')); + canceled = await current; + }); + it('should reject the pending current rather than deliver the late failure', () => { canceled.should.equal(true); }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_closing_before_a_scheduled_reread.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_closing_before_a_scheduled_reread.ts new file mode 100644 index 00000000..3e11a350 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_closing_before_a_scheduled_reread.ts @@ -0,0 +1,22 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../../given.js'; +import { a_gated_observation, flush } from './given/a_gated_observation.js'; + +describe('when closing an observation before a scheduled reread starts', given(a_gated_observation, context => { + let emissions: number[]; + beforeEach(async () => { + context.reads = 0; + emissions = []; + const observation = context.open(() => Promise.resolve(context.reads)); + observation.subscribe(value => emissions.push(value)); + await flush(); + context.notify(); + observation.close(); + await flush(); + }); + afterEach(() => context.close()); + it('should not start another read', () => { context.reads.should.equal(1); }); + it('should not emit after closing', () => { emissions.should.deep.equal([1]); }); + it('should release the listener', () => { context.listeners.should.equal(0); }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_deleting_and_recreating_a_row.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_deleting_and_recreating_a_row.ts new file mode 100644 index 00000000..b9936a73 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_deleting_and_recreating_a_row.ts @@ -0,0 +1,26 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { eq } from 'drizzle-orm'; +import { given } from '../../given.js'; +import { an_observation_scope, settle } from '../given/an_observation_scope.js'; +import { TaskRecord } from '../given/TaskRecord.js'; + +describe('when an observed row is deleted and recreated', given(an_observation_scope, context => { + let seen: (string | null)[]; + beforeEach(async () => { + await context.establish(); + seen = []; + const id = '00112233-4455-6677-8899-aabbccddeeff'; + const subscription = context.models.observeById(id).subscribe(row => seen.push(row?.title ?? null)); + await settle(); + context.fixture.database.delete(context.fixture.table).where(eq(context.fixture.table.title, 'z')).run(); + context.handle.notifyChanged(TaskRecord); + await settle(); + context.fixture.native.run(`insert into tasks values ('${id}', 'again')`); + context.handle.notifyChanged(TaskRecord); + await settle(); + subscription.unsubscribe(); + }); + afterEach(async () => { await context.dispose(); }); + it('should emit the missing row and the recreated value', () => { seen.should.deep.equal(['z', null, 'again']); }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_exceeding_the_small_list_limit.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_exceeding_the_small_list_limit.ts new file mode 100644 index 00000000..168fa416 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_exceeding_the_small_list_limit.ts @@ -0,0 +1,28 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../../given.js'; +import { an_observation_scope, settle } from '../given/an_observation_scope.js'; + +describe('when a SQL observation outgrows its complete-list limit', given(an_observation_scope, context => { + let error: string; + let emitted: number[]; + beforeEach(async () => { + await context.establish(2); + const observed = context.models.observe(); + emitted = []; + let fail!: (message: string) => void; + const failed = new Promise(resolve => { fail = resolve; }); + const subscription = observed.subscribe({ next: rows => emitted.push(rows.length), + error: (reason: Error) => fail(reason.message) }); + await settle(); + context.fixture.native.run("insert into tasks values ('31112233-4455-6677-8899-aabbccddeeff', 'new')"); + context.handle.notifyChanged(context.fixture.table); + error = await failed; + subscription.unsubscribe(); + }); + afterEach(async () => { await context.dispose(); }); + it('should fail with the paging error', () => { + error.should.equal('The result exceeds the maximum of 2 items; narrow the filter, or use observePage through defineObservableQuery for paged results'); + }); + it('should never emit a truncated later list', () => { emitted.should.deep.equal([2]); }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_not_enabled_through_withDrizzle.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_not_enabled_through_withDrizzle.ts new file mode 100644 index 00000000..31f298f4 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_not_enabled_through_withDrizzle.ts @@ -0,0 +1,36 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { ArcApplication, Severity } from '@cratis/arc.core'; +import { DrizzleDialect } from '../../DrizzleDialect.js'; +import { drizzleReadModel } from '../../drizzleToken.js'; +import '../../index.js'; +import { TaskRecord } from '../given/TaskRecord.js'; +import { a_sqlite_database } from '../given/a_sqlite_database.js'; + +describe('when starting an observation without enabling it in withDrizzle', () => { + let currentError: string; + let subscriptionError: string; + beforeEach(async () => { + const fixture = new a_sqlite_database(); + await fixture.establish(); + const builder = ArcApplication.createBuilder(); + builder.withDrizzle({ dialect: DrizzleDialect.SQLite, database: fixture.database, + readModels: [{ type: TaskRecord, table: fixture.table }] }); + const app = await builder.build(); + const scope = app.server.services.createScope({ tenantId: 'default', principal: undefined, + allowedSeverity: Severity.Warning, signal: new AbortController().signal, correlationId: crypto.randomUUID() }); + try { + const models = await scope.resolve(drizzleReadModel(TaskRecord)); + currentError = await models.observe().current().then(() => '', (error: Error) => error.message); + subscriptionError = await new Promise(resolve => models.observe().subscribe({ + error: (error: Error) => resolve(error.message) + })); + } finally { await scope.dispose(); await app.dispose(); fixture.close(); } + }); + it('should reject the current snapshot with the named error', () => { + currentError.should.equal('Drizzle observation is not enabled; set observation: DrizzleObservation.InProcess in withDrizzle'); + }); + it('should reject the subscription with the same error', () => { + subscriptionError.should.equal('Drizzle observation is not enabled; set observation: DrizzleObservation.InProcess in withDrizzle'); + }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_page_counts_disagree.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_page_counts_disagree.ts new file mode 100644 index 00000000..4f6817bc --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_page_counts_disagree.ts @@ -0,0 +1,32 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import type { DrizzleDatabase } from '../../DrizzleDatabase.js'; +import { DrizzleReadModels } from '../../DrizzleReadModels.js'; +import { TaskRecord } from '../given/TaskRecord.js'; +import { a_sqlite_database } from '../given/a_sqlite_database.js'; + +describe('when the count and rows of an observed page keep disagreeing', () => { + let error: string; + let counts: number; + let listeners: number; + beforeEach(async () => { + const fixture = new a_sqlite_database(); + await fixture.establish(); + counts = 0; + const fake = { select: fixture.database.select.bind(fixture.database), $count: () => { counts++; return 3; } }; + const bus = new DrizzleChangeNotifications(new Map([[TaskRecord, fixture.table]]), true); + const models = new DrizzleReadModels(fake as unknown as DrizzleDatabase, fixture.table, TaskRecord, 100, + undefined, bus, 'default'); + try { + error = await models.observePage(undefined, { paging: { page: 0, pageSize: 10 } }) + .current().then(() => '', (failure: Error) => failure.message); + listeners = bus.listenerCount('default'); + } finally { await models[Symbol.asyncDispose](); fixture.close(); } + }); + it('should retry no more than three times', () => { counts.should.equal(3); }); + it('should fail without emitting an inconsistent page', () => { + error.should.equal('Drizzle observation could not read a consistent page'); + }); + it('should release the listener', () => { listeners.should.equal(0); }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_publishing_outside_the_observer_context.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_publishing_outside_the_observer_context.ts new file mode 100644 index 00000000..4fa9b983 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_publishing_outside_the_observer_context.ts @@ -0,0 +1,38 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { AsyncLocalStorage } from 'node:async_hooks'; +import { sqliteTable, text } from 'drizzle-orm/sqlite-core'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import { DrizzleObservable } from '../../DrizzleObservable.js'; +import { DrizzleObservationSession } from '../../DrizzleObservationSession.js'; + +class Task {} +const table = sqliteTable('tasks', { id: text('id').primaryKey() }); + +describe('when a command publishes a change to a subscribed observer', () => { + let readsWhenNotifyReturns: number; + let contexts: (string | undefined)[]; + beforeEach(async () => { + const commandContext = new AsyncLocalStorage(); + const bus = new DrizzleChangeNotifications(new Map([[Task, table]]), true); + contexts = []; + let updated!: () => void; + const reread = new Promise(resolve => { updated = resolve; }); + const observable = new DrizzleObservable(closed => new DrizzleObservationSession(async () => { + contexts.push(commandContext.getStore()); + if (contexts.length === 2) updated(); + return contexts.length; + }, changed => bus.listen('tenant', table, changed), undefined, closed), () => {}); + const subscription = observable.subscribe(() => {}); + await new Promise(resolve => setImmediate(resolve)); + commandContext.run('command', () => { + bus.notify('tenant', [Task]); + readsWhenNotifyReturns = contexts.length; + }); + await reread; + subscription.unsubscribe(); + observable.close(); + }); + it('should return from notify before starting the observer reread', () => { readsWhenNotifyReturns.should.equal(1); }); + it('should not inherit the command async-local frame', () => { contexts.should.deep.equal([undefined, undefined]); }); +}); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_the_scope_closes.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_the_scope_closes.ts new file mode 100644 index 00000000..dd2fdd5b --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_the_scope_closes.ts @@ -0,0 +1,20 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { given } from '../../given.js'; +import { an_observation_scope } from '../given/an_observation_scope.js'; + +describe('when the read-model scope closes after priming', given(an_observation_scope, context => { + let closedPrime: boolean; + let newStart: boolean; + beforeEach(async () => { + await context.establish(); + const observed = context.models.observe(); + await observed.current(); + await context.scope.dispose(); + closedPrime = await observed.current().then(() => false, () => true); + newStart = await context.models.observe().current().then(() => false, () => true); + }); + afterEach(async () => { await context.dispose(); }); + it('should close the unused prime', () => { closedPrime.should.equal(true); }); + it('should reject further starts from the disposed owner', () => { newStart.should.equal(true); }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_observing/when_updating_a_sorted_page.ts b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_updating_a_sorted_page.ts new file mode 100644 index 00000000..bfa39163 --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_observing/when_updating_a_sorted_page.ts @@ -0,0 +1,25 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { SortDirection } from '@cratis/arc.core'; +import { given } from '../../given.js'; +import { an_observation_scope, settle } from '../given/an_observation_scope.js'; + +describe('when a new SQL row changes a sorted observed page', given(an_observation_scope, context => { + let pages: { title: string; total: number }[]; + beforeEach(async () => { + await context.establish(); + pages = []; + const subscription = context.models.observePage(undefined, { paging: { page: 0, pageSize: 1 }, + sorting: { field: 'title', direction: SortDirection.Ascending } }).subscribe(page => + pages.push({ title: page.items[0]!.title, total: page.totalItems })); + await settle(); + context.fixture.native.run("insert into tasks values ('41112233-4455-6677-8899-aabbccddeeff', '0')"); + context.handle.notifyChanged(context.fixture.table); + await settle(); + subscription.unsubscribe(); + }); + afterEach(async () => { await context.dispose(); }); + it('should recompute the SQL sort and total', () => { + pages.should.deep.equal([{ title: 'a', total: 2 }, { title: '0', total: 3 }]); + }); +})); diff --git a/Source/Drizzle/for_DrizzleReadModels/when_serving_a_sqlite_observation/with_each_http_adapter.ts b/Source/Drizzle/for_DrizzleReadModels/when_serving_a_sqlite_observation/with_each_http_adapter.ts new file mode 100644 index 00000000..12da6b9d --- /dev/null +++ b/Source/Drizzle/for_DrizzleReadModels/when_serving_a_sqlite_observation/with_each_http_adapter.ts @@ -0,0 +1,110 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { createServer, type Server } from 'node:http'; +import express from 'express'; +import fastify from 'fastify'; +import { Hono } from 'hono'; +import { serve } from '@hono/node-server'; +import { describe, it, should } from 'vitest'; +import { ArcApplication, Severity } from '@cratis/arc.core'; +import { cratisArc as expressArc } from '../../../Express/index.js'; +import { cratisArc as fastifyArc } from '../../../Fastify/index.js'; +import { cratisArc as honoArc } from '../../../Hono/index.js'; +import { DrizzleDialect, DrizzleObservation } from '../../index.js'; +import { DrizzleChangeNotifications } from '../../DrizzleChangeNotifications.js'; +import { TaskRecord } from '../given/TaskRecord.js'; +import { a_sqlite_database } from '../given/a_sqlite_database.js'; +import { AddObservedTask } from '../given/AddObservedTask.js'; +import { ObservedTaskQueries } from '../given/ObservedTaskQueries.js'; + +should(); +for (const adapter of ['Express', 'Fastify', 'Hono'] as const) { + describe(`when serving SQLite observation through ${adapter}`, () => { + it('should serve a current GET and stream a command change until disconnect', async () => { + const fixture = new a_sqlite_database(); + await fixture.establish(); + const builder = ArcApplication.createBuilder({ tenancy: { resolve: () => 'default' } }); + builder.add(AddObservedTask, ObservedTaskQueries).withDrizzle({ dialect: DrizzleDialect.SQLite, + database: fixture.database, readModels: [{ type: TaskRecord, table: fixture.table }], + observation: DrizzleObservation.InProcess }); + const application = await builder.build(); + let listener: Server | undefined; + let stop: () => Promise = async () => {}; + const controller = new AbortController(); + try { + if (adapter === 'Express') { + const host = express(); + host.use(expressArc(application)); + listener = createServer(host); + await new Promise(resolve => listener!.listen(0, '127.0.0.1', resolve)); + } else if (adapter === 'Fastify') { + const host = fastify(); + host.register(fastifyArc, { arc: application }); + await host.listen({ port: 0, host: '127.0.0.1' }); + listener = host.server as Server; + stop = () => host.close(); + } else { + const host = new Hono(); + host.use(honoArc(application)); + listener = serve({ fetch: host.fetch, port: 0, hostname: '127.0.0.1' }) as Server; + if (!listener.listening) await new Promise(resolve => listener!.once('listening', resolve)); + } + const address = listener.address(); + if (!address || typeof address === 'string') throw Error('No HTTP port'); + const routes = [...application.server.endpoints.keys()]; + const page = routes.find(route => route.endsWith('/page')); + const command = routes.find(route => route.endsWith('/add-observed-task')); + if (!page || !command) throw Error(`Observation routes missing: ${routes.join(', ')}`); + const base = `http://127.0.0.1:${address.port}`; + const url = `${base}${page}?page=0&pageSize=10`; + const snapshot = await fetch(url); + snapshot.status.should.equal(200); + const initial = await snapshot.json() as { paging: { totalItems: number } }; + initial.paging.totalItems.should.equal(2); + const response = await fetch(url, { headers: { Accept: 'text/event-stream' }, signal: controller.signal }); + response.status.should.equal(200); + const reader = response.body!.getReader(); + const decoder = new TextDecoder(); + let buffer = ''; + const nextFrame = async (): Promise<{ paging: { totalItems: number } }> => { + while (true) { + const boundary = buffer.indexOf('\n\n'); + if (boundary >= 0) { + const frame = buffer.slice(0, boundary); + buffer = buffer.slice(boundary + 2); + const data = frame.split('\n').filter(line => line.startsWith('data:')).map(line => line.slice(5)).join(''); + if (data) return JSON.parse(data) as { paging: { totalItems: number } }; + continue; + } + const part = await reader.read(); + if (part.done) throw Error('SSE closed before result'); + buffer += decoder.decode(part.value, { stream: true }).replaceAll('\r\n', '\n'); + } + }; + (await nextFrame()).paging.totalItems.should.equal(2); + const written = await fetch(`${base}${command}`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ id: '21112233-4455-6677-8899-aabbccddeeff', title: 'new' }) }); + written.status.should.equal(200); + (await nextFrame()).paging.totalItems.should.equal(3); + const busScope = application.server.services.createScope({ tenantId: 'default', principal: undefined, + allowedSeverity: Severity.Warning, correlationId: crypto.randomUUID(), signal: new AbortController().signal }); + const notifications = await busScope.resolve(DrizzleChangeNotifications); + notifications.listenerCount('default').should.equal(1); + controller.abort(); + await reader.cancel().catch(() => {}); + const deadline = Date.now() + 2000; + while (notifications.listenerCount('default') && Date.now() < deadline) + await new Promise(resolve => setTimeout(resolve, 10)); + notifications.listenerCount('default').should.equal(0); + await busScope.dispose(); + } finally { + controller.abort(); + if (listener && adapter !== 'Fastify') await new Promise((resolve, reject) => + listener!.close(error => error ? reject(error) : resolve())); + await stop(); + await application.dispose(); + fixture.close(); + } + }); + }); +} diff --git a/Source/Drizzle/index.ts b/Source/Drizzle/index.ts index a25d8a05..bd0ff789 100644 --- a/Source/Drizzle/index.ts +++ b/Source/Drizzle/index.ts @@ -7,6 +7,8 @@ export { DrizzleReadModels } from './DrizzleReadModels.js'; export { DrizzleReadModelForCommandResolver } from './DrizzleReadModelForCommandResolver.js'; export type { DrizzleOptions } from './DrizzleOptions.js'; export { DrizzleDialect } from './DrizzleDialect.js'; +export { DrizzleObservation } from './DrizzleObservation.js'; +export { DrizzleObservable } from './DrizzleObservable.js'; export type { DrizzleDatabase } from './DrizzleDatabase.js'; export type { DrizzleFilter } from './DrizzleFilter.js'; export type { ColumnCodec } from './ColumnCodec.js'; diff --git a/Source/Drizzle/package.json b/Source/Drizzle/package.json index f9b4db9b..5a1fa237 100644 --- a/Source/Drizzle/package.json +++ b/Source/Drizzle/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.drizzle", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { @@ -29,9 +29,10 @@ "README.md" ], "peerDependencies": { - "@cratis/arc.core": "^0.37.0", + "@cratis/arc.core": "^0.38.0", "@cratis/fundamentals": "^7.19.6", - "drizzle-orm": "^0.45.0" + "drizzle-orm": "^0.45.0", + "rxjs": "^7.8.2" }, "devDependencies": { "@cratis/arc.core": "workspace:^", @@ -43,6 +44,7 @@ "mysql2": "^3.24.4", "pg": "^8.0.0", "postgres": "^3.4.9", + "rxjs": "^7.8.2", "sql.js": "^1.14.2" } } diff --git a/Source/Drizzle/withDrizzle.ts b/Source/Drizzle/withDrizzle.ts index 92f19f1d..849175f7 100644 --- a/Source/Drizzle/withDrizzle.ts +++ b/Source/Drizzle/withDrizzle.ts @@ -11,6 +11,9 @@ import { drizzleDatabase, drizzleReadModel } from './drizzleToken.js'; import { DrizzleModelCodec } from './DrizzleModelCodec.js'; import { DrizzleReadModelForCommandResolver } from './DrizzleReadModelForCommandResolver.js'; import { getTableColumns } from 'drizzle-orm'; +import type { Table } from 'drizzle-orm'; +import { DrizzleObservation } from './DrizzleObservation.js'; +import { DrizzleChangeNotifications } from './DrizzleChangeNotifications.js'; /** An application retains ownership of its connections, pools, and migrations. */ export function withDrizzle(builder: ArcApplicationBuilder, options: DrizzleOptions): ArcApplicationBuilder { @@ -20,6 +23,10 @@ export function withDrizzle(builder: ArcApplicationBuilder, options: DrizzleOpti if (options.maxPageSize !== undefined && (!Number.isSafeInteger(options.maxPageSize) || options.maxPageSize <= 0 || options.maxPageSize > 10000)) throw new RangeError('maxPageSize must be between 1 and 10000'); + if (options.observation !== undefined && options.observation !== DrizzleObservation.InProcess) + throw new Error('Unsupported Drizzle observation mode'); + const tables = new Map object, Table>(); + const notifications = new DrizzleChangeNotifications(tables, options.observation === DrizzleObservation.InProcess); const registered = new Set object>(); // Check registrations before a request opens a tenant scope; reuse one codec per type. const codecs = new Map object, DrizzleModelCodec>(); @@ -27,6 +34,7 @@ export function withDrizzle(builder: ArcApplicationBuilder, options: DrizzleOpti for (const { type, table } of options.readModels ?? []) { if (registered.has(type)) throw new Error(`Duplicate Drizzle read model: ${type.name}`); registered.add(type); + tables.set(type, table); const columns = getTableColumns(table); const codec = new DrizzleModelCodec(type, columns); const keys = Object.entries(columns).filter(([, column]) => column.primary); @@ -34,7 +42,10 @@ export function withDrizzle(builder: ArcApplicationBuilder, options: DrizzleOpti new DrizzleReadModels({}, table, type, options.maxPageSize, codec); codecs.set(type, codec); } + builder.services.addSingleton(DrizzleChangeNotifications, () => notifications); builder.services.addScoped(DrizzleReadModelForCommandResolver, () => new DrizzleReadModelForCommandResolver(commandModels)); + builder.addCommandExecutionRunner((context, execute) => context.tenantId ? + notifications.run(context.tenantId.toLowerCase(), execute) : execute()); builder.addReadModelForCommandResolver(DrizzleReadModelForCommandResolver); builder.services.addScoped(drizzleDatabase(), async scope => { const context: ExecutionContext | undefined = scope.identity; @@ -43,11 +54,12 @@ export function withDrizzle(builder: ArcApplicationBuilder, options: DrizzleOpti if (options.database && tenant !== 'default') throw new Error('Drizzle database is only available for the default tenant'); const database = options.databaseFactory ? await options.databaseFactory(tenant, context) : options.database; if (!database) throw new Error('Drizzle tenant database resolver returned no database'); - return new DrizzleHandle(database); + return new DrizzleHandle(database, notifications, tenant); }); for (const { type, table } of options.readModels ?? []) { builder.services.addScoped(drizzleReadModel(type), async scope => - new DrizzleReadModels((await scope.resolve(drizzleDatabase())).native, table, type, options.maxPageSize, codecs.get(type)!)); + new DrizzleReadModels((await scope.resolve(drizzleDatabase())).native, table, type, options.maxPageSize, codecs.get(type)!, + notifications, scope.identity?.tenantId?.toLowerCase(), scope.identity?.signal)); } return builder; } diff --git a/Source/Express/package.json b/Source/Express/package.json index a40a2589..07299cf3 100644 --- a/Source/Express/package.json +++ b/Source/Express/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.express", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Fastify/package.json b/Source/Fastify/package.json index 1f3e8c29..df951dbf 100644 --- a/Source/Fastify/package.json +++ b/Source/Fastify/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.fastify", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Hono/package.json b/Source/Hono/package.json index 1de0fb35..7b431484 100644 --- a/Source/Hono/package.json +++ b/Source/Hono/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.hono", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/MongoDB/package.json b/Source/MongoDB/package.json index 120fee69..c902491c 100644 --- a/Source/MongoDB/package.json +++ b/Source/MongoDB/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.mongodb", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { @@ -29,7 +29,7 @@ "README.md" ], "peerDependencies": { - "@cratis/arc.core": "^0.37.0", + "@cratis/arc.core": "^0.38.0", "@cratis/fundamentals": "^7.19.6", "@opentelemetry/api": "^1.9.0", "mongodb": "^6.21.0", diff --git a/Source/Testing/package.json b/Source/Testing/package.json index bd66d388..aa0007d4 100644 --- a/Source/Testing/package.json +++ b/Source/Testing/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.testing", - "version": "0.37.0", + "version": "0.38.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Tools/ProxyGenerator/collectSourceQueries.ts b/Source/Tools/ProxyGenerator/collectSourceQueries.ts index 230b28b2..e29365d4 100644 --- a/Source/Tools/ProxyGenerator/collectSourceQueries.ts +++ b/Source/Tools/ProxyGenerator/collectSourceQueries.ts @@ -3,7 +3,7 @@ import ts from 'typescript'; import { ClientOperationKind } from '@cratis/arc.core'; import { annotation, roles, stringArgument } from './sourceAnnotations.js'; -import { identifier, isPackageSymbol, originalSymbol } from './sourceSymbols.js'; +import { identifier, isPackageSymbol, isTypeFrom, originalSymbol } from './sourceSymbols.js'; import { queryResult } from './queryResult.js'; import { warningOption, httpMethodOption } from './sourceOperationOptions.js'; import type { SourceField } from './SourceField.js'; @@ -27,9 +27,18 @@ function boundParameters(member: ts.MethodDeclaration, explicit: readonly ts.Exp isPackageSymbol(checker, item.expression, 'argument', '@cratis/arc.core') && item.arguments[0] && ts.isStringLiteral(item.arguments[0]) && item.arguments[0].text === parameterName); if (!binding) { - const serviceBinding = explicit.some(item => ts.isCallExpression(item) && - isPackageSymbol(checker, item.expression, 'service', '@cratis/arc.core') && item.arguments[0] && - originalSymbol(checker, item.arguments[0]) === checker.getTypeAtLocation(parameter).symbol); + const serviceBinding = explicit.some(item => { + if (!ts.isCallExpression(item) || !isPackageSymbol(checker, item.expression, 'service', '@cratis/arc.core') || + !item.arguments[0]) return false; + const binding = item.arguments[0]; + const parameterType = checker.getTypeAtLocation(parameter); + if (originalSymbol(checker, binding) === parameterType.symbol) return true; + const token = checker.getTypeAtLocation(binding); + if (!isTypeFrom(checker, token, 'ServiceToken', '@cratis/arc.core')) return false; + const serviceType = checker.getTypeArguments(token as ts.TypeReference)[0]; + return !!serviceType && checker.isTypeAssignableTo(serviceType, parameterType) && + checker.isTypeAssignableTo(parameterType, serviceType); + }); const type = checker.getTypeAtLocation(parameter); const actual = type.isUnion() ? type.types.find(part => !(part.flags & (ts.TypeFlags.Null | ts.TypeFlags.Undefined))) ?? type : type; diff --git a/Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_queries.ts b/Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_queries.ts new file mode 100644 index 00000000..a0af0826 --- /dev/null +++ b/Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_queries.ts @@ -0,0 +1,20 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { key, query, readModel, service } from '@cratis/arc.core'; +import { field } from '@cratis/fundamentals'; +import type { Observable } from 'rxjs'; +import { drizzleReadModel } from '../../../../Drizzle/drizzleToken.js'; +import type { DrizzleReadModels } from '../../../../Drizzle/DrizzleReadModels.js'; +class TaskRecord { + @field(String) @key() id!: string; + @field(String) title!: string; +} + +@readModel() +export class DerivedObservableQueries { + @query({ observable: true }, service(drizzleReadModel(TaskRecord))) + static unannotated(tasks: DrizzleReadModels) { return tasks.observe(); } + + @query({ observable: true }, service(drizzleReadModel(TaskRecord))) + static annotated(tasks: DrizzleReadModels): Observable { return tasks.observe(); } +} diff --git a/Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_returns.ts b/Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_returns.ts new file mode 100644 index 00000000..2184dda2 --- /dev/null +++ b/Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_returns.ts @@ -0,0 +1,17 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { Observable, Subject } from 'rxjs'; +import type { DrizzleObservable } from '../../../../Drizzle/DrizzleObservable.js'; +import type { MongoObservable } from '../../../../MongoDB/MongoObservable.js'; + +export declare function drizzle(): DrizzleObservable; +export declare function mongo(): MongoObservable; +export declare function annotated(): Observable; +export declare function subject(): Subject; + +class Mid extends Observable {} +class Leaf extends Mid {} +class Wrap extends Observable {} + +export declare function leaf(): Leaf; +export declare function wrapped(): Wrap; diff --git a/Source/Tools/ProxyGenerator/for_queryResult/when_generating_derived_observable_queries.ts b/Source/Tools/ProxyGenerator/for_queryResult/when_generating_derived_observable_queries.ts new file mode 100644 index 00000000..7821baa0 --- /dev/null +++ b/Source/Tools/ProxyGenerator/for_queryResult/when_generating_derived_observable_queries.ts @@ -0,0 +1,26 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { resolve } from 'node:path'; +import { ClientOperationKind } from '@cratis/arc.core'; +import { analyzeSource } from '../analyzeSource.js'; + +describe('when generating proxy operations for Drizzle observable queries', () => { + let operations: ReturnType['operations']; + beforeEach(() => { + const root = resolve(import.meta.dirname, '../../../..'); + operations = analyzeSource(resolve(root, 'tsconfig.specs.json'), + resolve(root, 'Source/Tools/ProxyGenerator/for_queryResult/given'), '', true).operations; + }); + it('should recognize the inferred tasks.observe return as an observable array', () => { + const result = operations.find(operation => operation.name === 'unannotated')!; + result.kind.should.equal(ClientOperationKind.Observable); + result.result.enumerable.should.equal(true); + result.result.text.should.include('TaskRecord'); + }); + it('should recognize an annotated Observable as the same shape', () => { + const result = operations.find(operation => operation.name === 'annotated')!; + result.kind.should.equal(ClientOperationKind.Observable); + result.result.enumerable.should.equal(true); + result.result.text.should.include('TaskRecord'); + }); +}); diff --git a/Source/Tools/ProxyGenerator/for_queryResult/when_resolving_derived_observables.ts b/Source/Tools/ProxyGenerator/for_queryResult/when_resolving_derived_observables.ts new file mode 100644 index 00000000..9f9d34f7 --- /dev/null +++ b/Source/Tools/ProxyGenerator/for_queryResult/when_resolving_derived_observables.ts @@ -0,0 +1,24 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { resolve } from 'node:path'; +import ts from 'typescript'; +import { queryResult } from '../queryResult.js'; + +describe('when resolving derived observable returns', () => { + let results: { name: string; observable: boolean; item: string }[]; + beforeEach(() => { + const file = resolve('Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_returns.ts'); + const program = ts.createProgram([file], { module: ts.ModuleKind.NodeNext, moduleResolution: ts.ModuleResolutionKind.NodeNext }); + const checker = program.getTypeChecker(); + results = program.getSourceFile(file)!.statements.filter(ts.isFunctionDeclaration) + .filter(node => ['drizzle', 'mongo', 'annotated', 'subject'].includes(node.name!.text)).map(node => { + const type = checker.getReturnTypeOfSignature(checker.getSignatureFromDeclaration(node)!); + const result = queryResult(type, checker, node); + return { name: node.name!.text, observable: result.observable, item: checker.typeToString(result.type) }; + }); + }); + it('should unwrap both unannotated derived and annotated RxJS array results', () => { + results.should.deep.equal(['drizzle', 'mongo', 'annotated', 'subject'].map(name => + ({ name, observable: true, item: 'string[]' }))); + }); +}); diff --git a/Source/Tools/ProxyGenerator/for_queryResult/when_resolving_nested_derived_observables.ts b/Source/Tools/ProxyGenerator/for_queryResult/when_resolving_nested_derived_observables.ts new file mode 100644 index 00000000..757b273b --- /dev/null +++ b/Source/Tools/ProxyGenerator/for_queryResult/when_resolving_nested_derived_observables.ts @@ -0,0 +1,40 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { resolve } from 'node:path'; +import ts from 'typescript'; +import { queryResult } from '../queryResult.js'; + +function fixture(name: string) { + const file = resolve('Source/Tools/ProxyGenerator/for_queryResult/given/derived_observable_returns.ts'); + const program = ts.createProgram([file], { module: ts.ModuleKind.NodeNext, moduleResolution: ts.ModuleResolutionKind.NodeNext }); + const checker = program.getTypeChecker(); + const node = program.getSourceFile(file)!.statements.filter(ts.isFunctionDeclaration).find(node => node.name!.text === name)!; + const type = checker.getReturnTypeOfSignature(checker.getSignatureFromDeclaration(node)!); + return { checker, node, type }; +} + +describe('when resolving an observable through two generic base classes', () => { + let item: string; + let observable: boolean; + beforeEach(() => { + const { checker, node, type } = fixture('leaf'); + const result = queryResult(type, checker, node); + observable = result.observable; + item = checker.typeToString(result.type); + }); + it('should classify the result as observable', () => { observable.should.equal(true); }); + it('should resolve the concrete item type', () => { item.should.equal('string[]'); }); +}); + +describe('when resolving an observable with a nested unbound type parameter', () => { + let failure: Error | undefined; + beforeEach(() => { + const { checker, node, type } = fixture('wrapped'); + try { queryResult(type, checker, node); } + catch (error) { failure = error as Error; } + }); + it('should require an explicit Observable return annotation', () => { + (failure instanceof Error).should.equal(true); + failure!.message.should.include('annotate the return as Observable<...>'); + }); +}); diff --git a/Source/Tools/ProxyGenerator/package.json b/Source/Tools/ProxyGenerator/package.json index fe6599ae..36672dc2 100644 --- a/Source/Tools/ProxyGenerator/package.json +++ b/Source/Tools/ProxyGenerator/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.proxygenerator", - "version": "0.37.0", + "version": "0.38.0", "description": "TypeScript source analyzer and deterministic Arc client proxy generator", "repository": { "type": "git", diff --git a/Source/Tools/ProxyGenerator/queryResult.ts b/Source/Tools/ProxyGenerator/queryResult.ts index 1cc90067..4f33c488 100644 --- a/Source/Tools/ProxyGenerator/queryResult.ts +++ b/Source/Tools/ProxyGenerator/queryResult.ts @@ -3,16 +3,46 @@ import ts from 'typescript'; import { isStandardType, isTypeFrom } from './sourceSymbols.js'; +function observableBase(type: ts.Type, checker: ts.TypeChecker, visited = new Set()): ts.Type[] | undefined { + if (visited.has(type)) return undefined; + visited.add(type); + if (isTypeFrom(checker, type, 'ObservableSource', '@cratis/arc.core') || + ['Observable', 'Subject', 'BehaviorSubject', 'ReplaySubject'].some(name => isTypeFrom(checker, type, name, 'rxjs')) || + isStandardType(type, 'AsyncIterable') || isStandardType(type, 'AsyncGenerator')) return [type]; + for (const base of checker.getBaseTypes(type as ts.InterfaceType) ?? []) { + const matched = observableBase(base, checker, visited); + if (matched) return [type, ...matched]; + } + return undefined; +} + +function hasTypeParameter(type: ts.Type, checker: ts.TypeChecker, visited = new Set()): boolean { + if (type.flags & ts.TypeFlags.TypeParameter) return true; + if (visited.has(type)) return false; + visited.add(type); + const parts = type.isUnionOrIntersection() ? type.types : + type.aliasTypeArguments ?? (type.flags & ts.TypeFlags.Object && (type as ts.ObjectType).objectFlags & ts.ObjectFlags.Reference ? + checker.getTypeArguments(type as ts.TypeReference) : []); + return parts.some(part => hasTypeParameter(part, checker, visited)); +} + export function queryResult(type: ts.Type, checker: ts.TypeChecker, node: ts.Node): { type: ts.Type; observable: boolean; paged: boolean } { const current = checker.getAwaitedType(type) ?? type; - const observable = isTypeFrom(checker, current, 'ObservableSource', '@cratis/arc.core') || - ['Observable', 'Subject', 'BehaviorSubject', 'ReplaySubject'].some(name => isTypeFrom(checker, current, name, 'rxjs')) || - isStandardType(current, 'AsyncIterable') || isStandardType(current, 'AsyncGenerator'); + const source = observableBase(current, checker); const paged = isTypeFrom(checker, current, 'QueryPage', '@cratis/arc.core'); - if (observable || paged) { - const argument = current.aliasTypeArguments?.[0] ?? checker.getTypeArguments(current as ts.TypeReference)[0]; + if (source || paged) { + const wrapped = source?.at(-1) ?? current; + let argument = wrapped.aliasTypeArguments?.[0] ?? checker.getTypeArguments(wrapped as ts.TypeReference)[0]; if (!argument) throw new Error(`${node.getSourceFile().fileName}: missing query result type`); - return { type: argument, observable, paged }; + for (const parent of source?.slice(0, -1).reverse() ?? []) { + if (!(argument.flags & ts.TypeFlags.TypeParameter)) break; + const reference = parent as ts.TypeReference; + const index: number = reference.target?.typeParameters?.indexOf(argument) ?? -1; + if (index >= 0) argument = checker.getTypeArguments(reference)[index] ?? argument; + } + if (source && hasTypeParameter(argument, checker)) + throw new Error(`${node.getSourceFile().fileName}: cannot resolve observable query result type; annotate the return as Observable<...>`); + return { type: argument, observable: !!source, paged }; } return { type: current, observable: false, paged: false }; } diff --git a/yarn.lock b/yarn.lock index 4dc0cef7..ab6ca7fe 100644 --- a/yarn.lock +++ b/yarn.lock @@ -45,8 +45,8 @@ __metadata: rxjs: "npm:^7.8.2" zod: "npm:^4.1.0" peerDependencies: - "@cratis/arc.core": ^0.37.0 - "@cratis/arc.testing": ^0.37.0 + "@cratis/arc.core": ^0.38.0 + "@cratis/arc.testing": ^0.38.0 "@cratis/chronicle": ^6.7.0 "@cratis/fundamentals": ^7.19.6 rxjs: ^7.8.2 @@ -148,11 +148,13 @@ __metadata: mysql2: "npm:^3.24.4" pg: "npm:^8.0.0" postgres: "npm:^3.4.9" + rxjs: "npm:^7.8.2" sql.js: "npm:^1.14.2" peerDependencies: - "@cratis/arc.core": ^0.37.0 + "@cratis/arc.core": ^0.38.0 "@cratis/fundamentals": ^7.19.6 drizzle-orm: ^0.45.0 + rxjs: ^7.8.2 languageName: unknown linkType: soft @@ -215,7 +217,7 @@ __metadata: mongodb: "npm:^6.21.0" rxjs: "npm:^7.8.2" peerDependencies: - "@cratis/arc.core": ^0.37.0 + "@cratis/arc.core": ^0.38.0 "@cratis/fundamentals": ^7.19.6 "@opentelemetry/api": ^1.9.0 mongodb: ^6.21.0