QueryBackend

Aggregate-bound storage SPI: four primitives that only check natively, translate and execute.

The core derives every query shape from them (see single, list, paged, cursor, aggregate in BackendQueries.kt): single is page(Offset(0, 1, withTotal = false)), list is stream, paged is page(Offset(..., withTotal = true)), cursor is page(Keyset). Windows come only from the core; the core also owns the cursor token shell, the residual aggregation operators the storage declares RESIDUAL and the empty summary row.

Every input is an AdmittedQuery: validated against its schema, normalized, with each field reference resolved. The gateway admits through a QueryAdmission instance, whose step 3 appends model defaults such as Snapshot ACTIVE. Direct callers admit through QueryAdmission.Trusted, which skips steps 0 to 3, so they supply any required deletion or access predicates themselves.

Every subscription to a returned publisher, including subscriptions created by retry, repeat, or concurrent callers, must own fresh mutable ObjectNode instances. Implementations must not cache or share nodes across subscriptions, publish cached nodes, mutate emitted nodes asynchronously, or mutate them after delivery.

Results must contain only standard JSON-tree values: the Backend converts its driver's values (Map/Document, BSON values, decimals) to them. The core checks every returned row, so no Backend repeats the check: a row with a NaN, an infinity, a POJONode, a binary or missing node fails the query in the core.

Inheritors

Properties

Link copied to clipboard

Encodes and decodes the native CursorPositions page returns for a PageWindow.Keyset.

Functions

Link copied to clipboard
abstract fun aggregate(query: AdmittedQuery<AggregationQuery>, window: GroupWindow): Flux<ObjectNode>

Streams the groups of query in its effective sort order, at most GroupWindow.First.limit of them, or every group for GroupWindow.All. The core has already removed from query what it computes itself.

Link copied to clipboard
fun QueryBackend.aggregate(query: AdmittedQuery<AggregationQuery>, budget: QueryBudget? = null): Flux<ObjectNode>

The groups of query. The core plans the residual operators the storage declares RESIDUAL: it removes them from the query it sends down, asks for every group when an operator needs them all, and applies dense fill, HAVING and top-N (or the limit) to the rows that come back. An aggregation without groups always has its summary row: when the backend emits none, because no record matched, the core emits the empty summary (EmptyAggregationValues). When it reads every group, budget's QueryBudget.maxResidualGroups bounds how many it processes, dense fill rows included: HAVING can discard every fill row, so without counting them a sparse fine-grained dense histogram would generate rows until the idle timeout.

Link copied to clipboard
abstract fun count(query: AdmittedQuery<FilterExpression>): Mono<Long>
Link copied to clipboard
fun QueryBackend.cursor(query: AdmittedQuery<ICursorQuery>): Mono<CursorPage<ObjectNode>>

One cursor page of query: page(Keyset) for one row more than the page, the look-ahead that decides whether a next page exists. The next token encodes the native position of the page's last row, never a value from the rows, so masking cannot leak into it. A token that does not decode for this model and effective sort is rejected as Invalid cursor. before any I/O.

Link copied to clipboard
fun QueryBackend.list(query: AdmittedQuery<IListQuery>): Flux<ObjectNode>

The records of query, streamed: at most its limit, or all of them when the limit is 0.

Link copied to clipboard
abstract fun page(query: AdmittedQuery<Queryable<*>>, window: PageWindow): Mono<BackendPage>

Returns one window of query's records: for PageWindow.Offset the rows and, when asked, the total; for PageWindow.Keyset the rows after the position and each row's native position, taken before any field added for the position is stripped from the row.

Link copied to clipboard
fun QueryBackend.paged(query: AdmittedQuery<IPagedQuery>): Mono<PagedList<ObjectNode>>

One page of query with the total: page(Offset(offset, size, withTotal = true)).

Link copied to clipboard
fun QueryBackend.single(query: AdmittedQuery<ISingleQuery>): Mono<ObjectNode>

The first record of query: page(Offset(0, 1, withTotal = false)).

Link copied to clipboard
abstract fun stream(query: AdmittedQuery<IListQuery>): Flux<ObjectNode>

Streams the records of query, at most query.limit of them, or all when the limit is 0.