Observable storage snapshots
Return observable storage snapshots
Section titled “Return observable storage snapshots”Both Spring Data integrations publish injectable query services whose methods return Kotlin Flow. A Flow is cold by default: it reads the initial snapshot when collected, owns its subscription, and closes it on cancellation. observeShared is the explicit shared alternative and replays one complete snapshot while subscribers exist. Java callers can use the demand-aware observePublisher methods or callback overloads, whose returned AutoCloseable cancels collection.
Observe MongoDB
Section titled “Observe MongoDB”Inject MongoObservableQuery into a model-bound query. observe, observeSingle, and observeById re-run the Spring Data query after insert, update, replace, delete, or invalidate notifications:
@JvmStaticfun observeTasks( @FromServices queries: MongoObservableQuery): Flow<List<TaskView>> = queries.observe(TaskView::class.java)
@JvmStaticfun observeTask( id: String, @FromServices queries: MongoObservableQuery): Flow<TaskView> = queries.observeById(TaskView::class.java, id)MongoChangeStreamWatcher is the change-stream SPI. The default ReconnectingMongoChangeStreamWatcher uses SpringDataMongoChangeStreamSource, resumes from the last token after transient failures, applies capped exponential backoff, and closes its cursor when the collector is canceled. Cursor operations run on Dispatchers.IO without a hidden flowOn channel. MongoObservationOptions.bufferCapacity is exact: zero keeps the producer and collector in rendezvous, while a positive value permits only that many queued changes. MongoDB change streams require a replica set or sharded cluster.
A custom MongoOperationsResolver can select operations from the captured tenantId; it must not consult thread-local request state. Pass the tenant explicitly to the observation method. The default resolver uses the application’s single MongoOperations bean. A Query supplies filtering, and observeById additionally narrows change notifications by document key. Updates and replacements always trigger a fresh snapshot, so a document that stops matching a filter is removed correctly.
Kotlin type-inference extensions for MongoDB
Section titled “Kotlin type-inference extensions for MongoDB”The arc-spring-data-mongodb artifact ships Kotlin extension functions on MongoObservableQuery that infer the document type from the generic parameter, eliminating the explicit ::class.java argument:
@JvmStaticfun observeTasks( @FromServices queries: MongoObservableQuery): Flow<List<TaskView>> = queries.observe<TaskView>()
@JvmStaticfun observeTask( id: String, @FromServices queries: MongoObservableQuery): Flow<TaskView> = queries.observeById<TaskView>(id)Criteria-based filtering is also supported without manually wrapping in a Query:
fun observeActiveTasks( @FromServices queries: MongoObservableQuery): Flow<List<TaskView>> = queries.observe<TaskView>(Criteria.where("active").`is`(true))The same extensions cover observeList, observeSingle, observeById, and observeShared; each delegates to the corresponding MongoObservableQuery method.
Limitation: extensions cannot be placed on Spring Data’s MongoCollection<T>, repository interfaces, or any other foreign type. All observe* extensions are on Arc’s own MongoObservableQuery.
Java convenience access for MongoDB
Section titled “Java convenience access for MongoDB”Kotlin extensions are invisible from Java. The MongoObservations object provides @JvmStatic Criteria-based helpers that Java callers can use without manually constructing a Query:
Flow.Publisher<List<TaskView>> publisher = MongoObservations.observe(queries, TaskView.class, Criteria.where("active").is(true));
Flow.Publisher<TaskView> single = MongoObservations.observeSingle(queries, TaskView.class, Criteria.where("id").is(taskId));Each MongoObservations helper wraps the Criteria into a Query and delegates to the matching MongoObservableQuery method. For Query-based and callback access, call the @JvmOverloads methods on MongoObservableQuery directly.
The tenant overload is also available:
MongoObservations.observe(queries, TaskView.class, Criteria.where("active").is(true), tenantId);Observe JPA
Section titled “Observe JPA”Inject JpaObservableQuery and return observe, observeList, observeSingle, or observeById from the read-model query. The default list query uses the mapped JPA entity name; JpaSnapshotQuery provides a Java-friendly customization seam for predicates, ordering, and fetch joins.
@JvmStaticfun observeTasks( @FromServices queries: JpaObservableQuery): Flow<List<TaskView>> = queries.observe(TaskView::class.java)Limitation: JPA observation is in-process only. It is backed by TransactionAwareDatabaseChangeNotifier, an explicit in-process publisher: call its DatabaseChangePublisher.publish(TaskView::class.java, tenantId) contract from the write side. Unlike MongoDB change streams there is no cross-process change mechanism; applications needing cross-process notifications must replace the DatabaseChangeNotifier bean with a database-native implementation and expose a corresponding DatabaseChangePublisher where local writes also need publishing. Arc does not silently poll either store; polling must be an application-owned, explicitly configured notifier.
Notifications are coalesced per transaction, discarded on rollback, and emitted only after commit. JpaObservationOptions.bufferCapacity must be positive and bounds pending invalidations; when it is full, the oldest invalidation is replaced because every notification causes a complete snapshot read. Snapshot delivery itself uses a rendezvous handoff rather than flowOn’s implicit buffer. The Flow performs its initial snapshot immediately, then debounce-coalesces committed changes into bounded replacement snapshots.
Kotlin type-inference extensions for JPA
Section titled “Kotlin type-inference extensions for JPA”The arc-spring-data-jpa artifact ships Kotlin extension functions on JpaObservableQuery that infer the entity type from the generic parameter:
@JvmStaticfun observeTasks( @FromServices queries: JpaObservableQuery): Flow<List<TaskView>> = queries.observe<TaskView>()
@JvmStaticfun observeTask( id: String, @FromServices queries: JpaObservableQuery): Flow<TaskView> = queries.observeById<TaskView>(id)A custom snapshot query can be supplied as a trailing lambda:
fun observeActiveTasks( @FromServices queries: JpaObservableQuery): Flow<List<TaskView>> = queries.observe<TaskView>( query = JpaSnapshotQuery { em -> em.createQuery( "select t from TaskView t where t.active = true", TaskView::class.java ).resultList } )The same extensions cover observeList, observeSingle, observeById, and observeShared.
Limitation: extensions cannot be placed on Spring Data’s JpaRepository or any other foreign interface. All observe* extensions are on Arc’s own JpaObservableQuery.
Java convenience access for JPA
Section titled “Java convenience access for JPA”Kotlin extensions are invisible from Java. The JpaObservations object provides @JvmStatic helpers for the two entry points that have no default query in the base API:
JpaObservations.observePublisher(queries, entityType)— returns a demand-awareFlow.Publisher<List<T>>using the default select-all JPQL query.JpaObservations.observeShared(scope, queries, entityType)— returns aSharedFlow<List<T>>that replays the latest snapshot while subscribers exist.
Flow.Publisher<List<TaskView>> publisher = JpaObservations.observePublisher(queries, TaskView.class);
SharedFlow<List<TaskView>> shared = JpaObservations.observeShared(scope, queries, TaskView.class);Both helpers supply the same default JPQL query (select entity from <EntityName> entity) as the Kotlin extensions. For a custom snapshot query or callback access, call the @JvmOverloads methods on JpaObservableQuery directly.
Current public-seam boundary
Section titled “Current public-seam boundary”The integrations deliberately use public Arc and Spring contracts. Contextual command read-model parameters, deterministic ownership, generated QueryRequest/QueryContext injection, exact Spring Data Commons Pageable/Sort parameters, and exact Page<T> response normalization are implemented. The store-specific request and page adapters remain compatibility utilities rather than request-scoped beans.
Repository injection remains ordinary Spring dependency injection and is not automatically tenant-routed by Arc. JPA observable snapshots still use their configured entity-manager factory; tenant labels on notifications are not storage isolation. Mongo observable paths use MongoOperationsResolver, including the certified adapter when a tenant-aware resolver is supplied. Imperative transaction scopes remain explicit fixed-store opt-ins and are not coroutine-safe.