diff --git a/ContractTests/Client/package.json b/ContractTests/Client/package.json index 5c7e4e95..72f63943 100644 --- a/ContractTests/Client/package.json +++ b/ContractTests/Client/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.core-client-contract", - "version": "0.25.0", + "version": "0.26.0", "private": true, "type": "module", "dependencies": { diff --git a/Documentation/index.md b/Documentation/index.md index a0fc1bdc..8b62a468 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.25.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.26.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/mongodb/change-stream-watcher.md b/Documentation/mongodb/change-stream-watcher.md new file mode 100644 index 00000000..712a9aeb --- /dev/null +++ b/Documentation/mongodb/change-stream-watcher.md @@ -0,0 +1,36 @@ +--- +title: Watch MongoDB changes across collections +description: React to tenant-scoped database changes with a shared stream and bounded observation lifetimes. +--- + +When several parts of your application need to respond to MongoDB changes, opening a separate change stream for each one is wasteful. Resolve `mongoDBWatcher` in the same Arc scope as your collections. The watcher shares one **database-level** change stream among subscriptions in that scope, and routes collection changes to each subscriber. + +This is a source-preview API. It is not a durable event consumer or a replacement for Chronicle. + +## React to changes + +The example assumes a configured Arc application, an active tenant scope, and registered `Book` and `Author` models: + +```typescript +import { mongoCollection, mongoDBWatcher } from '@cratis/arc.mongodb'; + +const books = await scope.resolve(mongoCollection(Book)); +const watcher = await scope.resolve(mongoDBWatcher); +const subscription = watcher.changes(books).subscribe(change => { + console.log(change.operationType, change.documentKey); +}); + +// On shutdown or when this work ends: +subscription.unsubscribe(); +await scope.dispose(); +``` + +`changes` returns an RxJS `Observable>`. Updates use MongoDB's `updateLookup` full-document option; deletes have a document key but no full document. Do not persist a resume token from this API or use it as an exactly-once feed. + +## Scope and failure + +MongoDB change streams require a replica set or sharded cluster and database-level watch permissions. A standalone server fails before the initial read. Join only collections resolved from the **same tenant and Arc scope** as the watcher; a cross-scope or cross-tenant collection is rejected. Disposing the scope, aborting its signal, or unsubscribing the last listener closes the cursor. A later subscription opens a new stream and starts at a new operation time. + +The driver resumes errors it recognizes as resumable. Other errors terminate all subscribers with `error`; the watcher does **not** silently reconnect and skip an unknown interval. Resubscribe with a fresh scope and re-read authoritative state if you need recovery. The watcher holds no durable checkpoint. + +For a live combined result, use [joined observation](joined-observe.md). For a single collection's complete snapshots, [observe on the collection](observing-collections.md) instead. diff --git a/Documentation/mongodb/geospatial.md b/Documentation/mongodb/geospatial.md new file mode 100644 index 00000000..3e686869 --- /dev/null +++ b/Documentation/mongodb/geospatial.md @@ -0,0 +1,53 @@ +--- +title: Store GeoJSON with Fundamentals types +description: Encode Point, LineString, and Polygon fields as MongoDB GeoJSON and restore typed geometry. +--- + +Use the geospatial types from `@cratis/fundamentals/geospatial` when a MongoDB read model stores locations, routes, or areas. `MongoDocumentCodec` writes their geometry as GeoJSON without changing driver-wide BSON serialization. + +## Point: one location + +Declare `@field(Point)` on the model. Longitude comes **first**, then latitude: + +```typescript +import { field } from '@cratis/fundamentals'; +import { Point } from '@cratis/fundamentals/geospatial'; +import { key } from '@cratis/arc.core'; + +class Place { + @field(String) @key() id!: string; + @field(Point) location!: Point; +} + +const place = Object.assign(new Place(), { id: 'one', location: new Point(10, 20) }); +// collection.codec.serialize(place).location -> { type: 'Point', coordinates: [10, 20] } +``` + +MongoDB can index the stored `location` field with a `2dsphere` index. Creating indexes and composing `$near` filters are application responsibilities; the Fundamentals `Point` is **not** a driver `GeoJSON` filter-builder argument. + +## LineString: a route + +A `LineString` holds two or more `Point`s. The codec writes `{ type: 'LineString', coordinates: [[10, 20], [30, 40]] }` for: + +```typescript +import { LineString, Point } from '@cratis/fundamentals/geospatial'; + +const route = new LineString([new Point(10, 20), new Point(30, 40)]); +``` + +Declare the field with `@field(LineString)`. Typed reads restore `LineString` and `Point` instances, not plain arrays. + +## Polygon: an area, optionally with holes + +A polygon stores the outer shell first and then its interior rings. Each `LinearRing` needs at least four points; repeat the first point as the last: + +```typescript +import { LinearRing, Point, Polygon } from '@cratis/fundamentals/geospatial'; + +const shell = new LinearRing([ + new Point(10, 20), new Point(30, 20), new Point(30, 40), new Point(10, 20) +]); +const area = new Polygon(shell); +``` + +Declare `@field(Polygon)` on the model. The BSON field contains `{ type: 'Polygon', coordinates: [/* shell coordinates, then holes */] }`; reads reconstruct `Polygon` and its `LinearRing`s. The codec rejects an unclosed ring, non-finite coordinates, a mismatched GeoJSON type, or missing coordinates rather than returning a misleading geometry. It does not enforce winding order, non-intersection, or MongoDB's full spatial-index rules. Validate domain geometry before writing, and use stored MongoDB field names in your own spatial filters. See [serializers](serializers.md) for the other model types and [naming policies](naming-policies.md) for stored names. diff --git a/Documentation/mongodb/index.md b/Documentation/mongodb/index.md index e38ad1ae..99c28a58 100644 --- a/Documentation/mongodb/index.md +++ b/Documentation/mongodb/index.md @@ -19,6 +19,9 @@ Your read models live in MongoDB, and every tenant has its own database. Wiring | Match Arc on .NET's property and collection naming | [Naming policies](naming-policies.md) | | Count, sort, and page in the database | [Paging](paging.md) | | Turn a change stream into an observable query | [Observing collections](observing-collections.md) | +| Combine two or three live collections | [Joined observation](joined-observe.md) | +| React to raw collection changes | [Change-stream watcher](change-stream-watcher.md) | +| Store GeoJSON Point, LineString, and Polygon fields | [Geospatial types](geospatial.md) | | Load a read model by command key | [Command context](../commands/command-context.md#load-a-read-model-by-key) | `withMongoDB` registers a read-model resolver for the models you list in `readModels`, so a command can declare `@inject(commandReadModel(TaskRecord))` and receive the document whose identity equals the command key. Do not also register another integration, such as Chronicle, as the owner of the same type; `build()` fails when two claim one type. @@ -33,6 +36,6 @@ The original `MongoReadModels` remains for low-level `defineQuery` users. ## Current boundaries -This integration does not supply transactions, a shared watcher or reconnect policy, joined observations, geospatial serializers, resilience middleware, or driver metrics. Do not infer any of those from Arc on .NET. The [capability reference](../reference/capabilities.md#persistence-and-chronicle) has the parity details. +This integration does not supply cross-store transactions or a durable change-stream checkpoint. The watcher shares a stream **within a tenant scope**, not across the process. Recognized transient reads retry at most twice; writes are not retried. Arc-owned clients expose OpenTelemetry MongoDB metrics, but caller-owned clients are not instrumented. Do not infer .NET's process-wide watcher or general-purpose resilience interceptors from these narrower guarantees. The [capability reference](../reference/capabilities.md#persistence-and-chronicle) has the parity details. Start with [Get started](getting-started.md). diff --git a/Documentation/mongodb/joined-observe.md b/Documentation/mongodb/joined-observe.md new file mode 100644 index 00000000..b529c30b --- /dev/null +++ b/Documentation/mongodb/joined-observe.md @@ -0,0 +1,27 @@ +--- +title: Observe joined MongoDB collections +description: Combine current tenant-scoped collections and re-emit the result when either changes. +--- + +A catalog can depend on books **and** their authors. Observing only books leaves the catalog stale when an author changes. Resolve the scoped watcher, join the two collections, and select the result you want to publish: + +```typescript +import { mongoCollection, mongoDBWatcher } from '@cratis/arc.mongodb'; + +const books = await scope.resolve(mongoCollection(Book)); +const authors = await scope.resolve(mongoCollection(Author)); +const watcher = await scope.resolve(mongoDBWatcher); +const catalog = watcher.observe(books, { published: true }) + .join(authors, { active: true }) + .select((publishedBooks, activeAuthors) => ({ publishedBooks, activeAuthors })); + +const subscription = catalog.subscribe(value => console.log(value)); +// When the consumer is done: +subscription.unsubscribe(); +``` + +The example is a service fragment: register both models with `withMongoDB({ readModels: [Book, Author], ... })`, and resolve the scope under a trusted tenant identity. The filters are **MongoDB filters in stored field names**, not predicates; use trusted values. See [naming policies](naming-policies.md#names-in-your-own-filters). + +`select` returns an RxJS `Observable`, rather than .NET's `ISubject`. It first emits a combined snapshot, then recomputes it after a change to either collection. Call `.join(thirdCollection, filter?)` before `select` to combine three collections. Every change in a watched collection triggers a refetch, even when the changed document does not match a filter. Consecutive changes while a read is in progress are coalesced into another complete snapshot, not buffered without a bound. Do not treat emissions as an audit trail or a transactionally consistent view across collections. + +The per-collection `maxObservableItems` cap applies to **each** joined snapshot. If a read, selector, or change stream fails, the observable errors instead of sending a partial list. The watcher shares one database-level stream per scope; subscriptions close on unsubscribe or scope disposal. Never keep a tenant-scoped collection in a singleton. See [change-stream watcher](change-stream-watcher.md) for recovery limits and [single-collection observation](observing-collections.md) when joins are unnecessary. diff --git a/Documentation/mongodb/observing-collections.md b/Documentation/mongodb/observing-collections.md index 35ceeb7e..6139a20a 100644 --- a/Documentation/mongodb/observing-collections.md +++ b/Documentation/mongodb/observing-collections.md @@ -66,7 +66,7 @@ Each emission then holds the first ten tasks by title, and `paging.totalItems` c - A full list is capped at 1,000 documents by default; `maxObservableItems` raises the cap to at most 10,000. A list that exceeds the cap, on the first read or on any later one, fails the subscription instead of sending a partial list. - A standalone MongoDB server is rejected with `MongoDB observe requires a replica set with change streams`. Replica sets and sharded clusters support change streams. -- Joined observation across collections is not available. +- For several collections, use [joined observation](joined-observe.md); each collection has its own snapshot cap. ## When the stream fails @@ -81,3 +81,4 @@ To recover, subscribe again. A new subscription opens a new change stream and re - [Observable queries](../queries/observable-queries.md) - [Paging](paging.md) - [MongoDB](index.md) +- [Change-stream watcher](change-stream-watcher.md) diff --git a/Documentation/mongodb/toc.yml b/Documentation/mongodb/toc.yml index 837f0285..81ecaa4a 100644 --- a/Documentation/mongodb/toc.yml +++ b/Documentation/mongodb/toc.yml @@ -12,3 +12,9 @@ href: paging.md - name: Observing collections href: observing-collections.md +- name: Joined observation + href: joined-observe.md +- name: Change-stream watcher + href: change-stream-watcher.md +- name: Geospatial types + href: geospatial.md diff --git a/Documentation/reference/capabilities.md b/Documentation/reference/capabilities.md index 54e5ab98..167600da 100644 --- a/Documentation/reference/capabilities.md +++ b/Documentation/reference/capabilities.md @@ -103,7 +103,7 @@ Evidence paths are relative to the repository root. Spec folders follow `for_ host.toString()).join(',')); + this.#clients.set(uri, client); + } return client; } async [Symbol.asyncDispose](): Promise { diff --git a/Source/MongoDB/MongoCollection.ts b/Source/MongoDB/MongoCollection.ts index 7ae19931..bdf8b942 100644 --- a/Source/MongoDB/MongoCollection.ts +++ b/Source/MongoDB/MongoCollection.ts @@ -8,6 +8,7 @@ import type { MongoCollectionOptions } from './MongoCollectionOptions.js'; import { MongoDocumentCodec } from './MongoDocumentCodec.js'; import { MongoObservation } from './MongoObservation.js'; import { MongoObservable } from './MongoObservable.js'; +import { retryRead } from './retryRead.js'; /** Tenant-bound model collection. The underlying driver collection remains available for writes. */ export class MongoCollection { @@ -27,12 +28,13 @@ export class MongoCollection { } /** Find models matching an application-owned filter. */ async find(filter: Filter = {}, options?: FindOptions): Promise { - const documents = await this.native.find(filter, { ...options, signal: this.context.signal }).toArray(); + const documents = await retryRead(() => this.native.find(filter, { ...options, signal: this.context.signal }).toArray(), this.context.signal); return documents.map(document => this.codec.deserialize(document)); } /** Read a model by its declared key. */ async findById(id: unknown): Promise { - const document = await this.native.findOne({ _id: this.codec.id(id) } as Filter, { signal: this.context.signal }); + const document = await retryRead(() => this.native.findOne({ _id: this.codec.id(id) } as Filter, + { signal: this.context.signal }), this.context.signal); return document ? this.codec.deserialize(document) : null; } /** Count and page in MongoDB, using only fields declared in the model for client sorting. */ @@ -50,14 +52,18 @@ export class MongoCollection { ...Object.fromEntries(Object.entries(findOptions?.sort ?? {}).filter(([name]) => name !== field)), ...(field === '_id' ? {} : { _id: 1 as const }) } : { ...findOptions?.sort, ...(!findOptions?.sort || !Object.hasOwn(findOptions.sort, '_id') ? { _id: 1 as const } : {}) }; - const total = await this.native.countDocuments(filter, { collation: findOptions?.collation, session: findOptions?.session, - signal: this.context.signal } as Parameters[1]); - const documents = await this.native.find(filter, { ...findOptions, sort, signal: this.context.signal }) - .skip(page * pageSize).limit(pageSize).toArray(); + const total = await retryRead(() => this.native.countDocuments(filter, { collation: findOptions?.collation, + session: findOptions?.session, signal: this.context.signal } as Parameters[1]), + this.context.signal); + const documents = await retryRead(() => this.native.find(filter, { ...findOptions, sort, signal: this.context.signal }) + .skip(page * pageSize).limit(pageSize).toArray(), this.context.signal); return queryPage(documents.map(document => this.codec.deserialize(document)), total, sorting); } - private async readObservable(filter: Filter): Promise { - const documents = await this.native.find(filter, { signal: this.context.signal }).limit(this.#maxObservableItems + 1).toArray(); + /** Read a bounded snapshot for joined observation. */ + async readForObservation(filter: Filter = {}): Promise { + if (this.context.signal.aborted) throw new DOMException('Aborted', 'AbortError'); + const documents = await retryRead(() => this.native.find(filter, { signal: this.context.signal }) + .limit(this.#maxObservableItems + 1).toArray(), this.context.signal); if (documents.length > this.#maxObservableItems) throw new RangeError('MongoDB observation exceeds maxObservableItems'); return documents.map(document => this.codec.deserialize(document)); } @@ -87,15 +93,19 @@ export class MongoCollection { } /** Open an async-iterable observation for consumers that do not use RxJS. */ observeIterable(filter: Filter = {}): Promise> { - return this.openObservation(() => this.readObservable(filter), []); + return this.openObservation(() => this.readForObservation(filter), []); } /** Open an async-iterable keyed observation for consumers that do not use RxJS. */ observeByIdIterable(id: unknown): Promise> { const filter = { _id: this.codec.id(id) } as Filter; - const read = () => this.native.findOne(filter, { signal: this.context.signal }) + const read = () => retryRead(() => this.native.findOne(filter, { signal: this.context.signal }), this.context.signal) .then(document => document ? this.codec.deserialize(document) : null); return this.openObservation(read, [{ $match: { 'documentKey._id': filter._id } }]); } + /** Refuse to join a collection from a different tenant or execution scope. */ + belongsTo(databaseName: string, context: ExecutionContext): boolean { + return this.database.databaseName === databaseName && this.context === context; + } /** End all observations with the owning Arc execution scope. */ async [Symbol.asyncDispose](): Promise { await Promise.all([...this.#observations].map(observation => observation.close())); diff --git a/Source/MongoDB/MongoDBClientMetrics.ts b/Source/MongoDB/MongoDBClientMetrics.ts new file mode 100644 index 00000000..46c177b1 --- /dev/null +++ b/Source/MongoDB/MongoDBClientMetrics.ts @@ -0,0 +1,23 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { metrics } from '@opentelemetry/api'; +import type { MongoClient } from 'mongodb'; + +/** Attach driver event instruments to Arc-owned MongoDB clients before connecting. */ +export function monitorMongoDBClient(client: MongoClient, server: string): void { + const meter = metrics.getMeter('Cratis.Arc.MongoDB'); + const open = meter.createUpDownCounter('mongodb-open-connections'); + const pool = meter.createUpDownCounter('mongodb-connections-in-pool'); + const commands = meter.createUpDownCounter('mongodb-commands'); + const failures = meter.createCounter('mongodb-failed-connections'); + const aggregate = meter.createCounter('mongodb-aggregated-commands'); + const labels = { Server: server }; + client.on('connectionCreated', () => { open.add(1, { ...labels, state: 'open' }); pool.add(1, labels); }); + client.on('connectionClosed', () => { open.add(-1, { ...labels, state: 'open' }); pool.add(-1, labels); }); + client.on('connectionCheckedOut', () => open.add(1, { ...labels, state: 'checked_out' })); + client.on('connectionCheckedIn', () => open.add(-1, { ...labels, state: 'checked_out' })); + client.on('connectionCheckOutFailed', () => failures.add(1, labels)); + client.on('commandStarted', () => { commands.add(1, labels); aggregate.add(1, labels); }); + client.on('commandSucceeded', () => commands.add(-1, labels)); + client.on('commandFailed', () => commands.add(-1, labels)); +} diff --git a/Source/MongoDB/MongoDBWatcher.ts b/Source/MongoDB/MongoDBWatcher.ts new file mode 100644 index 00000000..33154982 --- /dev/null +++ b/Source/MongoDB/MongoDBWatcher.ts @@ -0,0 +1,185 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import type { ExecutionContext } from '@cratis/arc.core'; +import { Observable } from 'rxjs'; +import type { ChangeStream, ChangeStreamDocument, Db, Document, Filter, Timestamp } from 'mongodb'; +import type { MongoCollection } from './MongoCollection.js'; + +function collectionName(change: ChangeStreamDocument): string | undefined { + return 'ns' in change && change.ns && 'coll' in change.ns ? change.ns.coll : undefined; +} + +/** Tenant- and scope-bound database watcher. Its subscribers share one database change stream. */ +export class MongoDBWatcher { + readonly #listeners = new Set<{ change: (change: ChangeStreamDocument) => void; error: (error: unknown) => void }>(); + #stream?: ChangeStream; + #opening?: Promise; + #disposed = false; + #generation = 0; + readonly #abort = () => { void this[Symbol.asyncDispose](); }; + + constructor(private readonly database: Db, private readonly context: ExecutionContext) { + context.signal.addEventListener('abort', this.#abort, { once: true }); + if (context.signal.aborted) this.#abort(); + } + + /** Start a typed observation. Join another collection before selecting the result. */ + observe(collection: MongoCollection, filter: Filter = {}): MongoDBObserveBuilder { + this.assertCollection(collection); + return new MongoDBObserveBuilder(this, collection, filter); + } + + /** React to raw changes on one collection, including deletes (which have no full document). */ + changes(collection: MongoCollection): Observable> { + this.assertCollection(collection); + return new Observable(subscriber => { + void this.listen(change => { + if (collectionName(change) === collection.native.collectionName) subscriber.next(change); + }, error => subscriber.error(error)).then(release => { + if (subscriber.closed) release(); + else subscriber.add(release); + }).catch(error => subscriber.error(error)); + }); + } + + /** Build a complete joined snapshot; notifications are coalesced, never buffered without a bound. */ + select(collections: readonly MongoCollection[], filters: readonly Filter[], + selector: (...documents: object[][]) => R): Observable { + collections.forEach(collection => this.assertCollection(collection)); + const names = new Set(collections.map(collection => collection.native.collectionName)); + return new Observable(subscriber => { + let dirty = false; + let running = false; + let release: (() => void) | undefined; + const read = async () => { + if (running || subscriber.closed) return; + running = true; + try { + do { + dirty = false; + const values = await Promise.all(collections.map((collection, index) => + collection.readForObservation(filters[index]))); + if (!subscriber.closed && !this.#disposed) subscriber.next(selector(...values)); + } while (dirty && !subscriber.closed && !this.#disposed); + } catch (error) { if (!subscriber.closed) subscriber.error(error); } + finally { running = false; } + }; + void this.listen(change => { + if (names.has(collectionName(change) ?? '')) { dirty = true; void read(); } + }, error => subscriber.error(error)).then(stop => { + release = stop; + if (subscriber.closed) stop(); + else void read(); + }).catch(error => { if (!subscriber.closed) subscriber.error(error); }); + return () => { release?.(); }; + }); + } + + private assertCollection(collection: MongoCollection): void { + if (this.#disposed || !collection.belongsTo(this.database.databaseName, this.context)) + throw new Error('MongoDB watcher requires a collection from the same tenant and scope'); + } + + private async listen(change: (event: ChangeStreamDocument) => void, error: (error: unknown) => void): Promise<() => void> { + if (this.#disposed) throw new Error('MongoDB watcher is disposed'); + const listener = { change, error }; + this.#listeners.add(listener); + try { await (this.#opening ??= this.start()); } + catch (failure) { this.#listeners.delete(listener); throw failure; } + if (this.#disposed) { this.#listeners.delete(listener); throw new Error('MongoDB watcher is disposed'); } + return () => { + this.#listeners.delete(listener); + if (!this.#listeners.size) void this.stop(); + }; + } + + private async start(generation = this.#generation): Promise { + try { + const hello = await this.database.command({ hello: 1 }, { signal: this.context.signal }); + if (typeof hello.setName !== 'string' && hello.msg !== 'isdbgrid') + throw new Error('MongoDB observe requires a replica set with change streams'); + const operationTime = (hello.operationTime ?? hello.$clusterTime?.clusterTime) as Timestamp | undefined; + if (!operationTime) throw new Error('MongoDB observe requires an operation time from the server'); + if (this.#disposed || generation !== this.#generation) return; + if (!this.#listeners.size) { this.#opening = undefined; return; } + const stream = this.database.watch([], { startAtOperationTime: operationTime, fullDocument: 'updateLookup' }); + this.#stream = stream; + void this.pump(stream); + } catch (error) { + if (generation === this.#generation) this.#opening = undefined; + throw error; + } + } + + private async pump(stream: ChangeStream): Promise { + try { + while (this.#stream === stream && !this.#disposed) { + const event = await stream.next(); + if (!event) throw new Error('MongoDB change stream ended'); + if (this.#stream !== stream || this.#disposed) break; + for (const listener of [...this.#listeners]) listener.change(event); + } + } catch (error) { + if (this.#stream === stream && !this.#disposed) { + const listeners = [...this.#listeners]; + this.#listeners.clear(); + try { await this.stop(); } + finally { for (const listener of listeners) listener.error(error); } + } + } finally { + if (this.#stream === stream) await this.stop(); + } + } + + private async stop(): Promise { + this.#generation++; + const stream = this.#stream; + this.#stream = undefined; + this.#opening = undefined; + if (stream) await stream.close(); + } + + async [Symbol.asyncDispose](): Promise { + if (this.#disposed) return; + this.#disposed = true; + this.context.signal.removeEventListener('abort', this.#abort); + for (const listener of [...this.#listeners]) listener.error(new DOMException('Aborted', 'AbortError')); + this.#listeners.clear(); + await this.stop(); + } +} + +/** Join two or three tenant-scoped collections and select a live result. */ +export class MongoDBObserveBuilder { + constructor(private readonly watcher: MongoDBWatcher, private readonly primary: MongoCollection, + private readonly filter: Filter) {} + join(collection: MongoCollection, filter: Filter = {}): MongoDBJoinedObserveBuilder { + return new MongoDBJoinedObserveBuilder(this.watcher, this.primary, this.filter, collection, filter); + } +} + +export class MongoDBJoinedObserveBuilder { + constructor(private readonly watcher: MongoDBWatcher, private readonly primary: MongoCollection, + private readonly primaryFilter: Filter, private readonly related: MongoCollection, + private readonly relatedFilter: Filter) {} + join(collection: MongoCollection, filter: Filter = {}): MongoDBThreeWayObserveBuilder { + return new MongoDBThreeWayObserveBuilder(this.watcher, this.primary, this.primaryFilter, this.related, + this.relatedFilter, collection, filter); + } + select(selector: (primary: T[], related: U[]) => R): Observable { + return this.watcher.select([this.primary, this.related], [this.primaryFilter, this.relatedFilter], + (primary, related) => selector(primary as T[], related as U[])); + } +} + +export class MongoDBThreeWayObserveBuilder { + constructor(private readonly watcher: MongoDBWatcher, private readonly primary: MongoCollection, + private readonly primaryFilter: Filter, private readonly related: MongoCollection, + private readonly relatedFilter: Filter, private readonly third: MongoCollection, + private readonly thirdFilter: Filter) {} + select(selector: (primary: T[], related: U[], third: V[]) => R): Observable { + return this.watcher.select([this.primary, this.related, this.third], + [this.primaryFilter, this.relatedFilter, this.thirdFilter], + (primary, related, third) => selector(primary as T[], related as U[], third as V[])); + } +} diff --git a/Source/MongoDB/MongoDocumentCodec.ts b/Source/MongoDB/MongoDocumentCodec.ts index 67b38ad3..c3697a7f 100644 --- a/Source/MongoDB/MongoDocumentCodec.ts +++ b/Source/MongoDB/MongoDocumentCodec.ts @@ -1,10 +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 { ConceptAs, DateOnly, DerivedType, Guid, TimeOnly, TimeSpan } from '@cratis/fundamentals'; +import { LineString, Point, Polygon } from '@cratis/fundamentals/geospatial'; import { fieldsFor, InvalidQuerySort, wireName } from '@cratis/arc.core'; import type { WireField } from '@cratis/arc.core'; import { Binary, Decimal128, Long, ObjectId } from 'mongodb'; import { defaultMongoNamingPolicy } from './MongoNamingPolicy.js'; +import { decodeGeometry, encodeGeometry } from './MongoGeoJSON.js'; import type { MongoNamingPolicy } from './MongoNamingPolicy.js'; import type { Document } from 'mongodb'; @@ -84,6 +86,10 @@ export class MongoDocumentCodec { if (value instanceof DateOnly) return new Date(`${(value as DateOnly).toString()}T12:00:00.000Z`); if (value instanceof TimeOnly) return new Date(`1970-01-01T${(value as TimeOnly).toString()}Z`); if (value instanceof TimeSpan) return (value as TimeSpan).toString(); + if (type === Point || type === LineString || type === Polygon) { + if (!(value instanceof type)) throw new TypeError(`Expected MongoDB ${type.name}`); + return encodeGeometry(value as Point | LineString | Polygon); + } if (type === Date || type === String || type === Number || type === Boolean || type === ObjectId || type === Binary) return value; if (fieldsFor(type).length) { const runtime = (value as object).constructor as WireField['type']; @@ -123,6 +129,7 @@ export class MongoDocumentCodec { if (type === DateOnly) return DateOnly.parse((value as Date).toISOString().slice(0, 10)); if (type === TimeOnly) return TimeOnly.parse((value as Date).toISOString().slice(11, 23)); if (type === TimeSpan) return TimeSpan.parse(value as string); + if (type === Point || type === LineString || type === Polygon) return decodeGeometry(type as typeof Point | typeof LineString | typeof Polygon, value); if (type === Number) { if (value instanceof Decimal128 || value instanceof Long) { const number = Number(value.toString()); diff --git a/Source/MongoDB/MongoGeoJSON.ts b/Source/MongoDB/MongoGeoJSON.ts new file mode 100644 index 00000000..1bd61829 --- /dev/null +++ b/Source/MongoDB/MongoGeoJSON.ts @@ -0,0 +1,57 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { LineString, LinearRing, Point, Polygon } from '@cratis/fundamentals/geospatial'; +import type { Document } from 'mongodb'; + +export type MongoGeometry = Point | LineString | Polygon; + +/** Encode Fundamentals geometry as MongoDB GeoJSON, with longitude before latitude. */ +export function encodeGeometry(value: MongoGeometry): Document { + if (value instanceof Point) return { type: 'Point', coordinates: point(value) }; + if (value instanceof LineString) { + if (value.coordinates.length < 2) throw new TypeError('GeoJSON LineString requires at least two points'); + return { type: 'LineString', coordinates: value.coordinates.map(point) }; + } + const rings = [value.shell, ...value.holes]; + return { type: 'Polygon', coordinates: rings.map(ring => { + const coordinates = ring.coordinates.map(point); + if (coordinates.length < 4 || coordinates[0]![0] !== coordinates.at(-1)![0] || + coordinates[0]![1] !== coordinates.at(-1)![1]) throw new TypeError('GeoJSON Polygon rings must be closed'); + return coordinates; + }) }; +} + +/** Reject malformed or mismatched geometry rather than returning a plausible but incorrect location. */ +export function decodeGeometry(type: typeof Point | typeof LineString | typeof Polygon, document: unknown): MongoGeometry { + if (!document || typeof document !== 'object') throw new TypeError('Invalid MongoDB GeoJSON geometry'); + const value = document as Document; + if (value.type !== type.name || !Array.isArray(value.coordinates)) throw new TypeError(`Expected GeoJSON ${type.name}`); + if (type === Point) return readPoint(value.coordinates); + if (type === LineString) { + if (value.coordinates.length < 2) throw new TypeError('GeoJSON LineString requires at least two points'); + return new LineString(value.coordinates.map(readPoint)); + } + if (!value.coordinates.length) throw new TypeError('GeoJSON Polygon requires a shell'); + const rings = value.coordinates.map((coordinates: unknown) => { + if (!Array.isArray(coordinates) || coordinates.length < 4) throw new TypeError('GeoJSON Polygon rings require four points'); + const points = coordinates.map(readPoint); + if (points[0]!.longitude !== points.at(-1)!.longitude || points[0]!.latitude !== points.at(-1)!.latitude) + throw new TypeError('GeoJSON Polygon rings must be closed'); + return new LinearRing(points); + }); + return new Polygon(rings[0]!, rings.slice(1)); +} + +function point(value: Point): number[] { + if (!Number.isFinite(value.longitude) || !Number.isFinite(value.latitude)) throw new TypeError('Invalid GeoJSON coordinate'); + return [value.longitude, value.latitude]; +} + +function readPoint(coordinates: unknown): Point { + if (!Array.isArray(coordinates) || coordinates.length !== 2 || + typeof coordinates[0] !== 'number' || typeof coordinates[1] !== 'number') + throw new TypeError('Invalid GeoJSON coordinate'); + if (!Number.isFinite(coordinates[0]) || !Number.isFinite(coordinates[1])) + throw new TypeError('Invalid GeoJSON coordinate'); + return new Point(coordinates[0], coordinates[1]); +} diff --git a/Source/MongoDB/for_MongoClientFactory/when_instrumenting_an_owned_client.ts b/Source/MongoDB/for_MongoClientFactory/when_instrumenting_an_owned_client.ts new file mode 100644 index 00000000..09d74a40 --- /dev/null +++ b/Source/MongoDB/for_MongoClientFactory/when_instrumenting_an_owned_client.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 { afterEach, beforeEach, describe, it, should } from 'vitest'; +import { metrics } from '@opentelemetry/api'; +import { AggregationTemporality, InMemoryMetricExporter, MeterProvider, PeriodicExportingMetricReader } from '@opentelemetry/sdk-metrics'; +import { MongoClient } from 'mongodb'; +import { given } from '../given.js'; +import { MongoClientFactory } from '../MongoClientFactory.js'; +import type { ExecutionContext } from '@cratis/arc.core'; + +should(); +class an_owned_client { + readonly exporter = new InMemoryMetricExporter(AggregationTemporality.CUMULATIVE); + readonly reader = new PeriodicExportingMetricReader({ exporter: this.exporter, exportIntervalMillis: 60000 }); + readonly provider = new MeterProvider({ readers: [this.reader] }); + readonly factory = new MongoClientFactory({ server: 'mongodb://user:secret@localhost:27017', readModels: [] }); + readonly context: ExecutionContext = { tenantId: 'default', signal: new AbortController().signal, + allowedSeverity: 3, correlationId: 'test', principal: undefined }; +} +describe('when instrumenting an owned client', given(an_owned_client, context => { + let client: MongoClient; + beforeEach(() => { + metrics.setGlobalMeterProvider(context.provider); + client = context.factory.get(context.context); + client.emit('commandStarted', { commandName: 'find' } as never); + client.emit('commandSucceeded', { commandName: 'find' } as never); + }); + afterEach(async () => { + await context.factory[Symbol.asyncDispose](); + await context.provider.shutdown(); + metrics.disable(); + }); + it('should count commands with a server label', async () => { + await context.reader.forceFlush(); + const instruments = context.exporter.getMetrics().flatMap(result => result.scopeMetrics.flatMap(scope => scope.metrics)); + const aggregate = instruments.find(instrument => instrument.descriptor.name === 'mongodb-aggregated-commands'); + aggregate!.dataPoints[0]!.value.should.equal(1); + (aggregate!.dataPoints[0]!.attributes.Server as string).should.equal('localhost:27017'); + }); +})); diff --git a/Source/MongoDB/for_MongoCollection/given/RelatedRecord.ts b/Source/MongoDB/for_MongoCollection/given/RelatedRecord.ts new file mode 100644 index 00000000..44d51fee --- /dev/null +++ b/Source/MongoDB/for_MongoCollection/given/RelatedRecord.ts @@ -0,0 +1,9 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { field } from '@cratis/fundamentals'; +import { key } from '@cratis/arc.core'; + +export class RelatedRecord { + @field(String) @key() id!: string; + @field(String) name!: string; +} diff --git a/Source/MongoDB/for_MongoCollection/when_observing_joined_collections/with_a_replica_set.integration.ts b/Source/MongoDB/for_MongoCollection/when_observing_joined_collections/with_a_replica_set.integration.ts new file mode 100644 index 00000000..9d0d8dbd --- /dev/null +++ b/Source/MongoDB/for_MongoCollection/when_observing_joined_collections/with_a_replica_set.integration.ts @@ -0,0 +1,58 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { afterEach, beforeEach, describe, it, should } from 'vitest'; +import { ArcApplication } from '@cratis/arc.core'; +import { Guid } from '@cratis/fundamentals'; +import { firstValueFrom, ReplaySubject, skip, take, timeout } from 'rxjs'; +import type { Subscription } from 'rxjs'; +import { given } from '../../given.js'; +import { mongoCollection, mongoDBWatcher } from '../../index.js'; +import { TaskRecord } from '../given/TaskRecord.js'; +import { RelatedRecord } from '../given/RelatedRecord.js'; +import { a_replica_set } from '../given/a_replica_set.js'; + +should(); +describe('when observing joined collections with a replica set', given(a_replica_set, context => { + let application: ArcApplication; + beforeEach(async () => { + if (!process.env.ARC_MONGO_TEST_URI) throw new Error('ARC_MONGO_TEST_URI is required'); + await context.client.connect(); + application = await ArcApplication.createBuilder().withMongoDB({ client: context.client, + databaseNameResolver: tenant => `${context.name}_${tenant}`, readModels: [TaskRecord, RelatedRecord] }).build(); + }); + afterEach(async () => { + await application.dispose(); + await context.client.db(`${context.name}_a`).dropDatabase(); + await context.client.db(`${context.name}_b`).dropDatabase(); + await context.client.close(); + }); + it('should emit snapshots after either collection changes without leaking another tenant', async () => { + const scope = application.server.services.createScope(context.context('a')); + const other = application.server.services.createScope(context.context('b')); + let subscription: Subscription | undefined; + try { + const primary = await scope.resolve(mongoCollection(TaskRecord)); + const related = await scope.resolve(mongoCollection(RelatedRecord)); + const watcher = await scope.resolve(mongoDBWatcher); + const source = watcher.observe(primary).join(related).select((tasks, relations) => + [tasks.map(task => task.title), relations.map(relation => relation.name)] as [string[], string[]]); + const updates = new ReplaySubject<[string[], string[]]>(1); + subscription = source.subscribe(updates); + (await firstValueFrom(updates.pipe(take(1), timeout(10000)))).should.deep.equal([[], []]); + const first = firstValueFrom(updates.pipe(skip(1), take(1), timeout(10000))); + await primary.native.insertOne(primary.codec.serialize(Object.assign(new TaskRecord(), { + id: Guid.parse('00112233-4455-6677-8899-aabbccddeeff'), title: 'task' + }))); + (await first)[0].should.deep.equal(['task']); + const second = firstValueFrom(updates.pipe(skip(1), take(1), timeout(10000))); + await related.native.insertOne(related.codec.serialize(Object.assign(new RelatedRecord(), { id: 'one', name: 'related' }))); + (await second)[1].should.deep.equal(['related']); + const otherTask = await other.resolve(mongoCollection(TaskRecord)); + (() => watcher.observe(otherTask)).should.throw('same tenant and scope'); + } finally { + subscription?.unsubscribe(); + await scope.dispose(); + await other.dispose(); + } + }, 30000); +})); diff --git a/Source/MongoDB/for_MongoDBWatcher/when_reacting_to_changes/with_a_replica_set.integration.ts b/Source/MongoDB/for_MongoDBWatcher/when_reacting_to_changes/with_a_replica_set.integration.ts new file mode 100644 index 00000000..67a2fc13 --- /dev/null +++ b/Source/MongoDB/for_MongoDBWatcher/when_reacting_to_changes/with_a_replica_set.integration.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 { afterEach, beforeEach, describe, it, should } from 'vitest'; +import { ArcApplication } from '@cratis/arc.core'; +import { Guid } from '@cratis/fundamentals'; +import { firstValueFrom, take, timeout } from 'rxjs'; +import { given } from '../../given.js'; +import { mongoCollection, mongoDBWatcher } from '../../index.js'; +import { TaskRecord } from '../../for_MongoCollection/given/TaskRecord.js'; +import { a_replica_set } from '../../for_MongoCollection/given/a_replica_set.js'; + +should(); +describe('when reacting to changes with a replica set', given(a_replica_set, context => { + let application: ArcApplication; + beforeEach(async () => { + if (!process.env.ARC_MONGO_TEST_URI) throw new Error('ARC_MONGO_TEST_URI is required'); + await context.client.connect(); + application = await ArcApplication.createBuilder().withMongoDB({ client: context.client, + databaseNameResolver: tenant => `${context.name}_${tenant}`, readModels: [TaskRecord] }).build(); + }); + afterEach(async () => { + await application.dispose(); + await context.client.db(`${context.name}_a`).dropDatabase(); + await context.client.close(); + }); + it('should deliver the inserted collection change and close with the scope', async () => { + const scope = application.server.services.createScope(context.context('a')); + try { + const collection = await scope.resolve(mongoCollection(TaskRecord)); + const watcher = await scope.resolve(mongoDBWatcher); + const change = firstValueFrom(watcher.changes(collection).pipe(take(1), timeout(10000))); + // The watcher opens lazily: wait for its initial snapshot before writing. + const ready = firstValueFrom(watcher.observe(collection).join(collection).select((items) => items) + .pipe(take(1), timeout(10000))); + await ready; + await collection.native.insertOne(collection.codec.serialize(Object.assign(new TaskRecord(), { + id: Guid.parse('00112233-4455-6677-8899-aabbccddeeff'), title: 'task' + }))); + (await change).operationType.should.equal('insert'); + } finally { await scope.dispose(); } + }, 30000); +})); diff --git a/Source/MongoDB/for_MongoDBWatcher/when_stream_fails/with_nonresumable_error.ts b/Source/MongoDB/for_MongoDBWatcher/when_stream_fails/with_nonresumable_error.ts new file mode 100644 index 00000000..7188301d --- /dev/null +++ b/Source/MongoDB/for_MongoDBWatcher/when_stream_fails/with_nonresumable_error.ts @@ -0,0 +1,44 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. +import { beforeEach, describe, it, should } from 'vitest'; +import { firstValueFrom } from 'rxjs'; +import { Timestamp } from 'mongodb'; +import type { Collection, Db, Document } from 'mongodb'; +import type { ExecutionContext } from '@cratis/arc.core'; +import { given } from '../../given.js'; +import { MongoCollection } from '../../MongoCollection.js'; +import { MongoDBWatcher } from '../../MongoDBWatcher.js'; +import { TaskRecord } from '../../for_MongoCollection/given/TaskRecord.js'; + +should(); +class a_failing_watcher { + readonly error = new Error('stream lost'); + readonly context: ExecutionContext = { tenantId: 'a', signal: new AbortController().signal, + allowedSeverity: 3, correlationId: 'test', principal: undefined }; + watches = 0; + readonly database = { + databaseName: 'tenant', + command: async () => ({ setName: 'rs0', operationTime: new Timestamp({ t: 1, i: 1 }) }), + watch: () => { + this.watches++; + return { next: async () => { throw this.error; }, close: async () => {} }; + } + } as unknown as Db; + readonly collection = new MongoCollection({ collectionName: 'Tasks' } as Collection, this.database, + TaskRecord, this.context); + readonly watcher = new MongoDBWatcher(this.database, this.context); +} +describe('when a stream fails with a nonresumable error', given(a_failing_watcher, context => { + let failure: unknown; + beforeEach(async () => { + try { await firstValueFrom(context.watcher.changes(context.collection)); } + catch (error) { failure = error; } + }); + it('should fail subscribers and allow a fresh stream on resubscription', async () => { + (failure as Error).should.equal(context.error); + try { await firstValueFrom(context.watcher.changes(context.collection)); } + catch (error) { (error as Error).should.equal(context.error); } + context.watches.should.equal(2); + await context.watcher[Symbol.asyncDispose](); + }); +})); diff --git a/Source/MongoDB/for_MongoDocumentCodec/given/a_geospatial_model.ts b/Source/MongoDB/for_MongoDocumentCodec/given/a_geospatial_model.ts new file mode 100644 index 00000000..4a055959 --- /dev/null +++ b/Source/MongoDB/for_MongoDocumentCodec/given/a_geospatial_model.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 { field } from '@cratis/fundamentals'; +import { LineString, Point, Polygon } from '@cratis/fundamentals/geospatial'; +import { key } from '@cratis/arc.core'; +import { MongoDocumentCodec } from '../../MongoDocumentCodec.js'; + +export class GeoRecord { + @field(String) @key() id!: string; + @field(Point) location!: Point; + @field(LineString) route!: LineString; + @field(Polygon) region!: Polygon; +} + +export class a_geospatial_model { + readonly codec = new MongoDocumentCodec(GeoRecord); +} diff --git a/Source/MongoDB/for_MongoDocumentCodec/when_rejecting_invalid_geometry.ts b/Source/MongoDB/for_MongoDocumentCodec/when_rejecting_invalid_geometry.ts new file mode 100644 index 00000000..02e1efe8 --- /dev/null +++ b/Source/MongoDB/for_MongoDocumentCodec/when_rejecting_invalid_geometry.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 { describe, it, should } from 'vitest'; +import { given } from '../given.js'; +import { a_geospatial_model } from './given/a_geospatial_model.js'; + +should(); +describe('when rejecting invalid geometry', given(a_geospatial_model, context => { + it('should reject a mismatched type', () => { + (() => context.codec.deserialize({ _id: 'one', location: { type: 'Polygon', coordinates: [1, 2] } })) + .should.throw(TypeError, 'Expected GeoJSON Point'); + }); + it('should reject an unclosed polygon ring', () => { + (() => context.codec.deserialize({ _id: 'one', region: { type: 'Polygon', coordinates: [ + [[1, 1], [2, 1], [2, 2], [3, 3]] + ] } })).should.throw(TypeError, 'GeoJSON Polygon rings must be closed'); + }); +})); diff --git a/Source/MongoDB/for_MongoDocumentCodec/when_round_tripping_geospatial_fields.ts b/Source/MongoDB/for_MongoDocumentCodec/when_round_tripping_geospatial_fields.ts new file mode 100644 index 00000000..a9ad2b29 --- /dev/null +++ b/Source/MongoDB/for_MongoDocumentCodec/when_round_tripping_geospatial_fields.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 { beforeEach, describe, it, should } from 'vitest'; +import { LineString, LinearRing, Point, Polygon } from '@cratis/fundamentals/geospatial'; +import { given } from '../given.js'; +import { a_geospatial_model, GeoRecord } from './given/a_geospatial_model.js'; + +should(); +describe('when round tripping geospatial fields', given(a_geospatial_model, context => { + let document: ReturnType; + let restored: GeoRecord; + beforeEach(() => { + const origin = new Point(10, 20); + const destination = new Point(30, 40); + const shell = new LinearRing([origin, destination, new Point(30, 20), origin]); + const hole = new LinearRing([new Point(11, 21), new Point(12, 22), new Point(12, 21), new Point(11, 21)]); + document = context.codec.serialize(Object.assign(new GeoRecord(), { id: 'one', location: origin, + route: new LineString([origin, destination]), region: new Polygon(shell, [hole]) })); + restored = context.codec.deserialize(document); + }); + it('should store GeoJSON with longitude before latitude and polygon holes', () => { + document.location.should.deep.equal({ type: 'Point', coordinates: [10, 20] }); + document.route.should.deep.equal({ type: 'LineString', coordinates: [[10, 20], [30, 40]] }); + document.region.coordinates.should.have.lengthOf(2); + }); + it('should restore fundamentals geospatial instances', () => { + restored.location.should.be.instanceOf(Point); + restored.route.should.be.instanceOf(LineString); + restored.region.should.be.instanceOf(Polygon); + restored.region.holes[0]!.coordinates[1]!.latitude.should.equal(22); + }); +})); diff --git a/Source/MongoDB/for_retryRead/when_exhausting_read_retries.ts b/Source/MongoDB/for_retryRead/when_exhausting_read_retries.ts new file mode 100644 index 00000000..ed21febb --- /dev/null +++ b/Source/MongoDB/for_retryRead/when_exhausting_read_retries.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 { beforeEach, describe, it, should } from 'vitest'; +import { MongoNetworkError } from 'mongodb'; +import { given } from '../given.js'; +import { retryRead } from '../retryRead.js'; + +should(); +class a_failing_read { + readonly signal = new AbortController().signal; + calls = 0; + readonly error = new MongoNetworkError('unavailable'); +} +describe('when exhausting read retries', given(a_failing_read, context => { + let failure: unknown; + beforeEach(async () => { + try { await retryRead(async () => { context.calls++; throw context.error; }, context.signal); } + catch (error) { failure = error; } + }); + it('should report the original error without a partial result', () => { + (failure as Error).should.equal(context.error); + context.calls.should.equal(3); + }); +})); diff --git a/Source/MongoDB/for_retryRead/when_retrying_a_transient_read.ts b/Source/MongoDB/for_retryRead/when_retrying_a_transient_read.ts new file mode 100644 index 00000000..58799ac5 --- /dev/null +++ b/Source/MongoDB/for_retryRead/when_retrying_a_transient_read.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 { beforeEach, describe, it, should } from 'vitest'; +import { MongoNetworkError } from 'mongodb'; +import { given } from '../given.js'; +import { retryRead } from '../retryRead.js'; + +should(); +class a_transient_read { + readonly signal = new AbortController().signal; + calls = 0; +} +describe('when retrying a transient read', given(a_transient_read, context => { + let value: number; + beforeEach(async () => { + value = await retryRead(async () => { + context.calls++; + if (context.calls < 3) throw new MongoNetworkError('transient'); + return 42; + }, context.signal); + }); + it('should return the recovered value after two retries', () => { + value.should.equal(42); + context.calls.should.equal(3); + }); +})); diff --git a/Source/MongoDB/index.ts b/Source/MongoDB/index.ts index 475feacc..a9d319c0 100644 --- a/Source/MongoDB/index.ts +++ b/Source/MongoDB/index.ts @@ -1,7 +1,10 @@ // Copyright (c) Cratis. All rights reserved. // Licensed under the MIT license. See LICENSE file in the project root for full license information. export { mongoCollection } from './collectionToken.js'; -export { withMongoDB, mongoClientFactory } from './withMongoDB.js'; +export { withMongoDB, mongoClientFactory, mongoDBWatcher } from './withMongoDB.js'; +export { MongoDBWatcher, MongoDBObserveBuilder, MongoDBJoinedObserveBuilder, MongoDBThreeWayObserveBuilder } from './MongoDBWatcher.js'; +export { encodeGeometry, decodeGeometry } from './MongoGeoJSON.js'; +export type { MongoGeometry } from './MongoGeoJSON.js'; export { defaultMongoNamingPolicy, camelCaseMongoNamingPolicy } from './MongoNamingPolicy.js'; export type { MongoNamingPolicy } from './MongoNamingPolicy.js'; export { MongoClientFactory } from './MongoClientFactory.js'; diff --git a/Source/MongoDB/package.json b/Source/MongoDB/package.json index 010f7fa0..6e88d7aa 100644 --- a/Source/MongoDB/package.json +++ b/Source/MongoDB/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.mongodb", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { @@ -29,8 +29,9 @@ "README.md" ], "peerDependencies": { - "@cratis/arc.core": "^0.25.0", + "@cratis/arc.core": "^0.26.0", "@cratis/fundamentals": "^7.19.6", + "@opentelemetry/api": "^1.9.0", "mongodb": "^6.21.0", "rxjs": "^7.8.2" }, @@ -38,6 +39,7 @@ "@cratis/arc.core": "workspace:^", "@cratis/arc.testing": "workspace:^", "@cratis/fundamentals": "7.19.6", + "@opentelemetry/api": "^1.9.0", "mongodb": "^6.21.0", "rxjs": "^7.8.2" } diff --git a/Source/MongoDB/retryRead.ts b/Source/MongoDB/retryRead.ts new file mode 100644 index 00000000..1694e764 --- /dev/null +++ b/Source/MongoDB/retryRead.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 { MongoError, MongoNetworkError, MongoServerSelectionError } from 'mongodb'; + +/** Retry only idempotent reads with recognized transient driver errors; never turn a failed read into an empty result. */ +export async function retryRead(read: () => Promise, signal: AbortSignal): Promise { + for (let attempt = 0; ; attempt++) { + if (signal.aborted) throw signal.reason ?? new DOMException('Aborted', 'AbortError'); + try { return await read(); } + catch (error) { + if (attempt >= 2 || !(error instanceof MongoNetworkError || error instanceof MongoServerSelectionError || + error instanceof MongoError && error.hasErrorLabel('RetryableReadError')) || signal.aborted) throw error; + await new Promise((resolve, reject) => { + const timer = setTimeout(() => { signal.removeEventListener('abort', abort); resolve(); }, 100 * (attempt + 1)); + const abort = () => { clearTimeout(timer); reject(signal.reason ?? new DOMException('Aborted', 'AbortError')); }; + signal.addEventListener('abort', abort, { once: true }); + if (signal.aborted) abort(); + }); + } + } +} diff --git a/Source/MongoDB/withMongoDB.ts b/Source/MongoDB/withMongoDB.ts index 104ed317..73977ba9 100644 --- a/Source/MongoDB/withMongoDB.ts +++ b/Source/MongoDB/withMongoDB.ts @@ -1,10 +1,11 @@ // Copyright (c) Cratis. All rights reserved. // Licensed under the MIT license. See LICENSE file in the project root for full license information. import { ArcApplicationBuilder, serviceToken } from '@cratis/arc.core'; -import type { ExecutionContext } from '@cratis/arc.core'; +import type { ExecutionContext, ServiceScope } from '@cratis/arc.core'; import { ArcApplicationBuilder as FetchArcApplicationBuilder } from '@cratis/arc.core/fetch'; import { MongoClientFactory } from './MongoClientFactory.js'; import { MongoCollection } from './MongoCollection.js'; +import { MongoDBWatcher } from './MongoDBWatcher.js'; import { MongoReadModelForCommandResolver } from './MongoReadModelForCommandResolver.js'; import type { MongoDBOptions } from './MongoDBOptions.js'; import { mongoCollection } from './collectionToken.js'; @@ -12,6 +13,8 @@ import { defaultMongoNamingPolicy } from './MongoNamingPolicy.js'; /** Public token for applications needing to resolve clients explicitly. */ export const mongoClientFactory = serviceToken('MongoClientFactory'); +/** Scoped watcher: resolve it in the same tenant scope as the collections it observes. */ +export const mongoDBWatcher = serviceToken('MongoDBWatcher'); export function withMongoDB(builder: ArcApplicationBuilder, configured: MongoDBOptions): ArcApplicationBuilder { const settings = { ...builder.configuration.Cratis?.MongoDB }; @@ -22,17 +25,24 @@ export function withMongoDB(builder: ArcApplicationBuilder, configured: MongoDBO builder.services.addSingleton(mongoClientFactory, () => factory); builder.services.addScoped(MongoReadModelForCommandResolver, () => new MongoReadModelForCommandResolver(options)); builder.addReadModelForCommandResolver(MongoReadModelForCommandResolver); + const resolveDatabase = async (scope: ServiceScope) => { + const context: ExecutionContext | undefined = scope.identity; + if (!context?.tenantId) throw new Error('A tenant is required for MongoDB access'); + const tenantId = context.tenantId.toLowerCase(); + const name = options.databaseNameResolver ? options.databaseNameResolver(tenantId, context) : + tenantId === 'default' ? options.database : `${options.database}+${tenantId}`; + if (!name) throw new Error('MongoDB database resolver returned no database'); + const client = (await scope.resolve(mongoClientFactory)).get(context); + return { context, database: client.db(name) }; + }; + builder.services.addScoped(mongoDBWatcher, async scope => { + const { context, database } = await resolveDatabase(scope); + return new MongoDBWatcher(database, context); + }); for (const type of options.readModels) { const token = mongoCollection(type); builder.services.addScoped(token, async scope => { - const context: ExecutionContext | undefined = scope.identity; - if (!context?.tenantId) throw new Error('A tenant is required for MongoDB access'); - const tenantId = context.tenantId.toLowerCase(); - const name = options.databaseNameResolver ? options.databaseNameResolver(tenantId, context) : - tenantId === 'default' ? options.database : `${options.database}+${tenantId}`; - if (!name) throw new Error('MongoDB database resolver returned no database'); - const client = (await scope.resolve(mongoClientFactory)).get(context); - const database = client.db(name); + const { context, database } = await resolveDatabase(scope); const namingPolicy = options.namingPolicy ?? defaultMongoNamingPolicy; const collectionName = options.collectionName?.(type) ?? namingPolicy.collectionName(type); if (!collectionName) throw new Error('MongoDB collection name is required'); diff --git a/Source/Testing/package.json b/Source/Testing/package.json index 5d2282f5..2da8c596 100644 --- a/Source/Testing/package.json +++ b/Source/Testing/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.testing", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Tools/ProxyGenerator/package.json b/Source/Tools/ProxyGenerator/package.json index 563ca862..da24f5cb 100644 --- a/Source/Tools/ProxyGenerator/package.json +++ b/Source/Tools/ProxyGenerator/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.proxygenerator", - "version": "0.25.0", + "version": "0.26.0", "description": "TypeScript source analyzer and deterministic Arc client proxy generator", "repository": { "type": "git", diff --git a/yarn.lock b/yarn.lock index 392daf06..5d93cc52 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.25.0 - "@cratis/arc.testing": ^0.25.0 + "@cratis/arc.core": ^0.26.0 + "@cratis/arc.testing": ^0.26.0 "@cratis/chronicle": ^6.7.0 "@cratis/fundamentals": ^7.19.6 rxjs: ^7.8.2 @@ -144,7 +144,7 @@ __metadata: postgres: "npm:^3.4.9" sql.js: "npm:^1.14.2" peerDependencies: - "@cratis/arc.core": ^0.25.0 + "@cratis/arc.core": ^0.26.0 "@cratis/fundamentals": ^7.19.6 drizzle-orm: ^0.45.0 languageName: unknown @@ -205,11 +205,13 @@ __metadata: "@cratis/arc.core": "workspace:^" "@cratis/arc.testing": "workspace:^" "@cratis/fundamentals": "npm:7.19.6" + "@opentelemetry/api": "npm:^1.9.0" mongodb: "npm:^6.21.0" rxjs: "npm:^7.8.2" peerDependencies: - "@cratis/arc.core": ^0.25.0 + "@cratis/arc.core": ^0.26.0 "@cratis/fundamentals": ^7.19.6 + "@opentelemetry/api": ^1.9.0 mongodb: ^6.21.0 rxjs: ^7.8.2 languageName: unknown