Observe joined MongoDB collections
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:
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.
select returns an RxJS Observable<TResult>, rather than .NET’s ISubject<TResult>. 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.
If a joined collection holds a Chronicle read model with personal or encrypted data that is registered with withChronicle, what happens depends on the selected shape. Decoded instances placed unmapped inside another object, as in the example, fail the query, because Arc does not release instances nested in the selected object. A selector that returns an array of the instances themselves has each one released. Instances the selector maps to DTOs or copies are served as stored, as ciphertext; release those yourself with the tenant store’s readModels.release. See What Arc does not release.
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 for recovery limits and single-collection observation when joins are unnecessary.