From eb77dc447bb304b5c7465f2dd9983903ef643b45 Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:39:08 +0200 Subject: [PATCH 1/9] Add GeoJSON mapping for Fundamentals geometry --- Source/MongoDB/MongoDocumentCodec.ts | 7 +++ Source/MongoDB/MongoGeoJSON.ts | 57 +++++++++++++++++++ .../given/a_geospatial_model.ts | 17 ++++++ .../when_rejecting_invalid_geometry.ts | 18 ++++++ .../when_round_tripping_geospatial_fields.ts | 32 +++++++++++ 5 files changed, 131 insertions(+) create mode 100644 Source/MongoDB/MongoGeoJSON.ts create mode 100644 Source/MongoDB/for_MongoDocumentCodec/given/a_geospatial_model.ts create mode 100644 Source/MongoDB/for_MongoDocumentCodec/when_rejecting_invalid_geometry.ts create mode 100644 Source/MongoDB/for_MongoDocumentCodec/when_round_tripping_geospatial_fields.ts 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_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); + }); +})); From b6a3b0e6eeef1f837167f8d250c0294ad23b2945 Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:39:13 +0200 Subject: [PATCH 2/9] Retry transient MongoDB reads within a bounded window --- Source/MongoDB/MongoCollection.ts | 30 ++++++++++++------- .../when_exhausting_read_retries.ts | 24 +++++++++++++++ .../when_retrying_a_transient_read.ts | 26 ++++++++++++++++ Source/MongoDB/retryRead.ts | 21 +++++++++++++ 4 files changed, 91 insertions(+), 10 deletions(-) create mode 100644 Source/MongoDB/for_MongoReadModels/when_exhausting_read_retries.ts create mode 100644 Source/MongoDB/for_MongoReadModels/when_retrying_a_transient_read.ts create mode 100644 Source/MongoDB/retryRead.ts 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/for_MongoReadModels/when_exhausting_read_retries.ts b/Source/MongoDB/for_MongoReadModels/when_exhausting_read_retries.ts new file mode 100644 index 00000000..ed21febb --- /dev/null +++ b/Source/MongoDB/for_MongoReadModels/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_MongoReadModels/when_retrying_a_transient_read.ts b/Source/MongoDB/for_MongoReadModels/when_retrying_a_transient_read.ts new file mode 100644 index 00000000..58799ac5 --- /dev/null +++ b/Source/MongoDB/for_MongoReadModels/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/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(); + }); + } + } +} From 1d1b63e7db0120696474e0d5dcc9406ca31910b8 Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:39:19 +0200 Subject: [PATCH 3/9] Instrument Arc-owned MongoDB clients with OpenTelemetry --- Source/MongoDB/MongoClientFactory.ts | 7 +++- Source/MongoDB/MongoDBClientMetrics.ts | 23 +++++++++++ .../when_instrumenting_an_owned_client.ts | 40 +++++++++++++++++++ Source/MongoDB/package.json | 2 + 4 files changed, 71 insertions(+), 1 deletion(-) create mode 100644 Source/MongoDB/MongoDBClientMetrics.ts create mode 100644 Source/MongoDB/for_MongoClientFactory/when_instrumenting_an_owned_client.ts diff --git a/Source/MongoDB/MongoClientFactory.ts b/Source/MongoDB/MongoClientFactory.ts index 3ee573ec..9e8e2f62 100644 --- a/Source/MongoDB/MongoClientFactory.ts +++ b/Source/MongoDB/MongoClientFactory.ts @@ -3,6 +3,7 @@ import type { ExecutionContext } from '@cratis/arc.core'; import { MongoClient } from 'mongodb'; import type { MongoDBOptions } from './MongoDBOptions.js'; +import { monitorMongoDBClient } from './MongoDBClientMetrics.js'; /** Owns URI-created clients only; never closes a caller-supplied client. */ export class MongoClientFactory { @@ -18,7 +19,11 @@ export class MongoClientFactory { const uri = this.options.serverResolver?.(context.tenantId.toLowerCase(), context) ?? this.options.server; if (!uri) throw new Error('MongoDB server resolver returned no server'); let client = this.#clients.get(uri); - if (!client) { client = new MongoClient(uri); this.#clients.set(uri, client); } + if (!client) { + client = new MongoClient(uri, { monitorCommands: true }); + monitorMongoDBClient(client, client.options.hosts.map(host => host.toString()).join(',')); + this.#clients.set(uri, client); + } return client; } async [Symbol.asyncDispose](): Promise { 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/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/package.json b/Source/MongoDB/package.json index 8a3865c2..75295afc 100644 --- a/Source/MongoDB/package.json +++ b/Source/MongoDB/package.json @@ -31,6 +31,7 @@ "peerDependencies": { "@cratis/arc.core": "^0.24.1", "@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" } From fc91a74a8f89a2019c1d66c786ed05633dc3fd3d Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:39:25 +0200 Subject: [PATCH 4/9] Add tenant-scoped shared MongoDB watcher and joined observations --- Source/MongoDB/MongoDBWatcher.ts | 185 ++++++++++++++++++ .../given/RelatedRecord.ts | 9 + .../with_a_replica_set.integration.ts | 58 ++++++ .../with_a_replica_set.integration.ts | 42 ++++ .../with_nonresumable_error.ts | 44 +++++ Source/MongoDB/index.ts | 5 +- Source/MongoDB/withMongoDB.ts | 28 ++- 7 files changed, 361 insertions(+), 10 deletions(-) create mode 100644 Source/MongoDB/MongoDBWatcher.ts create mode 100644 Source/MongoDB/for_MongoCollection/given/RelatedRecord.ts create mode 100644 Source/MongoDB/for_MongoCollection/when_observing_joined_collections/with_a_replica_set.integration.ts create mode 100644 Source/MongoDB/for_MongoDBWatcher/when_reacting_to_changes/with_a_replica_set.integration.ts create mode 100644 Source/MongoDB/for_MongoDBWatcher/when_stream_fails/with_nonresumable_error.ts 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/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/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/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'); From 655e2675e0de09264c39aec62543438da23c4fa3 Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:39:30 +0200 Subject: [PATCH 5/9] Document MongoDB watcher, joined results and GeoJSON support --- .../mongodb/change-stream-watcher.md | 36 +++++++++++++ Documentation/mongodb/geospatial.md | 53 +++++++++++++++++++ Documentation/mongodb/index.md | 5 +- Documentation/mongodb/joined-observe.md | 27 ++++++++++ .../mongodb/observing-collections.md | 3 +- Documentation/mongodb/toc.yml | 6 +++ Documentation/reference/capabilities.md | 4 +- 7 files changed, 130 insertions(+), 4 deletions(-) create mode 100644 Documentation/mongodb/change-stream-watcher.md create mode 100644 Documentation/mongodb/geospatial.md create mode 100644 Documentation/mongodb/joined-observe.md 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..b786b2df 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_ Date: Fri, 25 Sep 2026 08:39:51 +0200 Subject: [PATCH 6/9] Keep capability detail confined to the MongoDB row --- Documentation/reference/capabilities.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Documentation/reference/capabilities.md b/Documentation/reference/capabilities.md index b786b2df..167600da 100644 --- a/Documentation/reference/capabilities.md +++ b/Documentation/reference/capabilities.md @@ -165,7 +165,7 @@ The suite needs Docker. Its image, `cratis/chronicle:latest-development`, is mut ### MongoDB checks -`bash Source/MongoDB/run-integration.sh` starts a task-owned MongoDB 7 replica set in Docker and removes it afterward. The [integration spec](https://github.com/Cratis/Arc.TypeScript/blob/main/Source/MongoDB/for_MongoCollection/when_observing_changes/with_a_replica_set.integration.ts) exercises initial snapshots, insertion, deletion, tenant isolation, dependency injection, and provider paging, and a [second spec](https://github.com/Cratis/Arc.TypeScript/blob/main/Source/MongoDB/for_MongoCollection/when_serving_a_paged_query/with_each_http_adapter.integration.ts) serves sorted pages through Express, Fastify, and Hono. The script exits with 2 when Docker is not available, which means the check did not run. The live suite also checks joined snapshots and tenant isolation. Unit specs cover GeoJSON codec round trips and malformed geometry, bounded read retries, naming policies, burst coalescing, the item cap, and stream failures. The `for_MongoClientFactory` unit spec checks command-count export and verifies that the server label excludes URI credentials. Driver metrics are enabled for Arc-owned clients only; live pool event export is not checked. +`bash Source/MongoDB/run-integration.sh` starts a task-owned MongoDB 7 replica set in Docker and removes it afterward. The [integration spec](https://github.com/Cratis/Arc.TypeScript/blob/main/Source/MongoDB/for_MongoCollection/when_observing_changes/with_a_replica_set.integration.ts) exercises initial snapshots, insertion, deletion, tenant isolation, dependency injection, and provider paging, and a [second spec](https://github.com/Cratis/Arc.TypeScript/blob/main/Source/MongoDB/for_MongoCollection/when_serving_a_paged_query/with_each_http_adapter.integration.ts) serves sorted pages through Express, Fastify, and Hono. The script exits with 2 when Docker is not available, which means the check did not run. Unit specs cover the codec against .NET-shaped documents, naming policies, burst coalescing, the item cap, and stream failures. ### SQL checks From a2fe00105333c6408023bd36324164dc94be3699 Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:40:04 +0200 Subject: [PATCH 7/9] Place read retry specifications beside their subject --- .../when_exhausting_read_retries.ts | 0 .../when_retrying_a_transient_read.ts | 0 2 files changed, 0 insertions(+), 0 deletions(-) rename Source/MongoDB/{for_MongoReadModels => for_retryRead}/when_exhausting_read_retries.ts (100%) rename Source/MongoDB/{for_MongoReadModels => for_retryRead}/when_retrying_a_transient_read.ts (100%) diff --git a/Source/MongoDB/for_MongoReadModels/when_exhausting_read_retries.ts b/Source/MongoDB/for_retryRead/when_exhausting_read_retries.ts similarity index 100% rename from Source/MongoDB/for_MongoReadModels/when_exhausting_read_retries.ts rename to Source/MongoDB/for_retryRead/when_exhausting_read_retries.ts diff --git a/Source/MongoDB/for_MongoReadModels/when_retrying_a_transient_read.ts b/Source/MongoDB/for_retryRead/when_retrying_a_transient_read.ts similarity index 100% rename from Source/MongoDB/for_MongoReadModels/when_retrying_a_transient_read.ts rename to Source/MongoDB/for_retryRead/when_retrying_a_transient_read.ts From c8ece7a6cf1292f757cb43db1d81c6f8e4b33535 Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:52:10 +0200 Subject: [PATCH 8/9] Lock the MongoDB OpenTelemetry dependency --- yarn.lock | 2 ++ 1 file changed, 2 insertions(+) diff --git a/yarn.lock b/yarn.lock index 392daf06..24007382 100644 --- a/yarn.lock +++ b/yarn.lock @@ -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/fundamentals": ^7.19.6 + "@opentelemetry/api": ^1.9.0 mongodb: ^6.21.0 rxjs: ^7.8.2 languageName: unknown From 6c48c2286f816cfab436848bde09dee974095b7c Mon Sep 17 00:00:00 2001 From: woksin Date: Fri, 25 Sep 2026 08:52:22 +0200 Subject: [PATCH 9/9] Prepare the v0.26 source preview --- ContractTests/Client/package.json | 2 +- Documentation/index.md | 2 +- Documentation/reference/packages.md | 2 +- README.md | 2 +- Source/Chronicle/package.json | 6 +++--- Source/CodeAnalysis/package.json | 2 +- Source/Core/package.json | 2 +- Source/Cratis/package.json | 2 +- Source/Drizzle/package.json | 4 ++-- Source/Express/package.json | 2 +- Source/Fastify/package.json | 2 +- Source/Hono/package.json | 2 +- Source/MongoDB/package.json | 4 ++-- Source/Testing/package.json | 2 +- Source/Tools/ProxyGenerator/package.json | 2 +- yarn.lock | 8 ++++---- 16 files changed, 23 insertions(+), 23 deletions(-) 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/reference/packages.md b/Documentation/reference/packages.md index 8c7f0289..7856b6a0 100644 --- a/Documentation/reference/packages.md +++ b/Documentation/reference/packages.md @@ -3,7 +3,7 @@ title: Packages description: The packages this repository builds, what each exports, their peer dependencies and Node.js requirements, and how they relate to the published @cratis/arc client. --- -Every package in this repository is at version 0.25.0, the version of the source preview. **None is published to npm**; reference them from a clone with the `workspace:^` protocol. They ship ES modules only. +Every package in this repository is at version 0.26.0, the version of the source preview. **None is published to npm**; reference them from a clone with the `workspace:^` protocol. They ship ES modules only. ## Server packages diff --git a/README.md b/README.md index e376bcc9..ea1e1252 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.7.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.25.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 ran in Deno 2.9.7, while Bun, Cloudflare Workers, and Next.js deployments remain unverified. 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.26.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 ran in Deno 2.9.7, while Bun, Cloudflare Workers, and Next.js deployments remain unverified. 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 622ee750..d0e69729 100644 --- a/Source/Chronicle/package.json +++ b/Source/Chronicle/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.chronicle", - "version": "0.25.0", + "version": "0.26.0", "publishConfig": { "access": "public" }, @@ -34,8 +34,8 @@ "README.md" ], "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", diff --git a/Source/CodeAnalysis/package.json b/Source/CodeAnalysis/package.json index b2d42bc2..c38ade5c 100644 --- a/Source/CodeAnalysis/package.json +++ b/Source/CodeAnalysis/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/eslint-plugin-arc-core", - "version": "0.25.0", + "version": "0.26.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 b1713787..b8e037df 100644 --- a/Source/Core/package.json +++ b/Source/Core/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.core", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Cratis/package.json b/Source/Cratis/package.json index ba3a7b23..3208ae20 100644 --- a/Source/Cratis/package.json +++ b/Source/Cratis/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/cratis", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "description": "Arc and experimental Chronicle composition for Node.js", diff --git a/Source/Drizzle/package.json b/Source/Drizzle/package.json index 46ed7f90..9e0be0c8 100644 --- a/Source/Drizzle/package.json +++ b/Source/Drizzle/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.drizzle", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { @@ -29,7 +29,7 @@ "README.md" ], "peerDependencies": { - "@cratis/arc.core": "^0.25.0", + "@cratis/arc.core": "^0.26.0", "@cratis/fundamentals": "^7.19.6", "drizzle-orm": "^0.45.0" }, diff --git a/Source/Express/package.json b/Source/Express/package.json index 744b9cd4..e564f044 100644 --- a/Source/Express/package.json +++ b/Source/Express/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.express", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Fastify/package.json b/Source/Fastify/package.json index 78fbb5a9..e1bf951c 100644 --- a/Source/Fastify/package.json +++ b/Source/Fastify/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.fastify", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/Hono/package.json b/Source/Hono/package.json index b51cd9d0..5f22a07e 100644 --- a/Source/Hono/package.json +++ b/Source/Hono/package.json @@ -1,6 +1,6 @@ { "name": "@cratis/arc.hono", - "version": "0.25.0", + "version": "0.26.0", "type": "module", "license": "MIT", "publishConfig": { diff --git a/Source/MongoDB/package.json b/Source/MongoDB/package.json index 39c4573c..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,7 +29,7 @@ "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", 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 24007382..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 @@ -209,7 +209,7 @@ __metadata: 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