Skip to content

Commit 0ca4588

Browse files
committed
fix(files): close recovery and seed review gaps
1 parent fa478b2 commit 0ca4588

18 files changed

Lines changed: 725 additions & 152 deletions

File tree

apps/realtime/src/handlers/file-doc-app.ts

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -18,10 +18,8 @@ function postToApp(path: string, payload: unknown, timeoutMs: number): Promise<R
1818
}
1919

2020
/**
21-
* Ask the app to build a server-authoritative seed (markdown → Yjs) for a file's collaborative
22-
* document. Returns the Yjs update to apply, or `null` for a genuinely empty/missing file (an empty
23-
* document is correct). THROWS on a transport failure (non-2xx / network / timeout / malformed body)
24-
* so the caller can tell a real empty from a failure it should be allowed to retry.
21+
* Existing empty files have named, versioned seeds; only missing files return null.
22+
* Transport and malformed-response failures throw so callers retry instead of creating empty rooms.
2523
*/
2624
export async function fetchFileDocSeed(
2725
workspaceId: string,

apps/realtime/src/handlers/file-doc-store.test.ts

Lines changed: 87 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
/**
22
* @vitest-environment node
33
*/
4+
import { FILE_DOC_SEED } from '@sim/realtime-protocol/file-doc'
45
import { sleep } from '@sim/utils/helpers'
56
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
67
import * as Y from 'yjs'
@@ -311,6 +312,18 @@ function updateFor(text: string): Uint8Array {
311312
return update
312313
}
313314

315+
function seedFor(text: string): Uint8Array {
316+
const doc = docWithText(text)
317+
const config = doc.getMap(FILE_DOC_SEED.configMap)
318+
config.set(FILE_DOC_SEED.flag, true)
319+
config.set(FILE_DOC_SEED.docIdKey, `doc-${text}`)
320+
try {
321+
return Y.encodeStateAsUpdate(doc)
322+
} finally {
323+
doc.destroy()
324+
}
325+
}
326+
314327
let stores: FileDocStore[] = []
315328

316329
/** An existing stream from a relay predating generation markers; modern seeds use seedIfEmpty. */
@@ -445,7 +458,7 @@ describe('FileDocStore', () => {
445458
const token = await a.shouldSeed(NAME)
446459
expect(token).toBeTruthy()
447460
// A seeds and releases its lock.
448-
await a.seedIfEmpty(NAME, updateFor('hello'))
461+
await a.seedIfEmpty(NAME, seedFor('hello'))
449462
await vi.waitFor(async () => expect(await a.getStreamState(NAME)).not.toBeNull())
450463
await a.releaseSeedLock(NAME, token as string)
451464
// A different task must NOT seed again — the lock is free but the stream is non-empty.
@@ -455,7 +468,7 @@ describe('FileDocStore', () => {
455468

456469
it('fences stale publishers after invalidation and lets the next authoritative seed start fresh', async () => {
457470
const store = await newStore()
458-
const original = updateFor('old generation')
471+
const original = seedFor('old generation')
459472
await store.seedIfEmpty(NAME, original)
460473
await store.invalidateDocument(NAME, 10)
461474

@@ -467,7 +480,7 @@ describe('FileDocStore', () => {
467480
store.publishClientUpdateAndWait(NAME, 'stale-update', updateFor('stale acknowledged write'))
468481
).rejects.toThrow('replaced by a newer durable version')
469482

470-
const fresh = updateFor('fresh generation')
483+
const fresh = seedFor('fresh generation')
471484
await expect(store.seedIfEmpty(NAME, fresh, 11)).resolves.toBe(true)
472485
const recovered = new Y.Doc()
473486
Y.applyUpdate(recovered, (await store.getStreamState(NAME))!)
@@ -477,7 +490,7 @@ describe('FileDocStore', () => {
477490

478491
it('getStreamState reconstructs the shared document from the stream', async () => {
479492
const a = await newStore()
480-
await a.seedIfEmpty(NAME, updateFor('shared content'))
493+
await a.seedIfEmpty(NAME, seedFor('shared content'))
481494
let state: Uint8Array | null = null
482495
await vi.waitFor(async () => {
483496
state = await a.getStreamState(NAME)
@@ -491,7 +504,7 @@ describe('FileDocStore', () => {
491504

492505
it('lets a headless replica append against the generation of its shared base', async () => {
493506
const seeded = await newStore()
494-
await seeded.seedIfEmpty(NAME, updateFor('shared'), 20)
507+
await seeded.seedIfEmpty(NAME, seedFor('shared'), 20)
495508
const headless = await newStore()
496509
const generation = await headless.getDocumentGeneration(NAME)
497510
const doc = new Y.Doc()
@@ -508,7 +521,7 @@ describe('FileDocStore', () => {
508521

509522
it('keeps a newer seeded generation when an older invalidation arrives', async () => {
510523
const store = await newStore()
511-
await store.seedIfEmpty(NAME, updateFor('newest'), 20)
524+
await store.seedIfEmpty(NAME, seedFor('newest'), 20)
512525
const generation = await store.getDocumentGeneration(NAME)
513526
await expect(store.invalidateDocument(NAME, 10)).resolves.toEqual({ status: 'stale' })
514527
expect(await store.getDocumentGeneration(NAME)).toBe(generation)
@@ -518,20 +531,20 @@ describe('FileDocStore', () => {
518531

519532
it('rejects old seeds and version callbacks after an invalidation', async () => {
520533
const store = await newStore()
521-
await store.seedIfEmpty(NAME, updateFor('old'), 10)
534+
await store.seedIfEmpty(NAME, seedFor('old'), 10)
522535
const generation = await store.getDocumentGeneration(NAME)
523536
await store.invalidateDocument(NAME, 20)
524-
await expect(store.seedIfEmpty(NAME, updateFor('late stale seed'), 10)).resolves.toBe(false)
537+
await expect(store.seedIfEmpty(NAME, seedFor('late stale seed'), 10)).resolves.toBe(false)
525538
await store.setSyncedVersion(NAME, 30, generation)
526539
expect(await store.getSyncedVersion(NAME)).toBe(20)
527540
await expect(store.getStreamState(NAME)).resolves.toBeNull()
528541
})
529542

530543
it('does not repeat an invalidation after the same durable version is reseeded', async () => {
531544
const store = await newStore()
532-
await store.seedIfEmpty(NAME, updateFor('old'), 10)
545+
await store.seedIfEmpty(NAME, seedFor('old'), 10)
533546
await expect(store.invalidateDocument(NAME, 20)).resolves.toMatchObject({ status: 'applied' })
534-
await expect(store.seedIfEmpty(NAME, updateFor('replacement'), 20)).resolves.toBe(true)
547+
await expect(store.seedIfEmpty(NAME, seedFor('replacement'), 20)).resolves.toBe(true)
535548
const generation = await store.getDocumentGeneration(NAME)
536549
const doc = new Y.Doc()
537550
Y.applyUpdate(doc, (await store.getStreamState(NAME))!)
@@ -556,7 +569,7 @@ describe('FileDocStore', () => {
556569

557570
it('applies the first invalidation even when its durable version was already seeded', async () => {
558571
const store = await newStore()
559-
await store.seedIfEmpty(NAME, updateFor('same content, changed eligibility'), 20)
572+
await store.seedIfEmpty(NAME, seedFor('same content, changed eligibility'), 20)
560573
const docId = await store.getDocumentGeneration(NAME)
561574
await expect(store.invalidateDocument(NAME, 20)).resolves.toEqual({ status: 'applied', docId })
562575
await expect(store.invalidateDocument(NAME, 20)).resolves.toEqual({ status: 'stale' })
@@ -565,14 +578,14 @@ describe('FileDocStore', () => {
565578

566579
it('returns the removed generation and qualifies consecutive unsupported replacements', async () => {
567580
const store = await newStore()
568-
await store.seedIfEmpty(NAME, updateFor('old'), 10)
581+
await store.seedIfEmpty(NAME, seedFor('old'), 10)
569582
const oldId = await store.getDocumentGeneration(NAME)
570583
await expect(store.invalidateDocument(NAME, 20)).resolves.toEqual({
571584
status: 'applied',
572585
docId: oldId,
573586
})
574587
await expect(store.invalidateDocument(NAME, 30)).resolves.toEqual({ status: 'applied' })
575-
await store.seedIfEmpty(NAME, updateFor('replacement'), 30)
588+
await store.seedIfEmpty(NAME, seedFor('replacement'), 30)
576589
const replacementId = await store.getDocumentGeneration(NAME)
577590
expect(replacementId).not.toBe(oldId)
578591
await expect(store.invalidateDocument(NAME, 30)).resolves.toEqual({ status: 'stale' })
@@ -584,7 +597,7 @@ describe('FileDocStore', () => {
584597

585598
it('does not resurrect a tracked stream with a dependency-only update after Redis loses it', async () => {
586599
const store = await newStore()
587-
await store.seedIfEmpty(NAME, updateFor('base'), 10)
600+
await store.seedIfEmpty(NAME, seedFor('base'), 10)
588601
const generation = await store.getDocumentGeneration(NAME)
589602
state.backing!.streams.delete(`filedoc:stream:${NAME}`)
590603
state.backing!.kv.delete(`filedoc:generation:${NAME}`)
@@ -599,7 +612,7 @@ describe('FileDocStore', () => {
599612

600613
it('rejects appends and duplicate acknowledgements when only the stream is lost', async () => {
601614
const store = await newStore()
602-
await store.seedIfEmpty(NAME, updateFor('base'), 10)
615+
await store.seedIfEmpty(NAME, seedFor('base'), 10)
603616
const generation = await store.getDocumentGeneration(NAME)
604617
const delta = updateFor('edit')
605618
await store.publishClientUpdateAndWait(NAME, 'accepted-update', delta, generation)
@@ -642,7 +655,7 @@ describe('FileDocStore', () => {
642655

643656
it('rejects a shared replay if the document generation changes between pages', async () => {
644657
const store = await newStore()
645-
await store.seedIfEmpty(NAME, updateFor('old generation'), 10)
658+
await store.seedIfEmpty(NAME, seedFor('old generation'), 10)
646659
state.backing!.onRange = () => {
647660
state.backing!.kv.set(`filedoc:generation:${NAME}`, 'new generation')
648661
}
@@ -885,7 +898,7 @@ describe('FileDocStore', () => {
885898

886899
it('attachRoom catches a fresh task up to the current shared state', async () => {
887900
const a = await newStore()
888-
await a.seedIfEmpty(NAME, updateFor('already here'))
901+
await a.seedIfEmpty(NAME, seedFor('already here'))
889902
await vi.waitFor(async () => expect(await a.getStreamState(NAME)).not.toBeNull())
890903

891904
// A second task opens the same file: its doc must load the existing content, not start empty.
@@ -1271,18 +1284,18 @@ describe('FileDocStore', () => {
12711284
const staleDoc = new Y.Doc()
12721285
await store.attachRoom(NAME, staleDoc)
12731286
await store.invalidateDocument(NAME, 20)
1274-
await expect(store.seedIfEmpty(NAME, updateFor('stale fetched seed'), 10)).resolves.toBe(false)
1287+
await expect(store.seedIfEmpty(NAME, seedFor('stale fetched seed'), 10)).resolves.toBe(false)
12751288
await expect(store.isDocumentGenerationCurrent(NAME)).resolves.toBe(false)
12761289
expect(storeInternals(store).localInvalidations.size).toBe(1)
1277-
await expect(store.seedIfEmpty(NAME, updateFor('same-version seed'), 20)).resolves.toBe(true)
1290+
await expect(store.seedIfEmpty(NAME, seedFor('same-version seed'), 20)).resolves.toBe(true)
12781291
await expect(store.invalidateDocument(NAME, 20)).resolves.toEqual({ status: 'stale' })
12791292
await expect(store.isDocumentGenerationCurrent(NAME)).resolves.toBe(true)
12801293

12811294
store.detachRoom(NAME)
12821295
expect(storeInternals(store).localInvalidations.size).toBe(1)
12831296
const freshDoc = new Y.Doc()
12841297
await store.attachRoom(NAME, freshDoc)
1285-
await expect(store.seedIfEmpty(NAME, updateFor('fresh authoritative seed'), 20)).resolves.toBe(
1298+
await expect(store.seedIfEmpty(NAME, seedFor('fresh authoritative seed'), 20)).resolves.toBe(
12861299
true
12871300
)
12881301
await expect(store.invalidateDocument(NAME, 20)).resolves.toEqual({ status: 'stale' })
@@ -1302,7 +1315,7 @@ describe('FileDocStore', () => {
13021315
await expect(
13031316
store.publishClientUpdateAndWait(NAME, 'update-1', updateFor('x'))
13041317
).rejects.toThrow('not initialized')
1305-
await expect(store.seedIfEmpty(NAME, updateFor('seed'))).rejects.toThrow('not initialized')
1318+
await expect(store.seedIfEmpty(NAME, seedFor('seed'))).rejects.toThrow('not initialized')
13061319
await expect(store.getStreamState(NAME)).rejects.toThrow('not initialized')
13071320
expect(await store.acquireMergeSlot(NAME, 1_000)).toBeNull()
13081321
doc.destroy()
@@ -1311,7 +1324,7 @@ describe('FileDocStore', () => {
13111324
it('streamHasContent fences a seed apply against an already-seeded stream', async () => {
13121325
const a = await newStore()
13131326
expect(await a.streamHasContent(NAME)).toBe(false)
1314-
await a.seedIfEmpty(NAME, updateFor('seeded'))
1327+
await a.seedIfEmpty(NAME, seedFor('seeded'))
13151328
await vi.waitFor(async () => expect(await a.streamHasContent(NAME)).toBe(true))
13161329
})
13171330

@@ -1344,12 +1357,60 @@ describe('FileDocStore', () => {
13441357
doc.destroy()
13451358
})
13461359

1360+
it.each([
1361+
{ redis: true, docId: undefined },
1362+
{ redis: true, docId: '' },
1363+
{ redis: false, docId: undefined },
1364+
{ redis: false, docId: '' },
1365+
])('rejects an unnamed seed before publication (%j)', async ({ redis, docId }) => {
1366+
const store = redis ? await newStore() : new FileDocStore(undefined)
1367+
const doc = new Y.Doc()
1368+
const config = doc.getMap(FILE_DOC_SEED.configMap)
1369+
config.set(FILE_DOC_SEED.flag, true)
1370+
if (docId !== undefined) config.set(FILE_DOC_SEED.docIdKey, docId)
1371+
try {
1372+
await expect(store.seedIfEmpty(NAME, Y.encodeStateAsUpdate(doc), 1)).rejects.toThrow(
1373+
'missing its accepted document identity'
1374+
)
1375+
await expect(store.getStreamState(NAME)).resolves.toBeNull()
1376+
} finally {
1377+
doc.destroy()
1378+
if (!redis) await store.shutdown()
1379+
}
1380+
})
1381+
1382+
it('uses the same accepted identity for the seed owner and a replaying peer', async () => {
1383+
const owner = await newStore()
1384+
const peer = await newStore()
1385+
const ownerDoc = new Y.Doc()
1386+
const peerDoc = new Y.Doc()
1387+
try {
1388+
await owner.attachRoom(NAME, ownerDoc)
1389+
const seed = seedFor('shared identity')
1390+
expect(await owner.seedIfEmpty(NAME, seed, 1)).toBe(true)
1391+
Y.applyUpdate(ownerDoc, seed)
1392+
await peer.attachRoom(NAME, peerDoc)
1393+
for (const [store, doc] of [
1394+
[owner, ownerDoc],
1395+
[peer, peerDoc],
1396+
] as const) {
1397+
const docId = doc.getMap(FILE_DOC_SEED.configMap).get(FILE_DOC_SEED.docIdKey)
1398+
expect(docId).toBe('doc-shared identity')
1399+
expect(await store.getDocumentGeneration(NAME)).toBe(docId)
1400+
expect(await store.isDocumentGenerationCurrent(NAME, 'doc-shared identity')).toBe(true)
1401+
}
1402+
} finally {
1403+
ownerDoc.destroy()
1404+
peerDoc.destroy()
1405+
}
1406+
})
1407+
13471408
it('seedIfEmpty writes the seed once and reports it, then refuses a non-empty stream', async () => {
13481409
const a = await newStore()
1349-
expect(await a.seedIfEmpty(NAME, updateFor('first'))).toBe(true)
1410+
expect(await a.seedIfEmpty(NAME, seedFor('first'))).toBe(true)
13501411
// A second seed attempt (any task) must be refused — the stream already holds content.
13511412
const b = await newStore()
1352-
expect(await b.seedIfEmpty(NAME, updateFor('second'))).toBe(false)
1413+
expect(await b.seedIfEmpty(NAME, seedFor('second'))).toBe(false)
13531414
const doc = new Y.Doc()
13541415
Y.applyUpdate(doc, (await a.getStreamState(NAME))!)
13551416
expect(doc.getText('body').toString()).toBe('first')
@@ -1371,8 +1432,8 @@ describe('FileDocStore', () => {
13711432
expect(tokenB).toBeTruthy()
13721433
// Both tasks now race to seed with distinct client ids.
13731434
const [seededA, seededB] = await Promise.all([
1374-
a.seedIfEmpty(NAME, updateFor('SEED-A')),
1375-
b.seedIfEmpty(NAME, updateFor('SEED-B')),
1435+
a.seedIfEmpty(NAME, seedFor('SEED-A')),
1436+
b.seedIfEmpty(NAME, seedFor('SEED-B')),
13761437
])
13771438
expect([seededA, seededB].filter(Boolean)).toHaveLength(1)
13781439
// Exactly one seed is in the stream — the reconstructed text is a single seed, never a duplicated

apps/realtime/src/handlers/file-doc-store.ts

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -245,9 +245,10 @@ function generationOfSeed(update: Uint8Array): string {
245245
try {
246246
Y.applyUpdate(doc, update)
247247
const docId = doc.getMap(FILE_DOC_SEED.configMap).get(FILE_DOC_SEED.docIdKey)
248-
return typeof docId === 'string'
249-
? docId
250-
: `seed:${createHash('sha256').update(update).digest('hex')}`
248+
if (typeof docId !== 'string' || docId.length === 0) {
249+
throw new Error('File document seed is missing its accepted document identity')
250+
}
251+
return docId
251252
} finally {
252253
doc.destroy()
253254
}
@@ -623,6 +624,8 @@ export class FileDocStore {
623624
* true (single-replica: seed locally, no stream).
624625
*/
625626
async seedIfEmpty(name: string, update: Uint8Array, version = 0): Promise<boolean> {
627+
assertUpdateWithinLimit(update)
628+
const generation = generationOfSeed(update)
626629
if (!this.enabled) {
627630
const invalidation = this.localInvalidations.get(name)
628631
if (invalidation && invalidation.expiresAt > Date.now() && invalidation.version > version)
@@ -632,9 +635,7 @@ export class FileDocStore {
632635
return true
633636
}
634637
if (!this.write) throw new Error('FileDocStore is not initialized')
635-
assertUpdateWithinLimit(update)
636638
const encoded = Buffer.from(update).toString('base64')
637-
const generation = generationOfSeed(update)
638639
for (let attempt = 0; attempt <= PUBLISH_MAX_RETRIES; attempt++) {
639640
try {
640641
const wrote = await this.write.eval(SEED_IF_EMPTY_SCRIPT, {

0 commit comments

Comments
 (0)