diff --git a/package-lock.json b/package-lock.json index f4b555f..9308bec 100644 --- a/package-lock.json +++ b/package-lock.json @@ -45,7 +45,7 @@ "typescript": "^5.8.3" }, "engines": { - "node": ">=18.15.0" + "node": ">=22.0.0" } }, "node_modules/@ampproject/remapping": { diff --git a/readme.md b/readme.md index 8e5a711..a064ed7 100644 --- a/readme.md +++ b/readme.md @@ -72,8 +72,47 @@ The `osc_eyevinn_intercom_manager` resource requires these variables: | `ENDPOINT_IDLE_TIMEOUT_S` | Idle timeout in seconds for SMB endpoints (default: `60`) | | `OSC_ACCESS_TOKEN` | Personal Access Token from OSC for link sharing and reauthenticating (optional) | | `ICE_SERVERS` | Comma-separated list of ICE servers in the format: `turn:username:password@turn.example.com,stun:stun.example.com`. If no STUN server is provided, and WHIP endpoints are used, Google's default STUN server (`stun:stun.l.google.com:19302`) will be used. | +| `WHIP_GATEWAY_URL` | URL of the [SRT-WHIP Gateway](https://github.com/Eyevinn/srt-whip-gateway) for IO bridge transmitters (optional). Enables the transmitter bridge API when set | +| `WHIP_GATEWAY_API_KEY` | API key for the SRT-WHIP Gateway (optional) | +| `WHEP_GATEWAY_URL` | URL of the [WHEP-SRT Gateway](https://github.com/Eyevinn/whep-srt-gateway) for IO bridge receivers (optional). Enables the receiver bridge API when set | +| `WHEP_GATEWAY_API_KEY` | API key for the WHEP-SRT Gateway (optional) | +| `DEBUG_BRIDGE` | Set to `true` to enable verbose logging from the bridge manager reconcile loop (optional) | | `MONGODB_CONNECTION_STRING` | DEPRECATED: Use `DB_CONNECTION_STRING` instead | +## IO Bridge + +The IO bridge enables SRT-to-WebRTC and WebRTC-to-SRT bridging, allowing external SRT streams to be ingested into intercom production lines (transmitters) and intercom audio to be sent out as SRT streams (receivers). + +- **Transmitters** (SRT to WebRTC): An SRT source is received by the [SRT-WHIP Gateway](https://github.com/Eyevinn/srt-whip-gateway) and ingested into a production line via WHIP. +- **Receivers** (WebRTC to SRT): Audio from a production line is received via WHEP from the [WHEP-SRT Gateway](https://github.com/Eyevinn/whep-srt-gateway) and output as an SRT stream. + +The bridge is enabled by setting `WHIP_GATEWAY_URL` and/or `WHEP_GATEWAY_URL`. A bridge manager runs a sync loop (1s interval) that reconciles the desired state in the database with the actual state on the gateways. + +The bridge API is available at `/api/v1/bridge/transmitters` and `/api/v1/bridge/receivers`, and a configuration endpoint at `/api/v1/bridge/config` reports which gateways are enabled. + +### Local development with gateways + +To run the gateways locally for development: + +```sh +# SRT-WHIP Gateway (transmitters) — requires Node.js +git clone https://github.com/Eyevinn/srt-whip-gateway.git +cd srt-whip-gateway && npm install && npm run dev +# Runs on port 3000 + +# WHEP-SRT Gateway (receivers) — requires Node.js +git clone https://github.com/Eyevinn/whep-srt-gateway.git +cd whep-srt-gateway && npm install && npm run dev +# Runs on port 3001 +``` + +Then set the environment variables: + +```sh +WHIP_GATEWAY_URL=http://localhost:3000 +WHEP_GATEWAY_URL=http://localhost:3001 +``` + ## Installation / Usage Start an Intercom Manager instance: diff --git a/src/api.test.ts b/src/api.test.ts index 1b81ea6..d5af59f 100644 --- a/src/api.test.ts +++ b/src/api.test.ts @@ -43,7 +43,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(undefined), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(undefined) + updatePreset: jest.fn().mockResolvedValue(undefined), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const mockProductionManager = { diff --git a/src/api.ts b/src/api.ts index 6535f7b..b58d74f 100644 --- a/src/api.ts +++ b/src/api.ts @@ -14,6 +14,8 @@ import apiReAuth from './api_re_auth'; import apiShare from './api_share'; import apiWhip, { ApiWhipOptions } from './api_whip'; import apiWhep, { ApiWhepOptions } from './api_whep'; +import apiBridgeTx, { ApiBridgeTxOptions } from './api_bridge_tx'; +import apiBridgeRx, { ApiBridgeRxOptions } from './api_bridge_rx'; import { DbManager } from './db/interface'; import { IngestManager } from './ingest_manager'; import { ProductionManager } from './production_manager'; @@ -54,15 +56,22 @@ export interface ApiGeneralOptions { endpointIdleTimeout: string; smbServerApiKey?: string; publicHost: string; + whipAuthKey?: string; dbManager: DbManager; productionManager: ProductionManager; ingestManager: IngestManager; + whipGatewayUrl?: string; + whipGatewayApiKey?: string; + whepGatewayUrl?: string; + whepGatewayApiKey?: string; } export type ApiOptions = ApiGeneralOptions & ApiProductionsOptions & ApiWhipOptions & - ApiWhepOptions; + ApiWhepOptions & + ApiBridgeTxOptions & + ApiBridgeRxOptions; export default async (opts: ApiOptions) => { const api = fastify({ @@ -128,8 +137,7 @@ export default async (opts: ApiOptions) => { smbServerApiKey: opts.smbServerApiKey, dbManager: opts.dbManager, productionManager: opts.productionManager, - coreFunctions: opts.coreFunctions, - smb: opts.smb + coreFunctions: opts.coreFunctions }); api.register(apiWhip, { prefix: 'api/v1', @@ -140,7 +148,7 @@ export default async (opts: ApiOptions) => { productionManager: opts.productionManager, dbManager: opts.dbManager, whipAuthKey: opts.whipAuthKey, - smb: opts.smb + whipGatewayUrl: opts.whipGatewayUrl }); api.register(apiWhep, { prefix: 'api/v1', @@ -151,12 +159,54 @@ export default async (opts: ApiOptions) => { productionManager: opts.productionManager, dbManager: opts.dbManager, whipAuthKey: opts.whipAuthKey, - smb: opts.smb + whepGatewayUrl: opts.whepGatewayUrl }); api.register(apiShare, { publicHost: opts.publicHost, prefix: 'api/v1' }); api.register(apiReAuth, { prefix: 'api/v1' }); api.register(apiGroups, { prefix: 'api/v1', dbManager: opts.dbManager }); + // Bridge configuration endpoint + const BridgeConfig = Type.Object({ + whipGatewayEnabled: Type.Boolean(), + whepGatewayEnabled: Type.Boolean() + }); + + api.get<{ Reply: Static }>( + '/api/v1/bridge/config', + { + schema: { + description: 'Get bridge gateway configuration', + response: { + 200: BridgeConfig + } + } + }, + async (_, reply) => { + reply.send({ + whipGatewayEnabled: !!opts.whipGatewayUrl, + whepGatewayEnabled: !!opts.whepGatewayUrl + }); + } + ); + + // Register bridge IO endpoints (only if gateways are configured) + if (opts.whipGatewayUrl) { + api.register(apiBridgeTx, { + prefix: 'api/v1', + dbManager: opts.dbManager, + whipGatewayUrl: opts.whipGatewayUrl, + whipGatewayApiKey: opts.whipGatewayApiKey + }); + } + if (opts.whepGatewayUrl) { + api.register(apiBridgeRx, { + prefix: 'api/v1', + dbManager: opts.dbManager, + whepGatewayUrl: opts.whepGatewayUrl, + whepGatewayApiKey: opts.whepGatewayApiKey + }); + } + api.all('/whip/:productionId/:lineId', async (request, reply) => { if (request.method !== 'POST' && request.method !== 'OPTIONS') { return reply diff --git a/src/api_bridge_rx.ts b/src/api_bridge_rx.ts new file mode 100644 index 0000000..97eb0d4 --- /dev/null +++ b/src/api_bridge_rx.ts @@ -0,0 +1,484 @@ +import { FastifyPluginCallback } from 'fastify'; +import { Type } from '@sinclair/typebox'; +import { DbManager } from './db/interface'; +import { + NewReceiver, + Receiver, + ReceiverListResponse, + ReceiverStateChange, + PatchReceiver, + BridgeStatus +} from './models'; +import { Log } from './log'; +import { encodeSrtStreamId } from './utils'; + +export interface ApiBridgeRxOptions { + dbManager: DbManager; + whepGatewayUrl?: string; + whepGatewayApiKey?: string; +} + +const ParamsId = Type.Object({ + id: Type.String({ + description: 'Receiver ID' + }) +}); + +const apiBridgeRx: FastifyPluginCallback = ( + fastify, + opts, + next +) => { + const { dbManager, whepGatewayUrl, whepGatewayApiKey } = opts; + + // Helper function to call gateway API + const callGateway = async ( + method: string, + path: string, + body?: any + ): Promise => { + const url = `${whepGatewayUrl}${path}`; + const headers: Record = { + 'Content-Type': 'application/json' + }; + if (whepGatewayApiKey) { + headers['x-api-key'] = whepGatewayApiKey; + } + + const options: RequestInit = { + method, + headers + }; + + if (body) { + options.body = JSON.stringify(body); + } + + const response = await fetch(url, options); + + if (!response.ok) { + const errorText = await response.text(); + throw new Error( + `Gateway request failed: ${response.status} ${errorText}` + ); + } + + if (response.status === 204) { + return null; + } + + const text = await response.text(); + if (!text) { + return null; + } + + // Try to parse as JSON, if it fails, return the text as-is + try { + return JSON.parse(text); + } catch (e) { + // If it's a plain boolean string, convert it + if (text.toLowerCase() === 'true') return true; + if (text.toLowerCase() === 'false') return false; + // Otherwise return the text + return text; + } + }; + + // List all receivers + fastify.get<{ + Querystring: { limit?: string; offset?: string }; + Reply: ReceiverListResponse | { error: string }; + }>( + '/bridge/rx', + { + schema: { + description: 'List all receivers', + querystring: Type.Object({ + limit: Type.Optional(Type.String()), + offset: Type.Optional(Type.String()) + }), + response: { + 200: ReceiverListResponse, + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const limit = parseInt(request.query.limit || '100', 10); + const offset = parseInt(request.query.offset || '0', 10); + + const receivers = await dbManager.getReceivers(limit, offset); + const totalItems = await dbManager.getReceiversLength(); + + // Clean up undefined fields by using JSON.parse/stringify + // This removes undefined values which Fast JSON Stringify can't handle + const cleanedReceivers = JSON.parse(JSON.stringify(receivers)); + + reply.code(200).send({ + receivers: cleanedReceivers, + limit, + offset, + totalItems + }); + } catch (error) { + console.error('Failed to list receivers - Full error:', error); + console.error('Error stack:', (error as Error).stack); + Log().error('Failed to list receivers:', error); + reply.code(500).send({ error: 'Failed to list receivers' }); + } + } + ); + + // Get a specific receiver + fastify.get<{ + Params: { id: string }; + Reply: Receiver | { error: string }; + }>( + '/bridge/rx/:id', + { + schema: { + description: 'Get a receiver by ID', + params: ParamsId, + response: { + 200: Receiver, + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const receiver = await dbManager.getReceiver(request.params.id); + + if (!receiver) { + reply.code(404).send({ error: 'Receiver not found' }); + return; + } + + // Clean up undefined fields + const cleanedReceiver = JSON.parse(JSON.stringify(receiver)); + reply.code(200).send(cleanedReceiver); + } catch (error) { + Log().error('Failed to get receiver:', error); + reply.code(500).send({ error: 'Failed to get receiver' }); + } + } + ); + + // Create a new receiver + fastify.post<{ + Body: NewReceiver; + Reply: Receiver | { error: string }; + }>( + '/bridge/rx', + { + schema: { + description: 'Create a new receiver', + body: NewReceiver, + response: { + 201: Receiver, + 400: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + // Save to database first + const receiver = await dbManager.addReceiver(request.body); + + // Try to create on gateway + try { + await callGateway('POST', '/api/v1/rx', { + id: receiver._id, + whepUrl: receiver.whepUrl, + srtUrl: encodeSrtStreamId(receiver.srtUrl), + status: BridgeStatus.IDLE + }); + + // Update status to idle (gateway created successfully) + receiver.status = BridgeStatus.IDLE; + await dbManager.updateReceiver(receiver); + } catch (gatewayError) { + Log().error('Failed to create receiver on gateway:', gatewayError); + // Mark as failed but keep in database + receiver.status = BridgeStatus.FAILED; + await dbManager.updateReceiver(receiver); + } + + reply.code(201).send(receiver); + } catch (error) { + Log().error('Failed to create receiver:', error); + reply.code(500).send({ error: 'Failed to create receiver' }); + } + } + ); + + // Update receiver state + fastify.put<{ + Params: { id: string }; + Body: ReceiverStateChange; + Reply: Receiver | { error: string }; + }>( + '/bridge/rx/:id/state', + { + schema: { + description: 'Update receiver state', + params: ParamsId, + body: ReceiverStateChange, + response: { + 200: Receiver, + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const receiver = await dbManager.getReceiver(request.params.id); + + if (!receiver) { + reply.code(404).send({ error: 'Receiver not found' }); + return; + } + + // Save desired state in database + receiver.desiredStatus = request.body.desired; + await dbManager.updateReceiver(receiver); + + // Update gateway state + try { + await callGateway('PUT', `/api/v1/rx/${request.params.id}/state`, { + desired: request.body.desired + }); + + // Update actual status + receiver.status = request.body.desired; + await dbManager.updateReceiver(receiver); + + reply.code(200).send(receiver); + } catch (gatewayError) { + Log().error( + 'Failed to update receiver state on gateway:', + gatewayError + ); + // Desired state is saved, sync will retry + reply.code(500).send({ error: 'Failed to update receiver state' }); + } + } catch (error) { + Log().error('Failed to update receiver:', error); + reply.code(500).send({ error: 'Failed to update receiver' }); + } + } + ); + + // Update receiver metadata + fastify.patch<{ + Params: { id: string }; + Body: PatchReceiver; + Reply: Receiver | { error: string }; + }>( + '/bridge/rx/:id', + { + schema: { + description: 'Update receiver metadata (label, productionId, lineId)', + params: ParamsId, + body: PatchReceiver, + response: { + 200: Receiver, + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const receiver = await dbManager.getReceiver(request.params.id); + + if (!receiver) { + reply.code(404).send({ error: 'Receiver not found' }); + return; + } + + // Check if productionId or lineId are changing + const productionChanged = + request.body.productionId !== undefined && + request.body.productionId !== receiver.productionId; + const lineChanged = + request.body.lineId !== undefined && + request.body.lineId !== receiver.lineId; + + // If only label is changing, simple update + if ( + !productionChanged && + !lineChanged && + request.body.label !== undefined + ) { + receiver.label = request.body.label; + receiver.updatedAt = new Date().toISOString(); + await dbManager.updateReceiver(receiver); + const cleanedReceiver = JSON.parse(JSON.stringify(receiver)); + reply.code(200).send(cleanedReceiver); + return; + } + + // If production or line changed, need to recreate gateway object + if (productionChanged || lineChanged) { + // Save the current state to restore after recreation + const previousStatus = receiver.status; + const previousDesiredStatus = receiver.desiredStatus; + + // Update the receiver data + if (request.body.label !== undefined) { + receiver.label = request.body.label; + } + if (request.body.productionId !== undefined) { + receiver.productionId = request.body.productionId; + } + if (request.body.lineId !== undefined) { + receiver.lineId = request.body.lineId; + } + + // Reconstruct WHEP URL with new production/line IDs + // URL format: ${backendBaseUrl}/api/v1/whep/${productionId}/${lineId}/${whepUsername} + // Extract username from existing URL + const urlParts = receiver.whepUrl.split('/'); + const whepUsername = urlParts[urlParts.length - 1]; + const backendBaseUrl = receiver.whepUrl.split('/api/v1/')[0]; + receiver.whepUrl = `${backendBaseUrl}/api/v1/whep/${receiver.productionId}/${receiver.lineId}/${whepUsername}`; + + // Update timestamp + receiver.updatedAt = new Date().toISOString(); + + // Set desired state to STOPPED in database FIRST to prevent state enforcer from restarting + receiver.status = BridgeStatus.STOPPED; + receiver.desiredStatus = BridgeStatus.STOPPED; + await dbManager.updateReceiver(receiver); + + try { + // Stop the gateway first before deleting + try { + await callGateway( + 'PUT', + `/api/v1/rx/${request.params.id}/state`, + { + desired: BridgeStatus.STOPPED + } + ); + } catch (stopError) { + Log().warn('Failed to stop receiver before deletion:', stopError); + } + + // Delete from gateway + try { + await callGateway('DELETE', `/api/v1/rx/${request.params.id}`); + } catch (deleteError) { + Log().warn( + 'Failed to delete receiver from gateway:', + deleteError + ); + } + + // Create new gateway object with updated URL (gateway requires initial status) + await callGateway('POST', '/api/v1/rx', { + id: receiver._id, + whepUrl: receiver.whepUrl, + srtUrl: encodeSrtStreamId(receiver.srtUrl), + status: BridgeStatus.IDLE + }); + + // Restore previous state if it was running + if ( + previousStatus === BridgeStatus.RUNNING || + previousDesiredStatus === BridgeStatus.RUNNING + ) { + try { + await callGateway( + 'PUT', + `/api/v1/rx/${request.params.id}/state`, + { + desired: BridgeStatus.RUNNING + } + ); + receiver.status = BridgeStatus.RUNNING; + receiver.desiredStatus = BridgeStatus.RUNNING; + } catch (stateError) { + Log().warn('Failed to restore receiver state:', stateError); + receiver.status = BridgeStatus.IDLE; + } + } else { + receiver.status = BridgeStatus.IDLE; + receiver.desiredStatus = BridgeStatus.IDLE; + } + + await dbManager.updateReceiver(receiver); + } catch (gatewayError) { + Log().error( + 'Failed to recreate receiver on gateway:', + gatewayError + ); + receiver.status = BridgeStatus.FAILED; + await dbManager.updateReceiver(receiver); + } + } + + // Clean up undefined fields + const cleanedReceiver = JSON.parse(JSON.stringify(receiver)); + reply.code(200).send(cleanedReceiver); + } catch (error) { + Log().error('Failed to update receiver:', error); + reply.code(500).send({ error: 'Failed to update receiver' }); + } + } + ); + + // Delete a receiver + fastify.delete<{ + Params: { id: string }; + Reply: { success: boolean } | { error: string }; + }>( + '/bridge/rx/:id', + { + schema: { + description: 'Delete a receiver', + params: ParamsId, + response: { + 200: Type.Object({ success: Type.Boolean() }), + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const receiver = await dbManager.getReceiver(request.params.id); + + if (!receiver) { + reply.code(404).send({ error: 'Receiver not found' }); + return; + } + + // Delete from gateway first + try { + await callGateway('DELETE', `/api/v1/rx/${request.params.id}`); + } catch (gatewayError) { + Log().warn('Failed to delete receiver from gateway:', gatewayError); + // Continue with database deletion even if gateway fails + } + + // Delete from database + await dbManager.deleteReceiver(request.params.id); + + reply.code(200).send({ success: true }); + } catch (error) { + Log().error('Failed to delete receiver:', error); + reply.code(500).send({ error: 'Failed to delete receiver' }); + } + } + ); + + next(); +}; + +export default apiBridgeRx; diff --git a/src/api_bridge_tx.ts b/src/api_bridge_tx.ts new file mode 100644 index 0000000..ad2553b --- /dev/null +++ b/src/api_bridge_tx.ts @@ -0,0 +1,507 @@ +import { FastifyPluginCallback } from 'fastify'; +import { Type } from '@sinclair/typebox'; +import { DbManager } from './db/interface'; +import { + NewTransmitter, + Transmitter, + TransmitterListResponse, + TransmitterStateChange, + PatchTransmitter, + BridgeStatus +} from './models'; +import { Log } from './log'; + +export interface ApiBridgeTxOptions { + dbManager: DbManager; + whipGatewayUrl?: string; + whipGatewayApiKey?: string; +} + +const ParamsId = Type.Object({ + id: Type.String({ + description: 'Transmitter ID' + }) +}); + +const apiBridgeTx: FastifyPluginCallback = ( + fastify, + opts, + next +) => { + const { dbManager, whipGatewayUrl, whipGatewayApiKey } = opts; + + // Helper function to call gateway API + const callGateway = async ( + method: string, + path: string, + body?: any + ): Promise => { + const url = `${whipGatewayUrl}${path}`; + const headers: Record = { + 'Content-Type': 'application/json' + }; + if (whipGatewayApiKey) { + headers['x-api-key'] = whipGatewayApiKey; + } + + const options: RequestInit = { + method, + headers + }; + + if (body) { + options.body = JSON.stringify(body); + } + + const response = await fetch(url, options); + + if (!response.ok) { + const errorText = await response.text(); + throw new Error( + `Gateway request failed: ${response.status} ${errorText}` + ); + } + + if (response.status === 204) { + return null; + } + + const text = await response.text(); + if (!text) { + return null; + } + + // Try to parse as JSON, if it fails, return the text as-is + try { + return JSON.parse(text); + } catch (e) { + // If it's a plain boolean string, convert it + if (text.toLowerCase() === 'true') return true; + if (text.toLowerCase() === 'false') return false; + // Otherwise return the text + return text; + } + }; + + // List all transmitters + fastify.get<{ + Querystring: { limit?: string; offset?: string }; + Reply: TransmitterListResponse | { error: string }; + }>( + '/bridge/tx', + { + schema: { + description: 'List all transmitters', + querystring: Type.Object({ + limit: Type.Optional(Type.String()), + offset: Type.Optional(Type.String()) + }), + response: { + 200: TransmitterListResponse, + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const limit = parseInt(request.query.limit || '100', 10); + const offset = parseInt(request.query.offset || '0', 10); + + const transmitters = await dbManager.getTransmitters(limit, offset); + const totalItems = await dbManager.getTransmittersLength(); + + // Clean up undefined fields by using JSON.parse/stringify + // This removes undefined values which Fast JSON Stringify can't handle + const cleanedTransmitters = JSON.parse(JSON.stringify(transmitters)); + + reply.code(200).send({ + transmitters: cleanedTransmitters, + limit, + offset, + totalItems + }); + } catch (error) { + Log().error('Failed to list transmitters:', error); + reply.code(500).send({ error: 'Failed to list transmitters' }); + } + } + ); + + // Get a specific transmitter + fastify.get<{ + Params: { id: string }; + Reply: Transmitter | { error: string }; + }>( + '/bridge/tx/:id', + { + schema: { + description: 'Get a transmitter by id', + params: ParamsId, + response: { + 200: Transmitter, + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const id = request.params.id; + const transmitter = await dbManager.getTransmitter(id); + + if (!transmitter) { + reply.code(404).send({ error: 'Transmitter not found' }); + return; + } + + // Clean up undefined fields + const cleanedTransmitter = JSON.parse(JSON.stringify(transmitter)); + reply.code(200).send(cleanedTransmitter); + } catch (error) { + Log().error('Failed to get transmitter:', error); + reply.code(500).send({ error: 'Failed to get transmitter' }); + } + } + ); + + // Create a new transmitter + fastify.post<{ + Body: NewTransmitter; + Reply: Transmitter | { error: string }; + }>( + '/bridge/tx', + { + schema: { + description: 'Create a new transmitter', + body: NewTransmitter, + response: { + 201: Transmitter, + 400: Type.Object({ error: Type.String() }), + 409: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + // Save to database first + const transmitter = await dbManager.addTransmitter(request.body); + + // Try to create on gateway + try { + await callGateway('POST', '/api/v1/tx/id', { + id: transmitter._id, + label: transmitter.label, + port: transmitter.port, + mode: transmitter.mode === 'caller' ? 1 : 2, + srtUrl: transmitter.srtUrl?.replace(/^srt:\/\//, ''), + whipUrl: transmitter.whipUrl, + passThroughUrl: transmitter.passThroughUrl, + noVideo: transmitter.noVideo ?? true, + vp8: transmitter.vp8 ?? false, + bypassVideo: transmitter.bypassVideo ?? false, + status: BridgeStatus.IDLE + }); + + // Update status to idle (gateway created successfully) + transmitter.status = BridgeStatus.IDLE; + await dbManager.updateTransmitter(transmitter); + } catch (gatewayError) { + Log().error('Failed to create transmitter on gateway:', gatewayError); + // Mark as failed but keep in database + transmitter.status = BridgeStatus.FAILED; + await dbManager.updateTransmitter(transmitter); + } + + reply.code(201).send(transmitter); + } catch (error) { + Log().error('Failed to create transmitter:', error); + reply.code(500).send({ error: 'Failed to create transmitter' }); + } + } + ); + + // Update transmitter state + fastify.put<{ + Params: { id: string }; + Body: TransmitterStateChange; + Reply: Transmitter | { error: string }; + }>( + '/bridge/tx/:id/state', + { + schema: { + description: 'Update transmitter state', + params: ParamsId, + body: TransmitterStateChange, + response: { + 200: Transmitter, + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const id = request.params.id; + const transmitter = await dbManager.getTransmitter(id); + + if (!transmitter) { + reply.code(404).send({ error: 'Transmitter not found' }); + return; + } + + // Save desired state in database + transmitter.desiredStatus = request.body.desired; + await dbManager.updateTransmitter(transmitter); + + // Update gateway state + try { + await callGateway('PUT', `/api/v1/tx/id/${transmitter._id}/state`, { + desired: request.body.desired + }); + + // Update actual status + transmitter.status = request.body.desired; + await dbManager.updateTransmitter(transmitter); + + reply.code(200).send(transmitter); + } catch (gatewayError) { + Log().error( + 'Failed to update transmitter state on gateway:', + gatewayError + ); + // Desired state is saved, sync will retry + reply.code(500).send({ error: 'Failed to update transmitter state' }); + } + } catch (error) { + Log().error('Failed to update transmitter:', error); + reply.code(500).send({ error: 'Failed to update transmitter' }); + } + } + ); + + // Update transmitter metadata + fastify.patch<{ + Params: { id: string }; + Body: PatchTransmitter; + Reply: Transmitter | { error: string }; + }>( + '/bridge/tx/:id', + { + schema: { + description: + 'Update transmitter metadata (label, productionId, lineId)', + params: ParamsId, + body: PatchTransmitter, + response: { + 200: Transmitter, + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const id = request.params.id; + const transmitter = await dbManager.getTransmitter(id); + + if (!transmitter) { + reply.code(404).send({ error: 'Transmitter not found' }); + return; + } + + // Check if productionId or lineId are changing + const productionChanged = + request.body.productionId !== undefined && + request.body.productionId !== transmitter.productionId; + const lineChanged = + request.body.lineId !== undefined && + request.body.lineId !== transmitter.lineId; + + // If only label is changing, simple update + if ( + !productionChanged && + !lineChanged && + request.body.label !== undefined + ) { + transmitter.label = request.body.label; + transmitter.updatedAt = new Date().toISOString(); + await dbManager.updateTransmitter(transmitter); + const cleanedTransmitter = JSON.parse(JSON.stringify(transmitter)); + reply.code(200).send(cleanedTransmitter); + return; + } + + // If production or line changed, need to recreate gateway object + if (productionChanged || lineChanged) { + // Save the current state to restore after recreation + const previousStatus = transmitter.status; + const previousDesiredStatus = transmitter.desiredStatus; + + // Update the transmitter data + if (request.body.label !== undefined) { + transmitter.label = request.body.label; + } + if (request.body.productionId !== undefined) { + transmitter.productionId = request.body.productionId; + } + if (request.body.lineId !== undefined) { + transmitter.lineId = request.body.lineId; + } + + // Reconstruct WHIP URL with new production/line IDs + // URL format: ${backendBaseUrl}/api/v1/whip/${productionId}/${lineId}/${whipUsername} + // Extract username from existing URL + const urlParts = transmitter.whipUrl.split('/'); + const whipUsername = urlParts[urlParts.length - 1]; + const backendBaseUrl = transmitter.whipUrl.split('/api/v1/')[0]; + transmitter.whipUrl = `${backendBaseUrl}/api/v1/whip/${transmitter.productionId}/${transmitter.lineId}/${whipUsername}`; + + // Update timestamp + transmitter.updatedAt = new Date().toISOString(); + + // Set desired state to STOPPED in database FIRST to prevent state enforcer from restarting + transmitter.status = BridgeStatus.STOPPED; + transmitter.desiredStatus = BridgeStatus.STOPPED; + await dbManager.updateTransmitter(transmitter); + + try { + // Stop the gateway first before deleting + try { + await callGateway( + 'PUT', + `/api/v1/tx/id/${transmitter._id}/state`, + { + desired: BridgeStatus.STOPPED + } + ); + } catch (stopError) { + Log().warn( + 'Failed to stop transmitter before deletion:', + stopError + ); + } + + // Delete from gateway + try { + await callGateway('DELETE', `/api/v1/tx/id/${transmitter._id}`); + } catch (deleteError) { + Log().warn( + 'Failed to delete transmitter from gateway:', + deleteError + ); + } + + // Create new gateway object with updated URL (gateway requires initial status) + await callGateway('POST', '/api/v1/tx/id', { + id: transmitter._id, + label: transmitter.label, + port: transmitter.port, + mode: transmitter.mode === 'caller' ? 1 : 2, + srtUrl: transmitter.srtUrl?.replace(/^srt:\/\//, ''), + whipUrl: transmitter.whipUrl, + passThroughUrl: transmitter.passThroughUrl, + noVideo: transmitter.noVideo ?? true, + vp8: transmitter.vp8 ?? false, + bypassVideo: transmitter.bypassVideo ?? false, + status: BridgeStatus.IDLE + }); + + // Restore previous state if it was running + if ( + previousStatus === BridgeStatus.RUNNING || + previousDesiredStatus === BridgeStatus.RUNNING + ) { + try { + await callGateway( + 'PUT', + `/api/v1/tx/id/${transmitter._id}/state`, + { + desired: BridgeStatus.RUNNING + } + ); + transmitter.status = BridgeStatus.RUNNING; + transmitter.desiredStatus = BridgeStatus.RUNNING; + } catch (stateError) { + Log().warn('Failed to restore transmitter state:', stateError); + transmitter.status = BridgeStatus.IDLE; + } + } else { + transmitter.status = BridgeStatus.IDLE; + transmitter.desiredStatus = BridgeStatus.IDLE; + } + + await dbManager.updateTransmitter(transmitter); + } catch (gatewayError) { + Log().error( + 'Failed to recreate transmitter on gateway:', + gatewayError + ); + transmitter.status = BridgeStatus.FAILED; + await dbManager.updateTransmitter(transmitter); + } + } + + // Clean up undefined fields + const cleanedTransmitter = JSON.parse(JSON.stringify(transmitter)); + reply.code(200).send(cleanedTransmitter); + } catch (error) { + Log().error('Failed to update transmitter:', error); + reply.code(500).send({ error: 'Failed to update transmitter' }); + } + } + ); + + // Delete a transmitter + fastify.delete<{ + Params: { id: string }; + Reply: { success: boolean } | { error: string }; + }>( + '/bridge/tx/:id', + { + schema: { + description: 'Delete a transmitter', + params: ParamsId, + response: { + 200: Type.Object({ success: Type.Boolean() }), + 404: Type.Object({ error: Type.String() }), + 500: Type.Object({ error: Type.String() }) + } + } + }, + async (request, reply) => { + try { + const id = request.params.id; + const transmitter = await dbManager.getTransmitter(id); + + if (!transmitter) { + reply.code(404).send({ error: 'Transmitter not found' }); + return; + } + + // Delete from gateway first + try { + await callGateway('DELETE', `/api/v1/tx/id/${transmitter._id}`); + } catch (gatewayError) { + Log().warn( + 'Failed to delete transmitter from gateway:', + gatewayError + ); + // Continue with database deletion even if gateway fails + } + + // Delete from database + await dbManager.deleteTransmitter(id); + + reply.code(200).send({ success: true }); + } catch (error) { + Log().error('Failed to delete transmitter:', error); + reply.code(500).send({ error: 'Failed to delete transmitter' }); + } + } + ); + + next(); +}; + +export default apiBridgeTx; diff --git a/src/api_groups.test.ts b/src/api_groups.test.ts index d5839cb..44e57e1 100644 --- a/src/api_groups.test.ts +++ b/src/api_groups.test.ts @@ -45,7 +45,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(mockPreset), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(mockPreset) + updatePreset: jest.fn().mockResolvedValue(mockPreset), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const mockIngestManager = { diff --git a/src/api_productions.test.ts b/src/api_productions.test.ts index 40d7273..c978e19 100644 --- a/src/api_productions.test.ts +++ b/src/api_productions.test.ts @@ -56,7 +56,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(undefined), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(undefined) + updatePreset: jest.fn().mockResolvedValue(undefined), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const mockIngestManager = { diff --git a/src/api_re_auth.test.ts b/src/api_re_auth.test.ts index 5cc2287..2572ed7 100644 --- a/src/api_re_auth.test.ts +++ b/src/api_re_auth.test.ts @@ -44,7 +44,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(undefined), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(undefined) + updatePreset: jest.fn().mockResolvedValue(undefined), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const mockProductionManager = { diff --git a/src/api_share.test.ts b/src/api_share.test.ts index 768d6b5..dbaad30 100644 --- a/src/api_share.test.ts +++ b/src/api_share.test.ts @@ -44,7 +44,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(undefined), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(undefined) + updatePreset: jest.fn().mockResolvedValue(undefined), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const mockProductionManager = { diff --git a/src/api_validation.test.ts b/src/api_validation.test.ts index d0120de..f5eaaf5 100644 --- a/src/api_validation.test.ts +++ b/src/api_validation.test.ts @@ -44,7 +44,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(undefined), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(undefined) + updatePreset: jest.fn().mockResolvedValue(undefined), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const mockIngestManager = { diff --git a/src/api_whep.ts b/src/api_whep.ts index 0057d98..8b57999 100644 --- a/src/api_whep.ts +++ b/src/api_whep.ts @@ -3,11 +3,12 @@ import { Static, Type } from '@sinclair/typebox'; import { FastifyPluginCallback } from 'fastify'; import sdpTransform, { parse } from 'sdp-transform'; import { v4 as uuidv4 } from 'uuid'; +import { promises as dns } from 'dns'; import { CoreFunctions } from './api_productions_core_functions'; import { Log } from './log'; import { Line, WhipWhepRequest, WhipWhepResponse } from './models'; import { ProductionManager } from './production_manager'; -import { ISmbProtocol, SmbProtocol } from './smb'; +import { SmbProtocol } from './smb'; import { getIceServers } from './utils'; import { DbManager } from './db/interface'; @@ -19,16 +20,50 @@ export interface ApiWhepOptions { productionManager: ProductionManager; dbManager: DbManager; whipAuthKey?: string; - smb?: ISmbProtocol; + whepGatewayUrl?: string; } -export const apiWhep: FastifyPluginCallback = ( +export const apiWhep: FastifyPluginCallback = async ( fastify, - opts, - next + opts ) => { const productionManager = opts.productionManager; + // Build allowList for rate limiting + const rateLimitAllowList = ['127.0.0.1', '::1', '::ffff:127.0.0.1']; + if (opts.whepGatewayUrl) { + try { + const gatewayHost = new URL(opts.whepGatewayUrl).hostname; + if (gatewayHost && !rateLimitAllowList.includes(gatewayHost)) { + rateLimitAllowList.push(gatewayHost); + + try { + const addresses = await dns.resolve(gatewayHost); + for (const ip of addresses) { + if (!rateLimitAllowList.includes(ip)) { + rateLimitAllowList.push(ip); + } + } + Log().info( + `WHEP rate limit allowList - resolved ${gatewayHost} to IPs: ${addresses.join( + ', ' + )}` + ); + } catch (resolveErr) { + Log().warn( + `Failed to resolve WHEP gateway hostname ${gatewayHost}: ${resolveErr}` + ); + } + } + } catch (err) { + Log().warn( + `Failed to parse WHEP gateway URL for rate limit allowList: ${err}` + ); + } + } + + Log().info(`WHEP rate limit allowList: ${rateLimitAllowList.join(', ')}`); + fastify.addContentTypeParser( 'application/sdp', { parseAs: 'string' }, @@ -50,14 +85,14 @@ export const apiWhep: FastifyPluginCallback = ( opts.smbServerBaseUrl ).toString(); - const smb = opts.smb || new SmbProtocol(); + const smb = new SmbProtocol(); const smbServerApiKey = opts.smbServerApiKey || ''; const coreFunctions = opts.coreFunctions; const whipAuthKey = opts.whipAuthKey?.trim(); async function requireWhepAuth(request: any, reply: any): Promise { if (!whipAuthKey) { - return true; // auth disabled + return true; } const authHeader = @@ -91,11 +126,6 @@ export const apiWhep: FastifyPluginCallback = ( { schema: { description: 'WHEP endpoint for Egress WebRTC streams', - params: Type.Object({ - productionId: Type.String({ maxLength: 200 }), - lineId: Type.String({ maxLength: 200 }), - username: Type.String({ maxLength: 200 }) - }), body: WhipWhepRequest, response: { 201: WhipWhepResponse, @@ -108,9 +138,15 @@ export const apiWhep: FastifyPluginCallback = ( }, config: { rateLimit: { - max: 10, + max: 100, timeWindow: '1 minute', hook: 'onRequest', + allowList: rateLimitAllowList, + onExceeded: (req) => { + Log().warn( + `Rate limit exceeded for WHEP endpoint - IP: ${req.ip}, URL: ${req.url}` + ); + }, errorResponseBuilder: (_req, context) => { return { statusCode: 429, @@ -150,7 +186,6 @@ export const apiWhep: FastifyPluginCallback = ( lineId ); - // Allocate endpoint with audio support const endpoint = await coreFunctions.createEndpoint( smb, smbServerUrl, @@ -226,7 +261,6 @@ export const apiWhep: FastifyPluginCallback = ( ); // Create the Location URL for the WHEP resource - // Location URL can be relative to Request URL, so this is OK. const locationUrl = `/api/v1/whep/${productionId}/${lineId}/${sessionId}`; // Set response headers @@ -234,13 +268,20 @@ export const apiWhep: FastifyPluginCallback = ( 'Content-Type': 'application/sdp', Location: locationUrl, ETag: sessionId, - Link: getIceServers().join(',') + Link: getIceServers().join(','), + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'GET, POST, DELETE, OPTIONS, PATCH', + 'Access-Control-Allow-Headers': + 'Content-Type, Authorization, ETag, If-Match, Link', + 'Access-Control-Expose-Headers': 'Location, ETag, Link' }); - reply.code(201).send(sdpAnswer); + await reply.code(201).send(sdpAnswer); } catch (err) { Log().error(err); - reply.code(500).send({ error: 'Failed to process WHEP request' }); + reply + .code(500) + .send({ error: `Failed to process WHEP request: ${err}` }); } } ); @@ -252,11 +293,6 @@ export const apiWhep: FastifyPluginCallback = ( { schema: { description: 'Terminate a WHEP connection', - params: Type.Object({ - productionId: Type.String({ maxLength: 200 }), - lineId: Type.String({ maxLength: 200 }), - sessionId: Type.String({ maxLength: 200 }) - }), response: { 200: Type.String({ description: 'OK' }), 404: Type.Object({ error: Type.String() }), @@ -266,8 +302,9 @@ export const apiWhep: FastifyPluginCallback = ( }, async (request, reply) => { if (!(await requireWhepAuth(request, reply))) return; - const { sessionId } = request.params; try { + const { sessionId } = request.params; + Log().info( `Received WHEP DELETE request - sessionId: ${sessionId}, IP: ${request.ip}` ); @@ -289,13 +326,21 @@ export const apiWhep: FastifyPluginCallback = ( `WHEP session deleted successfully - sessionId: ${sessionId}` ); + reply.headers({ + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'POST, DELETE, OPTIONS', + 'Access-Control-Allow-Headers': 'Content-Type, Authorization' + }); + reply.code(200).send('OK'); } catch (err) { Log().error( - `Failed to delete WHEP session - sessionId: ${sessionId}:`, + `Failed to delete WHEP session - sessionId: ${request.params.sessionId}:`, err ); - reply.code(500).send({ error: 'Failed to terminate WHEP connection' }); + reply + .code(500) + .send({ error: `Failed to terminate WHEP connection: ${err}` }); } } ); @@ -323,7 +368,6 @@ export const apiWhep: FastifyPluginCallback = ( try { const { productionId, lineId } = request.params; - // Check if production and line exist const productionIdNum = parseInt(productionId, 10); if (isNaN(productionIdNum)) { reply.code(400).send({ error: 'Invalid production ID' }); @@ -345,18 +389,23 @@ export const apiWhep: FastifyPluginCallback = ( } reply.headers({ + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'POST, DELETE, OPTIONS, PATCH', + 'Access-Control-Allow-Headers': 'Content-Type, Authorization, ETag', + 'Access-Control-Expose-Headers': 'Location, ETag, Link', + 'Access-Control-Max-Age': '86400', 'Accept-Post': 'application/sdp' }); reply.code(200).send('OK'); } catch (err) { Log().error(err); - reply.code(500).send({ error: 'Failed to process OPTIONS request' }); + reply + .code(500) + .send({ error: `Failed to process OPTIONS request: ${err}` }); } } ); - - next(); }; export default apiWhep; diff --git a/src/api_whip.test.ts b/src/api_whip.test.ts index 5dcb291..ce19a97 100644 --- a/src/api_whip.test.ts +++ b/src/api_whip.test.ts @@ -59,7 +59,19 @@ const mockDbManager = { getPreset: jest.fn().mockResolvedValue(undefined), getPresets: jest.fn().mockResolvedValue([]), deletePreset: jest.fn().mockResolvedValue(true), - updatePreset: jest.fn().mockResolvedValue(undefined) + updatePreset: jest.fn().mockResolvedValue(undefined), + addTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitter: jest.fn().mockResolvedValue(undefined), + getTransmitters: jest.fn().mockResolvedValue([]), + getTransmittersLength: jest.fn().mockResolvedValue(0), + updateTransmitter: jest.fn().mockResolvedValue(undefined), + deleteTransmitter: jest.fn().mockResolvedValue(true), + addReceiver: jest.fn().mockResolvedValue(undefined), + getReceiver: jest.fn().mockResolvedValue(undefined), + getReceivers: jest.fn().mockResolvedValue([]), + getReceiversLength: jest.fn().mockResolvedValue(0), + updateReceiver: jest.fn().mockResolvedValue(undefined), + deleteReceiver: jest.fn().mockResolvedValue(true) }; const coreFunctions = new CoreFunctions( @@ -202,12 +214,14 @@ describe('apiWhip', () => { expect(response.statusCode).toBe(415); }); - it('should return 429 when rate limit is exceeded', async () => { + it('should not rate-limit requests from localhost (allowlisted)', async () => { const fastify = await createTestServer(); - // Send 10 valid requests (these should succeed or at least not trigger 429) - for (let i = 0; i < 10; i++) { - await fastify.inject({ + // Send more than max requests from localhost — all should succeed + // because 127.0.0.1 is in the rate limit allowList + let lastResponse; + for (let i = 0; i < 110; i++) { + lastResponse = await fastify.inject({ method: 'POST', url: '/whip/prod1/line1/testuser', headers: { @@ -217,22 +231,7 @@ describe('apiWhip', () => { }); } - // The 11th request should exceed the rate limit - const response = await fastify.inject({ - method: 'POST', - url: '/whip/prod1/line1/testuser', - headers: { - 'content-type': 'application/sdp' - }, - payload: 'v=0\r\n' - }); - - expect(response.statusCode).toBe(429); - expect(JSON.parse(response.body)).toEqual( - expect.objectContaining({ - error: expect.stringMatching(/Too many/i) - }) - ); + expect(lastResponse!.statusCode).not.toBe(429); }); }); diff --git a/src/api_whip.ts b/src/api_whip.ts index 52d82ed..16efc40 100644 --- a/src/api_whip.ts +++ b/src/api_whip.ts @@ -3,11 +3,12 @@ import { Static, Type } from '@sinclair/typebox'; import { FastifyPluginCallback } from 'fastify'; import sdpTransform, { parse } from 'sdp-transform'; import { v4 as uuidv4 } from 'uuid'; +import { promises as dns } from 'dns'; import { CoreFunctions } from './api_productions_core_functions'; import { Log } from './log'; import { Line, WhipWhepRequest, WhipWhepResponse } from './models'; import { ProductionManager } from './production_manager'; -import { ISmbProtocol, SmbProtocol } from './smb'; +import { SmbProtocol } from './smb'; import { getIceServers } from './utils'; import { DbManager } from './db/interface'; @@ -19,16 +20,50 @@ export interface ApiWhipOptions { productionManager: ProductionManager; dbManager: DbManager; whipAuthKey?: string; - smb?: ISmbProtocol; + whipGatewayUrl?: string; } -export const apiWhip: FastifyPluginCallback = ( +export const apiWhip: FastifyPluginCallback = async ( fastify, - opts, - next + opts ) => { const productionManager = opts.productionManager; + // Build allowList for rate limiting + const rateLimitAllowList = ['127.0.0.1', '::1', '::ffff:127.0.0.1']; + if (opts.whipGatewayUrl) { + try { + const gatewayHost = new URL(opts.whipGatewayUrl).hostname; + if (gatewayHost && !rateLimitAllowList.includes(gatewayHost)) { + rateLimitAllowList.push(gatewayHost); + + try { + const addresses = await dns.resolve(gatewayHost); + for (const ip of addresses) { + if (!rateLimitAllowList.includes(ip)) { + rateLimitAllowList.push(ip); + } + } + Log().info( + `WHIP rate limit allowList - resolved ${gatewayHost} to IPs: ${addresses.join( + ', ' + )}` + ); + } catch (resolveErr) { + Log().warn( + `Failed to resolve WHIP gateway hostname ${gatewayHost}: ${resolveErr}` + ); + } + } + } catch (err) { + Log().warn( + `Failed to parse WHIP gateway URL for rate limit allowList: ${err}` + ); + } + } + + Log().info(`WHIP rate limit allowList: ${rateLimitAllowList.join(', ')}`); + fastify.addContentTypeParser( 'application/sdp', { parseAs: 'string' }, @@ -50,14 +85,14 @@ export const apiWhip: FastifyPluginCallback = ( opts.smbServerBaseUrl ).toString(); - const smb = opts.smb || new SmbProtocol(); + const smb = new SmbProtocol(); const smbServerApiKey = opts.smbServerApiKey || ''; const coreFunctions = opts.coreFunctions; const whipAuthKey = opts.whipAuthKey?.trim(); async function requireWhipAuth(request: any, reply: any): Promise { if (!whipAuthKey) { - return true; // auth disabled + return true; } const authHeader = @@ -91,11 +126,6 @@ export const apiWhip: FastifyPluginCallback = ( { schema: { description: 'WHIP endpoint for ingesting WebRTC streams', - params: Type.Object({ - productionId: Type.String({ maxLength: 200 }), - lineId: Type.String({ maxLength: 200 }), - username: Type.String({ maxLength: 200 }) - }), body: WhipWhepRequest, response: { 201: WhipWhepResponse, @@ -108,9 +138,15 @@ export const apiWhip: FastifyPluginCallback = ( }, config: { rateLimit: { - max: 10, + max: 100, timeWindow: '1 minute', hook: 'onRequest', + allowList: rateLimitAllowList, + onExceeded: (req) => { + Log().warn( + `Rate limit exceeded for WHIP endpoint - IP: ${req.ip}, URL: ${req.url}` + ); + }, errorResponseBuilder: (_req, context) => { return { statusCode: 429, @@ -150,7 +186,6 @@ export const apiWhip: FastifyPluginCallback = ( lineId ); - // Allocate endpoint with audio support const endpoint = await coreFunctions.createEndpoint( smb, smbServerUrl, @@ -227,7 +262,6 @@ export const apiWhip: FastifyPluginCallback = ( ); // Create the Location URL for the WHIP resource - // Location URL can be relative to Request URL, so this is OK. const locationUrl = `/api/v1/whip/${productionId}/${lineId}/${sessionId}`; // Set response headers @@ -235,13 +269,20 @@ export const apiWhip: FastifyPluginCallback = ( 'Content-Type': 'application/sdp', Location: locationUrl, ETag: sessionId, - Link: getIceServers().join(',') + Link: getIceServers().join(','), + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'GET, POST, DELETE, OPTIONS, PATCH', + 'Access-Control-Allow-Headers': + 'Content-Type, Authorization, ETag, If-Match, Link', + 'Access-Control-Expose-Headers': 'Location, ETag, Link' }); - reply.code(201).send(sdpAnswer); + await reply.code(201).send(sdpAnswer); } catch (err) { Log().error(err); - reply.code(500).send({ error: 'Failed to process WHIP request' }); + reply + .code(500) + .send({ error: `Failed to process WHIP request: ${err}` }); } } ); @@ -253,11 +294,6 @@ export const apiWhip: FastifyPluginCallback = ( { schema: { description: 'Terminate a WHIP connection', - params: Type.Object({ - productionId: Type.String({ maxLength: 200 }), - lineId: Type.String({ maxLength: 200 }), - sessionId: Type.String({ maxLength: 200 }) - }), response: { 200: Type.String({ description: 'OK' }), 404: Type.Object({ error: Type.String() }), @@ -267,8 +303,9 @@ export const apiWhip: FastifyPluginCallback = ( }, async (request, reply) => { if (!(await requireWhipAuth(request, reply))) return; - const { sessionId } = request.params; try { + const { sessionId } = request.params; + Log().info( `Received WHIP DELETE request - sessionId: ${sessionId}, IP: ${request.ip}` ); @@ -282,7 +319,6 @@ export const apiWhip: FastifyPluginCallback = ( return; } - // Remove the user session await opts.dbManager.deleteUserSession(sessionId); productionManager.removeUserSession(sessionId); productionManager.emit('users:change'); @@ -290,13 +326,22 @@ export const apiWhip: FastifyPluginCallback = ( Log().info( `WHIP session deleted successfully - sessionId: ${sessionId}` ); + + reply.headers({ + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'POST, DELETE, OPTIONS', + 'Access-Control-Allow-Headers': 'Content-Type, Authorization' + }); + reply.code(200).send('OK'); } catch (err) { Log().error( - `Failed to delete WHIP session - sessionId: ${sessionId}:`, + `Failed to delete WHIP session - sessionId: ${request.params.sessionId}:`, err ); - reply.code(500).send({ error: 'Failed to terminate WHIP connection' }); + reply + .code(500) + .send({ error: `Failed to terminate WHIP connection: ${err}` }); } } ); @@ -324,7 +369,6 @@ export const apiWhip: FastifyPluginCallback = ( try { const { productionId, lineId } = request.params; - // Check if production and line exist const productionIdNum = parseInt(productionId, 10); if (isNaN(productionIdNum)) { reply.code(400).send({ error: 'Invalid production ID' }); @@ -346,18 +390,23 @@ export const apiWhip: FastifyPluginCallback = ( } reply.headers({ + 'Access-Control-Allow-Origin': '*', + 'Access-Control-Allow-Methods': 'POST, DELETE, OPTIONS, PATCH', + 'Access-Control-Allow-Headers': 'Content-Type, Authorization, ETag', + 'Access-Control-Expose-Headers': 'Location, ETag, Link', + 'Access-Control-Max-Age': '86400', 'Accept-Post': 'application/sdp' }); reply.code(200).send('OK'); } catch (err) { Log().error(err); - reply.code(500).send({ error: 'Failed to process OPTIONS request' }); + reply + .code(500) + .send({ error: `Failed to process OPTIONS request: ${err}` }); } } ); - - next(); }; export default apiWhip; diff --git a/src/bridge_manager.ts b/src/bridge_manager.ts new file mode 100644 index 0000000..10ce8eb --- /dev/null +++ b/src/bridge_manager.ts @@ -0,0 +1,523 @@ +import { DbManager } from './db/interface'; +import { Log } from './log'; +import { BridgeStatus } from './models'; +import { encodeSrtStreamId } from './utils'; + +// The bridge reconcile loop runs every second and logs per transmitter / +// receiver, which floods the terminal. Gate that verbose output behind its own +// flag so it stays quiet even when DEBUG (used by other diagnostic logs) is on. +// Set DEBUG_BRIDGE=true to re-enable it. +const bridgeDebug = (...args: unknown[]): void => { + if (process.env.DEBUG_BRIDGE === 'true') Log().debug(...args); +}; + +export class BridgeManager { + private dbManager: DbManager; + private whipGatewayUrl: string; + private whipGatewayApiKey?: string; + private whepGatewayUrl: string; + private whepGatewayApiKey?: string; + private syncInterval?: NodeJS.Timeout; + private syncIntervalMs = 1000; // 1 second + + constructor( + dbManager: DbManager, + whipGatewayUrl: string | undefined, + whipGatewayApiKey: string | undefined, + whepGatewayUrl: string | undefined, + whepGatewayApiKey: string | undefined + ) { + this.dbManager = dbManager; + this.whipGatewayUrl = whipGatewayUrl || ''; + this.whipGatewayApiKey = whipGatewayApiKey; + this.whepGatewayUrl = whepGatewayUrl || ''; + this.whepGatewayApiKey = whepGatewayApiKey; + } + + // Helper function to call gateway API + private async callGateway( + gatewayUrl: string, + apiKey: string | undefined, + method: string, + path: string, + body?: any + ): Promise { + const url = `${gatewayUrl}${path}`; + const headers: Record = { + 'Content-Type': 'application/json' + }; + if (apiKey) { + headers['x-api-key'] = apiKey; + } + + const options: RequestInit = { + method, + headers + }; + + if (body) { + options.body = JSON.stringify(body); + } + + try { + const response = await fetch(url, options); + + if (!response.ok) { + if (response.status === 404) { + return null; + } + const errorText = await response.text(); + throw new Error( + `Gateway request failed: ${response.status} ${errorText}` + ); + } + + if (response.status === 204 || response.status === 201) { + return null; + } + + const contentType = response.headers.get('content-type'); + const text = await response.text(); + + if (!text) { + return null; + } + + // Only parse as JSON if content-type indicates JSON + if (contentType && contentType.includes('application/json')) { + return JSON.parse(text); + } + + // For text responses, just return null (we don't need the response body) + return null; + } catch (error) { + Log().error(`Failed to call gateway ${url}:`, error); + throw error; + } + } + + // Start the sync service + start() { + Log().info('Starting bridge manager sync service'); + this.syncInterval = setInterval(() => { + this.syncAll().catch((error) => { + Log().error('Error during bridge sync:', error); + }); + }, this.syncIntervalMs); + + // Run initial sync + this.syncAll().catch((error) => { + Log().error('Error during initial bridge sync:', error); + }); + } + + // Stop the sync service + stop() { + if (this.syncInterval) { + clearInterval(this.syncInterval); + this.syncInterval = undefined; + Log().info('Stopped bridge manager sync service'); + } + } + + // Sync all transmitters and receivers + async syncAll() { + const tasks: Promise[] = []; + if (this.whipGatewayUrl) { + tasks.push(this.syncTransmitters()); + } + if (this.whepGatewayUrl) { + tasks.push(this.syncReceivers()); + } + await Promise.all(tasks); + } + + // Sync transmitters with gateway + private async syncTransmitters() { + try { + // Get all transmitters from database + const dbTransmitters = await this.dbManager.getTransmitters(1000, 0); + + // Get all transmitters from gateway + let gatewayTransmitters: any[] = []; + try { + gatewayTransmitters = + (await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'GET', + '/api/v1/tx' + )) || []; + } catch (error) { + Log().warn('Failed to fetch transmitters from gateway:', error); + // Mark all as failed if gateway is unreachable + for (const tx of dbTransmitters) { + if (tx.status !== BridgeStatus.FAILED) { + tx.status = BridgeStatus.FAILED; + await this.dbManager.updateTransmitter(tx); + } + } + return; + } + + // Create a map of gateway transmitters by their ID for quick lookup + const gatewayTxMap = new Map( + gatewayTransmitters.map((t: any) => [t.id, t]) + ); + + // Sync each database transmitter + for (const dbTx of dbTransmitters) { + const existsOnGateway = gatewayTxMap.has(dbTx._id); + + if (!existsOnGateway) { + // Transmitter missing from gateway - recreate it + try { + // Use desiredStatus if set, otherwise use current status + const statusToUse = dbTx.desiredStatus || dbTx.status; + + await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'POST', + '/api/v1/tx/id', + { + id: dbTx._id, + label: dbTx.label, + port: dbTx.port, + mode: dbTx.mode === 'caller' ? 1 : 2, + srtUrl: dbTx.srtUrl?.replace(/^srt:\/\//, ''), + whipUrl: dbTx.whipUrl, + passThroughUrl: dbTx.passThroughUrl, + noVideo: dbTx.noVideo ?? true, + vp8: dbTx.vp8 ?? false, + status: statusToUse + } + ); + Log().info( + `Recreated transmitter on gateway: port ${dbTx.port} with status ${statusToUse}` + ); + + // Update status to match what was set on gateway + if (dbTx.status !== statusToUse) { + dbTx.status = statusToUse; + await this.dbManager.updateTransmitter(dbTx); + } + } catch (error) { + const errorMessage = (error as Error).message || String(error); + + // If transmitter already exists, try to delete and recreate + if (errorMessage.includes('already exists')) { + Log().warn( + `Transmitter ${dbTx._id} already exists on gateway from previous run, removing and recreating` + ); + try { + // Delete the stale transmitter + await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'DELETE', + `/api/v1/tx/id/${dbTx._id}` + ); + + // Retry creation + const statusToUse = dbTx.desiredStatus || dbTx.status; + + await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'POST', + '/api/v1/tx/id', + { + id: dbTx._id, + label: dbTx.label, + port: dbTx.port, + mode: dbTx.mode === 'caller' ? 1 : 2, + srtUrl: dbTx.srtUrl?.replace(/^srt:\/\//, ''), + whipUrl: dbTx.whipUrl, + passThroughUrl: dbTx.passThroughUrl, + noVideo: dbTx.noVideo ?? true, + vp8: dbTx.vp8 ?? false, + status: statusToUse + } + ); + Log().info( + `Successfully recreated transmitter after cleanup: port ${dbTx.port}` + ); + + // Update status + if (dbTx.status !== statusToUse) { + dbTx.status = statusToUse; + await this.dbManager.updateTransmitter(dbTx); + } + } catch (retryError) { + Log().error( + `Failed to recreate transmitter after cleanup: port ${dbTx.port}`, + retryError + ); + if (dbTx.status !== BridgeStatus.FAILED) { + dbTx.status = BridgeStatus.FAILED; + await this.dbManager.updateTransmitter(dbTx); + } + } + } else { + Log().error( + `Failed to recreate transmitter on gateway: port ${dbTx.port}`, + error + ); + if (dbTx.status !== BridgeStatus.FAILED) { + dbTx.status = BridgeStatus.FAILED; + await this.dbManager.updateTransmitter(dbTx); + } + } + } + } else { + // Transmitter exists - check if desired state differs from actual + const gatewayTx = gatewayTxMap.get(dbTx._id); + + if (gatewayTx) { + // If desiredStatus is set and differs from gateway state, enforce it + if (dbTx.desiredStatus && dbTx.desiredStatus !== gatewayTx.status) { + bridgeDebug( + `Transmitter port ${dbTx.port} - Enforcing state change: gateway="${gatewayTx.status}" db="${dbTx.status}" desired="${dbTx.desiredStatus}"` + ); + try { + await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'PUT', + `/api/v1/tx/id/${dbTx._id}/state`, + { + desired: dbTx.desiredStatus + } + ); + bridgeDebug( + `Transmitter port ${dbTx.port} - Successfully enforced desired state: "${dbTx.desiredStatus}"` + ); + + // Update actual status to match desired + dbTx.status = dbTx.desiredStatus; + await this.dbManager.updateTransmitter(dbTx); + } catch (error) { + Log().error( + `Transmitter port ${dbTx.port} - Failed to enforce desired state "${dbTx.desiredStatus}" (will retry):`, + error + ); + } + } else if (gatewayTx.status !== dbTx.status) { + // No desired state set, just sync with gateway + bridgeDebug( + `Transmitter port ${dbTx.port} - Syncing from gateway: "${dbTx.status}" -> "${gatewayTx.status}"` + ); + dbTx.status = gatewayTx.status; + await this.dbManager.updateTransmitter(dbTx); + } + } + } + } + + // Remove orphaned transmitters from gateway (not in database) + const dbIdSet = new Set(dbTransmitters.map((t) => t._id)); + for (const gatewayTx of gatewayTransmitters) { + if (!dbIdSet.has(gatewayTx.id)) { + try { + // Stop transmitter first before deleting + try { + await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'PUT', + `/api/v1/tx/id/${gatewayTx.id}/state`, + { desired: BridgeStatus.STOPPED } + ); + } catch (stopError) { + Log().warn( + `Failed to stop orphaned transmitter before deletion: id ${gatewayTx.id}`, + stopError + ); + } + + // Now delete the transmitter + await this.callGateway( + this.whipGatewayUrl, + this.whipGatewayApiKey, + 'DELETE', + `/api/v1/tx/id/${gatewayTx.id}` + ); + Log().info( + `Removed orphaned transmitter from gateway: id ${gatewayTx.id}` + ); + } catch (error) { + Log().warn( + `Failed to remove orphaned transmitter from gateway: id ${gatewayTx.id}`, + error + ); + } + } + } + } catch (error) { + Log().error('Error syncing transmitters:', error); + } + } + + // Sync receivers with gateway + private async syncReceivers() { + try { + // Get all receivers from database + const dbReceivers = await this.dbManager.getReceivers(1000, 0); + + // Get all receivers from gateway + let gatewayReceivers: any[] = []; + try { + gatewayReceivers = + (await this.callGateway( + this.whepGatewayUrl, + this.whepGatewayApiKey, + 'GET', + '/api/v1/rx' + )) || []; + } catch (error) { + Log().warn('Failed to fetch receivers from gateway:', error); + // Mark all as failed if gateway is unreachable + for (const rx of dbReceivers) { + if (rx.status !== BridgeStatus.FAILED) { + rx.status = BridgeStatus.FAILED; + await this.dbManager.updateReceiver(rx); + } + } + return; + } + + const gatewayIdSet = new Set(gatewayReceivers.map((r: any) => r.id)); + + // Sync each database receiver + for (const dbRx of dbReceivers) { + const existsOnGateway = gatewayIdSet.has(dbRx._id); + + if (!existsOnGateway) { + // Receiver missing from gateway - recreate it + try { + // Use desiredStatus if set, otherwise use current status + const statusToUse = dbRx.desiredStatus || dbRx.status; + + await this.callGateway( + this.whepGatewayUrl, + this.whepGatewayApiKey, + 'POST', + '/api/v1/rx', + { + id: dbRx._id, + whepUrl: dbRx.whepUrl, + srtUrl: encodeSrtStreamId(dbRx.srtUrl), + status: statusToUse + } + ); + Log().info( + `Recreated receiver on gateway: id ${dbRx._id} with status ${statusToUse}` + ); + + // Update status to match what was set on gateway + if (dbRx.status !== statusToUse) { + dbRx.status = statusToUse; + await this.dbManager.updateReceiver(dbRx); + } + } catch (error) { + Log().error( + `Failed to recreate receiver on gateway: id ${dbRx._id}`, + error + ); + if (dbRx.status !== BridgeStatus.FAILED) { + dbRx.status = BridgeStatus.FAILED; + await this.dbManager.updateReceiver(dbRx); + } + } + } else { + // Receiver exists - check if desired state differs from actual + const gatewayRx = gatewayReceivers.find( + (r: any) => r.id === dbRx._id + ); + + if (gatewayRx) { + // If desiredStatus is set and differs from gateway state, enforce it + if (dbRx.desiredStatus && dbRx.desiredStatus !== gatewayRx.status) { + bridgeDebug( + `Receiver ${dbRx._id} - Enforcing state change: gateway="${gatewayRx.status}" db="${dbRx.status}" desired="${dbRx.desiredStatus}"` + ); + try { + await this.callGateway( + this.whepGatewayUrl, + this.whepGatewayApiKey, + 'PUT', + `/api/v1/rx/${dbRx._id}/state`, + { + desired: dbRx.desiredStatus + } + ); + bridgeDebug( + `Receiver ${dbRx._id} - Successfully enforced desired state: "${dbRx.desiredStatus}"` + ); + + // Update actual status to match desired + dbRx.status = dbRx.desiredStatus; + await this.dbManager.updateReceiver(dbRx); + } catch (error) { + Log().error( + `Receiver ${dbRx._id} - Failed to enforce desired state "${dbRx.desiredStatus}" (will retry):`, + error + ); + } + } else if (gatewayRx.status !== dbRx.status) { + // No desired state set, just sync with gateway + bridgeDebug( + `Receiver ${dbRx._id} - Syncing from gateway: "${dbRx.status}" -> "${gatewayRx.status}"` + ); + dbRx.status = gatewayRx.status; + await this.dbManager.updateReceiver(dbRx); + } + } + } + } + + // Remove orphaned receivers from gateway (not in database) + const dbIdSet = new Set(dbReceivers.map((r) => r._id)); + for (const gatewayRx of gatewayReceivers) { + if (!dbIdSet.has(gatewayRx.id)) { + try { + // Stop receiver first before deleting + try { + await this.callGateway( + this.whepGatewayUrl, + this.whepGatewayApiKey, + 'PUT', + `/api/v1/rx/${gatewayRx.id}/state`, + { desired: BridgeStatus.STOPPED } + ); + } catch (stopError) { + Log().warn( + `Failed to stop orphaned receiver before deletion: id ${gatewayRx.id}`, + stopError + ); + } + + // Now delete the receiver + await this.callGateway( + this.whepGatewayUrl, + this.whepGatewayApiKey, + 'DELETE', + `/api/v1/rx/${gatewayRx.id}` + ); + Log().info( + `Removed orphaned receiver from gateway: id ${gatewayRx.id}` + ); + } catch (error) { + Log().warn( + `Failed to remove orphaned receiver from gateway: id ${gatewayRx.id}`, + error + ); + } + } + } + } catch (error) { + Log().error('Error syncing receivers:', error); + } + } +} diff --git a/src/db/couchdb.ts b/src/db/couchdb.ts index 0ade56d..8b641eb 100644 --- a/src/db/couchdb.ts +++ b/src/db/couchdb.ts @@ -5,7 +5,12 @@ import { Line, NewIngest, Production, - UserSession + UserSession, + BridgeStatus, + NewReceiver, + NewTransmitter, + Receiver, + Transmitter } from '../models'; import { assert } from '../utils'; import { DbManager } from './interface'; @@ -30,6 +35,231 @@ export class DbManagerCouchDb implements DbManager { }); } + async addTransmitter(newTransmitter: NewTransmitter): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const now = new Date().toISOString(); + // Generate a unique ID for the transmitter using sequential index + const index = await this.getNextSequence('transmitters'); + const id = `tx-${index}`; + const transmitter: Transmitter = { + _id: id, + ...newTransmitter, + status: BridgeStatus.IDLE, + createdAt: now, + updatedAt: now + }; + const response = await this.nanoDb.insert( + transmitter as unknown as nano.MaybeDocument + ); + if (!response.ok) throw new Error('Failed to insert transmitter'); + return transmitter; + } + + async getTransmitter(id: string): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + try { + const transmitter = await this.nanoDb.get(id); + return transmitter as any | undefined; + } catch (error) { + return undefined; + } + } + + async getTransmitters(limit: number, offset: number): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const transmitters: Transmitter[] = []; + const response = await this.nanoDb.list({ + include_docs: true + }); + response.rows.forEach((row: any) => { + if (row.doc._id.toLowerCase().startsWith('tx-')) { + transmitters.push(row.doc); + } + }); + + const result = transmitters.slice(offset, offset + limit); + return result as any as Transmitter[]; + } + + async getTransmittersLength(): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const response = await this.nanoDb.list({ include_docs: false }); + const filteredRows = response.rows.filter((row: any) => + row.id.toLowerCase().startsWith('tx-') + ); + return filteredRows.length; + } + + async updateTransmitter( + transmitter: Transmitter + ): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const now = new Date().toISOString(); + try { + const existingTransmitter = await this.nanoDb.get(transmitter._id); + const updatedTransmitter = { + ...existingTransmitter, + ...transmitter, + updatedAt: now + }; + const response = await this.nanoDb.insert(updatedTransmitter); + return response.ok ? { ...transmitter, updatedAt: now } : undefined; + } catch (error) { + return undefined; + } + } + + async deleteTransmitter(id: string): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + try { + const transmitter = await this.nanoDb.get(id); + const response = await this.nanoDb.destroy( + transmitter._id, + transmitter._rev + ); + return response.ok; + } catch (error) { + return false; + } + } + async addReceiver(newReceiver: NewReceiver): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const now = new Date().toISOString(); + const index = await this.getNextSequence('receivers'); + const id = `rx-${index}`; + const receiver: Receiver = { + _id: id, + ...newReceiver, + status: BridgeStatus.IDLE, + createdAt: now, + updatedAt: now + }; + const response = await this.nanoDb.insert( + receiver as unknown as nano.MaybeDocument + ); + if (!response.ok) throw new Error('Failed to insert receiver'); + return receiver; + } + + async getReceiver(id: string): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + try { + const receiver = await this.nanoDb.get(id); + return receiver as any | undefined; + } catch (error) { + return undefined; + } + } + + async getReceivers(limit: number, offset: number): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const receivers: Receiver[] = []; + const response = await this.nanoDb.list({ + include_docs: true + }); + response.rows.forEach((row: any) => { + if ( + row.doc._id.toLowerCase().indexOf('counter') === -1 && + row.doc._id.toLowerCase().indexOf('session_') === -1 && + row.doc._id.toLowerCase().startsWith('rx-') + ) { + receivers.push(row.doc); + } + }); + + const result = receivers.slice(offset, offset + limit); + return result as any as Receiver[]; + } + + async getReceiversLength(): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const response = await this.nanoDb.list({ include_docs: false }); + const filteredRows = response.rows.filter( + (row: any) => + row.id.toLowerCase().indexOf('counter') === -1 && + row.id.toLowerCase().indexOf('session_') === -1 && + row.id.toLowerCase().startsWith('rx-') + ); + return filteredRows.length; + } + + async updateReceiver(receiver: Receiver): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + const now = new Date().toISOString(); + try { + const existingReceiver = await this.nanoDb.get(receiver._id); + const updatedReceiver = { + ...existingReceiver, + ...receiver, + _id: receiver._id, + updatedAt: now + }; + const response = await this.nanoDb.insert(updatedReceiver); + return response.ok ? receiver : undefined; + } catch (error) { + return undefined; + } + } + + async deleteReceiver(id: string): Promise { + await this.connect(); + if (!this.nanoDb) { + throw new Error('Database not connected'); + } + + try { + const receiver = await this.nanoDb.get(id); + const response = await this.nanoDb.destroy(receiver._id, receiver._rev); + return response.ok; + } catch (error) { + return false; + } + } + async connect(): Promise { if (!this.nanoDb) { const maxRetries = 5; diff --git a/src/db/interface.ts b/src/db/interface.ts index 9c56726..65f2cc2 100644 --- a/src/db/interface.ts +++ b/src/db/interface.ts @@ -4,7 +4,11 @@ import { Line, NewIngest, Production, - UserSession + UserSession, + NewReceiver, + NewTransmitter, + Receiver, + Transmitter } from '../models'; export interface DbManager { @@ -35,6 +39,18 @@ export interface DbManager { updates: Partial ): Promise; getSessionsByQuery(q: Partial): Promise; + addTransmitter(transmitter: NewTransmitter): Promise; + getTransmitter(id: string): Promise; + getTransmitters(limit: number, offset: number): Promise; + getTransmittersLength(): Promise; + updateTransmitter(transmitter: Transmitter): Promise; + deleteTransmitter(id: string): Promise; + addReceiver(receiver: NewReceiver): Promise; + getReceiver(id: string): Promise; + getReceivers(limit: number, offset: number): Promise; + getReceiversLength(): Promise; + updateReceiver(receiver: Receiver): Promise; + deleteReceiver(id: string): Promise; addPreset(preset: Omit): Promise; getPreset(id: string): Promise; getPresets(): Promise; diff --git a/src/db/mongodb.ts b/src/db/mongodb.ts index bba6767..2ad126b 100644 --- a/src/db/mongodb.ts +++ b/src/db/mongodb.ts @@ -6,7 +6,12 @@ import { Line, NewIngest, Production, - UserSession + UserSession, + BridgeStatus, + NewReceiver, + NewTransmitter, + Receiver, + Transmitter } from '../models'; import { v4 as uuidv4 } from 'uuid'; import { assert } from '../utils'; @@ -16,8 +21,10 @@ import { Log } from '../log'; const SESSION_PRUNE_SECONDS = 7_200; export class DbManagerMongoDb implements DbManager { private client: MongoClient; + private readonly databaseName: string; constructor(dbConnectionUrl: URL) { + this.databaseName = process.env.MONGODB_DATABASE ?? 'intercom-manager'; this.client = new MongoClient(dbConnectionUrl.toString()); } @@ -360,4 +367,132 @@ export class DbManagerMongoDb implements DbManager { return sessions.find(mongoQuery).toArray(); } + + // Transmitter operations + async addTransmitter(newTransmitter: NewTransmitter): Promise { + const db = this.client.db(this.databaseName); + const now = new Date().toISOString(); + const index = await this.getNextSequence('transmitters'); + const id = `tx-${index}`; + const transmitter: Transmitter = { + _id: id, + ...newTransmitter, + status: BridgeStatus.IDLE, + createdAt: now, + updatedAt: now + }; + await db + .collection('transmitters') + .insertOne(transmitter as any); + return transmitter; + } + + async getTransmitter(id: string): Promise { + const db = this.client.db(this.databaseName); + return db + .collection('transmitters') + .findOne({ _id: id } as any) as any | undefined; + } + + async getTransmitters(limit: number, offset: number): Promise { + const db = this.client.db(this.databaseName); + const transmitters = await db + .collection('transmitters') + .find() + .sort({ label: 1, createdAt: -1 }) + .skip(offset) + .limit(limit) + .toArray(); + return transmitters as Transmitter[]; + } + + async getTransmittersLength(): Promise { + const db = this.client.db(this.databaseName); + return await db.collection('transmitters').countDocuments(); + } + + async updateTransmitter( + transmitter: Transmitter + ): Promise { + const db = this.client.db(this.databaseName); + const now = new Date().toISOString(); + const result = await db + .collection('transmitters') + .updateOne({ _id: transmitter._id } as any, { + $set: { ...transmitter, updatedAt: now } + }); + return result.modifiedCount === 1 + ? { ...transmitter, updatedAt: now } + : undefined; + } + + async deleteTransmitter(id: string): Promise { + const db = this.client.db(this.databaseName); + const result = await db + .collection('transmitters') + .deleteOne({ _id: id } as any); + return result.deletedCount === 1; + } + + // Receiver operations + async addReceiver(newReceiver: NewReceiver): Promise { + const db = this.client.db(this.databaseName); + const now = new Date().toISOString(); + const index = await this.getNextSequence('receivers'); + const id = `rx-${index}`; + const receiver: Receiver = { + _id: id, + ...newReceiver, + status: BridgeStatus.IDLE, + createdAt: now, + updatedAt: now + }; + await db.collection('receivers').insertOne(receiver as any); + return receiver; + } + + async getReceiver(id: string): Promise { + const db = this.client.db(this.databaseName); + return db.collection('receivers').findOne({ _id: id } as any) as + | any + | undefined; + } + + async getReceivers(limit: number, offset: number): Promise { + const db = this.client.db(this.databaseName); + const receivers = await db + .collection('receivers') + .find() + .sort({ label: 1, createdAt: -1 }) + .skip(offset) + .limit(limit) + .toArray(); + return receivers as Receiver[]; + } + + async getReceiversLength(): Promise { + const db = this.client.db(this.databaseName); + return await db.collection('receivers').countDocuments(); + } + + async updateReceiver(receiver: Receiver): Promise { + const db = this.client.db(this.databaseName); + const now = new Date().toISOString(); + const result = await db + .collection('receivers') + .updateOne({ _id: receiver._id } as any, { + $set: { ...receiver, updatedAt: now } + }); + return result.modifiedCount === 1 + ? { ...receiver, updatedAt: now } + : undefined; + } + + async deleteReceiver(id: string): Promise { + const db = this.client.db(this.databaseName); + const result = await db + .collection('receivers') + .deleteOne({ _id: id } as any); + return result.deletedCount === 1; + } } diff --git a/src/models.ts b/src/models.ts index b38f1ea..bafda27 100644 --- a/src/models.ts +++ b/src/models.ts @@ -401,6 +401,121 @@ export const PatchIngest = Type.Union([ export const PatchIngestResponse = Type.Omit(Ingest, ['ipAddress']); +// Bridge IO - Transmitters and Receivers + +export enum BridgeStatus { + IDLE = 'idle', + RUNNING = 'running', + STOPPED = 'stopped', + FAILED = 'failed' +} + +export const Transmitter = Type.Object({ + _id: Type.String({ description: 'Unique transmitter ID (port as string)' }), + label: Type.Optional(Type.String({ description: 'Human-readable label' })), + port: Type.Number({ description: 'SRT port' }), + productionId: Type.Number({ description: 'Production ID' }), + lineId: Type.Number({ description: 'Line ID' }), + whipUrl: Type.String({ description: 'WHIP URL' }), + passThroughUrl: Type.Optional( + Type.String({ description: 'Pass-through to SRT ingest URL' }) + ), + mode: Type.Union([Type.Literal('caller'), Type.Literal('listener')]), + srtUrl: Type.Optional( + Type.String({ description: 'SRT URL (for caller mode)' }) + ), + noVideo: Type.Optional(Type.Boolean({ description: 'Disable video stream' })), + vp8: Type.Optional(Type.Boolean({ description: 'Transcode video to VP8' })), + bypassVideo: Type.Optional( + Type.Boolean({ description: 'Pass H264 video through without transcoding' }) + ), + status: Type.Enum(BridgeStatus), + desiredStatus: Type.Optional(Type.Enum(BridgeStatus)), + createdAt: Type.Optional(Type.String({ format: 'date-time' })), + updatedAt: Type.Optional(Type.String({ format: 'date-time' })) +}); +export type Transmitter = Static; + +export const NewTransmitter = Type.Object({ + label: Type.Optional(Type.String()), + port: Type.Number(), + productionId: Type.Number({ description: 'Production ID' }), + lineId: Type.Number({ description: 'Line ID' }), + whipUrl: Type.String({ description: 'WHIP URL' }), + passThroughUrl: Type.Optional(Type.String()), + mode: Type.Union([Type.Literal('caller'), Type.Literal('listener')]), + srtUrl: Type.Optional(Type.String()), + noVideo: Type.Optional(Type.Boolean({ description: 'Disable video stream' })), + vp8: Type.Optional(Type.Boolean({ description: 'Transcode video to VP8' })), + bypassVideo: Type.Optional( + Type.Boolean({ description: 'Pass H264 video through without transcoding' }) + ) +}); +export type NewTransmitter = Static; + +export const TransmitterListResponse = Type.Object({ + transmitters: Type.Array(Transmitter), + offset: Type.Number(), + limit: Type.Number(), + totalItems: Type.Number() +}); +export type TransmitterListResponse = Static; + +export const TransmitterStateChange = Type.Object({ + desired: Type.Enum(BridgeStatus) +}); +export type TransmitterStateChange = Static; + +export const Receiver = Type.Object({ + _id: Type.String({ description: 'Unique receiver ID' }), + label: Type.Optional(Type.String({ description: 'Human-readable label' })), + productionId: Type.Number({ description: 'Production ID' }), + lineId: Type.Number({ description: 'Line ID' }), + whepUrl: Type.String({ description: 'WHEP URL' }), + srtUrl: Type.String({ description: 'SRT output URL' }), + status: Type.Enum(BridgeStatus), + desiredStatus: Type.Optional(Type.Enum(BridgeStatus)), + createdAt: Type.Optional(Type.String({ format: 'date-time' })), + updatedAt: Type.Optional(Type.String({ format: 'date-time' })) +}); +export type Receiver = Static; + +export const NewReceiver = Type.Object({ + label: Type.Optional(Type.String()), + productionId: Type.Number({ description: 'Production ID' }), + lineId: Type.Number({ description: 'Line ID' }), + whepUrl: Type.String({ description: 'WHEP URL' }), + srtUrl: Type.String() +}); +export type NewReceiver = Static; + +export const ReceiverListResponse = Type.Object({ + receivers: Type.Array(Receiver), + offset: Type.Number(), + limit: Type.Number(), + totalItems: Type.Number() +}); +export type ReceiverListResponse = Static; + +export const ReceiverStateChange = Type.Object({ + desired: Type.Enum(BridgeStatus) +}); +export type ReceiverStateChange = Static; + +export const PatchTransmitter = Type.Object({ + label: Type.Optional(Type.String()), + productionId: Type.Optional(Type.Number()), + lineId: Type.Optional(Type.Number()) +}); +export type PatchTransmitter = Static; + +export const PatchReceiver = Type.Object({ + label: Type.Optional(Type.String()), + productionId: Type.Optional(Type.Number()), + lineId: Type.Optional(Type.Number()) +}); +export type PatchReceiver = Static; + export const PresetCall = Type.Object({ productionId: Type.String({ minLength: 1 }), lineId: Type.String({ minLength: 1 }), diff --git a/src/server.ts b/src/server.ts index 92739ad..25e50c0 100644 --- a/src/server.ts +++ b/src/server.ts @@ -1,5 +1,6 @@ import api from './api'; import { CoreFunctions } from './api_productions_core_functions'; +import { BridgeManager } from './bridge_manager'; import { ConnectionQueue } from './connection_queue'; import { DbManagerCouchDb } from './db/couchdb'; import { DbManagerMongoDb } from './db/mongodb'; @@ -9,6 +10,8 @@ import { ProductionManager } from './production_manager'; const SMB_ADDRESS: string = process.env.SMB_ADDRESS ?? 'http://localhost:8080'; const PUBLIC_HOST: string = process.env.PUBLIC_HOST ?? 'http://localhost:8000'; +const WHIP_GATEWAY_URL: string | undefined = process.env.WHIP_GATEWAY_URL; +const WHEP_GATEWAY_URL: string | undefined = process.env.WHEP_GATEWAY_URL; if (!process.env.SMB_ADDRESS) { Log().warn('SMB_ADDRESS environment variable not set, using defaults'); @@ -42,6 +45,22 @@ if (dbUrl.protocol === 'mongodb:' || dbUrl.protocol === 'mongodb+srv:') { const ingestManager = new IngestManager(dbManager); await ingestManager.load(); + // Initialize bridge manager for gateway sync (only if gateways are configured) + let bridgeManager: BridgeManager | null = null; + if (WHIP_GATEWAY_URL || WHEP_GATEWAY_URL) { + bridgeManager = new BridgeManager( + dbManager, + WHIP_GATEWAY_URL || '', + process.env.WHIP_GATEWAY_API_KEY, + WHEP_GATEWAY_URL || '', + process.env.WHEP_GATEWAY_API_KEY + ); + bridgeManager.start(); + Log().info('Bridge manager started'); + } else { + Log().info('Bridge manager disabled (no gateway URLs configured)'); + } + const server = await api({ title: 'intercom-manager', smbServerBaseUrl: SMB_ADDRESS, @@ -52,7 +71,11 @@ if (dbUrl.protocol === 'mongodb:' || dbUrl.protocol === 'mongodb+srv:') { dbManager: dbManager, productionManager: productionManager, ingestManager: ingestManager, - coreFunctions: new CoreFunctions(productionManager, connectionQueue) + coreFunctions: new CoreFunctions(productionManager, connectionQueue), + whipGatewayUrl: WHIP_GATEWAY_URL || '', + whipGatewayApiKey: process.env.WHIP_GATEWAY_API_KEY, + whepGatewayUrl: WHEP_GATEWAY_URL || '', + whepGatewayApiKey: process.env.WHEP_GATEWAY_API_KEY }); server.listen({ port: PORT, host: '0.0.0.0' }, (err, address) => { @@ -63,10 +86,19 @@ if (dbUrl.protocol === 'mongodb:' || dbUrl.protocol === 'mongodb+srv:') { Log().info( `Media Bridge at ${SMB_ADDRESS} (${ENDPOINT_IDLE_TIMEOUT_S}s idle timeout)` ); + if (WHIP_GATEWAY_URL) { + Log().info(`WHIP Gateway at ${WHIP_GATEWAY_URL}`); + } + if (WHEP_GATEWAY_URL) { + Log().info(`WHEP Gateway at ${WHEP_GATEWAY_URL}`); + } }); const shutdown = async (signal: string) => { Log().info(`${signal} received, shutting down gracefully`); + if (bridgeManager) { + bridgeManager.stop(); + } await server.close(); process.exit(0); }; diff --git a/src/utils.ts b/src/utils.ts index e207584..6a873e2 100644 --- a/src/utils.ts +++ b/src/utils.ts @@ -52,3 +52,9 @@ export function getIceServers(): string[] { return links; } + +export function encodeSrtStreamId(srtUrl: string): string { + return srtUrl.replace(/([?&]streamid=)([^&]*)/i, (_, prefix, value) => { + return prefix + encodeURIComponent(decodeURIComponent(value)); + }); +}