Skip to content

Commit d9f721d

Browse files
icecrasher321claude
andcommitted
fix(realtime): close join/eviction race and make sweep cleanups independent
Re-authorize immediately before socket.join so an in-flight join cannot reverse a sweep eviction, and run the sweep's best-effort cleanups independently so a room-state failure cannot skip the presence broadcast. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1 parent 4a13540 commit d9f721d

4 files changed

Lines changed: 108 additions & 8 deletions

File tree

apps/realtime/src/access-revalidation.test.ts

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -163,6 +163,20 @@ describe('access-revalidation sweep', () => {
163163
expect(noRoom.leave).not.toHaveBeenCalled()
164164
})
165165

166+
it('still broadcasts presence when room-state removal fails', async () => {
167+
const socket = makeSocket('sock-1', 'user-1', 'wf-1')
168+
const manager = makeManager([socket], [{ socketId: 'sock-1', role: 'read' }])
169+
manager.removeUserFromRoom.mockRejectedValue(new Error('redis down'))
170+
mockResolveRole.mockResolvedValue(null)
171+
172+
const sweep = startAccessRevalidationSweep(manager)
173+
await sweep.runOnce()
174+
sweep.stop()
175+
176+
expect(socket.leave).toHaveBeenCalledWith('wf-1')
177+
expect(manager.broadcastPresenceUpdate).toHaveBeenCalledWith('wf-1')
178+
})
179+
166180
it('still evaluates access when presence lookup fails (falls back safely)', async () => {
167181
const socket = makeSocket('sock-1', 'user-1', 'wf-1')
168182
const manager = makeManager([socket])

apps/realtime/src/access-revalidation.ts

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -87,9 +87,24 @@ export function startAccessRevalidationSweep(roomManager: IRoomManager): AccessR
8787
`Revoked live access for user ${socket.userId} on workflow ${workflowId} (socket ${socket.id})`
8888
)
8989

90-
// Best-effort room-state cleanup; failure here does not restore access.
91-
await roomManager.removeUserFromRoom(socket.id, workflowId)
92-
await roomManager.broadcastPresenceUpdate(workflowId)
90+
// Best-effort cleanups; each is independent so one failure neither restores
91+
// access nor prevents the other from running.
92+
try {
93+
await roomManager.removeUserFromRoom(socket.id, workflowId)
94+
} catch (error) {
95+
logger.warn(
96+
`Failed to remove evicted socket ${socket.id} from room state for ${workflowId}`,
97+
error
98+
)
99+
}
100+
try {
101+
await roomManager.broadcastPresenceUpdate(workflowId)
102+
} catch (error) {
103+
logger.warn(
104+
`Failed to broadcast presence after evicting socket ${socket.id} from ${workflowId}`,
105+
error
106+
)
107+
}
93108
}
94109

95110
async function runOnce(): Promise<void> {

apps/realtime/src/handlers/workflow.test.ts

Lines changed: 54 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,12 @@
44
import { beforeEach, describe, expect, it, vi } from 'vitest'
55
import type { IRoomManager } from '@/rooms'
66

7-
const { mockGetWorkflowState, mockVerifyWorkflowAccess } = vi.hoisted(() => ({
8-
mockGetWorkflowState: vi.fn(),
9-
mockVerifyWorkflowAccess: vi.fn(),
10-
}))
7+
const { mockGetWorkflowState, mockVerifyWorkflowAccess, mockResolveCurrentWorkflowRole } =
8+
vi.hoisted(() => ({
9+
mockGetWorkflowState: vi.fn(),
10+
mockVerifyWorkflowAccess: vi.fn(),
11+
mockResolveCurrentWorkflowRole: vi.fn(),
12+
}))
1113

1214
vi.mock('@sim/db', () => ({
1315
db: { select: vi.fn() },
@@ -20,6 +22,7 @@ vi.mock('@/database/operations', () => ({
2022

2123
vi.mock('@/middleware/permissions', () => ({
2224
verifyWorkflowAccess: mockVerifyWorkflowAccess,
25+
resolveCurrentWorkflowRole: mockResolveCurrentWorkflowRole,
2326
}))
2427

2528
import { setupWorkflowHandlers } from '@/handlers/workflow'
@@ -86,6 +89,7 @@ describe('setupWorkflowHandlers', () => {
8689
vi.clearAllMocks()
8790
mockGetWorkflowState.mockResolvedValue({ id: 'workflow-1', state: {} })
8891
mockVerifyWorkflowAccess.mockResolvedValue({ hasAccess: true, role: 'admin' })
92+
mockResolveCurrentWorkflowRole.mockResolvedValue('admin')
8993
})
9094

9195
it('includes workflowId when authentication is missing', async () => {
@@ -149,6 +153,52 @@ describe('setupWorkflowHandlers', () => {
149153
})
150154
})
151155

156+
it('denies the join when access is revoked while the join is in flight', async () => {
157+
mockResolveCurrentWorkflowRole.mockResolvedValue(null)
158+
159+
const { socket, handlers } = createSocket()
160+
const roomManager = createRoomManager()
161+
162+
setupWorkflowHandlers(
163+
socket as unknown as Parameters<typeof setupWorkflowHandlers>[0],
164+
roomManager
165+
)
166+
167+
await handlers['join-workflow']({ workflowId: 'workflow-1', tabSessionId: 'tab-1' })
168+
169+
expect(socket.emit).toHaveBeenCalledWith('join-workflow-error', {
170+
workflowId: 'workflow-1',
171+
error: 'Access denied to workflow',
172+
code: 'ACCESS_DENIED',
173+
retryable: false,
174+
})
175+
expect(socket.join).not.toHaveBeenCalled()
176+
expect(roomManager.addUserToRoom).not.toHaveBeenCalled()
177+
})
178+
179+
it('joins with the re-validated role, passing the join-time role as fallback', async () => {
180+
mockVerifyWorkflowAccess.mockResolvedValue({ hasAccess: true, role: 'write' })
181+
mockResolveCurrentWorkflowRole.mockResolvedValue('read')
182+
183+
const { socket, handlers } = createSocket()
184+
const roomManager = createRoomManager()
185+
186+
setupWorkflowHandlers(
187+
socket as unknown as Parameters<typeof setupWorkflowHandlers>[0],
188+
roomManager
189+
)
190+
191+
await handlers['join-workflow']({ workflowId: 'workflow-1', tabSessionId: 'tab-1' })
192+
193+
expect(mockResolveCurrentWorkflowRole).toHaveBeenCalledWith('user-1', 'workflow-1', 'write')
194+
expect(socket.join).toHaveBeenCalledWith('workflow-1')
195+
expect(roomManager.addUserToRoom).toHaveBeenCalledWith(
196+
'workflow-1',
197+
'socket-1',
198+
expect.objectContaining({ role: 'read' })
199+
)
200+
})
201+
152202
it('marks workflow access verification failures as retryable', async () => {
153203
mockVerifyWorkflowAccess.mockRejectedValue(new Error('database unavailable'))
154204

apps/realtime/src/handlers/workflow.ts

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ import { createLogger } from '@sim/logger'
33
import { eq } from 'drizzle-orm'
44
import { getWorkflowState } from '@/database/operations'
55
import type { AuthenticatedSocket } from '@/middleware/auth'
6-
import { verifyWorkflowAccess } from '@/middleware/permissions'
6+
import { resolveCurrentWorkflowRole, verifyWorkflowAccess } from '@/middleware/permissions'
77
import type { IRoomManager, UserPresence } from '@/rooms'
88

99
const logger = createLogger('WorkflowHandlers')
@@ -131,6 +131,27 @@ export function setupWorkflowHandlers(socket: AuthenticatedSocket, roomManager:
131131
}
132132
}
133133

134+
// Re-authorize immediately before joining: the access-revalidation sweep
135+
// may have evicted this socket while the awaits above were in flight, and
136+
// its eviction is recorded in the shared role cache before it runs — so a
137+
// revoked user resolves to null here. No awaits sit between this check
138+
// and socket.join, so a sweep eviction cannot interleave after it and be
139+
// reversed by this join.
140+
const currentRole = await resolveCurrentWorkflowRole(userId, workflowId, userRole)
141+
if (currentRole === null) {
142+
logger.warn(
143+
`User ${userId} (${userName}) lost access to workflow ${workflowId} before join completed`
144+
)
145+
socket.emit('join-workflow-error', {
146+
workflowId,
147+
error: 'Access denied to workflow',
148+
code: 'ACCESS_DENIED',
149+
retryable: false,
150+
})
151+
return
152+
}
153+
userRole = currentRole
154+
134155
// Join the new room
135156
socket.join(workflowId)
136157

0 commit comments

Comments
 (0)