Replication
Offline-first bidirectional sync between local RxDB and remote APIs
Tally UI uses RxDB's replication protocol to keep local data in sync with your backend. Each connector ships with replication adapters that handle the specifics of its API -- pagination, checkpoints, and conflict resolution are all built in.
Quick Setup
import { createTallyDatabase, startReplication, getStorage } from '@tallyui/database';
import { createWooCommerceConnector } from '@tallyui/connector-woocommerce';
// One connector per store session: build a new one on each sign-in or store change
const connector = createWooCommerceConnector();
// Create a persistent database
const db = await createTallyDatabase({
connector,
storage: getStorage(), // in-memory in Node; on the web, pass getRxStorageSQLiteWasm explicitly -- see Storage below
});
// Start replicating products
const replicationState = startReplication({
collection: db.products,
adapter: connector.replication.products,
context: {
connectorId: connector.id,
baseUrl: 'https://mystore.com/wp-json/wcpos/v2',
headers: connector.auth.getHeaders(credentials),
},
});
// Monitor sync status
replicationState.active$.subscribe((active) => {
console.log('Syncing:', active);
});
replicationState.error$.subscribe((error) => {
console.error('Sync error:', error);
});How It Works
Each connector's ReplicationAdapter provides a pull handler, which fetches documents changed since the last checkpoint. The checkpoint format is connector-specific -- WooCommerce uses offset + modified date, Shopify uses cursor pagination, Medusa uses offset + date, Vendure uses GraphQL skip + date.
Product replication is pull-only. Catalogue data is owned by the server, so the POS never writes products back. Data the POS creates, such as orders, will travel through a separate command outbox with idempotent writes (see the programme plan in docs/plans/). ReplicationAdapter keeps an optional push for that work.
RxDB's replicateRxCollection orchestrates the cycle: pull remote changes, apply them locally, repeat.
Storage
Web (SQLite-wasm)
Breaking: getStorage() on the web no longer returns Dexie. It throws with
guidance — pass the storage explicitly, bundling the worker entry:
import { createTallyDatabase } from '@tallyui/database';
import { getRxStorageSQLiteWasm } from '@tallyui/storage-sqlite/web';
// Vite (and most bundlers): a factory the bundler can statically detect and bundle
const workerInput = () =>
new Worker(new URL('@tallyui/storage-sqlite/web-worker', import.meta.url), { type: 'module' });
const db = await createTallyDatabase({
connector: myConnector,
storage: getRxStorageSQLiteWasm({ workerInput }),
});On Expo web there is no single import path across bundlers: your app supplies
its own worker file, built from @tallyui/storage-sqlite/web-worker, that its
bundler can emit, and passes that as workerInput.
This is RxDB Premium's SQLite on sqlite-wasm's
opfs-sahpool VFS, so your app needs a Premium licence key (ADR-031). It runs
one live tab per store (see Live tab), and a watchdog
reports the storage worker's health. No storage call is ever failed on a timer,
because a timed-out write may still commit. A write pending for more than 10s is
flagged as stalled, and its promise stays pending until the worker answers.
Reads are watched too: two silent 30s windows, with reads pending and nothing
settling, mean the worker is dead, and the app's recovery is to reload.
import { getStorageHealth } from '@tallyui/database';
getStorageHealth(db)?.subscribe(({ status }) => {
if (status === 'stalled') showBanner('Saving is slow…');
if (status === 'dead') showBanner('Storage stopped. Reload the app.');
});Switching from the previous Dexie storage is a cold resync, not a data migration — the SQLite database starts empty and replicates fresh. Once the new database has synced, delete the old one (the name depends on your app's database name):
indexedDB.deleteDatabase('tally_' + connector.id); // the old Dexie databaseReact Native (SQLite)
On React Native, provide the SQLite adapter explicitly. It runs on
RxDB Premium's SQLite storage, so your app installs
rxdb-premium (the same version as rxdb) under its own RxDB Premium licence:
import { createTallyDatabase } from '@tallyui/database';
import { getRxStorageSQLite } from '@tallyui/storage-sqlite';
import { openDatabaseSync } from 'expo-sqlite';
const sqliteDb = openDatabaseSync('tally.db');
const db = await createTallyDatabase({
connector: myConnector,
storage: getRxStorageSQLite(sqliteDb),
});Use one SQLite handle per RxDB database. The storage never closes the handle; your app does.
Tests and SSR
In Node.js environments (tests, server-side rendering), getStorage() returns in-memory storage automatically. No configuration needed.
ReplicationAdapter Interface
If you're building a custom connector, implement ReplicationAdapter for each collection:
import type { ReplicationAdapter } from '@tallyui/core';
type MyCheckpoint = { offset: number; updatedAt: string };
const myProductReplication: ReplicationAdapter<any, MyCheckpoint> = {
pull: {
async handler(lastCheckpoint, batchSize, context) {
// Fetch products changed since lastCheckpoint from your API
const response = await fetch(
`${context.baseUrl}/products?since=${lastCheckpoint?.updatedAt ?? ''}&limit=${batchSize}`,
{ headers: context.headers },
);
const products = await response.json();
// Add _deleted: false (your API returns non-deleted products)
const documents = products.map((p) => ({ ...p, _deleted: false }));
// Return documents + new checkpoint
const last = products[products.length - 1];
return {
documents,
checkpoint: last
? { offset: 0, updatedAt: last.updatedAt }
: lastCheckpoint ?? { offset: 0, updatedAt: '' },
};
},
},
};Product adapters have no push: the catalogue is server-owned, and the POS never writes products.
startReplication Options
startReplication({
collection, // RxDB collection to replicate
adapter, // ReplicationAdapter for this collection
context, // SyncContext with connectorId, baseUrl, headers
live: true, // Keep syncing after initial pull (default: true)
retryTime: 5000, // Retry interval on error in ms (default: 5000)
autoStart: true, // Start immediately (default: true)
maxRetryTimeMs: 300000, // Backoff ceiling and store-error wait (default: 5 minutes)
});Returns a TallyReplicationState: RxDB's RxReplicationState with notice$ and resume() added.
| Property/Method | Description |
|---|---|
active$ | Observable: is replication currently running? |
error$ | Observable: replication errors |
received$ | Observable: documents received from remote |
sent$ | Observable: documents sent to remote |
cancel() | Stop replication |
awaitInSync() | Promise that resolves when fully synced |
reSync() | Trigger an immediate sync cycle |
notice$ | Observable: a SyncNotice while a till or store error stops the pull (see Errors), otherwise undefined |
resume() | Clear the notice, the backoff and the error gates, and pull again from the stored checkpoint |
Errors
A pull error falls in one of three classes, by who can fix it.
errorKind(error) from @tallyui/core returns 'till' or 'store' when
the error's class sets fixedBy to that value and a string code, and
'transient' for anything else.
| Class | Who fixes it | Errors | What the pull does |
|---|---|---|---|
| Till | the till: someone signs in | ConnectorUnauthorizedError with status: 401 (unauthorized) | Stops until resume() |
| Store | the store owner, elsewhere | ConnectorUnauthorizedError with status: 403 (forbidden), WooDateFilterError (unsupported_store), WooTokenRefusedError (store_misconfigured), WooMissingUuidError (missing_plugin), VendureTimezoneConfigError (store_misconfigured) | Retries every 5 minutes and recovers by itself |
| Transient | nobody: it passes | anything without the marker: network errors, timeouts, 5xx, 429 | Retries with a doubling delay |
- Till (including only a 401
ConnectorUnauthorizedError): the pull makes that one request, setsnotice$once to{ code, since, fixedBy: 'till' }, and pauses. The error does not appear onerror$. Untilresume(), every pull returns an empty page without a request, whatever restarts the loop:reSync(),stream$, or RxDB's ownstart()when the page becomes visible again. Callresume()after sign-in; an app restart starts afresh anyway. - Store: the pull sets
notice$once to{ code, since, fixedBy: 'store' }(a repeat keeps the firstsince), emits the error onerror$, and waits 5 minutes (maxRetryTimeMs) before the next request. It does not pause and does not wait longer each time. A pull that runs sooner, for whatever reason, makes no request: it fails again with the same error, and RxDB retries when the 5 minutes are up. The first successful pull clears the notice and the backoff, so the till recovers by itself within 5 minutes of the owner's fix. A transient error in between leaves the notice in place. - Transient: RxDB emits it on
error$and retries. The wait starts atretryTime, doubles after each consecutive failure up to 5 minutes (maxRetryTimeMs), and goes back toretryTimeafter a successful pull.
An error may carry retryAfterMs, the server's own wait. It counts only
as a finite number of at least 0; the wait is then at least that long, up
to one hour (MAX_RETRY_AFTER_MS). A store error waits the longer of 5
minutes and retryAfterMs, also up to one hour.
The error's own software, minVersion and fix strings, when it has
them, are copied onto the notice: WooDateFilterError names WooCommerce
and 5.8, WooTokenRefusedError gives the fix, checking the JWT Authentication plugin's settings, and
VendureTimezoneConfigError gives the fix, running the server in UTC. Show the notice with SyncStatus's pullNotice prop,
which turns the code into words for the cashier and never shows the code
itself.
Push errors are retried by RxDB at a fixed retryTime, as before.
replicationState.notice$.subscribe((notice) => setPullNotice(notice));
// after the till signs in again:
await replicationState.resume();Connector Checkpoint Formats
Each connector uses a checkpoint format that matches its API's pagination model:
| Connector | Checkpoint | Pull Strategy |
|---|---|---|
| WooCommerce | { modified, offset, pass_mark?, pass_count? } | Pass-based: ?modified_after= (−1 s, dates_are_gmt=true) from a high-water mark, orderby=id, offset pages, ends on X-WP-Total or a short page |
| Shopify | { updated_at } | ?updated_at_min= with Link header cursors |
| Medusa | { offset, updated_at, pass_mark?, pass_count? } | Pass-based: ?updated_at[$gte]= from a high-water mark, order=id, offset pages, ends on count |
| Vendure | { skip, updatedAt, passHighWater?, passTotal? } | Pass-based: updatedAt.after from a high-water mark (−1 ms), sort: { id: ASC }, take/skip, ends on totalItems |
Writing a Pull Adapter
These rules come from bugs that only showed up inside RxDB's real
replicateRxCollection loop, or against a real backend (TallyUI #45, #48;
ADR-049, ADR-060). A new adapter, such as a variant feed or a new backend,
should follow them.
How RxDB's pull loop behaves (RxDB 16, replication-protocol/downstream.js):
- It stops on a short page. RxDB keeps pulling only while a handler
returns at least
batchSizedocuments. A short page ends the cycle, even mid-pass, until the nextreSync()or stream event. An adapter whose pages can shrink (a variant feed that re-delivers parent products, say) must keep reading source pages inside the same handler call until it has at leastbatchSizedocuments, or its pass has ended. - It discards an empty page's checkpoint. If a pass ends on an empty
page, the cursor stays wherever the last full page left it. With a
catalogue that is an exact multiple of the batch size, the cursor sticks
at the end, and every later sync returns nothing. End a pass on the list
total (
skip + page >= total), never on an empty page. - It merges checkpoints. Each new checkpoint is
Object.assigned onto the previous one (stackCheckpoints). Fields you drop are not removed; clear per-pass state explicitly withundefinedwhen a pass completes.
Cursor design:
- Keep the pass window fixed. Keep the lower bound fixed for the whole pass, and page by a stable key (id order).
- Take the next lower bound from a high-water mark read at pass start.
Read the newest
updated_atwith one small request. Don't use the newest timestamp seen during the pass, or an early row updated mid-pass is lost. - Return early when nothing changed. If the mark equals the stored bound, return an empty page with the checkpoint unchanged. This ends the loop even when a full batch of rows shares the mark's timestamp, and it makes an idle poll one request.
- Restart the pass if rows vanish. If a later page is empty, or the total shrank (rows already read were deleted, shifting later rows left), restart from zero inside the same call. Re-reads are harmless; skips are not.
Documents and timestamps:
- Project documents onto the schema's top-level fields. RxDB rejects
undeclared top-level properties. One extra field (
createdAt) made every Vendure batch fail, so nothing was ever stored. - Allow ±1 ms around timestamps. Postgres stores microseconds, and APIs
return milliseconds, so
after(t)matches the row stampedt. - Know the backend's time zone behaviour. Vendure's
timestamp without time zonecolumns skewupdatedAtfilters by the server's UTC offset, unless the process and database session both run in UTC. The Vendure connector guards against it and offersupdatedAtSkewMs.
Local writes:
- Never write locally into a replicated collection. While a local
write to a document is in flight, RxDB's downstream skips the next
version it pulls for that document, so a concurrent server change
can be lost until the document changes again. Keep derived or
corrected data in a separate, non-replicated collection and merge
it when reading, as the stock reconcile does with
stock_levels(ADR-060). A reconcile that corrects replicated documents hands them to the collection's pull throughcreateReconcileFeedinstead; the catalogue reconcile keeps its own gate and cursor in RxDB local documents, which are never replicated (#248). - Key the reconcile feed by the primary key. The feed matches
fetched documents to queued entries by
doc.idunless it is givenkey. Where the primary key is not the backend id (WooCommerce'suuid), passkey: (doc) => doc.uuid:fetchByIdsthen receives the entries, withlocalandremote, to find each backend id. Unkeyed, no fetched document would match, and every entry would be tombstoned. A fetched document may carry_deleted: true(an unpublished product), and the feed keeps it. - One replication per collection. Two replications on one
collection each see the other's writes as local writes, so either
can skip a pulled version the same way. Give a collection several
feeds (products and variants, say) by combining their pull adapters
with
combinePullAdaptersinto the one adapter the collection replicates with.
Proving it:
- A handler unit test is not enough. Add a test that runs the real
replicateRxCollectionagainst an in-memory server with the backend's exact filter semantics: exact multiples of the batch size, mid-pass updates, mid-pass deletions, and ties. - Also add an env-gated live test against the backend's dev store.
- Prove every incremental filter live with a future-dated mark. A
fake server honours whatever it is sent, so a misspelt filter passes
every unit test. Medusa 2.21 silently ignored
updated_at[gte]and returned the whole catalogue; onlyupdated_at[$gte]filters. A live test asserts that a mark in the future returns nothing.