---------------------------- MODULE SubscribeHook ----------------------------
(* A query collection with a `subscribe` option. Two subset Queries (Q) hold *)
(* one shared public row. The provider changes that row and pushes an event *)
(* to every open subscription. Fetches read the provider's current value. *)
EXTENDS Naturals, Integers, Sequences, FiniteSets, TLC
CONSTANTS Q, MaxVer, MaxFetch, MaxGen,
\* Design choices
UseReady, \* subscribe returns `ready`; the first fetch waits for it
CacheBranch, \* event for a cached Query: "update" its data or mark it "stale"
FailBranch, \* subscribe failure: "reject" the load or "degrade" to no live updates
Versioned, \* events and fetch results carry comparable versions
\* Design rules (TRUE / "owners" in the proposed design)
EventStamp, \* which Queries an event invalidates older fetches for:
\* "owners" (every Query holding the row) | "active" | "own" | "none"
RowGuard, \* a fetch result never replaces one from a fetch that started later
SharedRefetch, \* a direct write to a row another in-use subset holds refetches it
\* Environment
Ordered, \* the app's socket delivers events in order
\* Fault switches (FALSE in the proposed design)
AcceptDeadGen, \* writes from an evicted subscription still apply
IgnoreCached, \* events for a cached Query are dropped
DoubleCleanup \* the app's cleanup can run twice
None == -1
Gens == 1..MaxGen
Max(a, b) == IF a > b THEN a ELSE b
VARIABLES
srv, \* provider's current version of the row
phase, \* per Query: "absent" | "active" (in use) | "cached" (released, before gcTime)
cache, \* per Query: its cached data (a version), or None
stale, \* per Query: an event marked its cached data stale
cev, \* per Query: its cached data came from an event, not a fetch
succ, \* per Query: start number of its latest successful fetch
req, \* per Query: required fetch start (fetches starting at or before it do not count)
want, \* per Query: a refetch was requested
ready, \* per Query: subset readiness
fc, \* next fetch start number
fetches, \* in-flight fetches: [q, start, rd] (rd = version read, or None)
gen, \* per Query: its current subscription, or 0
sub, \* per subscription: [q, st] with st "subscribing" | "registered" |
\* "failed" | "closing" (unsubscribe sent, not yet processed) | "gone"
ng, \* next subscription id
evs, \* events queued on the app's socket: [q, g, v]
row, \* the published row (a version), or None when no subset holds it
rowStart, \* start number of whatever last published the row
\* Observation-only variables for the properties
floor, evMark, ovw, bad, cleaned
vars == <<srv, phase, cache, stale, cev, succ, req, want, ready, fc, fetches,
gen, sub, ng, evs, row, rowStart, floor, evMark, ovw, bad, cleaned>>
Init ==
/\ srv = 0
/\ phase = [q \in Q |-> "absent"]
/\ cache = [q \in Q |-> None]
/\ stale = [q \in Q |-> FALSE]
/\ succ = [q \in Q |-> 0]
/\ req = [q \in Q |-> 0]
/\ fc = 1
/\ fetches = {}
/\ gen = [q \in Q |-> 0]
/\ sub = [g \in Gens |-> [q |-> CHOOSE q \in Q : TRUE, st |-> "none"]]
/\ evs = <<>>
/\ ready = [q \in Q |-> FALSE]
/\ row = None
/\ floor = None
/\ bad = 0
/\ cleaned = [g \in Gens |-> 0]
/\ ng = 1
/\ evMark = 0
/\ ovw = 0
/\ cev = [q \in Q |-> FALSE]
/\ rowStart = 0
/\ want = [q \in Q |-> FALSE]
Active == {q \in Q : phase[q] = "active"}
Auth(q) == succ[q] > req[q] \* the Query's data is authoritative
SubSt(q) == IF gen[q] = 0 THEN "none" ELSE sub[gen[q]].st
\* Publish the shared row. `floor` records the highest version published
\* while the row stayed visible; NoRegression compares against it.
Publish(v, st) == /\ row' = v
/\ floor' = IF floor = None THEN v ELSE Max(floor, v)
/\ rowStart' = Max(rowStart, st)
\* The Queries whose older in-flight fetches an event invalidates.
Stamped(q) ==
CASE EventStamp = "owners" -> {p \in Q : phase[p] # "absent"}
[] EventStamp = "active" -> Active
[] EventStamp = "own" -> {q}
[] OTHER -> {}
\* A fetch result, or cached data from one, that started at `start` is about
\* to be published. `ovw` counts results that started before a published
\* event write. RowGuard and Versioned can keep the current row.
PublishResult(start, v) ==
/\ ovw' = IF start > evMark THEN ovw ELSE ovw + 1
/\ IF (Versioned /\ row # None /\ v < row) \/ (RowGuard /\ start < rowStart)
THEN UNCHANGED <<row, floor, rowStart>>
ELSE Publish(v, start)
Stamp(S) == req' = [p \in Q |-> IF p \in S THEN fc - 1 ELSE req[p]]
------------------------------------------------------------------------------
(* A live query needs the subset. A new Query subscribes; a cached one
settles at once if its data is current and authoritative. *)
Acquire(q) ==
/\ phase[q] \in {"absent", "cached"}
/\ phase' = [phase EXCEPT ![q] = "active"]
/\ IF phase[q] = "absent"
THEN /\ ng <= MaxGen
/\ gen' = [gen EXCEPT ![q] = ng]
/\ sub' = [sub EXCEPT ![ng] = [q |-> q, st |-> "subscribing"]]
/\ ng' = ng + 1
/\ UNCHANGED <<ready, row, floor, rowStart>>
ELSE /\ IF cache[q] # None /\ ~stale[q] /\ Auth(q)
THEN /\ ready' = [ready EXCEPT ![q] = TRUE]
/\ IF cev[q] THEN Publish(cache[q], succ[q]) /\ UNCHANGED ovw
ELSE PublishResult(succ[q], cache[q])
ELSE UNCHANGED <<ready, row, floor, rowStart, ovw>>
/\ UNCHANGED <<gen, sub, ng>>
/\ IF phase[q] = "absent" THEN UNCHANGED ovw ELSE TRUE
/\ UNCHANGED <<want, srv, cache, stale, succ, req, fc, fetches, evs, bad, cleaned,
evMark, cev>>
(* No live query needs the subset. Its Query stays cached; so does its
subscription. The row disappears when no subset holds it. *)
Release(q) ==
/\ phase[q] = "active"
/\ phase' = [phase EXCEPT ![q] = "cached"]
/\ ready' = [ready EXCEPT ![q] = FALSE]
/\ IF Active = {q} THEN row' = None /\ floor' = None /\ rowStart' = 0
ELSE UNCHANGED <<row, floor, rowStart>>
/\ UNCHANGED <<want, evMark, ovw, cev, srv, cache, stale, succ, req, fc, fetches, gen, sub, evs,
bad, cleaned, ng>>
Retire(g) == IF sub[g].st \in {"subscribing", "registered"}
THEN "closing" ELSE sub[g].st
CleanupCount(g) == IF cleaned[g] = 0 \/ DoubleCleanup
THEN cleaned[g] + 1 ELSE cleaned[g]
(* TanStack Query evicts the Query at gcTime. Cleanup runs once. The app's
unsubscribe may not reach the server for a while ("closing"). *)
Remove(q) ==
/\ phase[q] = "cached"
/\ LET g == gen[q] IN
/\ sub' = [sub EXCEPT ![g].st = Retire(g)]
/\ cleaned' = [cleaned EXCEPT ![g] = CleanupCount(g)]
/\ phase' = [phase EXCEPT ![q] = "absent"]
/\ gen' = [gen EXCEPT ![q] = 0]
/\ cache' = [cache EXCEPT ![q] = None]
/\ stale' = [stale EXCEPT ![q] = FALSE]
/\ succ' = [succ EXCEPT ![q] = 0]
/\ req' = [req EXCEPT ![q] = 0]
/\ fetches' = {f \in fetches : f.q # q} \* they belong to the evicted Query
/\ cev' = [cev EXCEPT ![q] = FALSE]
/\ want' = [want EXCEPT ![q] = FALSE]
/\ UNCHANGED <<evMark, ovw, srv, fc, evs, ready, row, floor, rowStart, bad, ng>>
(* Collection cleanup ends every subscription. *)
CollectionCleanup ==
/\ \E q \in Q : phase[q] # "absent"
/\ LET live == {gen[q] : q \in {p \in Q : phase[p] # "absent"}} IN
/\ sub' = [g \in Gens |-> IF g \in live
THEN [sub[g] EXCEPT !.st = Retire(g)]
ELSE sub[g]]
/\ cleaned' = [g \in Gens |-> IF g \in live THEN CleanupCount(g)
ELSE cleaned[g]]
/\ phase' = [q \in Q |-> "absent"]
/\ gen' = [q \in Q |-> 0]
/\ cache' = [q \in Q |-> None]
/\ stale' = [q \in Q |-> FALSE]
/\ succ' = [q \in Q |-> 0]
/\ req' = [q \in Q |-> 0]
/\ ready' = [q \in Q |-> FALSE]
/\ fetches' = {}
/\ row' = None /\ floor' = None /\ rowStart' = 0
/\ cev' = [q \in Q |-> FALSE]
/\ want' = [q \in Q |-> FALSE]
/\ UNCHANGED <<evMark, ovw, srv, fc, evs, bad, ng>>
------------------------------------------------------------------------------
(* The server starts delivering (`ready` resolves), or subscribe fails. *)
Register(g) ==
/\ sub[g].st = "subscribing"
/\ sub' = [sub EXCEPT ![g].st = "registered"]
/\ UNCHANGED <<want, evMark, ovw, cev, srv, phase, cache, stale, succ, req, fc, fetches, gen, evs,
ready, row, floor, rowStart, bad, cleaned, ng>>
Fail(g) ==
/\ sub[g].st = "subscribing"
/\ sub' = [sub EXCEPT ![g].st = "failed"]
/\ cleaned' = [cleaned EXCEPT ![g] = CleanupCount(g)]
/\ UNCHANGED <<want, evMark, ovw, cev, srv, phase, cache, stale, succ, req, fc, fetches, gen, evs,
ready, row, floor, rowStart, bad, ng>>
\* The app's unsubscribe finally reaches the server. Nothing waited for it.
ServerUnsub(g) ==
/\ sub[g].st = "closing"
/\ sub' = [sub EXCEPT ![g].st = "gone"]
/\ UNCHANGED <<want, evMark, ovw, cev, srv, phase, cache, stale, succ, req, fc, fetches, gen, evs,
ready, row, floor, rowStart, bad, cleaned, ng>>
------------------------------------------------------------------------------
(* When a fetch may start. *)
MayFetch(q) ==
LET st == SubSt(q) IN
IF st = "failed" THEN FailBranch = "degrade"
ELSE IF UseReady THEN st = "registered" ELSE TRUE
StartFetch(q) ==
/\ phase[q] = "active"
/\ fc <= MaxFetch
/\ ~ready[q] \/ want[q] \/ stale[q]
/\ ~\E f \in fetches : f.q = q /\ f.start > req[q]
/\ MayFetch(q)
/\ fetches' = fetches \cup {[q |-> q, start |-> fc, rd |-> None]}
/\ fc' = fc + 1
/\ UNCHANGED <<want, evMark, ovw, cev, srv, phase, cache, stale, succ, req, gen, sub, evs, ready,
row, floor, rowStart, bad, cleaned, ng>>
ReadFetch(f) ==
/\ f \in fetches /\ f.rd = None
/\ fetches' = (fetches \ {f}) \cup {[f EXCEPT !.rd = srv]}
/\ UNCHANGED <<want, evMark, ovw, cev, srv, phase, cache, stale, succ, req, fc, gen, sub, evs,
ready, row, floor, rowStart, bad, cleaned, ng>>
(* A fetch counts only if it started after the Query's required start. *)
CompleteFetch(f) ==
/\ f \in fetches /\ f.rd # None
/\ fetches' = fetches \ {f}
/\ LET q == f.q
newer == f.start > succ[q]
authoritative == f.start > req[q]
IN /\ succ' = [succ EXCEPT ![q] = IF newer THEN f.start ELSE succ[q]]
/\ cache' = [cache EXCEPT ![q] = IF newer THEN f.rd ELSE cache[q]]
/\ cev' = [cev EXCEPT ![q] = IF newer THEN FALSE ELSE cev[q]]
/\ stale' = [stale EXCEPT ![q] = IF newer /\ authoritative
THEN FALSE ELSE stale[q]]
/\ IF phase[q] = "active" /\ newer /\ authoritative
THEN /\ ready' = [ready EXCEPT ![q] = TRUE]
/\ PublishResult(f.start, f.rd)
ELSE UNCHANGED <<ready, row, floor, rowStart, ovw>>
/\ want' = [want EXCEPT ![f.q] = IF phase[f.q] = "active" /\ f.start > succ[f.q] /\ f.start > req[f.q]
THEN FALSE ELSE want[f.q]]
/\ UNCHANGED <<srv, phase, req, fc, gen, sub, evs, bad, cleaned, ng, evMark>>
------------------------------------------------------------------------------
(* The provider commits a new version and sends one event to each open
subscription, in generation order. *)
IsOpen(g) == sub[g].st \in {"registered", "closing"}
OpenGens == SelectSeq([i \in 1..MaxGen |-> i], IsOpen)
ServerWrite ==
/\ srv < MaxVer
/\ srv' = srv + 1
/\ evs' = evs \o [i \in 1..Len(OpenGens) |->
[q |-> sub[OpenGens[i]].q, g |-> OpenGens[i],
v |-> srv + 1]]
/\ UNCHANGED <<want, evMark, ovw, cev, phase, cache, stale, succ, req, fc, fetches, gen, sub,
ready, row, floor, rowStart, bad, cleaned, ng>>
Drop(i) == SubSeq(evs, 1, i - 1) \o SubSeq(evs, i + 1, Len(evs))
(* The handler receives an event and calls write() or invalidate(). *)
Deliver(i) ==
/\ i \in 1..Len(evs)
/\ Ordered => i = 1
/\ evs' = Drop(i)
/\ LET e == evs[i]
q == e.q
current == gen[q] = e.g /\ phase[q] # "absent"
IN
IF ~current THEN
IF AcceptDeadGen /\ Active # {}
THEN /\ bad' = bad + 1
/\ Publish(e.v, rowStart)
/\ UNCHANGED <<want, cache, stale, req, evMark, cev>>
ELSE UNCHANGED <<want, cache, stale, req, row, floor, rowStart, bad, evMark, cev>>
ELSE IF phase[q] = "active" THEN
\E kind \in {"write", "invalidate"} :
/\ Stamp(Stamped(q))
/\ IF kind = "write" /\ ~(Versioned /\ row # None /\ e.v <= row)
THEN Publish(e.v, fc - 1) /\ evMark' = fc - 1
ELSE UNCHANGED <<row, floor, rowStart, evMark>>
\* The event may be older than another subset's data: refetch it.
/\ want' = IF kind = "invalidate" THEN [want EXCEPT ![q] = TRUE]
ELSE [p \in Q |-> IF p # q /\ phase[p] = "active" /\ SharedRefetch
THEN TRUE ELSE want[p]]
/\ UNCHANGED <<cache, stale, bad, cev>>
ELSE \* cached: the released-but-cached window
IF IgnoreCached THEN UNCHANGED <<want, cache, stale, req, row, floor, rowStart, bad, evMark, cev>>
ELSE \E kind \in {"write", "invalidate"} :
/\ Stamp(Stamped(q)) \* includes the receiving Query itself
/\ IF CacheBranch = "update" /\ kind = "write"
THEN /\ cache' = [cache EXCEPT ![q] = e.v]
/\ cev' = [cev EXCEPT ![q] = TRUE]
/\ UNCHANGED stale
ELSE /\ stale' = [stale EXCEPT ![q] = TRUE]
/\ UNCHANGED <<cache, cev>>
/\ UNCHANGED <<want, row, floor, rowStart, bad, evMark>>
/\ UNCHANGED <<srv, phase, succ, fc, fetches, gen, sub, ready, cleaned, ng, ovw>>
------------------------------------------------------------------------------
Next ==
\/ \E q \in Q : Acquire(q) \/ Release(q) \/ Remove(q) \/ StartFetch(q)
\/ CollectionCleanup
\/ \E g \in Gens : Register(g) \/ Fail(g) \/ ServerUnsub(g)
\/ \E f \in fetches : ReadFetch(f) \/ CompleteFetch(f)
\/ ServerWrite
\/ \E i \in 1..Len(evs) : Deliver(i)
Spec == Init /\ [][Next]_vars
------------------------------------------------------------------------------
(* Safety properties *)
\* Writes from an evicted subscription never change published rows.
NoDeadEffects == bad = 0
\* The app's cleanup runs at most once per subscription.
CleanupOnce == \A g \in Gens : cleaned[g] <= 1
\* No fetch result (or cached data) that started before a published event
\* write is published, for any Query that holds the row.
NoStaleOverwrite == ovw = 0
\* The row never moves backward. Not promised without versions: a queued
\* older event can be written after a newer fetch result.
NoRegression == row = None \/ row >= floor
\* Once nothing is in flight for the subsets in use, the published row
\* equals the provider's current value.
Quiescent ==
/\ \A i \in 1..Len(evs) : ~(gen[evs[i].q] = evs[i].g /\ phase[evs[i].q] # "absent")
/\ fetches = {}
/\ \A p \in Active : ready[p] /\ ~want[p] /\ ~stale[p] /\ SubSt(p) = "registered"
QuiescentCurrent == (Active # {} /\ Quiescent) => row = srv
\* While a subset is ready and subscribed, every newer provider version is
\* still on its way: a queued event for a live subscription, a fetch that
\* will count, or a refetch that is owed.
Covered(v) ==
\/ \E i \in 1..Len(evs) :
/\ evs[i].v >= v
/\ gen[evs[i].q] = evs[i].g
/\ phase[evs[i].q] # "absent"
\/ \E f \in fetches :
phase[f.q] = "active" /\ f.start > req[f.q] /\ (f.rd = None \/ f.rd >= v)
\/ \E p \in Active : want[p] \/ stale[p] \/ ~ready[p]
NoGap ==
\A q \in Q :
(phase[q] = "active" /\ ready[q] /\ SubSt(q) = "registered")
=> \A v \in (row + 1)..srv : Covered(v)
\* With FailBranch = "reject", a failed subscription never fetches.
RejectHolds ==
FailBranch = "reject" =>
\A q \in Q : (SubSt(q) = "failed" /\ phase[q] = "active") =>
~\E f \in fetches : f.q = q
Sym == Permutations(Q)
Bound == ng <= MaxGen + 1 /\ Len(evs) <= MaxVer * MaxGen
=============================================================================
Summary
Add a
subscribeoption toqueryCollectionOptions. The collection calls it when a Query it manages is created and runs its cleanup when TanStack Query evicts that Query. This works in both sync modes. Inside it, the app listens to its own push channel (WebSocket, SSE, a realtime SDK) and either writes rows into the collection or asks for a refetch.This gives scoped live updates without a hand-built subscription manager.
Background: sync modes and subsets
A query collection loads data in one of two modes:
syncMode: 'eager', the default). One Query loads the whole collection. Every live query reads from that one result.syncMode: 'on-demand'). The collection loads only what live queries ask for. Each live query'swhere,orderByandlimitbecome a subset, for example "comments whereissueId = 42". Each distinct subset gets its own Query, keyed by that subset, and thequeryFnreceives the subset so it can ask the server for just those rows. Live queries that ask for the same subset share its Query. A subset's rows stay in the collection while any live query needs them.Subsets let an app with a large dataset load only what is on screen. They are also the natural unit for live updates: "comments on issue 42" is both what the app fetched and what it wants to hear about.
The problem
Many apps pair fetched data with a push channel. A common setup:
Two requirements make this hard to do well:
Today there are two common patterns, and neither meets both requirements:
useEffectnext to eachuseQueryoruseLiveQuerycall. This meets (1) but not (2): each component subscribes separately unless the app adds its own deduplication. It also ties the subscription to components rather than to the data.Teams end up writing a manager on top of the QueryClient cache events to get both. These managers are subtle. Common failure modes:
The documented WebSocket integration example shows how to apply events to a collection. It is right for data that is always loaded. It does not cover subscriptions scoped to what is currently in use.
Proposed API
Eager mode
In eager mode there is one Query, so
subscribeis called once, when the collection starts syncing.subsetdescribes the whole collection (nowhere). The subscription lasts as long as that Query:Compared with a WebSocket opened at module scope, this ties the subscription to the collection's sync. It starts when the collection is first used and stops when the collection is cleaned up. It also gives the same ordering guarantees as on-demand mode.
Context passed to
subscribesubset: the subset being loaded (the sameLoadSubsetOptionsthequeryFnreceives). In eager mode it covers the whole collection.data: the Query's data, as returned byqueryFn. It isundefinedunlesssubscribeOn: 'data'is set (see "When the channel comes from the data").queryKeyandmeta: the Query's key andmeta, for handlers that need them.write(change): applies a change to the collection's rows.invalidate(): refetches this subset. Calls are coalesced and can be debounced (see below).Two ways to handle an event
Both are first-class, and a handler can mix them per event type:
invalidate()). The event only says "this subset changed". It is the simplest option: the server needs no event payloads, and the rows always come from a fetch. It also avoids the brief step back and the shared-row refetch described under "Design notes". It is the right default for most apps.write()). The event carries the changed row, and no fetch is needed. Use it when refetching the subset on every event costs too much, for example for large or paged subsets or busy topics.Refetch hints are coalesced. A burst of
invalidate()calls while a refetch is pending leads to at most one follow-up refetch, and that refetch starts after the latest hint. An optional debounce, for exampleinvalidateDebounce: 200, waits for a burst to settle before refetching.Neither changes the guarantees. Only a refetch that starts after the latest event counts, and the model already lets any refetch start arbitrarily late.
Return value
unsubscribe: required. May return a promise. The library never waits on it.ready: optional promise. See "Ordering the first fetch" below.Lifecycle
The subscription has exactly the lifetime of its Query in the QueryClient cache. In on-demand mode that is one subscription per subset:
subscribe.readywas returned, the first fetch starts after it resolves. Otherwise it starts at once.writeorinvalidate.subscribecall happens.gcTimeexpires. The subscription stays open during that time, so if the user comes back, the cached data is still current.unsubscribe. After this,writeandinvalidatecalls from that subscription are ignored, even if the app's unsubscribe has not finished yet.So "held" covers two states: in use by a live query, and cached for reuse until
gcTimeruns out. There is deliberately no separate lifetime setting for the subscription:gcTime, the query collection serves the cached rows at once. If the subscription had closed, those rows may have missed updates, and the user sees them until a refetch replaces them.An app that wants the subscription to end sooner after the last consumer leaves can lower
gcTimefor that collection. The subscription and the cached data then go together.Guarantees
writeorinvalidatecannot replace that event's result when it completes. The query collection already applies this rule to local writes. This proposal applies it to events too, for every Query that holds the affected rows: in use or cached, including the Query that received the event.readyis used, every change after a subset is ready is either already shown, or on its way as a queued event, a fetch, or an owed refetch.unsubscribecannot corrupt data. It costs only an open server subscription.These guarantees assume the app's socket delivers each subscription's events in order, and that
writesends full row values, not deltas. With versions (below), ordering is no longer needed. Each guarantee is checked in the formal model described below.Design notes
Ordering the first fetch
If the fetch reads data before the server starts delivering events, a change made in between is lost: the fetch did not see it, and no event arrives for it.
readycloses this window. The library starts the first fetch only after the server confirms it is delivering.When the channel comes from the data
Some servers return the channel to listen on as part of the response, for example a
liveChannelfield on the list. Then the app cannot subscribe until the first fetch has returned. SetsubscribeOn: 'data':The library calls
subscribeafter the first successful fetch and passes its data. A change made after that fetch read, and before the subscription began, sends no event. So oncereadyresolves, the library refetches the subset one more time. The subscription still ends when the Query is evicted.Events that race a fetch
Take this order of events on the server:
The event is now older than the fetch result. What happens depends on which reaches the client first.
invalidate(), the only cost is one extra fetch.Without versions, the library cannot tell an event that is older than the data from one that is newer. It chooses the safe side each time: an extra fetch rather than stale data, and a brief step back rather than a missed update. Handlers that write directly should send full row values, not deltas, so that applying an event twice is harmless.
Rows shared by several subsets
Two subsets can hold the same row, for example "open issues" and "issues assigned to me". Two rules keep a shared row correct.
writelands on a row that another in-use subset also holds, that subset refetches.invalidate()needs no such rule, and rows held by one subset pay nothing extra.Versions (optional)
If events and fetch results carry a version, such as a commit sequence or
updated_at, most of these costs go away. The library can skip an event the data already includes, and keep a fetch result that is newer than the event. In the model, versions make three things unnecessary:This could be an optional addition, for example a
getVersionoption.Events for data that is cached but not shown
After the last consumer leaves, an event can arrive while a fetch for that subset is still in flight. That older fetch must not complete and mark the cache current. So every event also raises the bar for fetches of the Query that received it: only a fetch that starts after the event can make the data current again.
The library marks the cached data stale. Writing the event into the cached data instead gains nothing: under this rule, the next use refetches either way.
Reconnect
The app's socket owns reconnection. After a reconnect, the app should call
invalidate()for the subsets it serves, because events may have been missed. The library adds no reconnect logic and no shared socket abstraction. Topic fan-out and deduplication across collections stay in the app's own pub/sub.Formal model
We checked this design with a TLA+ model before writing any code. The model is the spec for the implementation. It is also meant to become the reference model for an oracle test in
packages/query-db-collection/tests.What it models
gcTime. Also collection cleanup.ready, failure, and an unsubscribe that may take arbitrarily long to reach the server.writeorinvalidatefor each event.Every action can interleave with every other. The checker explores all orderings within small bounds: 2 versions, 3 fetches and 2 subscriptions per run. That is 1 to 10 million distinct states and under a minute per configuration on a laptop.
Properties
NoStaleOverwriteNoGapQuiescentCurrentNoDeadEffectsCleanupOnceRejectHoldsNoRegressionWhat it changed
The first draft of this proposal failed four times. Each counterexample became a rule above:
It also settled one open question. Updating cached data from an event and marking it stale behave the same, so the library only marks it stale.
Configurations checked
Each rule and assumption has a configuration that turns it off. The checker must then find the failure that the rule prevents. This shows each rule is needed, and that the properties can fail.
ready: fetch before the server confirmsNoGapfails: an update is missedQuiescentCurrentfailsQuiescentCurrentfailsNoStaleOverwritefailsNoStaleOverwritefailsQuiescentCurrentfailsQuiescentCurrentfailsNoDeadEffectsfailsCleanupOncefailsQuiescentCurrentfailsNoRegressionNoRegressionfailsRun with TLC 2.19:
java -cp tla2tools.jar tlc2.TLC -deadlock -config Base.cfg SubscribeHook.tla. The-deadlockflag is needed because the model has no liveness steps once the bounds run out.From model to oracle
The model maps onto the repository's oracle structure:
Nextactions.fast-checkcommands. Acquire, release, evict, server change, register, fail, deliverwriteorinvalidate, and fetch start, read and complete.QueryClientand query collection with a controllable fake socket andqueryFn, so the test can hold and release each fetch and each event.The fault-switch configurations become the hostile controls: each should fail the oracle.
SubscribeHook.tlaBase.cfg(the proposed design)Open questions
subscribethrows orreadyrejects, should the subset load reject, or succeed without live updates and report that? The model passes with either policy whenreadyis used, and cleanup runs exactly once in both. Withoutready, "reject" cannot be kept: the fetch may already be running when the failure arrives.subscribeOn: 'data', a later fetch may return a different channel. Should the library resubscribe when a chosen field changes, or leave that to the handler?Alternatives considered
gcTime, for example closing once no live query needs the subset and its data is stale bystaleTime, then refetching on return. This works well with plain TanStack Query. It fits query collections less well, because they serve cached subsets at once on return. Rejected for the reasons under "Lifecycle". LoweringgcTimegives a shorter lifetime without a second setting.