Skip to content

Commit ebfa1c3

Browse files
committed
fix(realtime): preserve concurrent owner invalidation subscriptions
1 parent 9aed0f6 commit ebfa1c3

3 files changed

Lines changed: 217 additions & 83 deletions

File tree

‎apps/realtime/src/handlers/workspace-invalidation-room.test.ts‎

Lines changed: 193 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
1-
import { WORKSPACE_LIST_ROOM_TYPES } from '@sim/realtime-protocol/rooms'
1+
import { INVALIDATION_ROOM_TYPES, invalidationRoomIdKey } from '@sim/realtime-protocol/rooms'
2+
import { createDeferred } from '@sim/testing'
23
import { databaseMock } from '@sim/testing/mocks/database.mock'
34
import { sleep } from '@sim/utils/helpers'
45
import { beforeEach, describe, expect, it, vi } from 'vitest'
@@ -14,17 +15,22 @@ vi.mock('@sim/platform-authz/rooms', () => ({
1415
authorizeRoom: mockAuthorizeRoom,
1516
}))
1617

18+
vi.mock('@/handlers/file-list-app', () => ({
19+
fetchProjectRoomAccess: mockAuthorizeRoom,
20+
}))
21+
1722
import { setupWorkspaceInvalidationRoom } from '@/handlers/workspace-invalidation-room'
1823
import { beginRoomPermissionRead, commitRoomPermission } from '@/middleware/permissions'
1924

20-
type Payload = { workspaceId?: string }
25+
type Payload = { workspaceId?: string; projectId?: string }
2126

2227
function createSocket(overrides?: Record<string, unknown>) {
2328
const handlers: Record<string, (payload?: Payload) => Promise<void> | void> = {}
2429
// Live Set so the handler's native `socket.rooms` membership tracking works in tests.
2530
const rooms = new Set<string>()
2631
const socket = {
2732
id: 'socket-1',
33+
disconnected: false,
2834
userId: 'user-1',
2935
userName: 'Test User',
3036
userImage: 'avatar.png',
@@ -68,14 +74,13 @@ function createRoomManager(overrides?: Partial<IRoomManager>): IRoomManager {
6874
} as unknown as IRoomManager
6975
}
7076

71-
// The presence-free live-list rooms share one implementation; run the whole suite against each
72-
// so they can never drift. Event names and room names derive from the room type.
73-
describe.each(WORKSPACE_LIST_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (roomType) => {
77+
/** All invalidation room types share authorization and cancellation behavior. */
78+
describe.each(INVALIDATION_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (roomType) => {
7479
const joinEvent = `join-${roomType}`
7580
const successEvent = `${joinEvent}-success`
7681
const errorEvent = `${joinEvent}-error`
7782
const leaveEvent = `leave-${roomType}`
78-
const _roomOf = (workspaceId: string) => `${roomType}:${workspaceId}`
83+
const idKey = invalidationRoomIdKey(roomType)
7984

8085
const setup = (socket: ReturnType<typeof createSocket>['socket'], roomManager: IRoomManager) =>
8186
setupWorkspaceInvalidationRoom(
@@ -103,7 +108,7 @@ describe.each(WORKSPACE_LIST_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (
103108
const { socket, handlers } = createSocket()
104109
setup(socket, createRoomManager())
105110

106-
await handlers[joinEvent]({ workspaceId: 'ws-1' })
111+
await handlers[joinEvent]({ [idKey]: 'ws-1' })
107112

108113
expect(socket.emit).toHaveBeenCalledWith(
109114
errorEvent,
@@ -112,11 +117,7 @@ describe.each(WORKSPACE_LIST_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (
112117
})
113118

114119
it('aborts a join superseded during the access re-check await', async () => {
115-
// The access re-resolve is an await like any other: a leave landing during it must
116-
// still cancel this join, or the stale join would leave the room the client
117-
// switched to and commit the abandoned one. Forced down the re-resolve's DB path
118-
// by expiring the cached decision mid-join, so the interleaving is deterministic
119-
// rather than dependent on microtask ordering.
120+
/** Expire the cache to exercise a leave while the permission re-check is pending. */
120121
vi.useFakeTimers()
121122
try {
122123
const { handlers, socket } = createSocket({ id: 'socket-sup', userId: 'user-sup' })
@@ -141,12 +142,12 @@ describe.each(WORKSPACE_LIST_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (
141142
await sleep(31_000)
142143
} else {
143144
// Second call is the re-check's re-resolve: the client leaves during it.
144-
handlers[leaveEvent]({ workspaceId: 'ws-sup' })
145+
handlers[leaveEvent]({ [idKey]: 'ws-sup' })
145146
}
146147
return { allowed: true, status: 200, workspaceId: 'ws-sup', workspacePermission: 'admin' }
147148
})
148149

149-
const joining = handlers[joinEvent]({ workspaceId: 'ws-sup' })
150+
const joining = handlers[joinEvent]({ [idKey]: 'ws-sup' })
150151
await vi.advanceTimersByTimeAsync(31_000)
151152
await joining
152153

@@ -178,7 +179,7 @@ describe.each(WORKSPACE_LIST_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (
178179
return { allowed: true, status: 200, workspaceId: 'ws-race', workspacePermission: 'admin' }
179180
})
180181

181-
await handlers[joinEvent]({ workspaceId: 'ws-race' })
182+
await handlers[joinEvent]({ [idKey]: 'ws-race' })
182183

183184
expect(socket.emit).toHaveBeenCalledWith(
184185
errorEvent,
@@ -187,3 +188,180 @@ describe.each(WORKSPACE_LIST_ROOM_TYPES)('setupWorkspaceInvalidationRoom(%s)', (
187188
expect(socket.join).not.toHaveBeenCalled()
188189
})
189190
})
191+
192+
/** Exercise every owner address through the real authorization and membership handler. */
193+
describe.each(INVALIDATION_ROOM_TYPES)('concurrent owner subscriptions (%s)', (roomType) => {
194+
const idKey = invalidationRoomIdKey(roomType)
195+
const payload = (id: string): Payload => ({ [idKey]: id })
196+
const joinEvent = `join-${roomType}`
197+
const leaveEvent = `leave-${roomType}`
198+
const successEvent = `${joinEvent}-success`
199+
const errorEvent = `${joinEvent}-error`
200+
const allowed = { allowed: true, status: 200, workspacePermission: 'admin' }
201+
202+
function setup() {
203+
const state = createSocket({ disconnected: false })
204+
setupWorkspaceInvalidationRoom(
205+
state.socket as unknown as Parameters<typeof setupWorkspaceInvalidationRoom>[0],
206+
createRoomManager(),
207+
roomType
208+
)
209+
return state
210+
}
211+
212+
function pendingAuthorization() {
213+
const pending = createDeferred<typeof allowed>()
214+
mockAuthorizeRoom.mockImplementationOnce(() => pending.promise)
215+
return pending
216+
}
217+
218+
beforeEach(() => {
219+
mockAuthorizeRoom.mockReset().mockResolvedValue(allowed)
220+
})
221+
222+
it('keeps both owners subscribed after sequential joins', async () => {
223+
const { handlers, rooms } = setup()
224+
await handlers[joinEvent](payload('owner-a'))
225+
await handlers[joinEvent](payload('owner-b'))
226+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`, `${roomType}:owner-b`]))
227+
})
228+
229+
it('allows independent joins to finish in reverse order', async () => {
230+
const { handlers, rooms } = setup()
231+
const first = pendingAuthorization()
232+
const joining = handlers[joinEvent](payload('owner-a'))
233+
await handlers[joinEvent](payload('owner-b'))
234+
first.resolve(allowed)
235+
await joining
236+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`, `${roomType}:owner-b`]))
237+
})
238+
239+
it('scoped leave cancels only its pending owner and retains other membership', async () => {
240+
const { handlers, rooms } = setup()
241+
await handlers[joinEvent](payload('owner-c'))
242+
const first = pendingAuthorization()
243+
const joiningA = handlers[joinEvent](payload('owner-a'))
244+
const second = pendingAuthorization()
245+
const joiningB = handlers[joinEvent](payload('owner-b'))
246+
handlers[leaveEvent](payload('owner-a'))
247+
second.resolve(allowed)
248+
first.resolve(allowed)
249+
await Promise.all([joiningA, joiningB])
250+
expect(rooms).toEqual(new Set([`${roomType}:owner-b`, `${roomType}:owner-c`]))
251+
})
252+
253+
it('scoped leave removes only the specified joined owner', async () => {
254+
const { handlers, rooms } = setup()
255+
await handlers[joinEvent](payload('owner-a'))
256+
await handlers[joinEvent](payload('owner-b'))
257+
handlers[leaveEvent](payload('owner-b'))
258+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`]))
259+
})
260+
261+
it('leave all cancels every pending join and removes only this room type', async () => {
262+
const { handlers, rooms } = setup()
263+
rooms.add('other:owner')
264+
await handlers[joinEvent](payload('owner-c'))
265+
const first = pendingAuthorization()
266+
const joiningA = handlers[joinEvent](payload('owner-a'))
267+
const second = pendingAuthorization()
268+
const joiningB = handlers[joinEvent](payload('owner-b'))
269+
handlers[leaveEvent]()
270+
first.resolve(allowed)
271+
second.resolve(allowed)
272+
await Promise.all([joiningA, joiningB])
273+
expect(rooms).toEqual(new Set(['other:owner']))
274+
})
275+
276+
it('does not revive an old attempt after leave and rejoin of the same owner', async () => {
277+
const { handlers, socket, rooms } = setup()
278+
const first = pendingAuthorization()
279+
const joining = handlers[joinEvent](payload('owner-a'))
280+
handlers[leaveEvent](payload('owner-a'))
281+
await handlers[joinEvent](payload('owner-a'))
282+
first.resolve(allowed)
283+
await joining
284+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`]))
285+
expect(socket.emit.mock.calls.filter(([event]) => event === successEvent)).toHaveLength(1)
286+
})
287+
288+
it('supersedes a duplicate pending join for the same owner only', async () => {
289+
const { handlers, socket, rooms } = setup()
290+
const first = pendingAuthorization()
291+
const joining = handlers[joinEvent](payload('owner-a'))
292+
await handlers[joinEvent](payload('owner-b'))
293+
await handlers[joinEvent](payload('owner-a'))
294+
first.resolve(allowed)
295+
await joining
296+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`, `${roomType}:owner-b`]))
297+
expect(socket.emit.mock.calls.filter(([event]) => event === successEvent)).toHaveLength(2)
298+
})
299+
300+
it('suppresses stale authorization errors after the owner has rejoined', async () => {
301+
const { handlers, socket, rooms } = setup()
302+
const first = pendingAuthorization()
303+
const joining = handlers[joinEvent](payload('owner-a'))
304+
handlers[leaveEvent](payload('owner-a'))
305+
await handlers[joinEvent](payload('owner-a'))
306+
first.reject(new Error('Delayed authorization failure'))
307+
await joining
308+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`]))
309+
expect(socket.emit).not.toHaveBeenCalledWith(errorEvent, expect.anything())
310+
})
311+
312+
it('does not clear a newer pending attempt when the old attempt finishes', async () => {
313+
const { handlers, socket, rooms } = setup()
314+
const first = pendingAuthorization()
315+
const joiningA = handlers[joinEvent](payload('owner-a'))
316+
const second = pendingAuthorization()
317+
const joiningAgain = handlers[joinEvent](payload('owner-a'))
318+
first.resolve(allowed)
319+
await joiningA
320+
expect(rooms.size).toBe(0)
321+
second.resolve(allowed)
322+
await joiningAgain
323+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`]))
324+
expect(socket.emit.mock.calls.filter(([event]) => event === successEvent)).toHaveLength(1)
325+
})
326+
327+
it('preserves another owner while rejecting a revoked pending join', async () => {
328+
const { handlers, socket, rooms } = setup()
329+
await handlers[joinEvent](payload('owner-b'))
330+
const first = pendingAuthorization()
331+
const joining = handlers[joinEvent](payload('owner-a'))
332+
commitRoomPermission(
333+
socket.userId,
334+
{ type: roomType, id: 'owner-a' },
335+
null,
336+
beginRoomPermissionRead()
337+
)
338+
first.resolve(allowed)
339+
await joining
340+
expect(rooms).toEqual(new Set([`${roomType}:owner-b`]))
341+
expect(socket.emit).toHaveBeenCalledWith(
342+
errorEvent,
343+
expect.objectContaining({ [idKey]: 'owner-a', code: 'ACCESS_DENIED' })
344+
)
345+
})
346+
347+
it('does not commit any pending joins after disconnect', async () => {
348+
const { handlers, socket, rooms } = setup()
349+
const first = pendingAuthorization()
350+
const joining = handlers[joinEvent](payload('owner-a'))
351+
socket.disconnected = true
352+
first.resolve(allowed)
353+
await joining
354+
expect(rooms.size).toBe(0)
355+
expect(socket.emit).not.toHaveBeenCalledWith(successEvent, expect.anything())
356+
})
357+
358+
it('rejects malformed joins without cancelling a valid pending owner', async () => {
359+
const { handlers, rooms } = setup()
360+
const first = pendingAuthorization()
361+
const joining = handlers[joinEvent](payload('owner-a'))
362+
await handlers[joinEvent](payload(''))
363+
first.resolve(allowed)
364+
await joining
365+
expect(rooms).toEqual(new Set([`${roomType}:owner-a`]))
366+
})
367+
})

0 commit comments

Comments
 (0)