Skip to content
Open
4 changes: 2 additions & 2 deletions apps/docs/openapi-v2-files-audit.json
Original file line number Diff line number Diff line change
Expand Up @@ -3167,7 +3167,7 @@
"uploadedByEmail": {
"type": "string",
"format": "email",
"pattern": "^(?!\\.)(?!.*\\.\\.)([A-Za-z0-9_'+\\-\\.]*)[A-Za-z0-9_+-]@([A-Za-z0-9][A-Za-z0-9\\-]*\\.)+[A-Za-z]{2,}$",
"pattern": "^[a-zA-Z0-9.!#$%&'*+/=?^_`{|}~-]+@[a-zA-Z0-9](?:[a-zA-Z0-9-]{0,61}[a-zA-Z0-9])?(?:\\.[a-zA-Z0-9](?:[a-zA-Z0-9-]{0,61}[a-zA-Z0-9])?)*$",
"description": "Current email address of the uploader.",
"examples": ["jane@example.com"]
},
Expand Down Expand Up @@ -4029,7 +4029,7 @@
"uploadedByEmail": {
"type": "string",
"format": "email",
"pattern": "^(?!\\.)(?!.*\\.\\.)([A-Za-z0-9_'+\\-\\.]*)[A-Za-z0-9_+-]@([A-Za-z0-9][A-Za-z0-9\\-]*\\.)+[A-Za-z]{2,}$",
"pattern": "^[a-zA-Z0-9.!#$%&'*+/=?^_`{|}~-]+@[a-zA-Z0-9](?:[a-zA-Z0-9-]{0,61}[a-zA-Z0-9])?(?:\\.[a-zA-Z0-9](?:[a-zA-Z0-9-]{0,61}[a-zA-Z0-9])?)*$",
"description": "Current email address of the uploader.",
"examples": ["jane@example.com"]
},
Expand Down
73 changes: 73 additions & 0 deletions apps/realtime/src/handlers/connection.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
/**
* @vitest-environment node
*/
import { createServer, type Server as HttpServer } from 'node:http'
import type { AddressInfo } from 'node:net'
import { Server } from 'socket.io'
import { io as connect, type Socket } from 'socket.io-client'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { setupConnectionHandlers, waitForConnectionCleanup } from '@/handlers/connection'
import type { AuthenticatedSocket } from '@/middleware/auth'
import { MemoryRoomManager } from '@/rooms'

vi.mock('@/handlers/file-doc', () => ({ cleanupFileDocForSocket: vi.fn() }))
vi.mock('@/handlers/subblocks', () => ({ cleanupPendingSubblocksForSocket: vi.fn() }))
vi.mock('@/handlers/variables', () => ({ cleanupPendingVariablesForSocket: vi.fn() }))

describe('server shutdown connection drain', () => {
let httpServer: HttpServer
let io: Server
let manager: MemoryRoomManager
let client: Socket

beforeEach(async () => {
httpServer = createServer()
io = new Server(httpServer, { transports: ['websocket'] })
manager = new MemoryRoomManager(io)
await manager.initialize()
io.on('connection', (socket) => setupConnectionHandlers(socket as AuthenticatedSocket, manager))
await new Promise<void>((resolve) => httpServer.listen(0, '127.0.0.1', resolve))
const port = (httpServer.address() as AddressInfo).port
client = connect(`http://127.0.0.1:${port}`, { transports: ['websocket'], autoConnect: false })
const connected = new Promise<void>((resolve) => client.once('connect', resolve))
client.connect()
await connected
})

afterEach(async () => {
client.disconnect()
await io.close()
await waitForConnectionCleanup()
await manager.shutdown()
vi.restoreAllMocks()
})

it('keeps automatic reconnection active after transport shutdown', async () => {
const disconnected = new Promise<string>((resolve) => client.once('disconnect', resolve))
await io.close()
expect(await disconnected).toBe('transport close')
expect(client.active).toBe(true)
await waitForConnectionCleanup()
})

it('waits for asynchronous presence cleanup before releasing its dependencies', async () => {
let finishRemoval: (() => void) | undefined
vi.spyOn(manager, 'removeSocketFromAllRooms').mockImplementation(
() =>
new Promise((resolve) => {
finishRemoval = () => resolve([])
})
)
await io.close()
let drained = false
const drain = waitForConnectionCleanup().then(() => {
drained = true
})
await Promise.resolve()
expect(drained).toBe(false)
expect(finishRemoval).toBeDefined()
finishRemoval?.()
await drain
expect(drained).toBe(true)
})
})
15 changes: 14 additions & 1 deletion apps/realtime/src/handlers/connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,13 @@ const logger = createLogger('ConnectionHandlers')
*/
const PRESENCE_BEARING_TYPES = new Set<RoomRef['type']>([ROOM_TYPES.WORKFLOW, ROOM_TYPES.TABLE])

const pendingDisconnects = new Set<Promise<void>>()

/** Keep Redis available until disconnect listeners finish removing presence. */
export async function waitForConnectionCleanup(): Promise<void> {
await Promise.all(pendingDisconnects)
}

export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager: IRoomManager) {
socket.on('error', (error) => {
logger.error(`Socket ${socket.id} error:`, error)
Expand All @@ -28,7 +35,7 @@ export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager
// `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) => {
const handleDisconnect = async (reason: string) => {
try {
// Snapshot the live Socket.IO room membership SYNCHRONOUSLY, before any
// await: Socket.IO clears `socket.rooms` via leaveAll() as soon as the
Expand Down Expand Up @@ -91,5 +98,11 @@ export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager
} catch (error) {
logger.error(`Error handling disconnect for socket ${socket.id}:`, error)
}
}

socket.on('disconnecting', (reason) => {
const cleanup = handleDisconnect(reason)
pendingDisconnects.add(cleanup)
void cleanup.finally(() => pendingDisconnects.delete(cleanup))
})
}
Loading
Loading