Networked access

Subscription resumption

Resume a WebSocket change feed after a disconnect with sequence cursors and epochs, and recover cleanly when the server forces a resync.

Table of Contents

Every change event a subscription delivers carries a sequence number. A client that reconnects can hand the server the last one it processed and get back exactly the changes it missed. The client SDK resumes this way automatically whenever autoReconnect restores a dropped connection. This page drives the protocol by hand so that you can resume from any WebSocket client and handle the one case where replay is impossible. It builds one file, resume.ts, step by step: you subscribe, track the cursor through live changes, replay a gap after a disconnect, then trigger a forced resync.

Serve an orders table

Change subscriptions read from the server's change data capture log, which needs a file-backed database; the server rejects a subscription on an in-memory database with CDC_UNSUPPORTED. Open an orders table on disk and serve it. Start resume.ts here.

resume.ts
import { Sirannon } from '@delali/sirannon-db'
import { betterSqlite3 } from '@delali/sirannon-db/driver/better-sqlite3'
import { createServer } from '@delali/sirannon-db/server'
 
const driver = betterSqlite3()
const sirannon = new Sirannon({ driver })
const shop = await sirannon.open('shop', './data/shop.db')
 
await shop.execute(
  'CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY, customer TEXT NOT NULL, total INTEGER NOT NULL)'
)
 
const server = createServer(sirannon, { port: 9876 })
await server.listen()

Subscribe and store the cursor

Open a raw WebSocket to /db/shop and send a subscribe message. Two helpers make the socket awaitable: openSocket resolves once the connection is up, and listen queues incoming messages so that no message from a replay burst is dropped between reads.

resume.ts
import { Sirannon } from '@delali/sirannon-db'
import { betterSqlite3 } from '@delali/sirannon-db/driver/better-sqlite3'
import { createServer } from '@delali/sirannon-db/server'
 
const driver = betterSqlite3()
const sirannon = new Sirannon({ driver })
const shop = await sirannon.open('shop', './data/shop.db')
 
await shop.execute(
  'CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY, customer TEXT NOT NULL, total INTEGER NOT NULL)'
)
 
const server = createServer(sirannon, { port: 9876 })
await server.listen()
 
type WSMessage = { 
  type: string
  seq?: string
  epoch?: string
  resync?: boolean
  event?: { type: string; seq: string; row: Record<string, unknown> }
  data?: { rows: Array<Record<string, unknown>> }
}
 
const openSocket = (): Promise<WebSocket> =>
  new Promise((resolve, reject) => {
    const socket = new WebSocket('ws://localhost:9876/db/shop')
    socket.addEventListener('open', () => resolve(socket))
    socket.addEventListener('error', () => reject(new Error('connection failed')))
  })
 
const listen = (socket: WebSocket): (() => Promise<WSMessage>) => { 
  const queue: WSMessage[] = []
  const waiters: Array<(message: WSMessage) => void> = []
  socket.addEventListener('message', (event) => {
    const message = JSON.parse(String(event.data)) as WSMessage
    const waiter = waiters.shift()
    if (waiter) waiter(message)
    else queue.push(message)
  })
  return () => {
    const queued = queue.shift()
    if (queued) return Promise.resolve(queued)
    return new Promise((resolve) => waiters.push(resolve))
  }
}
 
const first = await openSocket() 
const nextFromFirst = listen(first)
first.send(JSON.stringify({ type: 'subscribe', id: 'orders-feed', table: 'orders' }))
 
const subscribed = await nextFromFirst() 
let cursor = subscribed.seq ?? '0'
const epoch = subscribed.epoch ?? ''
console.log('live from seq', cursor)
live from seq 0

The subscribed reply carries the two values a resumable client stores. seq is the sequence number the subscription is live from, sent as a decimal string so that cursors beyond 2^53 - 1 survive JSON. A client that has not received any change yet adopts it as its resume cursor, so a reconnect during an idle spell still replays what it missed. epoch identifies the sequence space the cursor was issued from; store it and echo it back when you resume so that a cursor from one database is never replayed against another.

Advance the cursor with each change

Insert an order through the embedded handle. The committed change reaches the socket as a change message, and its event.seq becomes the new cursor.

resume.ts
import { Sirannon } from '@delali/sirannon-db'
import { betterSqlite3 } from '@delali/sirannon-db/driver/better-sqlite3'
import { createServer } from '@delali/sirannon-db/server'
 
const driver = betterSqlite3()
const sirannon = new Sirannon({ driver })
const shop = await sirannon.open('shop', './data/shop.db')
 
await shop.execute(
  'CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY, customer TEXT NOT NULL, total INTEGER NOT NULL)'
)
 
const server = createServer(sirannon, { port: 9876 })
await server.listen()
 
type WSMessage = {
  type: string
  seq?: string
  epoch?: string
  resync?: boolean
  event?: { type: string; seq: string; row: Record<string, unknown> }
  data?: { rows: Array<Record<string, unknown>> }
}
 
const openSocket = (): Promise<WebSocket> =>
  new Promise((resolve, reject) => {
    const socket = new WebSocket('ws://localhost:9876/db/shop')
    socket.addEventListener('open', () => resolve(socket))
    socket.addEventListener('error', () => reject(new Error('connection failed')))
  })
 
const listen = (socket: WebSocket): (() => Promise<WSMessage>) => {
  const queue: WSMessage[] = []
  const waiters: Array<(message: WSMessage) => void> = []
  socket.addEventListener('message', (event) => {
    const message = JSON.parse(String(event.data)) as WSMessage
    const waiter = waiters.shift()
    if (waiter) waiter(message)
    else queue.push(message)
  })
  return () => {
    const queued = queue.shift()
    if (queued) return Promise.resolve(queued)
    return new Promise((resolve) => waiters.push(resolve))
  }
}
 
const first = await openSocket()
const nextFromFirst = listen(first)
first.send(JSON.stringify({ type: 'subscribe', id: 'orders-feed', table: 'orders' }))
 
const subscribed = await nextFromFirst()
let cursor = subscribed.seq ?? '0'
const epoch = subscribed.epoch ?? ''
console.log('live from seq', cursor)
 
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Amara Okafor', 180]) 
 
const change = await nextFromFirst() 
if (change.event) {
  cursor = change.event.seq
  console.log('seq', change.event.seq, change.event.row)
}
seq 1 { id: 1, customer: 'Amara Okafor', total: 180 }

Keep the highest seq you have fully processed, and only advance it after your application has applied the event. A cursor moved before the event is applied turns a crash into a silent gap.

Resume after a disconnect

Close the socket, let two orders commit while nobody is connected, then reconnect and subscribe again with sinceSeq set to the stored cursor and epoch set to the stored epoch. The server replays every retained change with a greater sequence number before delivering live events. The subscribed reply arrives first, so check its resync field before trusting the replay.

resume.ts
import { Sirannon } from '@delali/sirannon-db'
import { betterSqlite3 } from '@delali/sirannon-db/driver/better-sqlite3'
import { createServer } from '@delali/sirannon-db/server'
 
const driver = betterSqlite3()
const sirannon = new Sirannon({ driver })
const shop = await sirannon.open('shop', './data/shop.db')
 
await shop.execute(
  'CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY, customer TEXT NOT NULL, total INTEGER NOT NULL)'
)
 
const server = createServer(sirannon, { port: 9876 })
await server.listen()
 
type WSMessage = {
  type: string
  seq?: string
  epoch?: string
  resync?: boolean
  event?: { type: string; seq: string; row: Record<string, unknown> }
  data?: { rows: Array<Record<string, unknown>> }
}
 
const openSocket = (): Promise<WebSocket> =>
  new Promise((resolve, reject) => {
    const socket = new WebSocket('ws://localhost:9876/db/shop')
    socket.addEventListener('open', () => resolve(socket))
    socket.addEventListener('error', () => reject(new Error('connection failed')))
  })
 
const listen = (socket: WebSocket): (() => Promise<WSMessage>) => {
  const queue: WSMessage[] = []
  const waiters: Array<(message: WSMessage) => void> = []
  socket.addEventListener('message', (event) => {
    const message = JSON.parse(String(event.data)) as WSMessage
    const waiter = waiters.shift()
    if (waiter) waiter(message)
    else queue.push(message)
  })
  return () => {
    const queued = queue.shift()
    if (queued) return Promise.resolve(queued)
    return new Promise((resolve) => waiters.push(resolve))
  }
}
 
const first = await openSocket()
const nextFromFirst = listen(first)
first.send(JSON.stringify({ type: 'subscribe', id: 'orders-feed', table: 'orders' }))
 
const subscribed = await nextFromFirst()
let cursor = subscribed.seq ?? '0'
const epoch = subscribed.epoch ?? ''
console.log('live from seq', cursor)
 
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Amara Okafor', 180])
 
const change = await nextFromFirst()
if (change.event) {
  cursor = change.event.seq
  console.log('seq', change.event.seq, change.event.row)
}
 
first.close() 
 
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Lucas Meyer', 320]) 
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Priya Sharma', 95])
 
const second = await openSocket() 
const nextFromSecond = listen(second)
second.send(
  JSON.stringify({ type: 'subscribe', id: 'orders-feed', table: 'orders', sinceSeq: cursor, epoch })
)
 
const resumed = await nextFromSecond() 
console.log('resync needed:', resumed.resync === true)
 
for (let received = 0; received < 2; received += 1) { 
  const replayed = await nextFromSecond()
  if (replayed.event) {
    cursor = replayed.event.seq
    console.log('seq', replayed.event.seq, replayed.event.row)
  }
}
resync needed: false
seq 2 { id: 2, customer: 'Lucas Meyer', total: 320 }
seq 3 { id: 3, customer: 'Priya Sharma', total: 95 }

Both missed orders replay in commit order, and the feed then continues live on the same socket. Replayed events are indistinguishable from live ones, so one handler covers both.

Force a resync and recover

Replay is impossible in two cases: the requested sinceSeq fell below the server's retained history, or the cursor arrived with a foreign epoch and belongs to a different database. The server then sets resync: true on the subscribed reply and starts the subscription live anyway. Your side of the contract is to treat all prior state as stale and re-read the table before applying further events. Trigger the epoch case deliberately by resuming with a made-up epoch, then re-read over the same socket with a query message.

resume.ts
import { Sirannon } from '@delali/sirannon-db'
import { betterSqlite3 } from '@delali/sirannon-db/driver/better-sqlite3'
import { createServer } from '@delali/sirannon-db/server'
 
const driver = betterSqlite3()
const sirannon = new Sirannon({ driver })
const shop = await sirannon.open('shop', './data/shop.db')
 
await shop.execute(
  'CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY, customer TEXT NOT NULL, total INTEGER NOT NULL)'
)
 
const server = createServer(sirannon, { port: 9876 })
await server.listen()
 
type WSMessage = {
  type: string
  seq?: string
  epoch?: string
  resync?: boolean
  event?: { type: string; seq: string; row: Record<string, unknown> }
  data?: { rows: Array<Record<string, unknown>> }
}
 
const openSocket = (): Promise<WebSocket> =>
  new Promise((resolve, reject) => {
    const socket = new WebSocket('ws://localhost:9876/db/shop')
    socket.addEventListener('open', () => resolve(socket))
    socket.addEventListener('error', () => reject(new Error('connection failed')))
  })
 
const listen = (socket: WebSocket): (() => Promise<WSMessage>) => {
  const queue: WSMessage[] = []
  const waiters: Array<(message: WSMessage) => void> = []
  socket.addEventListener('message', (event) => {
    const message = JSON.parse(String(event.data)) as WSMessage
    const waiter = waiters.shift()
    if (waiter) waiter(message)
    else queue.push(message)
  })
  return () => {
    const queued = queue.shift()
    if (queued) return Promise.resolve(queued)
    return new Promise((resolve) => waiters.push(resolve))
  }
}
 
const first = await openSocket()
const nextFromFirst = listen(first)
first.send(JSON.stringify({ type: 'subscribe', id: 'orders-feed', table: 'orders' }))
 
const subscribed = await nextFromFirst()
let cursor = subscribed.seq ?? '0'
const epoch = subscribed.epoch ?? ''
console.log('live from seq', cursor)
 
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Amara Okafor', 180])
 
const change = await nextFromFirst()
if (change.event) {
  cursor = change.event.seq
  console.log('seq', change.event.seq, change.event.row)
}
 
first.close()
 
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Lucas Meyer', 320])
await shop.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Priya Sharma', 95])
 
const second = await openSocket()
const nextFromSecond = listen(second)
second.send(
  JSON.stringify({ type: 'subscribe', id: 'orders-feed', table: 'orders', sinceSeq: cursor, epoch })
)
 
const resumed = await nextFromSecond()
console.log('resync needed:', resumed.resync === true)
 
for (let received = 0; received < 2; received += 1) {
  const replayed = await nextFromSecond()
  if (replayed.event) {
    cursor = replayed.event.seq
    console.log('seq', replayed.event.seq, replayed.event.row)
  }
}
 
second.close() 
 
const third = await openSocket() 
const nextFromThird = listen(third)
third.send(
  JSON.stringify({
    type: 'subscribe',
    id: 'orders-feed',
    table: 'orders',
    sinceSeq: cursor,
    epoch: 'cursor-from-another-database',
  })
)
 
const forced = await nextFromThird() 
console.log('resync needed:', forced.resync === true, '- live from seq', forced.seq)
 
third.send( 
  JSON.stringify({ type: 'query', id: 'reread', sql: 'SELECT id, customer, total FROM orders ORDER BY id' })
)
const snapshot = await nextFromThird()
console.log('rows re-read:', snapshot.data?.rows.length)
 
third.close() 
resync needed: true - live from seq 3
rows re-read: 3

The reply's seq still reports where the live feed starts, so adopt it as the new cursor once the re-read completes. How much history the server retains for replay is set with the cdcRetentionMs server option; a subscriber away for longer than the retention window follows this same resync path. The finished file keeps serving after the final line prints, and the next section connects a second client to it; stop it with Ctrl+C when you are done.

Keep the cursor where it is until the re-read finishes. Events keep arriving while you re-read, so queue them and apply the queue once the snapshot is in place.

Let the client SDK resume for you

Application code rarely drives this by hand. The WebSocket transport stores the epoch and the last processed sequence number for every subscription, re-sends each subscribe with sinceSeq and epoch when autoReconnect restores the connection, and surfaces the resync case as a callback. Subscribe with an onReset handler, put your re-read there, and let the same client write an order so the feed has an event to deliver. Run feed.ts while resume.ts is still serving.

feed.ts
import { SirannonClient } from '@delali/sirannon-db/client'
 
const client = new SirannonClient('http://localhost:9876', { transport: 'websocket' })
const db = client.database('shop')
 
let markReceived = (): void => {}
const received = new Promise<void>((resolve) => {
  markReceived = resolve
})
 
const sub = await db.on('orders').subscribe(
  (event) => {
    console.log('apply', event.type, event.row)
    markReceived()
  },
  {
    onReset: () => {
      console.log('gap in history: re-read the orders table before trusting local state')
    },
  }
)
 
await db.execute('INSERT INTO orders (customer, total) VALUES (?, ?)', ['Sofia Almeida', 240])
await received
 
sub.unsubscribe()
client.close()
apply insert { id: 4, customer: 'Sofia Almeida', total: 240 }

The inserted order comes straight back as a change event, with id 4 following the three orders resume.ts created on a fresh database. onReset fires only when the server could not replay a gap, so most reconnects pass without it. When it does fire, replace the log line with the same table re-read the earlier sections walked through. sub.unsubscribe() ends the feed and client.close() drops the connection, so the file exits once the event arrives.