Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion .github/workflows/build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ on:
push:
branches: [main]
pull_request:
branches: [main]

jobs:
test:
Expand Down
20 changes: 19 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -229,7 +229,8 @@ YSR_TEST_REDIS=redis://127.0.0.1:36379 npm run test:consistency
```

The tests use a unique key prefix and remove only their own keys. They include
three server instances sharing Redis and a real Socket.IO reconnect.
three server instances sharing Redis, a real Socket.IO reconnect, delayed and
failed saves, worker replacement, and a real persistence worker thread.

A sync-update ACK confirms Redis stream acceptance, not database persistence.
Retries are appended again and safely merged by Yjs. Persistence leases use a
Expand All @@ -242,6 +243,23 @@ lease expiry can still finish after a successor starts persisting. Replacement
storage adapters need storage-enforced commit ordering to prevent stale
overwrites; that contract is outside this change and the storage API is unchanged.

Server-side saves reload a document and its Redis stream position together,
adding a Redis and storage read per save. The saved document is independent of
the subscription cache. Stream trimming is bounded by that snapshot's position
and the minimum message lifetime. `trimRoomStream(room, docid, persistedId)`
requires a known committed position; omitting it retains the stream. Queue
workers use the same bound, including when a task is reclaimed.

Persistence worker requests and success/error replies carry a `requestId`.
Custom worker implementations must echo it with the room; use the matching
package version in the server and worker. Failed saves retry after the normal
persistence delay. Worker errors, exits, and health-check replacement reject
pending requests so later saves and namespace cleanup can proceed.

These guarantees assume that storage resolves `persistDoc` only after a durable
save and rejects on failure. Redis stream expiration and eviction are separate
from trimming; configure retention and Redis durability for your deployment.

## License

[The MIT License](./LICENSE) © Kevin Jahns
97 changes: 58 additions & 39 deletions src/api.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import * as array from 'lib0/array'
import * as random from 'lib0/random'
import * as number from 'lib0/number'
import * as promise from 'lib0/promise'
import * as math from 'lib0/math'
import * as protocol from './protocol.js'
import * as env from 'lib0/environment'
import * as logging from 'lib0/logging'
Expand Down Expand Up @@ -254,8 +253,7 @@ export class Api {
}
const now = performance.now()
if (docstate) { Y.applyUpdateV2(ydoc, docstate.doc) }
let changed = false
ydoc.once('afterTransaction', (tr) => { changed = tr.changed.size > 0 })
const previousState = Y.encodeStateAsUpdateV2(ydoc)
ydoc.transact(() => {
docMessages?.messages.forEach(m => {
const decoder = decoding.createDecoder(m)
Expand All @@ -281,7 +279,9 @@ export class Api {
awareness,
redisLastId: docMessages?.lastId.toString() || '0',
storeReferences: docstate?.references || null,
changed
// Pending structs and delete sets also need persistence, even when no
// shared type changed during the transaction.
changed: !Buffer.from(previousState).equals(Y.encodeStateAsUpdateV2(ydoc))
}
}

Expand All @@ -295,15 +295,30 @@ export class Api {
return docMessages?.lastId.toString() || '0'
}

/**
* @param {string=} persistedId
*/
getTrimBoundary (persistedId) {
if (!persistedId || persistedId === '0') return null
const [ms, sequence = '0'] = persistedId.split('-')
const snapshotBoundary = BigInt(sequence) === (1n << 64n) - 1n
? `${BigInt(ms) + 1n}-0`
: `${ms}-${BigInt(sequence) + 1n}`
const retentionBoundary = `${Math.max(0, Date.now() - this.redisMinMessageLifetime)}-0`
return isSmallerRedisId(snapshotBoundary, retentionBoundary) ? snapshotBoundary : retentionBoundary
}

/**
* @param {string} room
* @param {string} docid
* @param {string=} persistedId Last stream entry included in the committed snapshot.
*/
async trimRoomStream (room, docid) {
async trimRoomStream (room, docid, persistedId) {
const boundary = this.getTrimBoundary(persistedId)
if (!boundary) return
const roomName = computeRedisRoomStreamName(room, docid, this.prefix)
const redisLastId = await this.getRedisLastId(room, docid)
const lastId = number.parseInt(redisLastId.split('-')[0])
await this.redis.xTrim(roomName, 'MINID', lastId - this.redisMinMessageLifetime)
// node-redis types XTRIM thresholds as numbers, but MINID accepts full IDs.
await this.redis.sendCommand(['XTRIM', roomName, 'MINID', boundary])
}

/**
Expand Down Expand Up @@ -349,41 +364,45 @@ export class Api {
} else {
reclaimCounts++
const { room, docid } = decodeRedisRoomStreamName(task.stream, this.prefix)
const { ydoc, storeReferences, redisLastId, changed } = await this.getDoc(room, docid)
const lastId = math.max(number.parseInt(redisLastId.split('-')[0]), number.parseInt(task.id.split('-')[0]))
if (changed) {
logWorker(`persisting changes in room: ${room}`)
await this.store.persistDoc(room, docid, ydoc)
} else logWorker(`skip persisting room: ${room} due to no changes`)
await promise.all([
storeReferences && changed ? this.store.deleteReferences(room, docid, storeReferences) : promise.resolve(),
this.redis.multi()
.xTrim(task.stream, 'MINID', lastId - this.redisMinMessageLifetime)
.xAdd(this.redisWorkerStreamName, '*', { compact: task.stream })
.xAck(this.redisWorkerStreamName, this.redisWorkerGroupName, task.id)
.xDel(this.redisWorkerStreamName, task.id)
.sAdd(this.workerSetName, task.stream)
.exec()
])
logWorker('Compacted stream ', { stream: task.stream, taskId: task.id, newLastId: lastId - this.redisMinMessageLifetime })
const { ydoc, awareness, storeReferences, redisLastId, changed } = await this.getDoc(room, docid)
try {
if (ydocUpdateCallback != null) {
// call YDOC_UPDATE_CALLBACK here
const formData = new FormData()
// @todo only convert ydoc to updatev2 once
// @ts-ignore
formData.append('ydoc', new Blob([Y.encodeStateAsUpdateV2(ydoc)]))
// @todo should add a timeout to fetch (see fetch signal abortcontroller)
const res = await fetch(new URL(room, ydocUpdateCallback), { body: formData, method: 'PUT' })
if (!res.ok) {
console.error(`Issue sending data to YDOC_UPDATE_CALLBACK. status="${res.status}" statusText="${res.statusText}"`)
if (changed) {
logWorker(`persisting changes in room: ${room}`)
await this.store.persistDoc(room, docid, ydoc)
} else logWorker(`skip persisting room: ${room} due to no changes`)
const boundary = this.getTrimBoundary(redisLastId)
const transaction = this.redis.multi()
if (boundary) transaction.addCommand(['XTRIM', task.stream, 'MINID', boundary])
await promise.all([
storeReferences && changed ? this.store.deleteReferences(room, docid, storeReferences) : promise.resolve(),
transaction
.xAdd(this.redisWorkerStreamName, '*', { compact: task.stream })
.xAck(this.redisWorkerStreamName, this.redisWorkerGroupName, task.id)
.xDel(this.redisWorkerStreamName, task.id)
.sAdd(this.workerSetName, task.stream)
.exec()
])
logWorker('Compacted stream ', { stream: task.stream, taskId: task.id, boundary })
try {
if (ydocUpdateCallback != null) {
// call YDOC_UPDATE_CALLBACK here
const formData = new FormData()
// @todo only convert ydoc to updatev2 once
// @ts-ignore
formData.append('ydoc', new Blob([Y.encodeStateAsUpdateV2(ydoc)]))
// @todo should add a timeout to fetch (see fetch signal abortcontroller)
const res = await fetch(new URL(room, ydocUpdateCallback), { body: formData, method: 'PUT' })
if (!res.ok) {
console.error(`Issue sending data to YDOC_UPDATE_CALLBACK. status="${res.status}" statusText="${res.statusText}"`)
}
}
} catch (e) {
console.error(e)
}
} catch (e) {
console.error(e)
} finally {
awareness?.destroy()
ydoc.destroy()
}
// destroy ydoc after persisting
ydoc.destroy()
}
}))
return { tasks, reclaimCounts }
Expand Down
20 changes: 13 additions & 7 deletions src/persist-worker-thread.js
Original file line number Diff line number Diff line change
Expand Up @@ -21,21 +21,27 @@ export class PersistWorkerThread {
parentPort?.postMessage({ event: 'ready' })
parentPort?.on('message', ({ event, ...rest }) => {
if (event === 'ping') parentPort?.postMessage({ event: 'pong' })
else this.persist(rest)
else if (event === undefined || event === 'persist') this.persist(rest)
})
}

/**
* @param {{ room: string, docstate: SharedArrayBuffer }} props
* @param {{ room: string, requestId: string, docstate: Uint8Array }} props
*/
persist = async ({ room, docstate }) => {
persist = async ({ room, requestId, docstate }) => {
this.log(`persisting ${room} in worker`)
const state = new Uint8Array(docstate)
const doc = new Y.Doc()
Y.applyUpdateV2(doc, state)
await this.store?.persistDoc(room, 'index', doc)
doc.destroy()
parentPort?.postMessage({ event: 'persisted', room })
try {
Y.applyUpdateV2(doc, state)
if (!this.store) throw new Error('Persistence storage is unavailable')
await this.store.persistDoc(room, 'index', doc)
parentPort?.postMessage({ event: 'persisted', room, requestId })
} catch (error) {
parentPort?.postMessage({ event: 'persist-error', room, requestId, error: String(error) })
} finally {
doc.destroy()
}
}
}

Expand Down
Loading
Loading