@syncular/server
v0.30.16
Published
Syncular server: handleSyncRequest + storage/auth interfaces for the sync protocol
Maintainers
Readme
@syncular/server
Framework-free embeddable SSP2 protocol library. Core surface:
handleSyncRequest(bytes, ctx) → bytes over host-provided storage /
scope-resolution / segment-store interfaces, plus the transport-agnostic
realtime hub (§8), the direct segment download handler (§5.5), commit-log
pruning (§4.6), and signed-URL token issuance (§5.4). SPEC.md is
normative for everything on the wire; this README covers the host
surface — in particular the ops seam and the pruning runbook.
Application processes can expose generated named queries and transactional
commands through RemoteOperationRegistry. Queries stay in the server
registry, command mutations use the ordinary serialized push path, and
RemoteOperationWatchHub provides live replacement snapshots. The protocol is
specified in docs/REMOTE.md and the practical setup is
in the remote operations guide.
Remote query registration requires generated relationPlans for each selected
SQL statement. Run syncular generate before upgrading existing query modules.
The server uses those boundaries to bind every physical table occurrence to
the authenticated partition, including quoted self joins and CTE bodies.
Application intent belongs in immutable domain event rows written in the same
commit as the state change. SyncularServerEvents below remains operational
telemetry. See the domain event guide.
Deployment matrix (runtime adapters)
The server core is runtime-neutral TypeScript — handleSyncRequest and
the realtime session speak only Web Request/Response/fetch/Web-Crypto,
no Bun- or Node-only builtin (enforced by a static import-graph scan,
test/runtime-neutrality.test.ts). Adapters wire that core to a runtime.
The supported set, and what deliberately does not get an adapter:
| Runtime | Adapter | Transport | Storage | Status |
| --- | --- | --- | --- | --- |
| Bun / Node 22.13+ | @syncular/server-hono | HTTP (POST /sync, segments, blobs) + WS realtime (§8, host-driven upgrade) | SqliteServerStorage through @syncular/server/sqlite, Postgres, memory | Supported now — the reference deployment; runs the full conformance catalog on both bindings. |
| Cloudflare Workers | @syncular/server-workers | HTTP binding via Hono (Workers-native) + optional WS realtime (§8) | D1ServerStorage behind one per-partition Durable Object queue; R2-as-S3 for segments/blobs | Supported now — D1 sync writes always traverse the DO; WebSocket upgrades remain optional. |
| Raw Deno / edge-misc | — | — | — | Not adapted (policy below). |
The policy for "not adapted". Untested ≠ unsupported forever. The core is runtime-neutral TS, so Deno/edge would very likely run it — but an adapter is only supported where the conformance catalog can run against it. We ship adapters for the runtimes where we run conformance (Bun/Node fully; Workers HTTP via the fetch-handler round-trip tests), and we do not claim runtimes we do not test. Deno is a plausible future adapter the day someone runs the catalog on it; until then it is neutral-core-friendly, not supported.
Workers coordination and realtime — the Durable Object. HTTP-only is fully
conformant at the protocol level, but D1 still requires a per-partition Durable
Object queue for every /sync round that may push. WebSocket upgrades are
optional; the write coordinator is not. SyncularRealtimeDO hosts that FIFO
and, when realtime is enabled, the RealtimeHub. The DO id derives from the
partition, so unrelated partitions remain concurrent. WebSocket
hibernation keeps idle sockets from billing wall time (the existing
RealtimeSession is the
per-connection state machine, driven from the hibernation callbacks and
rehydrated from a minimal socket attachment + the D1 client record on wake);
storage via the same D1 binding so realtime rounds and POST /sync rounds
share one commit log and one segment store; commit fan-out (§8.2) runs in-DO
(no LISTEN/NOTIFY needed — writes and sockets are co-located). Full shape,
wiring, hibernation semantics, and the
manual real-workerd smoke recipe in @syncular/server-workers/README.md.
No relay (decision). There is deliberately no relay — no bridge that forwards realtime from a self-hosted server to a managed realtime service. Realtime is the second binding of the same handler (§8.7, Direction decision 1 — the WS-native loop), so any host that runs the core serves realtime directly; multi-instance fanout is covered by LISTEN/NOTIFY on Postgres (below), and the Workers case is covered by the DO design (writes and sockets co-located per partition). Every job a relay would do is done by a binding of the core or by in-database fanout — a relay would add a hop, a second protocol surface, and a managed dependency for zero capability the core lacks.
Startup schema readiness
Fail startup before accepting traffic when the generated schema cannot compile or the storage projection cannot migrate:
import {
ensureSyncServerReady,
type SyncServerConfig,
} from '@syncular/server';
import { SqliteServerStorage } from '@syncular/server/sqlite';
const config: SyncServerConfig = {
schema,
storage: new SqliteServerStorage('./sync.db'),
segments,
resolveScopes,
};
await ensureSyncServerReady(config);
Bun.serve({ fetch: app.fetch });@syncular/server/sqlite selects bun:sqlite on Bun and the built-in
node:sqlite module on Node. It covers server storage, segment storage, blob
storage, leases, and SQLite-image generation without an external SQLite
package. The runtime-specific database wrappers are BunSqliteDatabase and
NodeSqliteDatabase when a host needs direct access to the native handle.
The helper accepts the generated ServerSchema, compiles it, and calls the
storage backend's low-level ensureSchema(CompiledSchema). A failure is a
SyncServerReadinessError with stable code sync.schema_not_ready, a phase
of schema_compile or storage_migration, and the generated schema version.
Log the cause for operators and stop startup. Do not catch readiness errors in
authentication or convert them to a 401; request-time schema checks are only a
defensive fallback.
For D1, finish storage.migrateSchema(compileSchema(schema)) across separate
Worker invocations before admitting traffic. Each call returns complete and
statementsExecuted; the default budget is 50 statements. See
D1 schema migration.
After restoring an authoritative database, keep traffic stopped and call
rotatePartitionLogEpoch({ storage, partition }) for every restored
partition. The rotation clears stale client cursors and requires version 2
clients to reset their server-derived rows while preserving the outbox. Follow
the complete backup and restore runbook.
Write validators and recovery metadata
validators is the server-authoritative seam for row business rules that
scope grants cannot express. A validator runs after row decode and scope
authorization, inside the commit transaction, for HTTP and WebSocket sync
rounds alike. Throw ValidationRejection for a deliberate host rejection:
import {
ValidationRejection,
type SyncServerConfig,
} from '@syncular/server';
const config: SyncServerConfig = {
schema,
storage,
segments,
resolveScopes,
validators: {
surgeries: ({ row }) => {
if (typeof row?.duration_minutes === 'number' && row.duration_minutes < 5) {
throw new ValidationRejection(
'surgery.duration_too_short',
'diagnostic only',
{
fieldPaths: ['duration_minutes'],
reason: 'below_minimum',
requiredAction: 'edit_fields',
references: { minimum_minutes: '5' },
},
);
}
},
},
};The third argument is optional. When supplied, Syncular validates and
normalizes a bounded RejectionDetails object and persists it with the
idempotency result. Its values replicate to the authorized client, so include
only non-sensitive identifiers that the host explicitly approves for recovery
UI. Unknown members, free-form tokens, malformed paths, and over-limit data
fail at construction. Diagnostic prose stays in message; apps should map
the stable code/details to localized copy instead of displaying that message.
Validators are per-operation. A validator must not mutate the row it receives.
A validator that decides from other rows (a membership, a role, a grant) calls
context.queryAuthoritative with the request storage.queryAuthoritative
takes, built from a generated query descriptor. The storage binds every
relation to the commit's partition and runs the statement on the push
transaction's connection, so the read sees the operations staged before this
one and needs no second pool connection. Never call
storage.queryAuthoritative from a validator: on SQLite and PGlite it waits
for the push transaction and the push never completes, and on a pool it reads
committed state. The whole-commit reader, the reaction planner reader, and the
remote command context expose the same method. See
Read a generated query inside a push transaction.
For a multi-row or multi-table invariant, install commitValidator. It runs
once after every operation is staged in the same transaction, and can inspect
both the decoded sibling operations and final candidate state:
import {
CommitValidationRejection,
type SyncServerConfig,
} from '@syncular/server';
const config: SyncServerConfig = {
schema,
storage,
segments,
resolveScopes,
commitValidator: ({ operations }) => {
const transition = operations.find(
(operation) =>
operation.table === 'surgeries' &&
operation.row !== undefined &&
operation.stored !== undefined &&
operation.row.status !== operation.stored.status,
);
if (transition === undefined) return;
const hasEvent = operations.some(
(operation) =>
operation.table === 'surgery_status_events' &&
operation.row?.surgery_id === transition.rowId &&
operation.row?.status === transition.row?.status,
);
if (!hasEvent) {
throw new CommitValidationRejection(
transition.opIndex,
'surgery.status_event_required',
'diagnostic only',
{
fieldPaths: ['status'],
reason: 'missing_sibling_operation',
requiredAction: 'repair_aggregate',
},
);
}
},
};The callback also receives read.getRow(), bounded read.scanRows(), and
read.queryAuthoritative() APIs. Those reads see the final candidate state, including staged sibling upserts and
deletes. Syncular serializes the partition before any operation read so two
aggregate validators cannot both accept mutually invalid candidates. It also
re-checks idempotency after taking that lock and persists a rejected outcome
while the lock is retained, so overlapping duplicate deliveries cannot rerun
the callback.
SQLite and PostgreSQL provide that serialization directly for every push. D1
does not expose an interactive lock: D1ServerStorage fails closed unless it
is constructed inside an explicit per-partition coordinator with
{ pushApplySerialized: true } (normally the packaged Durable Object FIFO).
Do not set that assertion on a stateless Worker. Custom storages must implement
the pre-operation partition lock, locked idempotency re-check, atomic rejection
finalization, and—when commitValidator is used—candidate scans.
Whole-commit validation checks a client-proposed commit; it does not grant authority. Privileged operations such as connecting facilities still belong in explicit server-authoritative commands.
Seed idempotency and safe revisioning
seedMutations uses the real push path and a stable clientId/commitId.
Both applied and rejected outcomes are terminal for that key. A rejected call
throws SeedMutationError, whose structured code, opIndex, replayed,
recordedAtMs, and cacheIdentity fields distinguish a fresh policy failure
from replay of an older cached rejection.
After correcting a development seed definition, advance an explicit seed
revision (catalog-v1 to catalog-v2) and rerun it; do not delete the database
or unrelated rows. This does not apply to application commands. After an
unknown command outcome, reuse the original idempotency key because changing it
can execute the operation twice. The full inspection and recovery recipe is in
the public server guide.
One extra rule applies when the seed actor changes. A Syncular clientId is
bound to its first actor inside a partition. Keep the client ID stable for seed
revisions under that actor, but allocate both a new purpose-specific clientId
and a new commitId when moving a seed to another actor. Changing only the
commit ID correctly fails with sync.invalid_client_id; do not erase the
database to bypass that identity evidence. For example, move
catalog-import/seed-user/catalog-v1 to
catalog-import/server-authority/catalog-v2, not merely to
catalog-import/seed-user/catalog-v2.
The task-oriented concurrency and conflict-correction guide shows version projection, aggregate rollback, corrected replacement commits, explicit acknowledgement, and restart-safe recovery UI together.
Trusted relational-index lookups for authoritative commands
scanRows is a Syncular scope-index scan, never an unscoped administrative
query. Passing an empty or omitted scopeFilter throws the exported
StorageQueryError with code: 'sync.storage.scan_requires_scope' on
SQLite, PostgreSQL, and D1.
An authoritative command that needs an exact alternate lookup can instead use
the optional storage.scanRowsByIndex(partition, query) or transactional
tx.scanRowsByIndex(query) capability. The query names one declared
TableSchema.indexes entry, supplies one exact value per index column, uses an
exclusive afterRowId, and has a required limit from 1 through 1,000. All
shipped adapters implement it; transaction reads see staged writes and deletes.
It requires a materialized table.
values is always a complete, order-sensitive index tuple. A compound index
on (workspace_id, state, id) therefore requires exactly three values; passing
only [workspaceId] is rejected with
sync.storage.index_value_count_mismatch and is never interpreted as a
leading-prefix scan. Declare a separate (workspace_id) index when the command
must enumerate every row in one Workspace. Syncular does not currently expose
trusted prefix or range scans.
This is a trusted @syncular/server storage capability, not SSP2: it creates no
scope variable, named-query obligation, subscription descriptor, or client
authority. Never expose table/index/value selection through a client-controlled
route. Custom adapters may omit the additive method; authoritative commands
must check for it and fail closed. See the public
storage lookup guide
for a user-scoped key-grant table revoked through a Workspace index and for the
atomic reverse-index/queue fallback required by ordered or derived lookups.
Durable server reactions
reactionPlanner turns an accepted candidate commit into bounded work records.
It runs after operation and whole-commit validation, inside the authoritative
push transaction. It may use the candidate-state reader and must perform no
external side effects. A rejected or replayed commit does not run it.
import {
ReactionRunner,
type ReactionPlanner,
type SyncServerConfig,
} from '@syncular/server';
type AppReactions = {
'invoice.email': { invoiceId: string };
};
const reactionPlanner: ReactionPlanner<AppReactions> = ({ operations }) =>
operations.flatMap((operation) =>
operation.table === 'invoice_events' &&
operation.row?.kind === 'invoice_finalized' &&
typeof operation.row.invoice_id === 'string'
? [{
key: `invoice:${operation.row.invoice_id}`,
type: 'invoice.email',
version: 1,
payload: { invoiceId: operation.row.invoice_id },
maxAttempts: 8,
}]
: [],
);
const config: SyncServerConfig = {
schema, storage, segments, resolveScopes, reactionPlanner,
};App rows, commit metadata, reaction rows, and the push idempotency result land
in one transaction. The handler idempotency key is derived from the source
partition, clientId, clientCommitId, and planner key. Reaction rows use
their own partition-scoped table and survive pruneCommitLog.
Drive delivery from a host scheduler or queue wake:
const runner = new ReactionRunner<AppReactions>({
storage,
partition: 'main',
workerId: 'invoice-worker-1',
handlers: {
'invoice.email': async ({ payload, idempotencyKey, extendLease }) => {
await extendLease();
await emailProvider.send({
invoiceId: payload.invoiceId,
idempotencyKey,
});
},
},
});
await runner.runOnce();Claims and acknowledgements compare a lease owner. Expired leases can be
claimed by another worker, and long handlers can call extendLease(). Ordinary
throws and RetryableReactionError retry with bounded exponential backoff.
PermanentReactionError and exhausted retry limits enter dead-letter.
retryDeadLetterReaction resets one row for an explicit operator retry.
Each runOnce() uses a fresh lease token and rechecks ownership before
starting every handler in a claimed batch.
Schedule terminal retention separately from commit-log pruning:
import { pruneReactions } from '@syncular/server';
let result;
do {
result = await pruneReactions({
storage,
partition: 'main',
nowMs: Date.now(),
events,
});
} while (result.mayHaveMore);Defaults retain completed rows for 30 days, dead-lettered rows for 90 days,
and remove at most 1,000 rows per pass. Override them with
retention: { completedRetentionMs, deadLetterRetentionMs, batchSize }.
Only terminal rows strictly older than their cutoff are eligible. Pending and
leased work is preserved, including expired leases. Cleanup and manual retry
serialize at storage, so one transition wins a race. Each pass emits
reaction.prune_completed with both cutoffs, both removal counts, the limit,
and mayHaveMore.
Delivery is at least once. A crash after the handler's external call and before
acknowledgement can run that call again. Handlers receive the same stable
idempotencyKey on every attempt and should pass it to external providers.
Syncular does not claim exactly-once external effects.
Planned payloads are plain JSON, versioned, limited to 64 KiB and 16 levels, with at most 100 reactions per commit. Failure details are plain JSON limited to 8 KiB. The planner API cannot enforce purity in JavaScript; running an external effect from the planner violates the transaction contract.
SQLite and PostgreSQL use their existing push transactions. D1 appends the
reaction writes to the same atomic batch as the source commit and retains its
mandatory per-partition Durable Object coordination for pushes. D1 claims use
one UPDATE ... RETURNING statement, so there is no claim read/write gap.
Structured events (the ops seam)
One optional interface, SyncularServerEvents, carries every
operator-relevant signal as a typed, JSON-able, stable-shaped event:
import { consoleJsonEvents, type SyncServerConfig } from '@syncular/server';
const config: SyncServerConfig = {
schema, storage, segments, resolveScopes,
events: consoleJsonEvents(), // one JSON line per event on stdout
};There is no logger dependency and no formatting — emission only. The same
shapes feed one-line JSON logs, metrics counters, and error trackers; a
Sentry adapter is a ~20-line emit implementation over this seam. The
events sink rides on SyncServerConfig, so the Hono adapter (and any
other adapter that spreads the config into the request context) passes it
through with no extra wiring. The realtime hub and pruneCommitLog take
the same sink via their own config/options (they run outside the request
context). The demo server wires it behind SYNCULAR_DEMO_EVENTS=1.
Guarantees
- Never throws through. Emission is fire-and-forget: a throwing
emitis swallowed at the seam and cannot affect request processing, realtime delivery, or pruning. (Tested.) - Zero cost when off. With no sink configured, no event object is ever built — every call site checks the sink before constructing the event. The benches run with events unset.
- Stable, JSON-able shapes. Flat objects, no
undefinedvalues, no classes;JSON.stringifyround-trips every event. Shapes andtypestrings are append-only surface. - Virtual-clock clean. All timestamps and durations come from the ctx
clock (
clockon the config / hub;nowMsfor prune), so conformance and tests under a virtual clock stay deterministic. Wall clock is never read behind the host's back.
Event catalog
| Event | When | Key fields |
| --- | --- | --- |
| request.handled | Once per POST /sync, after the response bytes are fully produced (or the request was rejected up front) | kind (sync), partition, actorId, durationMs, bytesIn, bytesOut, outcome (ok | schema_floor | rejected | error), errorCode?, pushCommits, pulled, subscriptions |
| push.applied | A PUSH_COMMIT applied, or replayed from the idempotency cache (§2.3) | clientId, clientCommitId, operations, commitSeq?, replay |
| push.rejected | A commit rejected (§6.3) | clientId, clientCommitId, operations, code (§10.2), opIndex |
| push.conflicted | A commit terminated by a version conflict (§6.2) | clientId, clientCommitId, operations, opIndex |
| reaction.queued | A planned reaction committed with its source push | clientId, clientCommitId, commitSeq, idempotencyKey, reactionType, version |
| reaction.started | A worker claimed and began one attempt | workerId, idempotencyKey, reactionType, version, attempt |
| reaction.retried | A retryable attempt failed and was rescheduled | started fields plus nextAttemptAtMs, errorCode |
| reaction.completed | A handler finished and its owner acknowledged | started fields |
| reaction.dead_lettered | A permanent or exhausted failure was recorded | started fields plus errorCode |
| reaction.prune_completed | One bounded terminal-reaction retention pass finished | completedBeforeMs, deadLetterBeforeMs, limit, removedCompleted, removedDeadLetter, mayHaveMore |
| pull.served | Once per served pull half, after all sections streamed | clientId, subscriptions[]: {id, table, status, mode (bootstrap | incremental | none), fromCursor, nextCursor, commits, changes, segments[]}; each segment: {mediaType (rows | sqlite), delivery (inline | ref), origin (built | reused), bytes, rows} |
| segment.downloaded | Every direct segment download (§5.5), success or failure | segmentId, outcome (ok | error), errorCode?, mediaType?, bytes?, durationMs |
| blob.swept | Every sweepOrphanBlobs pass (§5.9.2 orphan GC) | partition, swept (deleted count), referenced (keep-set size), graceMs |
| realtime.opened | A socket registered with the hub and got hello (§8.1) | sessionId, clientId, registrations, cursor, latestSeq |
| realtime.closed | A session left the hub (once per session) | sessionId, durationMs |
| realtime.delta | A delta message pushed over the socket (§8.2) | sessionId, commitSeq, bytes, changes |
| realtime.wake | A sync wake-up sent (§8.3) | sessionId, reason (catchup-required | delta-too-large | reset-required) |
| prune.completed | Every pruneCommitLog pass, moved or not | partition, previousHorizonSeq, horizonSeq, advanced, removedCommits |
| scopes.resolve_failed | The host resolveScopes callback threw — the §3.2/§3.4 fail-loud path | phase (request | realtime | segment-download), message |
All events also carry type, atMs, and (where a request identity
exists) partition / actorId.
Admin / console surface (SyncularAdmin)
The operator-facing read surface over the server core lives in this package and adds zero wire protocol. Authorization for these reads is entirely the host's. It delivers the 80% operator value (who's connected, what's flowing, horizon health, the event tail) as a handful of read-only, partition-scoped, JSON-able queries.
The event ring (the "event stream")
RingBufferEvents is a SyncularServerEvents sink that retains the last N
events in memory (bounded — oldest dropped when full) with a
query({type?, sinceMs?, clientId?, actorId?, limit}) and a
subscribe(listener) hook for live tails (the SSE route below). It is the
event stream without any infrastructure dependency. Compose it with any
other sink so the console tail and your logs/metrics see the same
emissions:
import {
RingBufferEvents, composeEvents, consoleJsonEvents, SyncularAdmin,
} from '@syncular/server';
const ring = new RingBufferEvents({ capacity: 1000 });
const config: SyncServerConfig = {
schema, storage, segments, resolveScopes,
events: composeEvents(ring, consoleJsonEvents()), // both see every event
};
const admin = SyncularAdmin.fromConfig(config, { ring });Query surface
Every method is read-only. All are partition-scoped except the two fleet
reads (listPartitions / partitionsOverview), which enumerate partitions
by design:
| Method | Returns |
| --- | --- |
| listClients(partition) | Known clients: clientId, actorId, cursor, lag (commits not yet pulled: maxCommitSeq − max(cursor, 0)), updatedAtMs, subscriptions[], and an active flag (cursor touched within the §4.6 active window). |
| clientDetail(partition, clientId, {eventLimit?}) | One client's drill-down: {exists, client?, lease?, events} — the record (with lag), its §7.3 lease when a lease store is wired, and its slice of the event tail. Answers "why is this client stale" in one read. |
| listCommits(partition, {afterSeq?, limit?, table?}) | Commit-log metadata (never payloads), newest first: commitSeq, clientId, clientCommitId, actorId, createdAtMs, changeCount, tables[]. |
| listReactions(partition, {statuses?, types?, limit?}) | Durable reaction lifecycle rows, newest first, including source commit, attempts, lease, completion, and bounded failure information. |
| inspectRow(partition, table, rowId) | {exists, serverVersion?, scopes?} — current row version + stored scopes, payload not decoded. |
| scopeActivity(partition, {variable, value}, {limit?}) | Recent commits touching one scope key, via the §3.1 change-scope index (never a log scan). |
| horizonStatus(partition) | {maxCommitSeq, horizonSeq, retainedCommits, activeCursorFloor, recommendedHorizonSeq, recommendation} — the horizon a prune pass would reach now (§4.6) + a coarse up-to-date / prune-recommended. |
| segmentStats() / blobStats(partition) / stats(partition) | Counts/bytes where the stores expose them (segments split rows/sqlite). undefined when a store omits stats(). |
| metrics(partition, {windowMs?, buckets?}) | Ring-derived request/push health over a trailing window (default 5 min): request count/rate, error share, p50/p95 duration, push applied/rejected/conflicted, and per-bucket counts for a sparkline. Zero new server state; all zeros when no ring is wired. |
| listPartitionRegistry() / partitionsOverview() | Every authenticated partition; the registry carries log continuity and last-authenticated time, while the overview adds retained backlog, client counts, and the prune recommendation. |
| events({type?, sinceMs?, clientId?, actorId?, limit?}) | The ring tail, newest first. Empty when no ring is wired (hasEventStream reports which). |
| subscribeEvents(listener) | Live events as they land in the ring; returns the unsubscribe function, undefined when no ring is wired. |
The query surface leans on additive, optional storage/store methods
(ServerStorage.listClientRecords / listCommitMetadata / scopeActivity /
getRowScopes / listPartitions; SegmentStore.stats; BlobStore.stats)
— the established
optional-method pattern. SqliteServerStorage, PostgresServerStorage,
D1ServerStorage, and the memory/sqlite stores implement them; the shared
ServerStorage contract suite exercises them on all backends. A backend that
omits one makes the corresponding admin read fail loud (it never returns a
silently-empty console). The S3SegmentStore does report stats() — from
a LIST-free pointer-object accumulator (see "S3 stats" below) — but its
counters are marked approximate: true, an additive field the admin surface
carries through so the console can label them honestly. The exact in-process
stores (memory/sqlite) omit the marker.
HTTP routes + the single console page
@syncular/server-hono exports createSyncularAdminRoutes(admin, opts),
a mountable Hono sub-app. The auth seam is required: the factory throws
if you omit the authorize guard — there is no default-open admin. Every
endpoint (including the page) runs the guard first; a falsy result is a 401.
import { createSyncularAdminRoutes } from '@syncular/server-hono';
const routes = createSyncularAdminRoutes(admin, {
defaultPartition: 'main',
authorize: ({ request }) => isOperator(request), // YOUR check — mandatory
});
app.route('/admin', routes);| Route | Mirrors |
| --- | --- |
| GET / | The console page (see below). |
| GET /clients | listClients |
| GET /clients/:clientId?eventLimit | clientDetail |
| GET /commits?afterSeq&limit&table | listCommits |
| GET /reactions?status&type&limit | listReactions |
| GET /rows/:table/:rowId | inspectRow |
| GET /scope-activity?variable&value&limit | scopeActivity |
| GET /horizon | horizonStatus |
| GET /stats | stats |
| GET /metrics?windowMs | metrics |
| GET /partitions | partitionsOverview (fleet view) |
| GET /events?type&sinceMs&clientId&actorId&limit | events (ring tail) |
| GET /events/stream?type&clientId&actorId&limit | live SSE tail (backlog replay, then subscribeEvents) |
?partition= selects the partition (falls back to defaultPartition).
GET /partitions is the one cross-partition endpoint: the guard still runs
on it, with an empty partition in its context unless the query passes one
— a host that authorizes per partition should treat it accordingly.
GET / (or /admin) serves a single static HTML page — zero
framework, no build step, no React. It fetches the sibling JSON endpoints
(relative to its own mount path, so it works under any prefix and the same
guard covers its XHRs) and renders: a metrics statusbar (push rate,
conflicts, error share, p95, ASCII sparkline), the fleet view (which doubles
as the partition picker), horizon, store stats, clients with cursor lag and
a per-client drill-down (subscriptions with full scope sets, lease status,
the client's event slice), recent commits, the row inspector, scope
activity, and the event tail — with an auto-refresh toggle (2 s poll).
This is a deliberate one-file console: 5% of the code a full console app
would cost, the 80% operator value.
Live tail over SSE. GET /events/stream streams the ring as
server-sent events: a recent backlog replays first, then events arrive as
they land (a subscribe hook on the ring), with a comment ping every 5 s
(under Bun.serve's 10 s default idle timeout, so a quiet stream survives).
It is plain web-streams, so it serves identically on Bun, Node, and
Workers. The page upgrades its event panel to the stream automatically and
falls back to the 2-second poll when the stream is unavailable (no ring, or
EventSource cannot pass the host's auth headers — cookie-authorized
setups stream fine). The tail can carry sensitive identifiers (actorIds,
session ids), the same as GET /events; the admin guard gates both.
The demo server (apps/demo) mounts the admin behind a dev guard:
SYNCULAR_DEMO_ADMIN=1 enables /admin (optionally token-gated with
SYNCULAR_DEMO_ADMIN_TOKEN), so the console is inspectable live.
Docs-site coverage of the console is a follow-up: the docs app is owned by a concurrent workstream this round (the schema-bump page), so this README is the console's documentation home for now.
Segment storage on S3 / R2 (S3SegmentStore)
Three SegmentStore backends ship in-tree and pass one shared contract
suite (test/segment-store-contract.ts): MemorySegmentStore (tests,
single process), SqliteSegmentStore (single node), and
S3SegmentStore — the production backend for any S3-compatible object
store (AWS S3, Cloudflare R2, MinIO). It is dependency-free: SigV4 is
hand-rolled over fetch (sigv4.ts, pinned by the published AWS test
vectors).
import { S3SegmentStore, s3PresignedUrls } from '@syncular/server';
const segments = new S3SegmentStore({
endpoint: 'https://s3.eu-central-1.amazonaws.com', // origin only, no bucket
region: 'eu-central-1',
bucket: 'my-app-segments',
accessKeyId: process.env.AWS_ACCESS_KEY_ID!,
secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY!,
keyPrefix: 'syncular/', // optional namespace inside the bucket
ttlMs: 24 * 60 * 60 * 1000, // §5.1 default
});
const config: SyncServerConfig = {
schema, storage, resolveScopes,
segments,
signedUrls: s3PresignedUrls(segments, { ttlSeconds: 900 }), // §5.4 delegated presign
};Cloudflare R2 specifics. The endpoint is your account's S3 API host
and the region is always auto:
const segments = new S3SegmentStore({
endpoint: 'https://<account-id>.r2.cloudflarestorage.com',
region: 'auto',
bucket: 'my-app-segments',
accessKeyId: R2_ACCESS_KEY_ID, // R2 API token pair
secretAccessKey: R2_SECRET_ACCESS_KEY,
});MinIO works the same way (endpoint: 'http://127.0.0.1:9000', any
region string). Requests are path-style ({endpoint}/{bucket}/{key}),
which all three providers accept.
Key layout. Deterministic, so every lookup is a GET/HEAD — never a LIST:
{keyPrefix}seg/sha256/{hex}— the segment bytes, verbatim and immutable (the object body is exactly the content-addressed bytes, so presigned GETs serve them directly and the client's §5.1 hash check passes).{keyPrefix}rec/sha256/{hex}.json— the mutableSegmentRecord. It carries thepublicationsarray: one entry per (partition, log epoch, table, schema version, media type, scope digest, pin, cursors) context the content was published under, each with its own TTL.getis two GETs (record + bytes). A record written by 0.21 read the metadata from the bytes object's user metadata;getstill reads that shape and materializes one publication per recorded digest.{keyPrefix}find/{sha256(reuse key)}.json— the §5.3 whole-table reuse pointer, written only forrowCursor: nullsegments;findis one GET plus a HEAD to confirm the segment object still exists, and it selects the publication the key names.
A host that implements SegmentStore itself must return publications
and scopeDigests on every SegmentRecord (scopeDigests is the
compatibility view over publications). Download and find select a
publication matching the caller's partition, live log epoch, and digest;
they never union digests across contexts. The shared contract suite is
test/segment-store-contract.ts.
TTL and lifecycle. Expiry is store-side and authoritative:
expiresAtMs (put time + ttlMs, default 24 h) is recorded with the
record; get returns expired records so the §5.5 endpoint can answer
the precise, retryable sync.segment_expired, and find filters them
itself. Bucket lifecycle expiration is garbage collection only — set
it comfortably above ttlMs (e.g. 2 days for the 24 h default) and
never below it. After lifecycle deletes an object, clients see
sync.not_found instead of sync.segment_expired; both recover by
re-pulling, but the former loses the "just re-pull, this is normal"
signal, so keep the GC margin generous.
S3 stats — a LIST-free, approximate accumulator
S3SegmentStore.stats() reports store-wide {count, bytes, rowsSegments,
sqliteSegments} for the admin console without a bucket LIST — a LIST
would defeat the store's whole every-lookup-is-a-GET/HEAD design and cost
real money at scale. Instead the store keeps a tiny counter object under a
fixed key ({keyPrefix}stats/segments.json) and folds each new segment into
it read-modify-write on put. A HEAD before the segment PUT detects an
idempotent re-put (same content address ⇒ same key), so each distinct
segment is counted once.
The accumulator write is guarded by an ETag compare-and-swap — If-Match
against the ETag we read (or If-None-Match: * to create) — and the store
retries on a 412 PreconditionFailed, so two writers folding concurrently do
not silently clobber each other's increment. AWS S3 and Cloudflare R2 both
honor these conditional headers.
Even so, the counters are approximate, and stats() marks them
approximate: true (an additive field the admin surface carries through, so
admin.segmentStats() / admin.stats() expose it and the console labels the
numbers honestly). They can drift: a crash between the segment PUT and the
accumulator CAS, lifecycle GC deleting objects the accumulator still counts,
or enough concurrent writers to exhaust the CAS retry budget. They are a
health gauge, not an invoice — reconcile against a periodic inventory report
if you need an exact number. The exact in-process stores (memory / sqlite)
count on demand and omit the marker.
The blob store's stats() uses the same accumulator + approximate: true
marker (see "Blob bytes on S3 / R2" below).
Blob bytes on S3 / R2 (S3BlobStore)
File-attachment bytes (§5.9) get the same object-store backend as segments.
S3BlobStore is the blob twin of S3SegmentStore — same hand-rolled SigV4,
same content-addressed key layout, same LIST-free-on-the-hot-path stats
accumulator — for AWS S3, Cloudflare R2, or MinIO. It closes the
"attachments are SQLite-only" gap: a Workers/edge or horizontally-scaled
deployment can now store blobs durably in an object store instead of the
database.
import {
S3BlobStore,
s3PresignedBlobUploads,
s3PresignedBlobUrls,
} from '@syncular/server';
const blobs = new S3BlobStore({
endpoint: 'https://s3.us-east-1.amazonaws.com',
region: 'us-east-1',
bucket: 'my-attachments',
accessKeyId: AWS_ACCESS_KEY_ID,
secretAccessKey: AWS_SECRET_ACCESS_KEY,
keyPrefix: 'syncular/', // optional namespace inside the bucket
});
const config: SyncServerConfig = {
// …schema, storage, segments, resolveScopes…
blobs,
// Presigned DOWNLOAD (always-issue): serve blob downloads as provider
// presigned GET URLs (§5.9.5). The client fetches bytes straight from the
// object store — the sync server exits the download egress path.
blobSignedUrls: s3PresignedBlobUrls(blobs, { ttlSeconds: 900 }),
// Presigned UPLOAD (direct-to-storage): mint single presigned PUT URLs so
// clients upload straight to the object store, bypassing the server upload
// bandwidth path (§5.9.3). Optional — absent ⇒ clients stream through the
// direct `PUT /blobs/{blobId}` endpoint (a capability, never a fallback).
blobUploadUrls: s3PresignedBlobUploads(blobs, { ttlSeconds: 900 }),
};Presigned blobs, end to end (§5.9.3 / §5.9.5)
Two independent presign switches let the sync server step out of the blob byte path in both directions:
Download — blobSignedUrls (always-issue). After the §5.9.5 row-derived
authorization check passes, GET /blobs/{blobId} returns
{ url, urlExpiresAtMs } and no bytes; the client fetches the URL directly
(no host auth — the URL is the entire grant) and re-verifies the content
address. Always-issue, not accept-bit negotiation (the pinned decision,
§5.9.5): unlike segments — where a descriptor rides the pull stream to a client
that may be unable to fetch a bare URL, so issuance is gated on accept bit 3 —
a blob download is a plain request/response, so the response simply carries the
URL and a client that cannot consume it re-requests. Always-issue is harmless
(the authorized endpoint is the same route; nothing is stuck on a stream) and
simpler (no accept-bit plumbing on a non-pull path). Recovery mirrors §5.4:
a failed or expired URL fetch re-requests the endpoint (which re-authorizes
and mints a fresh URL) — never a fall-through.
Upload — blobUploadUrls (the grant flow). POST
/blobs/{blobId}/upload-grant (host-authenticated, body
{ byteLength, mediaType? }) mints a single presigned PUT; the client PUTs
bytes straight to the object store. Upload authz is host-authentication only
(any authenticated actor may obtain a grant within the size cap) — uploading
bytes is not a scope-bearing act, because the content address discloses nothing
and integrity is enforced at reference time (the §5.9.6 push existence check
- every download's content-address verify), never by a store-side hash
recompute. The size cap is enforced up front against the declared
byteLength, before any URL is minted (the object-store hop cannot re-check the streamed byte count). An already-present blob returns{ present: true }(skip the PUT, idempotent §5.9.3). A single PUT only — never a multipart or chunk protocol; resumable upload, when it lands, is provider multipart behind this same grant. Absent config ⇒ the client streams through the direct host-authenticatedPUT /blobs/{blobId}endpoint — a capability choice, not a fallback (that endpoint was always the other path).
Cloudflare R2. Identical to the segment store — point endpoint at
https://<account-id>.r2.cloudflarestorage.com, region: 'auto', and use an
R2 API-token access-key pair. R2 honors the SigV4 header/query auth and the
ListObjectsV2 the sweep uses.
Key layout. {keyPrefix}blob/{partition}/sha256/{hex} — content-addressed
and partition-scoped (the same bytes uploaded under two partitions are two
objects; a partition cannot read another's attachment by guessing a content
address). The object body is the blob bytes verbatim, so a presigned GET
serves exactly the content-addressed bytes and the client's §5.9.1 hash check
passes. byteLength + optional mediaType + createdAtMs ride along as
object user metadata (x-amz-meta-syncular-blob = base64url(JSON)), so get
is a single GET.
Durability, not TTL — the difference from segments
This is the honest interface difference from S3SegmentStore. Segments are
TTL cache entries; blobs are durable. A blob referenced by a live row must
stay downloadable indefinitely (§5.9.5 B3). So S3BlobStore writes no
expiresAtMs, has no ttlMs config, and maps to no S3 lifecycle-expiration
rule. Reclamation is reference-driven, not time-driven: the only thing
that deletes a blob is the orphan sweep, and it deletes only blobs no live
row references.
Do NOT put an S3/R2 lifecycle-expiration rule on the
blob/prefix. It would delete still-referenced attachments out from under live rows. This is the exact opposite of the segment guidance (where lifecycle expiration abovettlMsis encouraged as GC). Blobs are cleaned by the sweep below.
Orphan sweep (the GC runbook)
Nothing reclaims blobs automatically — the host schedules the sweep, the blob
analogue of pruneCommitLog. sweepOrphanBlobs(storage, blobStore, partition,
{ graceMs }) reads the live keep-set from the §5.9.4 reference index
(storage.listReferencedBlobIds) and deletes every blob that is both
unreferenced and older than the grace period. It emits one blob.swept
ops event ({ swept, referenced, graceMs }) and returns the deleted ids.
import { sweepOrphanBlobs } from '@syncular/server';
// A periodic per-partition GC job (hourly to daily is sensible).
const { swept } = await sweepOrphanBlobs(storage, blobs, partition, {
graceMs: 24 * 60 * 60 * 1000, // default; see the race note below
events,
});The grace period is not optional cleverness — it is the correctness
mechanism. Uploads are content-addressed and land before the referencing
push (§5.9.2 upload-before-reference): a client PUTs the bytes, then pushes
the row. Between those two steps the blob is legitimately unreferenced. If
the sweep ran with no grace it would delete that fresh upload before its push
arrived. So the grace period must comfortably exceed any sane upload→push
latency — the default is 24 h, deliberately far above any push window.
Lower it only if you fully understand your clients' outbox latency; there is
no upside to a tight grace and a real data-loss risk. The sweep compares
against the blob's upload time (createdAtMs from object metadata), and an
idempotent re-upload does not reset that clock.
sweepOrphanBlobs requires storage.listReferencedBlobIds (the §5.9.4
reference index — SQLite, D1, and Postgres all implement it). Against a
storage without it, the helper throws rather than sweep with an empty keep-set
(which would delete everything). The S3BlobStore sweep is the store's only
LISTing operation — it pages ListObjectsV2 over the partition's blob/
prefix, an admin/GC path off the hot path; every point lookup is still a
GET/HEAD.
Blob stats — the same approximate accumulator
S3BlobStore.stats() reports store-wide {count, bytes} from a fixed counter
object ({keyPrefix}stats/blobs.json) folded read-modify-write on put under
the same ETag compare-and-swap as segments, and marks the result
approximate: true (the marker already present on BlobStoreStats, carried
through admin.blobStats()). A HEAD before the PUT detects an idempotent
re-upload so distinct blobs count once; the sweep decrements on delete. Same
caveats as segment stats — a health gauge, not an invoice. The in-process
stores (memory / sqlite) count on demand and omit the marker.
Presigned blob downloads (§5.9.5)
Setting blobSignedUrls: s3PresignedBlobUrls(blobs) makes the server issue a
provider-presigned GET URL for a blob download — but only after the
row-derived authorization check passes (handleBlobDownload resolves the
actor's scopes and tests the referencing rows first; a blobId is never a
bearer capability minted from the id alone). The signed object key embeds the
blobId, TTL SHOULD be ≤ 15 min (default 900 s), and the URL is a short-lived
grant to exactly those immutable bytes. The issued url/urlExpiresAtMs ride
additively on BlobDownloadResult alongside the bytes; the mounted §5.9.5
direct-download endpoint stays the default serving path. Client consumption of
the presigned URL (following it instead of streaming bytes through the sync
server) is a later rung — the server-side issuance ships now.
Native HMAC vs delegated presign (§5.4)
SyncServerConfig.signedUrls accepts either scheme; the pull emits
SEGMENT_REF.url/urlExpiresAtMs identically for both (issuance always
happens inside the pull, immediately after scope resolution), and
clients cannot tell them apart.
- Native HMAC (
SignedUrlConfig) — you serve the segment bytes yourself (or from something that delegates auth to you, e.g. a CDN worker callingverifySegmentTokenat the edge). Thesttoken binds segment + scope digest + partition audience. Choose this when segments live inSqliteSegmentStoreor when you want claim-level binding at your own edge. - Delegated presign (
DelegatedPresignConfig, vias3PresignedUrls(store)) — the object store enforces the grant; the sync server never proxies segment bytes (zero egress through it — the bootstrap-storm answer). The §5.4 equivalence rule holds by construction: the signed object key embeds exactly onesegmentId, and the expiry obeys the same ≤ 15 min TTL guidance (default 900 s for both schemes).
Either way, keep the §5.5 direct-download endpoint mounted: it is the mandatory fallback for expired/failed URLs and for clients that never advertised accept bit 3.
CDN in front
Segment URLs are safe to cache by content: the object key is the
content address (seg/sha256/{hex}), the bytes are immutable for a
given key, and the client verifies the hash after download (§5.1) — so
a CDN can cache segment objects keyed on the path alone and can never
serve wrong bytes, only stale-but-correct ones. Two rules:
- Strip the query from the cache key, never from the auth check.
Presigned query parameters (or the native
sttoken) differ per client; the path is the content address. Configure the CDN to cache on the path while still forwarding the query for origin authorization (or validate at the edge:verifySegmentTokenfor native tokens). Never cache the authorization decision. - Align the CDN TTL with the store TTL. Cache lifetime at or below
ttlMskeeps the CDN from serving objects the store already declared expired (harmless — the client would still verify and apply — but it masks the §5.1 cache-entry semantics and can hide lifecycle GC). Content-addressing makes over-caching safe, not useful.
The §5.5 endpoint responses stay Cache-Control: private, max-age=0
— only segment-object URLs are CDN-cacheable, never the re-authorized
download path.
Horizon & pruning: operational guidance
The commit log grows forever unless you prune it. pruneCommitLog
(SPEC §4.6) advances the per-partition horizonSeq and deletes commits
at or below it. Nothing prunes automatically — the host schedules it.
When to run. A periodic job per partition — hourly to daily is the
sensible range; there is no benefit below the granularity of your
activeWindowMs. Pass events to get prune.completed per pass.
Pruning verifies the captured log epoch and updates the horizon together with
commit/change/scope deletion in one transaction. Concurrent passes cannot
lower the horizon. A retry cleans up eligible records even when the horizon
already covers them. A restore invalidates a pending pass with
sync.storage.prune_epoch_mismatch; recompute its retention inputs.
Unregistered partitions cannot be pruned. D1 maintenance enters the owning
Durable Object's existing write queue through
SyncularRealtimeHost.pruneCommitLog.
Custom storage adapters must add getPartitionLogEpoch(partition) and replace
pruneCommitsThrough(partition, seq) with
pruneCommitsThrough(partition, { logEpoch, throughSeq }), returning
{ previousHorizonSeq, horizonSeq, removedCommits } from the transaction.
setHorizonSeq remains monotonic and is no longer used by the pruning helper.
The retention floors (§4.6, encoded in RetentionPolicy). The
horizon never advances past min(cursor) of active clients — clients
whose cursor record was touched within activeWindowMs (default 14
days). Two escape hatches keep laggards from pinning the log forever:
commits older than ageForceMs (default 30 days) may be pruned
regardless, and at least the newest minRetainedCommits (default 1000)
commits are always kept. Raise the defaults freely, lower them with care.
What sync.cursor_expired means operationally. A client whose
cursor fell behind the horizon gets SUB_START.status = reset and
re-bootstraps from scratch (§4.7). That is correct behavior, not an
error — but its rate is your pruning health signal. A steady trickle
means devices returning from >30-day absences (expected). A spike means
you pruned faster than your fleet syncs: ageForceMs or
activeWindowMs is too tight for real usage, and you are paying for it
in bootstrap load (full re-scans + segment builds), not just in resets.
Observe it via pull.served subscriptions with status: "reset".
Segment TTL interplay. Segments are cache entries, not durable state
(§5.1; default TTL 24 h). A bootstrap that resumes past segment expiry
answers sync.segment_expired and the client re-pulls for fresh
descriptors — again correct, again a cost signal. Keep the segment TTL
comfortably longer than the slowest plausible bootstrap (a multi-page
bootstrap must finish while its segments live), and note that pruning
and segment expiry compound: a reset storm triggers a bootstrap storm,
which the §5.3 image-reuse rule absorbs only while images stay
unexpired. If you see origin: "built" dominating "reused" for the
same table+scope during a storm, your TTL is shorter than the storm.
What to alert on.
push.rejectedrate, bycode— a risingsync.forbiddenshare usually means an authorization regression, not misbehaving clients. (push.conflictedis normal offline-first traffic; alert only on gross shifts.)scopes.resolve_failed— any nonzero rate. This is the fail-loud path: every occurrence revokes subscriptions or rejects writes for a real request, and it is almost always a host bug or a dead dependency of the resolver.request.handledwithoutcome: "error"anderrorCode: "internal"— storage failures surfacing mid-stream.- Reset rate (
pull.served→status: "reset") — see above; alert on spikes relative to fleet size. - Prune backlog:
prune.completedwithadvanced: falsefor many consecutive passes while the log grows means one laggard cursor inside the active window is pinning retention — inspectlistClientCursorsfor the offender; the §4.6 floors bound the damage toageForceMs. realtime.wakewithreason: "delta-too-large"— sustained occurrences mean commits routinely exceedmaxDeltaBytesand clients are falling back to HTTP pulls; raise the limit or shrink commits.
Postgres storage (the production database path)
SqliteServerStorage from @syncular/server/sqlite uses bun:sqlite or
Node's built-in node:sqlite. For
production, PostgresServerStorage implements the same ServerStorage
contract against Postgres, with the inverted scope index carried through
as covering indexes so scope fanout is an index range scan, never a
scan-before-LIMIT. The
schema (POSTGRES_DDL) and its index design live in
src/postgres-storage.ts; storage.migrate() applies it idempotently
(every DDL is CREATE … IF NOT EXISTS, run statement-by-statement, so
migrate() is safe to call on every boot).
Blobs (§5.9.4) on Postgres. PostgresServerStorage implements the
optional blob-reference index — setBlobRefs (in the commit transaction)
plus listRowsReferencingBlob / listReferencedBlobIds — at full parity
with the SQLite and D1 storages. The sync_blob_refs table keys
(partition, tbl, row_id, blob_id) and carries a secondary
(partition, blob_id) index that drives the §5.9.5 download-authorization
candidate set as an index range (asserted in postgres-explain.test.ts, same
no-Seq Scan doctrine as the scope indexes). So a Bun/Node or Workers
deployment on Postgres supports file attachments end-to-end: push writes the
row's references atomically with the commit, and the blob-download handler
authorizes via the reference index. The shared ServerStorage contract runs
its blob section on pglite alongside sqlite and D1.
The PgExecutor seam (zero runtime deps)
The server library never imports a Postgres driver. PostgresServerStorage
is written against the minimal PgExecutor interface (query(text, params)
plus a transaction(fn) scope) — you wire your driver of choice:
Bun.sql (built into bun):
import {
PostgresServerStorage,
type PgExecutor,
type PgQueryable,
} from '@syncular/server';
function bunSqlExecutor(sql: import('bun').SQL): PgExecutor {
const over = (h: any): PgQueryable => ({
async query(text, params) {
const rows = await h.unsafe(text, params ? [...params] : []);
return { rows, rowCount: rows.length };
},
});
return {
query: over(sql).query,
transaction: (fn) => sql.begin((tx: any) => fn(over(tx))),
close: () => sql.end(),
};
}
const storage = new PostgresServerStorage(
bunSqlExecutor(new Bun.SQL(process.env.DATABASE_URL!)),
);
await storage.migrate();node-postgres (pg) — adapt a Pool:
import { Pool, type PoolClient } from 'pg';
import { PostgresServerStorage, type PgExecutor } from '@syncular/server';
function pgPoolExecutor(pool: Pool): PgExecutor {
const over = (c: Pool | PoolClient) => ({
query: (text: string, params?: readonly unknown[]) =>
c.query(text, params ? [...params] : []),
});
return {
query: over(pool).query,
async transaction(fn) {
const client = await pool.connect();
try {
await client.query('BEGIN');
const result = await fn(over(client));
await client.query('COMMIT');
return result;
} catch (error) {
await client.query('ROLLBACK');
throw error;
} finally {
client.release();
}
},
close: () => pool.end(),
};
}Type-parser note. commit_seq/server_version are int8. Drivers
decode int8 differently (node-postgres → string, Bun.sql → bigint,
pglite → number); the storage layer coerces every sequence read through
Number(...), so no driver-specific type-parser config is required.
bytea must decode to Uint8Array/Buffer (all three do).
Tests wire @electric-sql/pglite (embedded WASM Postgres, a
devDependency — hermetic, no docker) via pgliteExecutor from
@syncular/server/pglite. Both backends run the shared
ServerStorage contract (test/storage-contract.ts), and
test/postgres-explain.test.ts asserts via EXPLAIN that the fanout
candidate scans are index-driven so the scan-before-LIMIT regression
cannot silently return.
commitSeq allocation under concurrency
Per-partition commitSeq is dense and gap-free (§2.1). appendCommit
allocates it with an upsert that increments sync_partitions.max_commit_seq
and returns the allocated value. A common table expression feeds that value
into the commit metadata insert in the same SQL statement. The upsert takes a
row-level write lock on the partition for the transaction duration. Concurrent
pushes to that partition serialize; separate partitions use separate locks. A Postgres SEQUENCE is deliberately not used: it would leave
gaps on rollback, which the §4.5 pull-window arithmetic does not tolerate.
Each change insert expands its scope object into the inverted scope entries in the same statement. Serialized scopes bind as text before JSONB parsing; this avoids driver-specific JSON string encoding. Historical string-form scopes remain readable. Empty scopes still produce a change row, and repeated scopes within a commit retain one index entry.
Multi-instance fanout (LISTEN/NOTIFY)
Behind a load balancer, a commit applied on instance A fans out to A's
local realtime sessions in-memory, but a client whose socket lives on
instance B never sees it. PostgresFanout bridges the gap: after a commit
lands, the originating instance NOTIFYs syncular_commit with a
<partition>:<commitSeq> payload; every instance runs a listen() loop
that, on a notification, calls hub.wake(partition, 'catchup-required') —
remote sessions then pull the delta from the shared Postgres storage they
already read from (§8.3). NOTIFY payloads are capped (~8 KB) and are not an
ordered delta channel, so we wake rather than re-broadcast bytes; only
cross-instance delivery pays the re-pull. Single-instance deployments
install no fanout at all.
import { PostgresFanout, type PgNotificationConnection } from '@syncular/server';
// node-postgres: a dedicated Client for LISTEN + the pool for NOTIFY.
const conn: PgNotificationConnection = {
async listen(channel, handler) {
const client = await pool.connect(); // long-lived, NOT released
client.on('notification', (m) => m.payload && handler(m.payload));
await client.query(`LISTEN ${channel}`);
},
notify: (channel, payload) =>
pool.query('SELECT pg_notify($1, $2)', [channel, payload]).then(() => {}),
};
const fanout = new PostgresFanout(conn);
await fanout.install(hub); // start the LISTEN loop
// after a push commit lands:
await fanout.notifyCommit(partition, commitSeq);pglite is single-connection and cannot exercise cross-connection NOTIFY,
so the fanout integration test is env-gated on SYNCULAR_PG_URL (it wires
Bun.sql as a worked example) and skips cleanly; the payload encode/parse
and wake wiring are unit-tested hermetically.
Bench lane
bench has an env-gated Postgres lane measuring 100k bootstrap +
propagation on the production path. It runs only with SYNCULAR_PG_URL
set and is never part of bench:ci budgets (those stay on the
deterministic in-process sqlite loopback):
SYNCULAR_PG_URL=postgres://user:pass@localhost:5432/db bun run benchCustom storage adapters must implement
getActiveClientCursorFloor(partition, cutoffMs). Return the minimum cursor
whose updatedAtMs >= cutoffMs, or null when no client qualifies. Preserve
negative bootstrap cursors. Pruning and admin horizon status use this scalar
aggregate; listClientCursors remains the explicit listing interface.
SQLite image builders now return Promise<Uint8Array> and receive
rowBatches, an iterable or async iterable of row arrays. Replace custom
builders' input.rows loop with for await (const rows of input.rowBatches),
insert each batch into the dedicated image database, and count rows during
consumption. Write the final row count into _syncular_segment before
serialization. Await buildSqliteImage(input) when calling the built-in
Bun or Node builder directly.
The server shares in-flight builds for the same storage pair and artifact
identity after authorization. Sharing is local to one process. Signed URL
grants remain per request. The first eligibility probe has at most
limitSnapshotRows + 1 rows; subsequent builder batches have at most 5,000
rows. The image database and serialized output still consume memory.
Realtime acknowledgements call advanceClientCursor(partition, clientId,
actorId, logEpoch, cursor, updatedAtMs). Custom storage adapters must implement
this atomic update: advance the cursor and activity timestamp to their respective
maxima, preserve registration fields, and require a matching actor and current
partition log epoch. Leave missing records unchanged. SQLite, Postgres, and D1
perform one update without reading or serializing the subscription list. HTTP
registration keeps its existing cursor and subscription replacement rules.
Reviewed schema windows
Pass schemaWindow: [compileSchema(schema), compileSchema(previousSchema)] to
serve the prior codec through the current storage and validation rules. Entries
are newest first, with the compiled current schema first. Added tables and
appended nullable columns are supported. Removed/renamed columns, changed codecs,
primary keys, references or scopes reject the window before serving. Hosts must
exclude semantically incompatible prior versions. The schema upgrade guide
owns the pull, segment, push and realtime contract.
