From 9a397afaa908f92a96cfe41c695b32fcd860328f Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:53:58 -0700 Subject: [PATCH 1/2] Add tests for room-registry teardown data loss and zombie writes Reproduces two compounding bugs in destroyRoomIfEmpty/applyDocUpdate: a room's live Y.Doc is destroyed the instant it looks empty with no guarantee its pending content was flushed first, and applyDocUpdate silently auto-creates a room for any write, letting a delayed/zombie POST from an evicted client repopulate a freshly torn-down room and defeat seedOnJoin's "only seed an empty doc" guard. --- .../presence/src/artifact-persistence.test.ts | 60 +++++- packages/presence/src/room-registry.test.ts | 180 +++++++++++++++++- packages/presence/test/routes.test.ts | 26 +++ 3 files changed, 263 insertions(+), 3 deletions(-) diff --git a/packages/presence/src/artifact-persistence.test.ts b/packages/presence/src/artifact-persistence.test.ts index 304df7438..cc5c7395e 100644 --- a/packages/presence/src/artifact-persistence.test.ts +++ b/packages/presence/src/artifact-persistence.test.ts @@ -12,6 +12,22 @@ function docUpdateInserting(text: string): Uint8Array { return Y.encodeStateAsUpdate(doc); } +/** `applyDocUpdate` only accepts updates for a room that already exists — + * it never auto-creates one (see `room-registry.ts`) — so every test that + * posts a doc update joins first to open the room, exactly like the real + * join-then-edit flow the HTTP routes use. */ +function joinPrincipal( + registry: ReturnType, + key: { tenantId: string; surface: string }, + principalId: string, +): void { + registry.join(key, { + principalId, + displayName: principalId, + color: "hsl(0 65% 45%)", + }); +} + /** Drains pending microtasks — generous rather than an exact tick count, * since the write chain's depth (and therefore how many `.then()` hops a * test needs to wait out) is an implementation detail these tests @@ -87,6 +103,7 @@ describe("createArtifactDocPersistence", () => { registry.seedDocText(key, ""); registry.subscribeDocUpdates(key, () => undefined); // keep the room alive + joinPrincipal(registry, key, "prn_alice"); // Simulate three quick edits, each a real Yjs update. const edit = (text: string) => { @@ -131,14 +148,13 @@ describe("createArtifactDocPersistence", () => { }, }); - registry.applyDocUpdate(key, docUpdateInserting("final edit"), "prn_bob"); - const unsubscribe = registry.subscribe(key, () => undefined); registry.join(key, { principalId: "prn_bob", displayName: "Bob", color: "hsl(0 65% 45%)", }); + registry.applyDocUpdate(key, docUpdateInserting("final edit"), "prn_bob"); registry.leave(key, "prn_bob"); unsubscribe(); @@ -149,6 +165,42 @@ describe("createArtifactDocPersistence", () => { expect(writes).toEqual(["final edit"]); }); + test("a post-eviction zombie update cannot pre-empt or block seedOnJoin's content restoration", async () => { + const registry = createPresenceRoomRegistry(); + const persistence = createArtifactDocPersistence({ + registry, + loadArtifactContent: async () => "content from storage", + writeArtifactSnapshot: async () => ({ version: 1 }), + }); + + // Alice joins, edits, and leaves — the room empties and is torn down. + const unsubscribe = registry.subscribe(key, () => undefined); + joinPrincipal(registry, key, "prn_alice"); + registry.applyDocUpdate( + key, + docUpdateInserting("alice's edit"), + "prn_alice", + ); + registry.leave(key, "prn_alice"); + unsubscribe(); + // The room's teardown flush is asynchronous (see room-registry.ts's + // deferred destroy) — wait for it to actually settle and the room to + // be gone before simulating a POST arriving after that point. + await flushMicrotasks(); + + // A delayed POST from Alice's now-evicted client lands after the + // room was destroyed. It must be rejected outright, not silently + // recreate the room and populate it with stale content. + expect(() => + registry.applyDocUpdate(key, docUpdateInserting("zombie"), "prn_alice"), + ).toThrow(); + + // A legitimate rejoin must still restore the real stored content — + // the zombie write must not have made the doc look already-seeded. + await persistence.seedOnJoin(key); + expect(registry.docText(key)).toBe("content from storage"); + }); + test("never schedules a snapshot for a non-artifact surface", () => { const registry = createPresenceRoomRegistry(); const clock = fakeClock(); @@ -166,6 +218,7 @@ describe("createArtifactDocPersistence", () => { }); const channelKey = { tenantId: "tnt_a", surface: "channel:chn_1" }; + joinPrincipal(registry, channelKey, "prn_alice"); registry.applyDocUpdate( channelKey, docUpdateInserting("not an artifact"), @@ -234,6 +287,7 @@ describe("createArtifactDocPersistence", () => { writeArtifactSnapshot: async () => ({ version: 12 }), }); + joinPrincipal(registry, key, "prn_alice"); registry.applyDocUpdate(key, docUpdateInserting("x"), "prn_alice"); clock.advance(2_100); await Promise.resolve(); @@ -258,6 +312,7 @@ describe("createArtifactDocPersistence", () => { onSnapshotError: (_key, error) => errors.push(error), }); + joinPrincipal(registry, key, "prn_alice"); registry.applyDocUpdate(key, docUpdateInserting("x"), "prn_alice"); clock.advance(3_000); @@ -298,6 +353,7 @@ describe("createArtifactDocPersistence", () => { }, }); + joinPrincipal(registry, key, "prn_alice"); registry.applyDocUpdate(key, docUpdateInserting("a"), "prn_alice"); clock.advance(2_000); // schedules write #1 (slow — awaits its resolver) await Promise.resolve(); diff --git a/packages/presence/src/room-registry.test.ts b/packages/presence/src/room-registry.test.ts index 630ac1fc9..a9eadc6d7 100644 --- a/packages/presence/src/room-registry.test.ts +++ b/packages/presence/src/room-registry.test.ts @@ -136,6 +136,12 @@ function clientDoc(): Y.Doc { return new Y.Doc(); } +/** Drains pending microtasks — generous rather than an exact tick count, + * matching `artifact-persistence.test.ts`'s helper of the same name. */ +async function flushMicrotasks(times = 10): Promise { + for (let i = 0; i < times; i += 1) await Promise.resolve(); +} + describe("createPresenceRoomRegistry: doc sync", () => { const docKey: PresenceRoomKey = { tenantId: "tnt_a", @@ -144,6 +150,8 @@ describe("createPresenceRoomRegistry: doc sync", () => { test("concurrent inserts from two clients both land in the server's doc", () => { const registry = createPresenceRoomRegistry(); + registry.join(docKey, state("prn_alice")); + registry.join(docKey, state("prn_bob")); const alice = clientDoc(); const bob = clientDoc(); alice.getText("content").insert(0, "hello"); @@ -162,6 +170,7 @@ describe("createPresenceRoomRegistry: doc sync", () => { test("a late joiner catches up to the full doc state via docStateAsUpdate", () => { const registry = createPresenceRoomRegistry(); + registry.join(docKey, state("prn_alice")); const alice = clientDoc(); alice.getText("content").insert(0, "already written"); registry.applyDocUpdate(docKey, Y.encodeStateAsUpdate(alice), "prn_alice"); @@ -174,6 +183,7 @@ describe("createPresenceRoomRegistry: doc sync", () => { test("applyDocUpdate notifies both the room-scoped and global doc-change listeners", () => { const registry = createPresenceRoomRegistry(); + registry.join(docKey, state("prn_alice")); const roomUpdates: string[] = []; const globalChanges: { key: PresenceRoomKey; author: string }[] = []; registry.subscribeDocUpdates(docKey, (_update, author) => @@ -189,6 +199,54 @@ describe("createPresenceRoomRegistry: doc sync", () => { expect(globalChanges).toEqual([{ key: docKey, author: "prn_alice" }]); }); + test("applyDocUpdate rejects an update for a room nobody currently holds open", () => { + const registry = createPresenceRoomRegistry(); + const alice = clientDoc(); + alice.getText("content").insert(0, "never joined"); + + expect(() => + registry.applyDocUpdate( + docKey, + Y.encodeStateAsUpdate(alice), + "prn_alice", + ), + ).toThrow(); + expect(registry.docText(docKey)).toBe(""); + }); + + test("a zombie update arriving after the room emptied cannot repopulate the freshly recreated room, and a real rejoin still restores it", () => { + const registry = createPresenceRoomRegistry(); + const unsubscribe = registry.subscribe(docKey, () => undefined); + registry.join(docKey, state("prn_alice")); + registry.applyDocUpdate( + docKey, + Y.encodeStateAsUpdate(clientDoc()), + "prn_alice", + ); + registry.leave(docKey, "prn_alice"); + unsubscribe(); + + // The room is now torn down. A delayed POST from Alice's now-evicted + // client — built before she left, delivered after — must not be able + // to create a fresh, unseeded room out of nothing. + const zombie = clientDoc(); + zombie.getText("content").insert(0, "zombie content"); + expect(() => + registry.applyDocUpdate( + docKey, + Y.encodeStateAsUpdate(zombie), + "prn_alice", + ), + ).toThrow(); + expect(registry.docText(docKey)).toBe(""); + + // A legitimate rejoin still sees a genuinely empty doc, ready for + // `seedOnJoin` to restore real content — not silently defeated by the + // zombie write. + expect(registry.seedDocText(docKey, "real stored content")).toBe(true); + expect(registry.docText(docKey)).toBe("real stored content"); + }); + test("seedDocText only seeds an empty doc, never clobbering real content", () => { const registry = createPresenceRoomRegistry(); expect(registry.seedDocText(docKey, "seeded")).toBe(true); @@ -201,7 +259,9 @@ describe("createPresenceRoomRegistry: doc sync", () => { test("onEmpty fires when a room with doc content is torn down", () => { const registry = createPresenceRoomRegistry(); const emptied: PresenceRoomKey[] = []; - registry.onEmpty((key) => emptied.push(key)); + registry.onEmpty((key) => { + emptied.push(key); + }); registry.seedDocText(docKey, "content before anyone joins"); const unsubscribe = registry.subscribe(docKey, () => undefined); @@ -255,3 +315,121 @@ describe("createPresenceRoomRegistry: doc sync", () => { expect(seen).toEqual([]); }); }); + +describe("createPresenceRoomRegistry: deferred destroy", () => { + const docKey: PresenceRoomKey = { + tenantId: "tnt_a", + surface: "artifact:art_1", + }; + + test("a rejoin before a pending onEmpty flush settles cancels the deferred destroy without losing content", async () => { + const registry = createPresenceRoomRegistry(); + let resolveFlush: (() => void) | undefined; + registry.onEmpty( + () => + new Promise((resolve) => { + resolveFlush = resolve; + }), + ); + + registry.seedDocText(docKey, "unsaved content"); + const unsubscribe = registry.subscribe(docKey, () => undefined); + registry.join(docKey, state("prn_alice")); + registry.leave(docKey, "prn_alice"); + unsubscribe(); // room now empty; destroy deferred on the pending flush + + // Bob joins before the flush resolves. + registry.join(docKey, state("prn_bob")); + expect(registry.docText(docKey)).toBe("unsaved content"); + + resolveFlush?.(); + await flushMicrotasks(); + + // The doc must still be intact — the deferred destroy must not have + // fired against a room that came back to life while it waited. + expect(registry.docText(docKey)).toBe("unsaved content"); + }); + + test("a doc update landing on a room while its flush is pending gets its own flush, never silently destroyed with it", async () => { + const registry = createPresenceRoomRegistry(); + const flushedContents: string[] = []; + let resolveFirstFlush: (() => void) | undefined; + let flushCount = 0; + registry.onEmpty((key) => { + flushCount += 1; + flushedContents.push(registry.docText(key)); + if (flushCount === 1) { + return new Promise((resolve) => { + resolveFirstFlush = resolve; + }); + } + return undefined; + }); + + registry.seedDocText(docKey, "first content"); + const unsubscribe = registry.subscribe(docKey, () => undefined); + registry.join(docKey, state("prn_alice")); + registry.leave(docKey, "prn_alice"); + unsubscribe(); // flush #1 dispatched and pending + + // A doc update lands on the still-live (pending-destroy) room before + // the first flush resolves. + const alice = clientDoc(); + alice.getText("content").insert(0, "second content"); + registry.applyDocUpdate( + docKey, + Y.encodeStateAsUpdate(alice), + "prn_alice", + ); + + resolveFirstFlush?.(); + await flushMicrotasks(); + + // The second write must have triggered its own flush dispatch instead + // of being silently dropped once the room was eventually destroyed. + expect(flushCount).toBe(2); + expect(flushedContents[1]).toContain("second content"); + }); + + test("a rejoin/edit/leave cycle happening entirely inside a pending flush window is not lost", async () => { + const registry = createPresenceRoomRegistry(); + const flushedContents: string[] = []; + const resolvers: (() => void)[] = []; + registry.onEmpty((key) => { + flushedContents.push(registry.docText(key)); + return new Promise((resolve) => resolvers.push(resolve)); + }); + + registry.seedDocText(docKey, "first content"); + const unsubscribeAlice = registry.subscribe(docKey, () => undefined); + registry.join(docKey, state("prn_alice")); + registry.leave(docKey, "prn_alice"); + unsubscribeAlice(); // flush #1 dispatched and pending + + // Bob joins, edits, and leaves entirely within the first flush's + // pending window. + const unsubscribeBob = registry.subscribe(docKey, () => undefined); + registry.join(docKey, state("prn_bob")); + const bob = clientDoc(); + bob.getText("content").insert(0, "bob's edit"); + registry.applyDocUpdate(docKey, Y.encodeStateAsUpdate(bob), "prn_bob"); + registry.leave(docKey, "prn_bob"); + unsubscribeBob(); + + resolvers[0]?.(); + await flushMicrotasks(); + + // Settling the first flush must not destroy the doc out from under + // Bob's session — it must trigger a second flush that actually sees + // Bob's content, not silently discard it. And that second flush must + // see it on the *same* doc Alice's flush already captured — proven by + // both her content and Bob's showing up together, not a doc that got + // destroyed-and-recreated fresh (losing "first content") in between. + expect(flushedContents.length).toBe(2); + expect(flushedContents[1]).toContain("first content"); + expect(flushedContents[1]).toContain("bob's edit"); + + resolvers[1]?.(); + await flushMicrotasks(); + }); +}); diff --git a/packages/presence/test/routes.test.ts b/packages/presence/test/routes.test.ts index 5474d1951..ce565c328 100644 --- a/packages/presence/test/routes.test.ts +++ b/packages/presence/test/routes.test.ts @@ -290,6 +290,17 @@ describe("presence routes: doc sync", () => { }, ); + await alice.request(`/rooms/${ARTIFACT_SURFACE}/join`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: "{}", + }); + await bob.request(`/rooms/${ARTIFACT_SURFACE}/join`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: "{}", + }); + await alice.request(`/rooms/${ARTIFACT_SURFACE}/update`, { method: "POST", headers: { "content-type": "application/json" }, @@ -328,6 +339,11 @@ describe("presence routes: doc sync", () => { }, ); + await alice.request(`/rooms/${ARTIFACT_SURFACE}/join`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: "{}", + }); await alice.request(`/rooms/${ARTIFACT_SURFACE}/update`, { method: "POST", headers: { "content-type": "application/json" }, @@ -419,6 +435,11 @@ describe("presence routes: doc sync", () => { }, ); + await tenantA.request(`/rooms/${ARTIFACT_SURFACE}/join`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: "{}", + }); await tenantA.request(`/rooms/${ARTIFACT_SURFACE}/update`, { method: "POST", headers: { "content-type": "application/json" }, @@ -530,6 +551,11 @@ describe("presence routes: doc sync", () => { { tenantId: "tnt_a", principalId: "prn_alice" }, ); + await app.request(`/rooms/${ARTIFACT_SURFACE}/join`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: "{}", + }); const response = await app.request(`/rooms/${ARTIFACT_SURFACE}/update`, { method: "POST", headers: { "content-type": "application/json" }, From 21eaadcacc75c7b7f2b4b41cdde9916024f29701 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:54:21 -0700 Subject: [PATCH 2/2] Stop room teardown from discarding unflushed content or reviving empty MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit applyDocUpdate no longer auto-creates a room: it only ever writes into a room that's already open (via join or an active subscription), so a stale write arriving after a room has been torn down is rejected with PresenceRoomNotFoundError instead of silently recreating an empty, unseeded room for it to populate. destroyRoomIfEmpty now defers actually destroying a room's Y.Doc until every onEmpty listener's returned promise (persistence's teardown flush) has settled, and re-checks the room's identity and a monotonic per-room epoch before finalizing — a rejoin, edit, or reseed that lands while a flush is still in flight cancels the destroy or gets its own flush instead of being silently dropped. artifact-persistence.ts wires its flush's promise back through onEmpty so the registry can wait on it. This is a room-lifecycle guard, not an authorization check: doc writes stay decoupled from presence/awareness join by design, and grant-based write authorization is unchanged and lives upstream in routes.ts. --- packages/presence/src/artifact-persistence.ts | 64 ++++++--- packages/presence/src/index.ts | 1 + packages/presence/src/room-registry.test.ts | 11 +- packages/presence/src/room-registry.ts | 126 +++++++++++++++--- 4 files changed, 157 insertions(+), 45 deletions(-) diff --git a/packages/presence/src/artifact-persistence.ts b/packages/presence/src/artifact-persistence.ts index ff0ea2b16..73f5b83e5 100644 --- a/packages/presence/src/artifact-persistence.ts +++ b/packages/presence/src/artifact-persistence.ts @@ -93,7 +93,15 @@ export function createArtifactDocPersistence( clearTimeout(handle as ReturnType)); const timers = new Map(); - const pendingAuthor = new Map(); + // The single source of truth for "this room has a doc change that + // hasn't been captured into a write yet" — set on every `onDocChange` + // fire, cleared only at the moment `flush` reads the doc's current + // content into a queued write. Deliberately independent of `timers`: + // `flushNow` (called from `onEmpty`, right before the registry destroys + // the room's doc) must flush unflushed content whether or not a + // debounce timer happens to still be pending, or a room emptying + // between two scheduling events would silently drop its last edit. + const unflushedAuthor = new Map(); // Per-room write serialization: two snapshot writes for the same room // can otherwise both be in flight at once (a slow write #1 still // pending when write #2's own debounce elapses and starts). Since each @@ -144,8 +152,8 @@ export function createArtifactDocPersistence( function enqueueSnapshot( key: PresenceRoomKey, authorPrincipalId: string, - ): void { - if (artifactIdForSurface(key.surface) === null) return; + ): Promise { + if (artifactIdForSurface(key.surface) === null) return Promise.resolve(); const content = deps.registry.docText(key); const id = roomKeyId(key); const previous = writeChains.get(id) ?? Promise.resolve(); @@ -156,10 +164,24 @@ export function createArtifactDocPersistence( // swallows its own errors via `onSnapshotError`, but a defensive // catch here keeps one unexpected throw from permanently wedging // every later write queued behind it for this room. - writeChains.set( - id, - next.catch(() => undefined), - ); + const chained = next.catch(() => undefined); + writeChains.set(id, chained); + return chained; + } + + /** Reads whatever content is currently unflushed for `key` and hands it + * to `enqueueSnapshot`, clearing the room's dirty bit. Both the + * debounce timer's natural firing and `flushNow`'s bypass funnel + * through this one place so "is there unflushed content" has exactly + * one answer. Resolves once the write (if any was actually unflushed) + * has settled, so `flushNow` can hand that back to the registry as the + * promise it defers a pending room destroy on. */ + function flush(key: PresenceRoomKey): Promise { + const id = roomKeyId(key); + const author = unflushedAuthor.get(id); + if (author === undefined) return Promise.resolve(); + unflushedAuthor.delete(id); + return enqueueSnapshot(key, author); } function scheduleSnapshot( @@ -170,34 +192,32 @@ export function createArtifactDocPersistence( const id = roomKeyId(key); const existing = timers.get(id); if (existing !== undefined) clearTimeoutImpl(existing); - pendingAuthor.set(id, authorPrincipalId); + unflushedAuthor.set(id, authorPrincipalId); const handle = setTimeoutImpl(() => { timers.delete(id); - const author = pendingAuthor.get(id); - pendingAuthor.delete(id); - if (author !== undefined) enqueueSnapshot(key, author); + void flush(key); }, debounceMs); timers.set(id, handle); } - function flushNow(key: PresenceRoomKey): void { + function flushNow(key: PresenceRoomKey): Promise { const id = roomKeyId(key); const existing = timers.get(id); - const author = pendingAuthor.get(id); - if (existing === undefined || author === undefined) return; - clearTimeoutImpl(existing); - timers.delete(id); - pendingAuthor.delete(id); - enqueueSnapshot(key, author); + if (existing !== undefined) { + clearTimeoutImpl(existing); + timers.delete(id); + } + return flush(key); } deps.registry.onDocChange((key, authorPrincipalId) => { scheduleSnapshot(key, authorPrincipalId); }); - deps.registry.onEmpty((key) => { - flushNow(key); - }); + // Returning the flush's promise lets the registry defer actually + // destroying the room's `Y.Doc` until this settles — see + // `PresenceRoomRegistry.onEmpty`'s doc comment. + deps.registry.onEmpty((key) => flushNow(key)); return { async seedOnJoin(key) { @@ -211,7 +231,7 @@ export function createArtifactDocPersistence( dispose() { for (const handle of timers.values()) clearTimeoutImpl(handle); timers.clear(); - pendingAuthor.clear(); + unflushedAuthor.clear(); }, }; } diff --git a/packages/presence/src/index.ts b/packages/presence/src/index.ts index d93e5d4b4..aaaf1b6cf 100644 --- a/packages/presence/src/index.ts +++ b/packages/presence/src/index.ts @@ -1,6 +1,7 @@ export { createPresenceRoomRegistry, PRESENCE_DOC_TEXT_FIELD, + PresenceRoomNotFoundError, type PresenceRoomRegistry, type PresenceRoomKey, type PresenceRoomListener, diff --git a/packages/presence/src/room-registry.test.ts b/packages/presence/src/room-registry.test.ts index a9eadc6d7..41f7ea96f 100644 --- a/packages/presence/src/room-registry.test.ts +++ b/packages/presence/src/room-registry.test.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test"; import * as Y from "yjs"; import { createPresenceRoomRegistry, + PresenceRoomNotFoundError, type PresenceRoomKey, type PresenceState, } from "./room-registry"; @@ -210,7 +211,7 @@ describe("createPresenceRoomRegistry: doc sync", () => { Y.encodeStateAsUpdate(alice), "prn_alice", ), - ).toThrow(); + ).toThrow(PresenceRoomNotFoundError); expect(registry.docText(docKey)).toBe(""); }); @@ -237,7 +238,7 @@ describe("createPresenceRoomRegistry: doc sync", () => { Y.encodeStateAsUpdate(zombie), "prn_alice", ), - ).toThrow(); + ).toThrow(PresenceRoomNotFoundError); expect(registry.docText(docKey)).toBe(""); // A legitimate rejoin still sees a genuinely empty doc, ready for @@ -376,11 +377,7 @@ describe("createPresenceRoomRegistry: deferred destroy", () => { // the first flush resolves. const alice = clientDoc(); alice.getText("content").insert(0, "second content"); - registry.applyDocUpdate( - docKey, - Y.encodeStateAsUpdate(alice), - "prn_alice", - ); + registry.applyDocUpdate(docKey, Y.encodeStateAsUpdate(alice), "prn_alice"); resolveFirstFlush?.(); await flushMicrotasks(); diff --git a/packages/presence/src/room-registry.ts b/packages/presence/src/room-registry.ts index 84646b258..2775936c1 100644 --- a/packages/presence/src/room-registry.ts +++ b/packages/presence/src/room-registry.ts @@ -73,6 +73,16 @@ interface Room { readonly listeners: Set; readonly docListeners: Set; readonly snapshotListeners: Set; + /** + * Bumped on every `join` and `applyDocUpdate`/`seedDocText` call. A + * deferred destroy (see `destroyRoomIfEmpty`) captures this before + * dispatching `onEmpty` and compares it again once any pending flush + * settles — a mismatch means new membership or new content arrived + * during the flush, so the room isn't actually stale enough to destroy + * yet, and gets a fresh empty-check (and fresh flush) instead of being + * torn down out from under whatever just landed on it. + */ + epoch: number; } export interface PresenceRoomRegistry { @@ -106,10 +116,22 @@ export interface PresenceRoomRegistry { /** * Applies a Yjs update (already decoded from the wire's base64) to the * room's doc, attributing it to `authorPrincipalId` for persistence's - * "who wrote this snapshot" bookkeeping. Creates the room if it doesn't - * exist yet — mirrors `join`'s own "first call creates" behavior. - * Throws if `update` isn't a well-formed Yjs update; the caller (the - * HTTP route) is expected to turn that into a 400. + * "who wrote this snapshot" bookkeeping. Unlike `join`, this never + * creates a room: a room only exists once something has opened it + * (`join` or a live `subscribe*`), and a write against a room nobody + * currently holds open is exactly the "zombie" case (a stale POST + * arriving after the room emptied and was torn down) that would + * otherwise populate a freshly recreated, not-yet-seeded doc and + * permanently defeat `seedOnJoin`'s "only seed an empty doc" guard. + * Throws `PresenceRoomNotFoundError` for that case, and throws + * (unspecified) if `update` isn't a well-formed Yjs update — the caller + * (the HTTP route) is expected to turn either into an error response. + * + * This is a room-lifecycle guard, not an authorization check: it does + * not require `authorPrincipalId` to currently be a joined member (doc + * edits are deliberately decoupled from presence/awareness join in this + * design). Whether `authorPrincipalId` is allowed to write at all is the + * caller's `requireGrant` check, upstream of this call. */ applyDocUpdate( key: PresenceRoomKey, @@ -135,8 +157,16 @@ export interface PresenceRoomRegistry { onDocChange( listener: (key: PresenceRoomKey, authorPrincipalId: string) => void, ): () => void; - /** Fires when a room is torn down for having no members and no SSE subscribers left — persistence's cue to flush a pending snapshot immediately rather than wait out the debounce window on a doc nobody is looking at anymore. */ - onEmpty(listener: (key: PresenceRoomKey) => void): () => void; + /** + * Fires when a room is about to be torn down for having no members and + * no SSE subscribers left — persistence's cue to flush a pending + * snapshot immediately rather than wait out the debounce window on a + * doc nobody is looking at anymore. A listener may return a `Promise`; + * the registry defers the actual `Y.Doc` destruction until every + * listener's promise has settled, so a flush that's still in flight is + * guaranteed to see a live doc for its entire duration. + */ + onEmpty(listener: (key: PresenceRoomKey) => void | Promise): () => void; /** * Announces that a snapshot write for the room finished — persistence @@ -157,12 +187,33 @@ function roomKeyId(key: PresenceRoomKey): string { return `${key.tenantId}::${key.surface}`; } +/** Thrown by `applyDocUpdate` for a room nobody currently holds open — see + * that method's doc comment. Exported so a caller (e.g. the HTTP routes) + * could map it to a distinct response instead of folding it into + * "malformed update"; today's route still does the latter, since + * distinguishing the two response codes is out of this change's scope. */ +export class PresenceRoomNotFoundError extends Error { + constructor(key: PresenceRoomKey) { + super(`no open presence room for ${roomKeyId(key)}`); + this.name = "PresenceRoomNotFoundError"; + } +} + export function createPresenceRoomRegistry(): PresenceRoomRegistry { const rooms = new Map(); const docChangeListeners = new Set< (key: PresenceRoomKey, authorPrincipalId: string) => void >(); - const emptyListeners = new Set<(key: PresenceRoomKey) => void>(); + const emptyListeners = new Set< + (key: PresenceRoomKey) => void | Promise + >(); + // Room ids currently mid-teardown: their `onEmpty` listeners have fired + // and destruction is waiting on the returned promises to settle. Guards + // against a `leave`/`sweepStale`/unsubscribe that fires while that's + // still in flight re-dispatching `onEmpty` a second time for the same + // teardown — the eventual `finalize` re-checks emptiness (and, via + // `epoch`, freshness) on its own once the first round settles. + const pendingDestroys = new Set(); let nextClientId = 1; function ensureRoom(key: PresenceRoomKey): Room { @@ -185,6 +236,7 @@ export function createPresenceRoomRegistry(): PresenceRoomRegistry { listeners: new Set(), docListeners: new Set(), snapshotListeners: new Set(), + epoch: 0, }; rooms.set(id, room); } @@ -224,17 +276,55 @@ export function createPresenceRoomRegistry(): PresenceRoomRegistry { room.awareness.meta.set(clientId, { clock, lastUpdated: Date.now() }); } + function isRoomEmpty(room: Room): boolean { + return room.clientIdByPrincipal.size === 0 && room.listeners.size === 0; + } + function destroyRoomIfEmpty(key: PresenceRoomKey, room: Room): void { - if (room.clientIdByPrincipal.size > 0 || room.listeners.size > 0) return; + const id = roomKeyId(key); + // `room` may be a stale reference captured by a closure (e.g. + // `subscribe`'s returned unsubscribe) from before the room was last + // torn down and recreated under the same key — never act on anything + // but the room currently registered. + if (rooms.get(id) !== room) return; + if (!isRoomEmpty(room)) return; + if (pendingDestroys.has(id)) return; + pendingDestroys.add(id); + + const epochAtDispatch = room.epoch; // Fired synchronously, before the doc is destroyed: a listener (e.g. // persistence's flush-on-empty) that reads `docText`/`doc` needs it - // intact for the duration of this call — any `await` inside a - // listener naturally suspends past the destroy below, but every - // synchronous read within the listener still sees a live doc. - for (const listener of emptyListeners) listener(key); - room.awareness.destroy(); - room.doc.destroy(); - rooms.delete(roomKeyId(key)); + // intact for the duration of this call — every synchronous read + // within a listener sees a live doc regardless of whether it also + // returns a promise. A returned promise defers `finalize` below, but + // never the synchronous dispatch itself. + const pendingFlushes = [...emptyListeners] + .map((listener) => listener(key)) + .filter((result): result is Promise => result instanceof Promise); + + const finalize = (): void => { + pendingDestroys.delete(id); + if (rooms.get(id) !== room) return; + if (!isRoomEmpty(room)) return; + if (room.epoch !== epochAtDispatch) { + // A join or a doc/text change landed on the still-live room while + // its flush was in flight — the room is empty again now, but with + // content newer than what just got flushed. Re-run the whole + // empty check so that content gets its own `onEmpty` dispatch (and + // its own flush) instead of being silently destroyed unpersisted. + destroyRoomIfEmpty(key, room); + return; + } + room.awareness.destroy(); + room.doc.destroy(); + rooms.delete(id); + }; + + if (pendingFlushes.length === 0) { + finalize(); + } else { + void Promise.allSettled(pendingFlushes).then(finalize); + } } return { @@ -247,6 +337,7 @@ export function createPresenceRoomRegistry(): PresenceRoomRegistry { room.clientIdByPrincipal.set(state.principalId, clientId); } room.lastSeenAtByPrincipal.set(state.principalId, now); + room.epoch += 1; writeAwarenessState(room, clientId, state); broadcast(room); return currentStates(room); @@ -315,7 +406,9 @@ export function createPresenceRoomRegistry(): PresenceRoomRegistry { }, applyDocUpdate(key, update, authorPrincipalId) { - const room = ensureRoom(key); + const room = rooms.get(roomKeyId(key)); + if (room === undefined) throw new PresenceRoomNotFoundError(key); + room.epoch += 1; Y.applyUpdate(room.doc, update, "remote"); for (const listener of room.docListeners) { listener(update, authorPrincipalId); @@ -340,6 +433,7 @@ export function createPresenceRoomRegistry(): PresenceRoomRegistry { const room = ensureRoom(key); const yText = room.doc.getText(PRESENCE_DOC_TEXT_FIELD); if (yText.length > 0) return false; + room.epoch += 1; yText.insert(0, text); return true; },