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
47 changes: 39 additions & 8 deletions apps/realtime/src/handlers/connection.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { createLogger } from '@sim/logger'
import { parseRoomName, type RoomRef, roomName } from '@sim/realtime-protocol/rooms'
import { cleanupPendingSubblocksForSocket } from '@/handlers/subblocks'
import { cleanupPendingVariablesForSocket } from '@/handlers/variables'
import type { AuthenticatedSocket } from '@/middleware/auth'
Expand All @@ -15,20 +16,50 @@ export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager
logger.error(`Socket ${socket.id} connection error:`, error)
})

socket.on('disconnect', async (reason) => {
// `disconnecting` (not `disconnect`): here `socket.rooms` is still populated and
// authoritative, so presence is cleaned up even if the Redis room-set key was
// evicted or TTL-expired (which would leave the manager's stored rooms empty).
socket.on('disconnecting', async (reason) => {
try {
// Clean up pending debounce entries for this socket to prevent memory leaks
cleanupPendingSubblocksForSocket(socket.id)
cleanupPendingVariablesForSocket(socket.id)

const workflowIdHint = [...socket.rooms].find((roomId) => roomId !== socket.id)
const workflowId = await roomManager.removeUserFromRoom(socket.id, workflowIdHint)
// A socket may occupy multiple rooms (one per type). Remove it from every
// room the manager knows about.
const removedRooms = await roomManager.removeSocketFromAllRooms(socket.id)

if (workflowId) {
await roomManager.broadcastPresenceUpdate(workflowId)
logger.info(
`Socket ${socket.id} disconnected from workflow ${workflowId} (reason: ${reason})`
)
// Union with the live Socket.IO membership (authoritative here, and it
// survives a Redis eviction/TTL lapse that would leave the manager's tracked
// rooms empty). Attempt removal for any room the manager didn't already
// remove — best-effort, since a transient Redis error can't be recovered here.
const wasInRooms = new Map<string, RoomRef>()
for (const room of removedRooms) wasInRooms.set(roomName(room), room)
for (const name of socket.rooms) {
// `wasInRooms.has(name)` already excludes every room the manager removed
// (same room-name key via the roomName/parseRoomName bijection), so any
// room reaching here was NOT in `removedRooms` and needs a removal attempt.
if (name === socket.id || wasInRooms.has(name)) continue
const ref = parseRoomName(name)
if (!ref) continue
wasInRooms.set(name, ref)
await roomManager.removeUserFromRoom(ref, socket.id)
}
Comment thread
cursor[bot] marked this conversation as resolved.

// Broadcast a correction to every room this socket was in, EXCLUDING this
// socket — so it is never shown as a ghost collaborator even if its presence
// entry outlived a failed removal (transient Redis error; the hashes have no
// TTL). Any orphaned entry is additionally reclaimed by the next join's
// stale-presence sweep.
for (const room of wasInRooms.values()) {
await roomManager.broadcastPresenceUpdate(room, socket.id)
}
Comment thread
greptile-apps[bot] marked this conversation as resolved.

if (wasInRooms.size > 0) {
const rooms = Array.from(wasInRooms.values())
.map((room) => `${room.type}:${room.id}`)
.join(', ')
logger.info(`Socket ${socket.id} disconnected from [${rooms}] (reason: ${reason})`)
}
} catch (error) {
logger.error(`Error handling disconnect for socket ${socket.id}:`, error)
Expand Down
41 changes: 23 additions & 18 deletions apps/realtime/src/handlers/operations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,15 @@ import {
type VariableOperation,
WORKFLOW_OPERATIONS,
} from '@sim/realtime-protocol/constants'
import { ROOM_TYPES } from '@sim/realtime-protocol/rooms'
import { WorkflowOperationSchema } from '@sim/realtime-protocol/schemas'
import { getErrorMessage } from '@sim/utils/errors'
import { generateId } from '@sim/utils/id'
import { ZodError } from 'zod'
import { persistWorkflowOperation } from '@/database/operations'
import type { AuthenticatedSocket } from '@/middleware/auth'
import { checkWorkflowOperationPermission } from '@/middleware/permissions'
import type { IRoomManager, UserSession } from '@/rooms'
import { type IRoomManager, type UserSession, workflowRoom as wf } from '@/rooms'

const logger = createLogger('OperationsHandlers')

Expand Down Expand Up @@ -44,7 +45,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
let session: UserSession | null = null

try {
workflowId = await roomManager.getWorkflowIdForSocket(socket.id)
workflowId = (await roomManager.getRoomForSocket(socket.id, ROOM_TYPES.WORKFLOW))?.id ?? null
session = await roomManager.getUserSession(socket.id)
} catch (error) {
logger.error('Error loading session for workflow operation:', error)
Expand All @@ -65,7 +66,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager

let hasRoom = false
try {
hasRoom = await roomManager.hasWorkflowRoom(workflowId)
hasRoom = await roomManager.hasRoom(wf(workflowId))
} catch (error) {
logger.error('Error checking workflow room:', error)
emitOperationError(
Expand Down Expand Up @@ -98,14 +99,16 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
const operationTimestamp = isPositionUpdate ? timestamp : Date.now()

// Get user presence for permission checking
const users = await roomManager.getWorkflowUsers(workflowId)
const users = await roomManager.getRoomUsers(wf(workflowId))
const userPresence = users.find((u) => u.socketId === socket.id)

// Skip permission checks for non-committed position updates (broadcasts only, no persistence)
if (isPositionUpdate && !commitPositionUpdate) {
// Update last activity
if (userPresence) {
await roomManager.updateUserActivity(workflowId, socket.id, { lastActivity: Date.now() })
await roomManager.updateUserActivity(wf(workflowId), socket.id, {
lastActivity: Date.now(),
})
}
} else {
// Check permissions from cached role for all other operations
Expand All @@ -123,7 +126,9 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
return
}

await roomManager.updateUserActivity(workflowId, socket.id, { lastActivity: Date.now() })
await roomManager.updateUserActivity(wf(workflowId), socket.id, {
lastActivity: Date.now(),
})

// Re-validate the workspace role against the DB (cached per pod for a short
// window) so revoked or downgraded collaborators lose write access live.
Expand Down Expand Up @@ -198,7 +203,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
timestamp: operationTimestamp,
userId: session.userId,
})
await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

if (operationId) {
socket.emit('operation-confirmed', {
Expand Down Expand Up @@ -244,7 +249,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
timestamp: operationTimestamp,
userId: session.userId,
})
await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

if (operationId) {
socket.emit('operation-confirmed', { operationId, serverTimestamp: Date.now() })
Expand Down Expand Up @@ -277,7 +282,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

const broadcastData = {
operation,
Expand Down Expand Up @@ -317,7 +322,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

const broadcastData = {
operation,
Expand Down Expand Up @@ -354,7 +359,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -386,7 +391,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -415,7 +420,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -447,7 +452,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -479,7 +484,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -511,7 +516,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -540,7 +545,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

socket.to(workflowId).emit('workflow-operation', {
operation,
Expand Down Expand Up @@ -569,7 +574,7 @@ export function setupOperationsHandlers(socket: AuthenticatedSocket, roomManager
userId: session.userId,
})

await roomManager.updateRoomLastModified(workflowId)
await roomManager.updateRoomLastModified(wf(workflowId))

const broadcastData = {
operation,
Expand Down
21 changes: 11 additions & 10 deletions apps/realtime/src/handlers/presence.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { createLogger } from '@sim/logger'
import { ROOM_TYPES } from '@sim/realtime-protocol/rooms'
import type { AuthenticatedSocket } from '@/middleware/auth'
import type { IRoomManager } from '@/rooms'

Expand All @@ -7,16 +8,16 @@ const logger = createLogger('PresenceHandlers')
export function setupPresenceHandlers(socket: AuthenticatedSocket, roomManager: IRoomManager) {
socket.on('cursor-update', async ({ cursor }) => {
try {
const workflowId = await roomManager.getWorkflowIdForSocket(socket.id)
const room = await roomManager.getRoomForSocket(socket.id, ROOM_TYPES.WORKFLOW)
const session = await roomManager.getUserSession(socket.id)

if (!workflowId || !session) return
if (!room || !session) return

// Update cursor in room state
await roomManager.updateUserActivity(workflowId, socket.id, { cursor })
await roomManager.updateUserActivity(room, socket.id, { cursor })

// Broadcast to other users in the room
socket.to(workflowId).emit('cursor-update', {
// Broadcast to other users in the room (workflow room name is the bare id)
socket.to(room.id).emit('cursor-update', {
socketId: socket.id,
userId: session.userId,
userName: session.userName,
Expand All @@ -30,16 +31,16 @@ export function setupPresenceHandlers(socket: AuthenticatedSocket, roomManager:

socket.on('selection-update', async ({ selection }) => {
try {
const workflowId = await roomManager.getWorkflowIdForSocket(socket.id)
const room = await roomManager.getRoomForSocket(socket.id, ROOM_TYPES.WORKFLOW)
const session = await roomManager.getUserSession(socket.id)

if (!workflowId || !session) return
if (!room || !session) return

// Update selection in room state
await roomManager.updateUserActivity(workflowId, socket.id, { selection })
await roomManager.updateUserActivity(room, socket.id, { selection })

// Broadcast to other users in the room
socket.to(workflowId).emit('selection-update', {
// Broadcast to other users in the room (workflow room name is the bare id)
socket.to(room.id).emit('selection-update', {
socketId: socket.id,
userId: session.userId,
userName: session.userName,
Expand Down
12 changes: 7 additions & 5 deletions apps/realtime/src/handlers/subblocks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,13 @@ import { workflow, workflowBlocks } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { assertWorkflowMutable, WorkflowLockedError } from '@sim/platform-authz/workflow'
import { SUBBLOCK_OPERATIONS } from '@sim/realtime-protocol/constants'
import { ROOM_TYPES } from '@sim/realtime-protocol/rooms'
import { getErrorMessage } from '@sim/utils/errors'
import { isWorkflowBlockProtected } from '@sim/workflow-types/workflow'
import { and, eq } from 'drizzle-orm'
import type { AuthenticatedSocket } from '@/middleware/auth'
import { checkWorkflowOperationPermission } from '@/middleware/permissions'
import type { IRoomManager } from '@/rooms'
import { type IRoomManager, workflowRoom as wf } from '@/rooms'

const logger = createLogger('SubblocksHandlers')

Expand Down Expand Up @@ -69,7 +70,8 @@ export function setupSubblocksHandlers(socket: AuthenticatedSocket, roomManager:
}

try {
const sessionWorkflowId = await roomManager.getWorkflowIdForSocket(socket.id)
const sessionWorkflowId =
(await roomManager.getRoomForSocket(socket.id, ROOM_TYPES.WORKFLOW))?.id ?? null
const session = await roomManager.getUserSession(socket.id)

if (!sessionWorkflowId || !session) {
Expand Down Expand Up @@ -106,7 +108,7 @@ export function setupSubblocksHandlers(socket: AuthenticatedSocket, roomManager:
return
}

const hasRoom = await roomManager.hasWorkflowRoom(workflowId)
const hasRoom = await roomManager.hasRoom(wf(workflowId))
if (!hasRoom) {
logger.debug(`Ignoring subblock update: workflow room not found`, {
socketId: socket.id,
Expand All @@ -117,7 +119,7 @@ export function setupSubblocksHandlers(socket: AuthenticatedSocket, roomManager:
return
}

const users = await roomManager.getWorkflowUsers(workflowId)
const users = await roomManager.getRoomUsers(wf(workflowId))
const userPresence = users.find((user) => user.socketId === socket.id)
if (!userPresence) {
socket.emit('operation-forbidden', {
Expand Down Expand Up @@ -182,7 +184,7 @@ export function setupSubblocksHandlers(socket: AuthenticatedSocket, roomManager:
}

// Update user activity
await roomManager.updateUserActivity(workflowId, socket.id, { lastActivity: Date.now() })
await roomManager.updateUserActivity(wf(workflowId), socket.id, { lastActivity: Date.now() })

// Server-side debounce/coalesce by workflowId+blockId+subblockId
const debouncedKey = `${workflowId}:${blockId}:${subblockId}`
Expand Down
12 changes: 7 additions & 5 deletions apps/realtime/src/handlers/variables.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,12 @@ import { workflow } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { assertWorkflowMutable, WorkflowLockedError } from '@sim/platform-authz/workflow'
import { VARIABLE_OPERATIONS } from '@sim/realtime-protocol/constants'
import { ROOM_TYPES } from '@sim/realtime-protocol/rooms'
import { getErrorMessage } from '@sim/utils/errors'
import { eq } from 'drizzle-orm'
import type { AuthenticatedSocket } from '@/middleware/auth'
import { checkWorkflowOperationPermission } from '@/middleware/permissions'
import type { IRoomManager } from '@/rooms'
import { type IRoomManager, workflowRoom as wf } from '@/rooms'

const logger = createLogger('VariablesHandlers')

Expand Down Expand Up @@ -61,7 +62,8 @@ export function setupVariablesHandlers(socket: AuthenticatedSocket, roomManager:
}

try {
const sessionWorkflowId = await roomManager.getWorkflowIdForSocket(socket.id)
const sessionWorkflowId =
(await roomManager.getRoomForSocket(socket.id, ROOM_TYPES.WORKFLOW))?.id ?? null
const session = await roomManager.getUserSession(socket.id)

if (!sessionWorkflowId || !session) {
Expand Down Expand Up @@ -98,7 +100,7 @@ export function setupVariablesHandlers(socket: AuthenticatedSocket, roomManager:
return
}

const hasRoom = await roomManager.hasWorkflowRoom(workflowId)
const hasRoom = await roomManager.hasRoom(wf(workflowId))
if (!hasRoom) {
logger.debug(`Ignoring variable update: workflow room not found`, {
socketId: socket.id,
Expand All @@ -109,7 +111,7 @@ export function setupVariablesHandlers(socket: AuthenticatedSocket, roomManager:
return
}

const users = await roomManager.getWorkflowUsers(workflowId)
const users = await roomManager.getRoomUsers(wf(workflowId))
const userPresence = users.find((user) => user.socketId === socket.id)
if (!userPresence) {
socket.emit('operation-forbidden', {
Expand Down Expand Up @@ -174,7 +176,7 @@ export function setupVariablesHandlers(socket: AuthenticatedSocket, roomManager:
}

// Update user activity
await roomManager.updateUserActivity(workflowId, socket.id, { lastActivity: Date.now() })
await roomManager.updateUserActivity(wf(workflowId), socket.id, { lastActivity: Date.now() })

const debouncedKey = `${workflowId}:${variableId}:${field}`
const existing = pendingVariableUpdates.get(debouncedKey)
Expand Down
Loading