Skip to content

Proposal: a subscribe option for query collections #1935

Description

@KyleAMathews

Summary

Add a subscribe option to queryCollectionOptions. 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:

  • Eager (syncMode: 'eager', the default). One Query loads the whole collection. Every live query reads from that one result.
  • On-demand (syncMode: 'on-demand'). The collection loads only what live queries ask for. Each live query's where, orderBy and limit become a subset, for example "comments where issueId = 42". Each distinct subset gets its own Query, keyed by that subset, and the queryFn receives 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:

  • One app-wide WebSocket.
  • Each part of the UI asks the server for events on a topic, such as "comments on issue 42".
  • When an event arrives, the app refetches the affected data or writes the changed row directly.

Two requirements make this hard to do well:

  1. Subscribe only while the app holds the data. Once the app has dropped issue 42's comments, it should stop listening for them. Data that is cached for quick reuse counts as held (see "Lifecycle").
  2. Subscribe once. If three components show the same comments, the server should see one subscription, not three.

Today there are two common patterns, and neither meets both requirements:

  • A useEffect next to each useQuery or useLiveQuery call. 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.
  • One subscription at the route or app root. This meets (2) but not (1): the subscription outlives the data, or is missing when the same data appears on another route.

Teams end up writing a manager on top of the QueryClient cache events to get both. These managers are subtle. Common failure modes:

  • An update that happens between the first fetch and the start of the subscription is missed.
  • A resubscribe waits on a teardown that never finishes, because the socket is gone.
  • A fetch that started before an event lands afterwards and overwrites it.

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

const comments = createCollection(
  queryCollectionOptions({
    queryKey: ["comments"],
    queryFn: ({ meta }) => api.comments.list(meta.loadSubsetOptions),
    queryClient,
    getKey: (c) => c.id,
    syncMode: "on-demand",

    subscribe: ({ subset, write, invalidate }) => {
      const topic = `issue:${subset.where.issueId}:comments`;

      const unsubscribe = socket.listen(topic, (event) => {
        switch (event.type) {
          case "comment-upserted":
            write({ type: "upsert", value: event.comment });
            break;
          case "comment-deleted":
            write({ type: "delete", key: event.id });
            break;
          default:
            invalidate(); // too large to patch, or unknown: refetch this subset
        }
      });

      return {
        unsubscribe,
        // Optional. Resolves when the server has started delivering events.
        ready: socket.whenSubscribed(topic),
      };
    },
  }),
);

Eager mode

In eager mode there is one Query, so subscribe is called once, when the collection starts syncing. subset describes the whole collection (no where). The subscription lasts as long as that Query:

queryCollectionOptions({
  queryKey: ["projects"],
  queryFn: () => api.projects.list(),
  queryClient,
  getKey: (p) => p.id,
  subscribe: ({ write, invalidate }) => {
    const unsubscribe = socket.listen("projects", (event) =>
      event.type === "project-upserted"
        ? write({ type: "upsert", value: event.project })
        : invalidate(),
    );
    return { unsubscribe, ready: socket.whenSubscribed("projects") };
  },
});

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 subscribe

  • subset: the subset being loaded (the same LoadSubsetOptions the queryFn receives). In eager mode it covers the whole collection.
  • data: the Query's data, as returned by queryFn. It is undefined unless subscribeOn: 'data' is set (see "When the channel comes from the data").
  • queryKey and meta: the Query's key and meta, 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:

  • Refetch hints (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.
  • Direct writes (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 example invalidateDebounce: 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:

  1. A live query needs comments for issue 42. The collection creates the Query for that subset and calls subscribe.
  2. If ready was returned, the first fetch starts after it resolves. Otherwise it starts at once.
  3. Events arrive and the handler calls write or invalidate.
  4. A second component needs the same subset. It shares the Query, so no second subscribe call happens.
  5. No component needs the subset any more. Its rows leave the collection, but the Query stays cached until gcTime expires. The subscription stays open during that time, so if the user comes back, the cached data is still current.
  6. TanStack Query evicts the Query. The library calls unsubscribe. After this, write and invalidate calls 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 gcTime runs out. There is deliberately no separate lifetime setting for the subscription:

  • Closing it before eviction leaves cached rows unwatched. When a live query needs the subset again within 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.
  • Keeping it open after eviction would accept writes for data the collection no longer holds.

An app that wants the subscription to end sooner after the last consumer leaves can lower gcTime for that collection. The subscription and the cached data then go together.

Guarantees

  • One subscription per subset. All consumers of a subset share it.
  • A write is not overwritten by an older fetch. A fetch that started before an event's write or invalidate cannot 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.
  • A newer fetch is not overwritten by an older one. When two subsets hold the same row, a fetch result never replaces a result from a fetch that started later. (This also fixes an existing race between overlapping subsets. See "Rows shared by several subsets".)
  • No missed updates while ready. When ready is 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.
  • Rows end up current. Once nothing is in flight for the subsets in use, every row matches the server.
  • No writes after eviction. A slow or hung unsubscribe cannot corrupt data. It costs only an open server subscription.
  • Cleanup runs once. Collection cleanup ends every subscription, and the app's cleanup runs exactly once per subscription, including after a failed subscribe.

These guarantees assume the app's socket delivers each subscription's events in order, and that write sends 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. ready closes 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 liveChannel field on the list. Then the app cannot subscribe until the first fetch has returned. Set subscribeOn: 'data':

queryCollectionOptions({
  // ...
  subscribeOn: "data",
  subscribe: ({ data, write, invalidate }) => {
    const unsubscribe = socket.listen(data.liveChannel, (event) => {
      /* ... */
    });
    return { unsubscribe, ready: socket.whenSubscribed(data.liveChannel) };
  },
});

The library calls subscribe after the first successful fetch and passes its data. A change made after that fetch read, and before the subscription began, sends no event. So once ready resolves, 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:

  1. The client starts a fetch.
  2. A row changes, and the server emits an event.
  3. The backend runs the fetch's query, so the result already includes the change.

The event is now older than the fetch result. What happens depends on which reaches the client first.

  • The event arrives before the fetch completes. The fetch started before the event, so its result cannot count. If the subset still needs data, the library runs another fetch. This can cost one extra request, even though the discarded result was already current: the library cannot tell that the backend read after the change.
  • The event arrives after the fetch completes, for example because it was queued on the socket behind other messages. The row already shows the change.
    • If the handler calls invalidate(), the only cost is one extra fetch.
    • If the handler writes the event's row directly, the row is rewritten with the same value, or with an older one if later changes have also landed. In that case the row goes back briefly. It catches up when that subscription's next event arrives.

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.

  • A newer fetch wins. If subset A's fetch started before subset B's but completes after it, A's older result does not replace B's row. The query collection does not do this today, so an overlapping subset can briefly show older data even without subscriptions. The rule is worth fixing on its own.
  • A direct write refetches the row's other holders. Subset A's socket may still hold an old event when subset B has already fetched newer data. If A's handler writes that event, B would show the older value. If A is then evicted before its next event arrives, nothing would ever correct it. So when a write lands 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:

  • the brief step back;
  • the shared-row refetch;
  • the requirement that the socket deliver events in order.

This could be an optional addition, for example a getVersion option.

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

  • Two subset Queries that hold one shared row. Two are enough to reach every shared-row case above.
  • A provider that changes the row, and pushes an event to every open subscription.
  • Each subset Query's lifecycle: in use, released but cached, and evicted at gcTime. Also collection cleanup.
  • Subscriptions: subscribe, ready, failure, and an unsubscribe that may take arbitrarily long to reach the server.
  • Fetches in three steps: start, read the provider's value, complete. Any number of steps from other actions can fall in between.
  • Handlers that call write or invalidate for 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

Property Guarantee
NoStaleOverwrite A write is not overwritten by an older fetch
NoGap No missed updates while ready
QuiescentCurrent Rows end up current
NoDeadEffects No writes after eviction
CleanupOnce Cleanup runs once
RejectHolds With "reject on failure", a failed subscription never fetches
NoRegression The row never moves backward (with versions only)

What it changed

The first draft of this proposal failed four times. Each counterexample became a rule above:

  1. A queued older event can overwrite newer fetched data. So without versions, the design cannot promise that a row never moves backward. It promises that rows end up current.
  2. A subset released during a fetch. An event for the now-cached Query marked it stale, but the older fetch then completed and cleared the mark. Hence the rule that every event also invalidates older fetches for the Query that received it.
  3. Two subsets' fetches completing out of order left the shared row older than either subset's data, with nothing to correct it. Hence "a newer fetch wins". This race exists in query collections today.
  4. An old event on one subset's socket overwrote a shared row after another subset had fetched newer data, and the first subset was then evicted. Hence the shared-row refetch.

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.

Configuration Result
Proposed design (either failure policy; with or without versions) all properties hold
With versions, events delivered out of order all properties hold
No ready: fetch before the server confirms NoGap fails: an update is missed
No "newer fetch wins" QuiescentCurrent fails
No shared-row refetch QuiescentCurrent fails
Events invalidate older fetches only for their own Query NoStaleOverwrite fails
Events do not invalidate older fetches NoStaleOverwrite fails
Events invalidate older fetches only for in-use Queries QuiescentCurrent fails
Events for cached Queries are ignored QuiescentCurrent fails
Writes from an evicted subscription still apply NoDeadEffects fails
Cleanup not guarded against running twice CleanupOnce fails
Events delivered out of order, no versions QuiescentCurrent fails
No versions, checking NoRegression NoRegression fails

Run with TLC 2.19: java -cp tla2tools.jar tlc2.TLC -deadlock -config Base.cfg SubscribeHook.tla. The -deadlock flag 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:

  • Contract: the guarantees above.
  • Model: this spec, or a small TypeScript port of its Next actions.
  • History grammar: the same actions, generated with fast-check commands. Acquire, release, evict, server change, register, fail, deliver write or invalidate, and fetch start, read and complete.
  • Production driver: a real QueryClient and query collection with a controllable fake socket and queryFn, so the test can hold and release each fetch and each event.
  • Refinement check: the published rows compared at each publication, using the same properties.

The fault-switch configurations become the hostile controls: each should fail the oracle.

SubscribeHook.tla
---------------------------- 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
=============================================================================
Base.cfg (the proposed design)
SPECIFICATION Spec
CONSTANTS
  Q = {q1, q2}
  MaxVer = 2
  MaxFetch = 3
  MaxGen = 2
  UseReady = TRUE
  CacheBranch = "stale"
  FailBranch = "reject"
  EventStamp = "owners"
  Ordered = TRUE
  AcceptDeadGen = FALSE
  IgnoreCached = FALSE
  DoubleCleanup = FALSE
  Versioned = FALSE
  RowGuard = TRUE
  SharedRefetch = TRUE
INVARIANTS NoDeadEffects CleanupOnce NoStaleOverwrite QuiescentCurrent NoGap RejectHolds
CONSTRAINT Bound
SYMMETRY Sym

Open questions

  • Subscribe failure. When subscribe throws or ready rejects, should the subset load reject, or succeed without live updates and report that? The model passes with either policy when ready is used, and cleanup runs exactly once in both. Without ready, "reject" cannot be kept: the fetch may already be running when the failure arrives.
  • Channel changes in later data. With 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?
  • Status. Should subsets expose whether they are receiving live updates, so the UI can show data that is only cached?

Alternatives considered

  • A shared-socket adapter in the library (topic reference counting, reconnect, subscribe confirmation). Rejected. Apps already have a pub/sub layer, and topic routing is app-specific.
  • A subscription lifetime separate from gcTime, for example closing once no live query needs the subset and its data is stale by staleTime, 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". Lowering gcTime gives a shorter lifetime without a second setting.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions