diff --git a/backend/src/handlers/registerBroadcaster.handler.ts b/backend/src/handlers/registerBroadcaster.handler.ts index 94124a2..7c025aa 100644 --- a/backend/src/handlers/registerBroadcaster.handler.ts +++ b/backend/src/handlers/registerBroadcaster.handler.ts @@ -223,28 +223,120 @@ const registerBroadcasterHandler = async (socket: Socket) => { try { const room = getRoom(roomId); const hostUserId = socket.data.user?.id; - const router = room.router; + const socketId = socket.id; const roomKey = `room:${roomId}`; - const redisRoom = await getRedisRoom(roomKey) - //TODO: Implementation for different pod connections - if(redisRoom.nodeId !== config.instanceId){ - logger.error('Different pod') - throw new ApiError(409,"Room belongs to another node") + let redisRoom; + try { + redisRoom = await getRedisRoom(roomKey) + } catch (error) { + logger.error('Error finding redis room',{ + error: (error as Error).message, + stack: (error as Error).stack + }) + + ack({success: false, code: "TRANSPORT_CREATION_FAILED"}) + return; } - + if(redisRoom?.nodeId !== config.instanceId){ + const requestId = crypto.randomUUID() + const date = Date.now(); + const args = {roomId, socketId}; + const replyTo = `pod:${config.instanceId}:response`; + + const payLoad : PodCommandPayload = { + requestId, + type: 'createBroadcasterTransport', + args, + replyTo, + date + } + + const TIMEOUTMS = 5000; + const timeOuthandle = setTimeout(() => { + const entry = podRequestHandleMap.get(requestId); + if(!entry) return logger.warn(`Entry not found with requestId: ${requestId}`); + entry.onComplete({}, 'TRANSPORT_CREATION_FAILED'); + podRequestHandleMap.delete(requestId); + }, TIMEOUTMS) + + podRequestHandleMap.set(requestId, { + requestId, + socketId, + startDate: date, + requestType: 'createBroadcasterTransport', + status: 'pending', + replyTo, + + onComplete: (result, error) => { + clearTimeout(timeOuthandle) + + if(error){ + logger.error('[Transport] pod error', { + error : error + }) + + ack({success: false, code: 'TRANSPORT_CREATION_FAILED'}) + return; + } + + ack({success: true, data: result}) + + void Broadcaster.findOneAndUpdate( + { + broadcasterId: hostUserId, + roomId, + }, + { + $push: { + transportIds: result?.id, + }, + } + ).catch((error) => { + logger.error("Failed to update broadcaster transport", error); + }); + + } + + }) + + if(!redisRoom.nodeId){ + logger.error('Redis room nodeId not found'); + podRequestHandleMap.delete(requestId); + ack({success: false, code: 'TRANSPORT_CREATION_FAILED'}) + return; + } + + let receivers: number; + try { + receivers = await publishCommand(payLoad, redisRoom.nodeId) + } catch (error) { + clearTimeout(timeOuthandle); + podRequestHandleMap.delete(requestId); + throw error; + } + if(receivers === 0){ + logger.error('Cross [POD] connection failed'); + clearTimeout(timeOuthandle); + podRequestHandleMap.delete(requestId); + ack({success: false, code: 'TRANSPORT_CREATION_FAILED'}) + } + return; + } + + const router = room.router; const broadcasterTransport = await createWebRtcTransport( router, roomId, - socket.id, + socketId, "producer" ); await saveBroadcasterTransport( roomId, - socket.id, + socketId, broadcasterTransport ); diff --git a/backend/src/utils/podConnection.ts b/backend/src/utils/podConnection.ts index ccbf58d..30133bb 100644 --- a/backend/src/utils/podConnection.ts +++ b/backend/src/utils/podConnection.ts @@ -7,7 +7,8 @@ import { connectConsumerTransport, createConsumerTransport, joinAsViewer } from import { heartBeat } from "./roomCordinator"; import {consume} from '../handlers/viewer.handler' import { pauseConsumer, resumeConsumer } from "../consumer/consumer.handler"; - +import { createWebRtcTransport } from "../mediasoup/transport"; +import { addBroadcaster, saveBroadcasterTransport } from "./broadcaster.util"; export interface PodCommandPayload { type: string; @@ -76,7 +77,33 @@ export const handleIncomingRequest = async(payload: PodCommandPayload) => { break; } - case(type === ''): {} + case(type === 'createBroadcasterTransport'): { + const {roomId, socketId} = args as unknown as generalArgs; + const routerId = roomToRouter.get(roomId); + if(!routerId){ + error = 'RouterId not found' + throw new Error('RouterId not found') + } + const router = getRouter(routerId); + + const broadcasterTransport = await createWebRtcTransport(router,roomId, socketId, 'producer'); + + result = { + id: broadcasterTransport?.id, + iceParameters: broadcasterTransport?.iceParameters, + iceCandidates: broadcasterTransport?.iceCandidates, + dtlsParameters: broadcasterTransport?.dtlsParameters, + } + addBroadcaster(roomId, socketId) + await saveBroadcasterTransport(roomId, socketId, broadcasterTransport) + + const payLoad: PodResponsePayload = { + requestId, + result + } + await publishResponse(payLoad, replyTo) + break; + } case(type === ''): {} case(type === ''): {} @@ -266,15 +293,16 @@ export const handleIncomingResponse = async(payload: PodResponsePayload) => { podRequestHandleMap.delete(payload.requestId) } -export const publishCommand = async(payload: PodCommandPayload, targetNode:string) => { +export const publishCommand = async(payload: PodCommandPayload, targetNode:string):Promise => { const channel = `pod:${targetNode}:cmd` - const receivers = await redis.spublish(channel , JSON.stringify(payload)) + const receivers = await redis.spublish(channel , JSON.stringify(payload)) as number if(receivers === 0){ logger.error(`No subscribers for ${channel} — pod may be down`, { requestId: payload.requestId }); } logger.info("Redis delivered to", receivers, "subscribers"); + return receivers; } diff --git a/docs/MultiplePodSignaling/broadcasterTransport.md b/docs/MultiplePodSignaling/broadcasterTransport.md new file mode 100644 index 0000000..fbaa7f2 --- /dev/null +++ b/docs/MultiplePodSignaling/broadcasterTransport.md @@ -0,0 +1,263 @@ +```txt +# Cross-Pod Co-Broadcaster: Create Broadcaster Transport + +> IMPORTANT: +> +> This flow is specifically for a **co-broadcaster whose Socket.IO +> connection lands on a different pod from the pod that owns the room +> and Mediasoup router**. +> +> The original room creator/broadcaster normally already lives on the +> room-owning pod and does not need this cross-pod forwarding path. +> +> This cross-pod flow exists for: +> +> Co-Broadcaster Socket → POD B +> Room + Router Owner → POD A +> +> POD A must create and own the actual Mediasoup producer transport. + +--- + +POD B (Co-Broadcaster socket lives here) POD A (owns Room + Router) +───────────────────────────────────────── ────────────────────────────── + +socket.on("createBroadcasterTransport") + │ + ├─ roomId received from client + ├─ socketId = socket.id + ├─ hostUserId = socket.data.user?.id + │ + ├─ getRedisRoom(`room:${roomId}`) + │ + │ → nodeId = "A" + │ → config.instanceId = "B" + │ + ├─ nodeId !== instanceId + │ + ▼ +CROSS-POD PATH + │ + ├─ requestId = crypto.randomUUID() + │ + ├─ Create payload: + │ + │ { + │ type: "createBroadcasterTransport", + │ requestId, + │ args: { + │ roomId, + │ socketId + │ }, + │ replyTo: "pod:B:response" + │ } + │ + ├─ Start timeout (5000ms) + │ + ├─ podRequestHandleMap.set(requestId, { + │ + │ status: "pending", + │ requestType: "createBroadcasterTransport", + │ + │ onComplete: (result, error) => { + │ + │ clearTimeout(timeoutHandle) + │ + │ if (error) { + │ + │ ack({ + │ success: false, + │ code: "TRANSPORT_CREATION_FAILED" + │ }) + │ + │ return + │ } + │ + │ ack({ + │ success: true, + │ data: result + │ }) + │ + │ Broadcaster.findOneAndUpdate(...) + │ + │ → store result.id in transportIds + │ } + │ }) + │ + │ ▲ + │ │ + │ │ Stored LOCALLY on POD B + │ │ This entry waits for POD A's response + │ │ + ├─ publishCommand(payload, "A") + │ + │ spublish("pod:A:cmd", { + │ + │ type: "createBroadcasterTransport", + │ requestId, + │ + │ args: { + │ roomId, + │ socketId + │ }, + │ + │ replyTo: "pod:B:response" + │ }) + │ + │ + │ │ + │ │ Redis + │ ▼ + │ + │ POD A + │ │ + │ │ podConnectionSubscriber + │ │ receives "pod:A:cmd" + │ ▼ + │ + │ handleIncomingRequest(payload) + │ │ + │ ├─ type === + │ │ "createBroadcasterTransport" + │ │ + │ ▼ + │ + │ roomToRouter.get(roomId) + │ │ + │ ▼ + │ + │ Get routerId + │ │ + │ ▼ + │ + │ getRouter(routerId) + │ │ + │ ▼ + │ + │ createWebRtcTransport( + │ router, + │ roomId, + │ socketId, + │ "producer" + │ ) + │ │ + │ ▼ + │ + │ Mediasoup creates the + │ PRODUCER TRANSPORT + │ on POD A + │ │ + │ ▼ + │ + │ result = { + │ id, + │ iceParameters, + │ iceCandidates, + │ dtlsParameters + │ } + │ │ + │ ▼ + │ + │ addBroadcaster( + │ roomId, + │ socketId + │ ) + │ │ + │ ▼ + │ + │ saveBroadcasterTransport( + │ roomId, + │ socketId, + │ broadcasterTransport + │ ) + │ + │ │ + │ ▼ + │ + │ Room state on POD A: + │ + │ room.broadcasters + │ └── socketId + │ └── producer transport + │ + │ │ + │ ▼ + │ + │ publishResponse( + │ { + │ requestId, + │ result + │ }, + │ "pod:B:response" + │ ) + │ + │ ▲ + │ │ + │ │ Redis + │ │ + └────────────────────┘ + │ + ▼ + +POD B + │ + ├─ podConnectionSubscriber + │ + ├─ receives: + │ + │ "pod:B:response" + │ + ▼ +handleIncomingResponse(payload) + │ + ├─ podRequestHandleMap.get(requestId) + │ + ├─ Finds the SAME pending request + │ + ├─ Calls: + │ + │ entry.onComplete( + │ payload.result, + │ payload.error + │ ) + │ + ▼ +onComplete(result, error) + │ + ├─ clearTimeout(timeoutHandle) + │ + ├─ If error: + │ + │ ack({ + │ success: false, + │ code: "TRANSPORT_CREATION_FAILED" + │ }) + │ + └─ If success: + │ + ├─ ack({ + │ success: true, + │ data: { + │ id, + │ iceParameters, + │ iceCandidates, + │ dtlsParameters + │ } + │ }) + │ + ▼ + + Broadcaster.findOneAndUpdate(...) + + $push: + transportIds: result.id + + ▼ + + Co-broadcaster database record + now contains the transport ID + + ▼ + + podRequestHandleMap.delete(requestId) +``` \ No newline at end of file