Skip to content

Commit a9fd086

Browse files
committed
fix(files): require explicit owners in live documents
1 parent 845b9c5 commit a9fd086

14 files changed

Lines changed: 232 additions & 137 deletions

File tree

‎.github/workflows/test-build.yml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -128,7 +128,7 @@ jobs:
128128
http-e2e:
129129
name: End-to-end over real HTTP
130130
runs-on: ${{ (vars.CI_PROVIDER == '' || vars.CI_PROVIDER == 'blacksmith') && 'blacksmith-8vcpu-ubuntu-2404' || 'ubuntu-latest' }}
131-
timeout-minutes: 40
131+
timeout-minutes: 20
132132
services:
133133
redis:
134134
image: redis:8.2-alpine

‎apps/realtime/src/handlers/file-doc-owner.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -68,9 +68,9 @@ const OWNER_ADAPTERS: Record<FileDocOwner['entityType'], FileDocOwnerAdapter> =
6868
},
6969
}
7070

71-
/** Only registered scopes can use callbacks; legacy workspace addresses resolve their owner at join. */
72-
export function fileDocOwnerAdapter(owner?: FileDocOwner): FileDocOwnerAdapter {
73-
const type = owner?.entityType ?? 'workspace'
71+
/** Select callbacks only after the document's canonical owner has been resolved. */
72+
export function fileDocOwnerAdapter(owner: FileDocOwner): FileDocOwnerAdapter {
73+
const type = owner.entityType
7474
if (!Object.hasOwn(OWNER_ADAPTERS, type)) throw new Error('Unsupported document owner')
7575
return OWNER_ADAPTERS[type]
7676
}

‎apps/realtime/src/handlers/file-doc.multireplica.test.ts‎

Lines changed: 31 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -62,9 +62,11 @@ describe('applyMarkdownToLiveFileDoc — multi-replica (store-enabled) ordering'
6262

6363
it('drops a stale durable write against the SHARED synced version', async () => {
6464
// A durable write (e.g. a concurrent human save on another process) records the shared synced version.
65-
expect(await applyMarkdownToLiveFileDoc('file-1', '# durable', { version: 100 })).toBe(
66-
'applied'
67-
)
65+
expect(
66+
await applyMarkdownToLiveFileDoc({ type: 'workspace-file-doc', id: 'file-1' }, '# durable', {
67+
version: 100,
68+
})
69+
).toBe('applied')
6870
expect(fakeStore.setSyncedVersion).toHaveBeenCalledWith(ROOM_NAME, 100, 'shared-generation')
6971
expect(fakeStore.getStreamState).toHaveBeenCalledWith(ROOM_NAME, 'shared-generation')
7072
expect(fakeStore.publishAndWait).toHaveBeenCalledWith(
@@ -76,15 +78,23 @@ describe('applyMarkdownToLiveFileDoc — multi-replica (store-enabled) ordering'
7678

7779
// A durable write with an OLDER version than the SHARED synced version is stale — rejected under the
7880
// lock before any diff is built, so it can't regress the doc across replicas.
79-
expect(await applyMarkdownToLiveFileDoc('file-1', '# older durable', { version: 50 })).toBe(
80-
'stale'
81-
)
81+
expect(
82+
await applyMarkdownToLiveFileDoc(
83+
{ type: 'workspace-file-doc', id: 'file-1' },
84+
'# older durable',
85+
{ version: 50 }
86+
)
87+
).toBe('stale')
8288
expect(mockFetchFileDocMerge).not.toHaveBeenCalled()
8389

8490
// A newer durable write applies and advances the shared synced version.
85-
expect(await applyMarkdownToLiveFileDoc('file-1', '# durable again', { version: 150 })).toBe(
86-
'applied'
87-
)
91+
expect(
92+
await applyMarkdownToLiveFileDoc(
93+
{ type: 'workspace-file-doc', id: 'file-1' },
94+
'# durable again',
95+
{ version: 150 }
96+
)
97+
).toBe('applied')
8898
expect(fakeStore.setSyncedVersion).toHaveBeenCalledWith(ROOM_NAME, 150, 'shared-generation')
8999
// setSyncedVersion fired only for the two applied durable writes, never for the stale one.
90100
expect(fakeStore.setSyncedVersion).toHaveBeenCalledTimes(2)
@@ -98,17 +108,25 @@ describe('applyMarkdownToLiveFileDoc — multi-replica (store-enabled) ordering'
98108
fakeStore.isAgentStreaming.mockResolvedValue(true)
99109

100110
expect(
101-
await applyMarkdownToLiveFileDoc('file-1', '# streamed by a client', { version: 100 })
111+
await applyMarkdownToLiveFileDoc(
112+
{ type: 'workspace-file-doc', id: 'file-1' },
113+
'# streamed by a client',
114+
{ version: 100 }
115+
)
102116
).toBe('applied')
103117
expect(mockFetchFileDocMerge).not.toHaveBeenCalled() // content deferred to the client
104118
expect(fakeStore.publishAndWait).not.toHaveBeenCalled()
105119
expect(fakeStore.setSyncedVersion).toHaveBeenCalledWith(ROOM_NAME, 100, 'shared-generation') // version still recorded
106120

107121
// Once streaming stops the flag clears and the (now near-noop) durable merge resumes normally.
108122
fakeStore.isAgentStreaming.mockResolvedValue(false)
109-
expect(await applyMarkdownToLiveFileDoc('file-1', '# final durable', { version: 150 })).toBe(
110-
'applied'
111-
)
123+
expect(
124+
await applyMarkdownToLiveFileDoc(
125+
{ type: 'workspace-file-doc', id: 'file-1' },
126+
'# final durable',
127+
{ version: 150 }
128+
)
129+
).toBe('applied')
112130
expect(mockFetchFileDocMerge).toHaveBeenCalledTimes(1)
113131
})
114132
})

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

Lines changed: 24 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -376,7 +376,8 @@ describe('setupWorkspaceFileDocHandlers', () => {
376376
source.getText(FILE_DOC_FIELD).insert(0, 'Stale text')
377377
const acknowledge = vi.fn()
378378
try {
379-
if (timing === 'before append') await invalidateLiveFileDocument('file-1', 2)
379+
if (timing === 'before append')
380+
await invalidateLiveFileDocument({ type: 'workspace-file-doc', id: 'file-1' }, 2)
380381
sent.length = 0
381382
handlers[FILE_DOC_EVENTS.UPDATE](
382383
{
@@ -389,7 +390,7 @@ describe('setupWorkspaceFileDocHandlers', () => {
389390
)
390391
if (timing !== 'before append') {
391392
expect(publish).toHaveBeenCalledTimes(1)
392-
await invalidateLiveFileDocument('file-1', 2)
393+
await invalidateLiveFileDocument({ type: 'workspace-file-doc', id: 'file-1' }, 2)
393394
if (timing === 'during append and reseed') {
394395
mockFetchFileDocSeed.mockResolvedValue({
395396
...seedResult('Replacement', 'doc-new'),
@@ -949,19 +950,31 @@ describe('setupWorkspaceFileDocHandlers', () => {
949950
mockFetchFileDocMerge.mockResolvedValue(Y.encodeStateAsUpdate(new Y.Doc()))
950951

951952
// A newer durable version lands and is recorded as the synced version.
952-
expect(await applyMarkdownToLiveFileDoc('file-1', '# newer', { version: 100 })).toBe('applied')
953+
expect(
954+
await applyMarkdownToLiveFileDoc({ type: 'workspace-file-doc', id: 'file-1' }, '# newer', {
955+
version: 100,
956+
})
957+
).toBe('applied')
953958
mockFetchFileDocMerge.mockClear()
954959

955960
// An older durable version arriving out of order (e.g. a concurrent write on another process) is
956961
// stale: skipped before any diff is computed, so the live doc never regresses to older content and
957962
// no diff is published that a later persist could write back.
958-
expect(await applyMarkdownToLiveFileDoc('file-1', '# older, stale', { version: 50 })).toBe(
959-
'stale'
960-
)
963+
expect(
964+
await applyMarkdownToLiveFileDoc(
965+
{ type: 'workspace-file-doc', id: 'file-1' },
966+
'# older, stale',
967+
{ version: 50 }
968+
)
969+
).toBe('stale')
961970
// The same version is idempotent — also skipped.
962-
expect(await applyMarkdownToLiveFileDoc('file-1', '# same version', { version: 100 })).toBe(
963-
'stale'
964-
)
971+
expect(
972+
await applyMarkdownToLiveFileDoc(
973+
{ type: 'workspace-file-doc', id: 'file-1' },
974+
'# same version',
975+
{ version: 100 }
976+
)
977+
).toBe('stale')
965978
expect(mockFetchFileDocMerge).not.toHaveBeenCalled()
966979
})
967980

@@ -979,8 +992,8 @@ describe('setupWorkspaceFileDocHandlers', () => {
979992
.mockReturnValueOnce(new Promise((resolve) => (resolveFirst = resolve)))
980993
.mockResolvedValueOnce(noOpUpdate)
981994

982-
const first = applyMarkdownToLiveFileDoc('file-1', '# One')
983-
const second = applyMarkdownToLiveFileDoc('file-1', '# Two')
995+
const first = applyMarkdownToLiveFileDoc({ type: 'workspace-file-doc', id: 'file-1' }, '# One')
996+
const second = applyMarkdownToLiveFileDoc({ type: 'workspace-file-doc', id: 'file-1' }, '# Two')
984997
await flushMicrotasks(SEED_CHAIN_TICKS)
985998
expect(mockFetchFileDocMerge).toHaveBeenCalledTimes(1) // second is queued behind the first
986999

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

Lines changed: 36 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -144,7 +144,7 @@ interface FileDocPresenceOwner {
144144
interface FileDocRoom {
145145
/** The `workspace_files.id` this room edits. */
146146
fileId: string
147-
owner?: FileDocOwner
147+
owner: FileDocOwner
148148
lastEditorConnectionId: string | null
149149
doc: Y.Doc
150150
awareness: awarenessProtocol.Awareness
@@ -786,12 +786,14 @@ interface MergeOrder {
786786
* reconcile treats it distinctly (retry later) rather than as "nothing to reconcile into".
787787
*/
788788
export function applyMarkdownToLiveFileDoc(
789-
fileId: string,
789+
ref: RoomRef,
790790
markdown: string,
791-
order: MergeOrder = {},
792-
owner?: FileDocOwner
791+
order: MergeOrder = {}
793792
): Promise<'applied' | 'no-live-room' | 'merge-unavailable' | 'stale'> {
794-
const name = roomName(fileDocRoom({ fileId, owner }))
793+
const address = fileDocTargetFromRoom(ref)
794+
if (!address) throw new Error('Invalid file document room')
795+
const { fileId } = address
796+
const name = roomName(ref)
795797
return serializeFileDocMutation(name, () => mergeMarkdownIntoRoom(name, fileId, markdown, order))
796798
}
797799

@@ -819,11 +821,11 @@ async function acquireFileDocMergeSlot(name: string): Promise<string | null> {
819821

820822
/** Serializes and version-orders an unsupported durable replacement with live Markdown merges. */
821823
export function invalidateLiveFileDocument(
822-
fileId: string,
823-
version: number,
824-
owner?: FileDocOwner
824+
ref: RoomRef,
825+
version: number
825826
): Promise<{ status: 'applied'; docId?: string } | { status: 'stale' }> {
826-
const name = roomName(fileDocRoom({ fileId, owner }))
827+
if (!fileDocTargetFromRoom(ref)) throw new Error('Invalid file document room')
828+
const name = roomName(ref)
827829
return serializeFileDocMutation(name, async () => {
828830
const store = getFileDocStore()
829831
const token = await acquireFileDocMergeSlot(name)
@@ -970,19 +972,13 @@ async function mergeMarkdownIntoRoom(
970972
return 'applied'
971973
}
972974

973-
function requireFileDocTarget(ref: RoomRef): FileDocTarget {
974-
const target = fileDocTargetFromRoom(ref)
975-
if (!target) throw new Error('Invalid file document room')
976-
return target
977-
}
978-
979975
/**
980976
* Get (or lazily create) the authoritative document for a room, wiring the two
981977
* relay handlers exactly once: document updates and awareness changes are
982978
* broadcast to the room.
983979
*/
984-
function getOrCreateRoom(io: Server, ref: RoomRef): FileDocRoom {
985-
const name = roomName(ref)
980+
function getOrCreateRoom(io: Server, target: FileDocTarget): FileDocRoom {
981+
const name = roomName(fileDocRoom(target))
986982
const existing = fileDocRooms.get(name)
987983
if (existing) return existing
988984

@@ -994,7 +990,7 @@ function getOrCreateRoom(io: Server, ref: RoomRef): FileDocRoom {
994990
// Started BEFORE the room is registered so no join can observe a room without its hydration handle.
995991
const hydrated = getFileDocStore().attachRoom(name, doc)
996992
const room: FileDocRoom = {
997-
...requireFileDocTarget(ref),
993+
...target,
998994
lastEditorConnectionId: null,
999995
doc,
1000996
awareness,
@@ -1129,10 +1125,10 @@ function isFileDocWriteAllowed(
11291125
name: string
11301126
): boolean | Promise<boolean> {
11311127
const userId = socket.userId
1132-
const fileId = fileDocRooms.get(name)?.fileId
1133-
if (!userId || !fileId) return false
1128+
const document = fileDocRooms.get(name)
1129+
if (!userId || !document) return false
11341130

1135-
const owner = fileDocRooms.get(name)?.owner
1131+
const { fileId, owner } = document
11361132
const room = fileDocRoom({ fileId, owner })
11371133
if (fileDocOwnerAdapter(owner).requiresCurrentActor)
11381134
return (async () => {
@@ -1343,8 +1339,8 @@ async function handleClientUpdate(
13431339
if (
13441340
!room ||
13451341
(target.owner &&
1346-
(target.owner.entityType !== room.owner?.entityType ||
1347-
target.owner.entityId !== room.owner?.entityId))
1342+
(target.owner.entityType !== room.owner.entityType ||
1343+
target.owner.entityId !== room.owner.entityId))
13481344
) {
13491345
reject('NOT_JOINED', true, candidate.updateId)
13501346
return
@@ -1485,7 +1481,6 @@ export function setupWorkspaceFileDocHandlers(
14851481
socket.on(FILE_DOC_EVENTS.JOIN, async (payload: JoinFileDocPayload) => {
14861482
const { fileId, clientId } = payload
14871483
const target = parseFileDocTarget(payload)
1488-
const owner = target?.owner
14891484
// Hoisted so the catch can tell whether this join was superseded (a switch to another file)
14901485
// before surfacing a retryable error for the abandoned one.
14911486
let generation: number | undefined
@@ -1547,9 +1542,9 @@ export function setupWorkspaceFileDocHandlers(
15471542
generation = joinGeneration.get(socket.id) ?? 0
15481543
}
15491544

1550-
const room = fileDocRoom({ fileId, owner })
1545+
const room = fileDocRoom(target)
15511546
const name = roomName(room)
1552-
const admissionName = fileDocAdmissionRoom({ fileId, owner })
1547+
const admissionName = fileDocAdmissionRoom(target)
15531548

15541549
const authorizeJoin = () =>
15551550
resolveRoomJoinAuth({
@@ -1569,11 +1564,21 @@ export function setupWorkspaceFileDocHandlers(
15691564
})
15701565
const authorized = await authorizeJoin()
15711566
if (!authorized) return
1572-
if (owner?.entityType === 'workspace' && owner.entityId !== authorized.workspaceId) {
1567+
if (
1568+
target.owner?.entityType === 'workspace' &&
1569+
target.owner.entityId !== authorized.workspaceId
1570+
) {
15731571
emitJoinError(socket, fileId, clientId, 'File not found', 'NOT_FOUND', false)
15741572
return
15751573
}
15761574

1575+
const owner =
1576+
target.owner ??
1577+
(authorized.workspaceId
1578+
? { entityType: 'workspace' as const, entityId: authorized.workspaceId }
1579+
: null)
1580+
if (!owner) throw new Error('Document authorization did not resolve its owner')
1581+
15771582
// Server-authenticated identity for the presence roster (never trusts the client-set
15781583
// awareness). Resolved here so the generation guard below also covers this await.
15791584
const avatarUrl = await resolveAvatarUrl(socket, userId)
@@ -1606,11 +1611,7 @@ export function setupWorkspaceFileDocHandlers(
16061611
discardInvalidatedRoom(name, io)
16071612
}
16081613

1609-
const entry = getOrCreateRoom(io, room)
1610-
// The workspace the server-side persist writes back to — and what the seed is built from, so it
1611-
// must be captured BEFORE the room is prepared below.
1612-
if (authorized.workspaceId)
1613-
entry.owner = { entityType: 'workspace', entityId: authorized.workspaceId }
1614+
const entry = getOrCreateRoom(io, { fileId, owner })
16141615

16151616
// Hold the room open across the awaits below: it has no owner until this join commits, so a
16161617
// concurrent last-leave would otherwise tear down the very document being prepared.
@@ -1795,12 +1796,12 @@ export function setupWorkspaceFileDocHandlers(
17951796
docId: docIdOf(entry.doc),
17961797
version: joinedVersion,
17971798
schemaVersion: FILE_DOC_SCHEMA_VERSION,
1799+
...fileDocOwnerWireFields(entry.owner),
17981800
...(store.enabled || fileDocOwnerAdapter(owner).requiresCurrentActor
17991801
? { acknowledgedUpdates: true as const }
18001802
: {}),
18011803
...(fileDocOwnerAdapter(owner).requiresCurrentActor
18021804
? {
1803-
...fileDocOwnerWireFields(entry.owner),
18041805
canWrite: finalPermission === 'write' || finalPermission === 'admin',
18051806
}
18061807
: {}),
@@ -1846,7 +1847,8 @@ export function setupWorkspaceFileDocHandlers(
18461847
} catch (error) {
18471848
logger.error('Error joining file-doc room:', error)
18481849
try {
1849-
const name = roomName(fileDocRoom({ fileId, owner }))
1850+
const name = target ? roomName(fileDocRoom(target)) : null
1851+
if (!name) throw error
18501852
/**
18511853
* Roll back ownership only if this attempt committed it. A failed provisional admission must
18521854
* preserve a previous file's binding and any co-mounted provider already in the target room.

0 commit comments

Comments
 (0)