import { readFileSync } from 'node:fs'
import { timingSafeEqual } from 'node:crypto'
import https from 'node:https'
import { WebSocketServer } from 'ws'
const username = process.env.FORK_USERNAME ?? ''
const password = process.env.FORK_PASSWORD ?? ''
if (!username || !password) throw new Error('Set FORK_USERNAME and FORK_PASSWORD')
function sameSecret(received, expected) {
const a = Buffer.from(received)
const b = Buffer.from(expected)
return a.length === b.length && timingSafeEqual(a, b)
}
function hasValidBasicAuth(header) {
if (!header?.startsWith('Basic ')) return false
const decoded = Buffer.from(header.slice(6), 'base64').toString('utf8')
const separator = decoded.indexOf(':')
if (separator < 0) return false
return sameSecret(decoded.slice(0, separator), username) &&
sameSecret(decoded.slice(separator + 1), password)
}
const keyFile = process.env.TLS_KEY_FILE
const certFile = process.env.TLS_CERT_FILE
if (!keyFile || !certFile) throw new Error('Set TLS_KEY_FILE and TLS_CERT_FILE')
const server = https.createServer({
key: readFileSync(keyFile),
cert: readFileSync(certFile),
})
const wss = new WebSocketServer({ noServer: true })
server.on('upgrade', (request, socket, head) => {
const path = new URL(request.url ?? '/', 'https://receiver.invalid').pathname
if (path !== '/orbit-fork' || !hasValidBasicAuth(request.headers.authorization)) {
socket.write('HTTP/1.1 401 Unauthorized\r\nConnection: close\r\nWWW-Authenticate: Basic realm="media-fork"\r\n\r\n')
socket.destroy()
return
}
wss.handleUpgrade(request, socket, head, (ws) => {
wss.emit('connection', ws, request)
})
})
wss.on('connection', (ws) => {
let forkId
ws.on('message', (data, isBinary) => {
if (isBinary) {
const frame = Array.isArray(data) ? Buffer.concat(data) : Buffer.from(data)
// Write raw PCM to stdout for a downstream consumer; keep logs on stderr.
// Pause this socket if stdout is backpressured to avoid unbounded buffering.
if (!process.stdout.write(frame)) {
ws.pause()
process.stdout.once('drain', () => ws.resume())
}
return
}
let message
try {
message = JSON.parse(data.toString())
} catch {
ws.close(1003, 'Expected JSON control message')
return
}
if (message.type === 'metadata') {
forkId = message.fork_id
console.error('Media fork connected', {
forkId,
roomId: message.room_id,
participantId: message.participant_id,
sampleRate: message.sampleRate,
mixType: message.mixType,
channels: message.channels,
encoding: message.encoding,
frameMs: message.frame_ms,
metadata: message.metadata,
})
} else if (message.type === 'stopped') {
console.error('Media fork stopped', { forkId: message.fork_id, reason: message.reason })
}
})
ws.on('close', () => {
// Orbit may reconnect once after a customer-side socket loss. A new
// connection starts with metadata again, so initialize per connection.
console.error('Media fork socket closed', { forkId })
})
})
server.listen(443)