Skip to content

Commit e6032bd

Browse files
committed
fix(files): preserve collaboration through reconnects and peer edits
1 parent dbe426d commit e6032bd

31 files changed

Lines changed: 2164 additions & 272 deletions

apps/docs/openapi-v2-files-audit.json

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3167,7 +3167,7 @@
31673167
"uploadedByEmail": {
31683168
"type": "string",
31693169
"format": "email",
3170-
"pattern": "^(?!\\.)(?!.*\\.\\.)([A-Za-z0-9_'+\\-\\.]*)[A-Za-z0-9_+-]@([A-Za-z0-9][A-Za-z0-9\\-]*\\.)+[A-Za-z]{2,}$",
3170+
"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])?)*$",
31713171
"description": "Current email address of the uploader.",
31723172
"examples": ["jane@example.com"]
31733173
},
@@ -4029,7 +4029,7 @@
40294029
"uploadedByEmail": {
40304030
"type": "string",
40314031
"format": "email",
4032-
"pattern": "^(?!\\.)(?!.*\\.\\.)([A-Za-z0-9_'+\\-\\.]*)[A-Za-z0-9_+-]@([A-Za-z0-9][A-Za-z0-9\\-]*\\.)+[A-Za-z]{2,}$",
4032+
"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])?)*$",
40334033
"description": "Current email address of the uploader.",
40344034
"examples": ["jane@example.com"]
40354035
},
Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,73 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { createServer, type Server as HttpServer } from 'node:http'
5+
import type { AddressInfo } from 'node:net'
6+
import { Server } from 'socket.io'
7+
import { io as connect, type Socket } from 'socket.io-client'
8+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
9+
import { setupConnectionHandlers, waitForConnectionCleanup } from '@/handlers/connection'
10+
import type { AuthenticatedSocket } from '@/middleware/auth'
11+
import { MemoryRoomManager } from '@/rooms'
12+
13+
vi.mock('@/handlers/file-doc', () => ({ cleanupFileDocForSocket: vi.fn() }))
14+
vi.mock('@/handlers/subblocks', () => ({ cleanupPendingSubblocksForSocket: vi.fn() }))
15+
vi.mock('@/handlers/variables', () => ({ cleanupPendingVariablesForSocket: vi.fn() }))
16+
17+
describe('server shutdown connection drain', () => {
18+
let httpServer: HttpServer
19+
let io: Server
20+
let manager: MemoryRoomManager
21+
let client: Socket
22+
23+
beforeEach(async () => {
24+
httpServer = createServer()
25+
io = new Server(httpServer, { transports: ['websocket'] })
26+
manager = new MemoryRoomManager(io)
27+
await manager.initialize()
28+
io.on('connection', (socket) => setupConnectionHandlers(socket as AuthenticatedSocket, manager))
29+
await new Promise<void>((resolve) => httpServer.listen(0, '127.0.0.1', resolve))
30+
const port = (httpServer.address() as AddressInfo).port
31+
client = connect(`http://127.0.0.1:${port}`, { transports: ['websocket'], autoConnect: false })
32+
const connected = new Promise<void>((resolve) => client.once('connect', resolve))
33+
client.connect()
34+
await connected
35+
})
36+
37+
afterEach(async () => {
38+
client.disconnect()
39+
await io.close()
40+
await waitForConnectionCleanup()
41+
await manager.shutdown()
42+
vi.restoreAllMocks()
43+
})
44+
45+
it('keeps automatic reconnection active after transport shutdown', async () => {
46+
const disconnected = new Promise<string>((resolve) => client.once('disconnect', resolve))
47+
await io.close()
48+
expect(await disconnected).toBe('transport close')
49+
expect(client.active).toBe(true)
50+
await waitForConnectionCleanup()
51+
})
52+
53+
it('waits for asynchronous presence cleanup before releasing its dependencies', async () => {
54+
let finishRemoval: (() => void) | undefined
55+
vi.spyOn(manager, 'removeSocketFromAllRooms').mockImplementation(
56+
() =>
57+
new Promise((resolve) => {
58+
finishRemoval = () => resolve([])
59+
})
60+
)
61+
await io.close()
62+
let drained = false
63+
const drain = waitForConnectionCleanup().then(() => {
64+
drained = true
65+
})
66+
await Promise.resolve()
67+
expect(drained).toBe(false)
68+
expect(finishRemoval).toBeDefined()
69+
finishRemoval?.()
70+
await drain
71+
expect(drained).toBe(true)
72+
})
73+
})

apps/realtime/src/handlers/connection.ts

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,13 @@ const logger = createLogger('ConnectionHandlers')
1616
*/
1717
const PRESENCE_BEARING_TYPES = new Set<RoomRef['type']>([ROOM_TYPES.WORKFLOW, ROOM_TYPES.TABLE])
1818

19+
const pendingDisconnects = new Set<Promise<void>>()
20+
21+
/** Keep Redis available until disconnect listeners finish removing presence. */
22+
export async function waitForConnectionCleanup(): Promise<void> {
23+
await Promise.all(pendingDisconnects)
24+
}
25+
1926
export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager: IRoomManager) {
2027
socket.on('error', (error) => {
2128
logger.error(`Socket ${socket.id} error:`, error)
@@ -28,7 +35,7 @@ export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager
2835
// `disconnecting` (not `disconnect`): here `socket.rooms` is still populated and
2936
// authoritative, so presence is cleaned up even if the Redis room-set key was
3037
// evicted or TTL-expired (which would leave the manager's stored rooms empty).
31-
socket.on('disconnecting', async (reason) => {
38+
const handleDisconnect = async (reason: string) => {
3239
try {
3340
// Snapshot the live Socket.IO room membership SYNCHRONOUSLY, before any
3441
// await: Socket.IO clears `socket.rooms` via leaveAll() as soon as the
@@ -91,5 +98,11 @@ export function setupConnectionHandlers(socket: AuthenticatedSocket, roomManager
9198
} catch (error) {
9299
logger.error(`Error handling disconnect for socket ${socket.id}:`, error)
93100
}
101+
}
102+
103+
socket.on('disconnecting', (reason) => {
104+
const cleanup = handleDisconnect(reason)
105+
pendingDisconnects.add(cleanup)
106+
void cleanup.finally(() => pendingDisconnects.delete(cleanup))
94107
})
95108
}

0 commit comments

Comments
 (0)