Skip to content
Merged
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
60 changes: 58 additions & 2 deletions packages/presence/src/artifact-persistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof createPresenceRoomRegistry>,
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
Expand Down Expand Up @@ -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) => {
Expand Down Expand Up @@ -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();

Expand All @@ -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();
Expand All @@ -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"),
Expand Down Expand Up @@ -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();
Expand All @@ -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);
Expand Down Expand Up @@ -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();
Expand Down
64 changes: 42 additions & 22 deletions packages/presence/src/artifact-persistence.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,15 @@ export function createArtifactDocPersistence(
clearTimeout(handle as ReturnType<typeof setTimeout>));

const timers = new Map<string, unknown>();
const pendingAuthor = new Map<string, string>();
// 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<string, string>();
// 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
Expand Down Expand Up @@ -144,8 +152,8 @@ export function createArtifactDocPersistence(
function enqueueSnapshot(
key: PresenceRoomKey,
authorPrincipalId: string,
): void {
if (artifactIdForSurface(key.surface) === null) return;
): Promise<void> {
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();
Expand All @@ -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<void> {
const id = roomKeyId(key);
const author = unflushedAuthor.get(id);
if (author === undefined) return Promise.resolve();
unflushedAuthor.delete(id);
return enqueueSnapshot(key, author);
}

function scheduleSnapshot(
Expand All @@ -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<void> {
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) {
Expand All @@ -211,7 +231,7 @@ export function createArtifactDocPersistence(
dispose() {
for (const handle of timers.values()) clearTimeoutImpl(handle);
timers.clear();
pendingAuthor.clear();
unflushedAuthor.clear();
},
};
}
1 change: 1 addition & 0 deletions packages/presence/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
export {
createPresenceRoomRegistry,
PRESENCE_DOC_TEXT_FIELD,
PresenceRoomNotFoundError,
type PresenceRoomRegistry,
type PresenceRoomKey,
type PresenceRoomListener,
Expand Down
Loading
Loading