@codefusion-cc/realtime-media
v0.1.3
Published
Calls for realtime apps: peer-to-peer mesh and Cloudflare Realtime SFU sessions with simulcast, TURN credentials, the room broker and account budgets that bound them, the policy that decides how a call travels, camera profiles, device choice, speaking det
Downloads
1,295
Maintainers
Readme
@codefusion-cc/realtime-media
Calls for apps on @codefusion-cc/realtime: peer to peer first, through Cloudflare Realtime SFU with simulcast
when peer to peer fails or a room outgrows it, with TURN relays and per-account budgets that bound what one
account can allocate. Around them, the page's side of any WebRTC app: choosing and remembering cameras,
microphones and speakers, telling who is speaking, reading what a connection carries from getStats(), knowing
which videos anyone looks at, and an on-device screen that blurs explicit pictures and live video.
npm install @codefusion-cc/realtime-media| Entry | Runs in | What it has |
| --- | --- | --- |
| @codefusion-cc/realtime-media | Worker | a room's calls (RoomCalls), its SFU broker (RoomMediaBroker), the gateway's part (mediaGateway), TURN credentials, the SFU and TURN health check |
| @codefusion-cc/realtime-media/browser | page | capture (LocalCapture), the call sessions (MeshSession, SfuSession), TURN credentials, call stats, devices, speaking, video demand, the screen |
| @codefusion-cc/realtime-media/react | page | useSpeakingStreams |
| @codefusion-cc/realtime-media/policy | both | how a call travels (callTransport, CallSettings) and camera profiles (CameraProfile, simulcast layers) |
How a call travels
The room decides for everyone in its call and tells them as mediaTransport: { mode, relay, reason }.
callTransport() takes the app's settings (peer to peer first or not, when the SFU and TURN may be used, how many
people a peer-to-peer call holds) and the month's egress, and answers p2p, sfu or off. An account or room
with priority (a paying member, say) always has both paid paths. The call's rules (CallRules) say what kind of
call it is: growsToSfu moves it to the SFU past meshMax (a team's, a public room's), and without p2pRelay it
gets no TURN relay peer to peer (a public room's strangers). A two-person call has neither rule's limit. The page runs the session the mode asks for
and swaps it when the mode changes, keeping the camera and microphone.
In the Worker
import { combineGatewayApps, defineUserGateway } from '@codefusion-cc/realtime'
import { mediaGateway, RoomCalls, RoomMediaBroker, MEDIA_SCHEMA, ROOM_CALLS_SCHEMA } from '@codefusion-cc/realtime-media'
// The account's budget: 8 SFU allocations a minute, 12 live connections across its rooms.
export const UserGateway = defineUserGateway(combineGatewayApps(own, mediaGateway({ rooms: 'ROOM_CHANNELS' })))
export class RoomChannel extends DurableObject {
constructor(ctx, env) {
super(ctx, env)
ctx.storage.sql.exec(MEDIA_SCHEMA + ROOM_CALLS_SCHEMA)
this.media = new RoomMediaBroker(ctx, {
env: () => env, eligible, authorize, kind,
transport: () => this.calls.transport(),
schedule: () => this.scheduleAlarm(), // the room's one alarm, at the earliest of its deadlines and media.deadlines()
})
this.calls = new RoomCalls(ctx, {
env: () => env, broker: this.media, attachmentOf, update, sockets, policy,
rules: () => this.isPublic() ? { growsToSfu: true, p2pRelay: false } : {},
})
}
async webSocketMessage(socket, raw) {
const frame = JSON.parse(raw)
// Checks, rate-limits and answers it; tells everyone the call's state (`CallFrame`) when it changed.
if (frame.type === 'media') return this.calls.handleMediaFrame(socket, frame, { tell: call => this.broadcast({ type: 'call', ...call }) })
}
// The gateway asks a room whether a slot it holds is released.
mediaResourceState(input) { return this.media.resourceState(input) }
}- Members hear the call's state (
calls.callFrame()) in their snapshot, and again only when it changed:calls.tellIfChanged(tell)after anything that may change it (a member leaving, the alarm), andcalls.heard(call)when the room tells everyone itself.calls.forget(connectionId)once a socket is gone. - The broker owns every SFU operation: browsers send room commands, and SFU ids and secrets never reach them. Each provider mutation is persisted before it is sent, so an evicted room never replays one; it retires the connection, and the account's slot comes back once the room says it is quiet.
- A retired connection's SFU session is closed from the room's alarm: retried within the first minute, then at most
a minute apart while Cloudflare answers and up to an hour apart while it does not. A session Cloudflare never
reports ended is given up a day after its cleanup came due: counted in
media_cleanup_abandoned, its slot released. Aconsole.warnline says when the SFU first stops answering about a session, and when one is given up. - Bindings:
USER_GATEWAYSand the rooms' binding,MEDIA_RATE_LIMITER(TURN credential requests), andMEDIA_USAGE(the Analytics Engine dataset CodeFusion Console reads calls from). - Secrets:
SFU_APP_IDandSFU_APP_SECRETfor the SFU,TURN_KEY_IDandTURN_KEY_API_TOKENfor TURN. Without the SFU's the room never chooses it; without TURN's, browsers get Cloudflare's STUN only. - Health:
handleMediaHealthanswersGET MEDIA_HEALTH_PATH(an SFU session and TURN credentials, at most once a minute), andcheckMediaHealthfrom a cron (MEDIA_HEALTH_CRON) logs a failure when it happens.
A call in an authorized channel
channelCalls(host, { policy, priority, growsToSfu, kind }) is the extension (defineAuthorizedChannel({ extend })
of @codefusion-cc/realtime) that hosts a call in a workspace or conversation: everyone the channel admits may join,
peer to peer first and through the SFU when the policy allows, with the same rules as a room (RoomCalls,
RoomMediaBroker). priority says whether a member's account may use relays and the SFU (a paid plan); its
presence gives the call relays. The channel class answers its gateways with
mediaResourceState(input) { return this.extension.resourceState(input) }.
In the page
import { LocalCapture, MeshSession, SfuSession } from '@codefusion-cc/realtime-media/browser'
import { cameraConstraints, cameraEncodings } from '@codefusion-cc/realtime-media/policy'
const capture = new LocalCapture({ onChange, videoConstraints: () => cameraConstraints(profile) })
const session = transport.mode === 'sfu'
? new SfuSession({ userId, iceServers, request, onChange, videoEncodings: cameraEncodings(profile), demand })
: new MeshSession({ userId, connectionId, iceServers, request, onChange })request sends a call command over the room's channel (openDurableChannel from @codefusion-cc/realtime/browser),
and requestIceCredentials asks the room for TURN credentials before they expire.
Peer to peer, MeshSession keeps one PeerLink from @codefusion-cc/realtime-p2p for each person: the link
negotiates, batches candidates and restarts ICE once; the session adds the tracks and each person's share of the
bitrate, and tells the room when a link fails so it may move the call to the SFU. The room checks every signal with
the same parsePeerSignal (@codefusion-cc/realtime-p2p/protocol) the links answer them with.
An account this device blocked (isBlocked) is never connected to: the session declines each of its connections
(p2p-decline), and the room passes the decline on as a media-signal without data. That side then stops
connecting to this one, without a failure, and the room counts no failure between the two as a reason to move the
call to the SFU: a block hides media between those two people only.
A connection has an audio and a video section from its first offer, whatever either side sends, so a kind sent later only changes a direction, and a camera turned off or changed needs no new offer. When the first offers cross, the side that gives way answers on a new connection instead of rolling its own offer back: libwebrtc (Chromium, Safari) keeps a rolled-back offer's header extension ids, and could refuse that side's next offer. An offer the browser refuses all the same counts as peer to peer failing, which moves the call to the SFU.
CallController does all of that for a page: it takes the channel's frames and connection state, starts the
session the channel's transport asks for and swaps it when the call moves, renews TURN credentials, and keeps a
reconnect from replaying anything of the last socket. Framework-free; in React, read it with useSyncExternalStore.
import { CallController } from '@codefusion-cc/realtime-media/browser'
import { cameraVideo } from '@codefusion-cc/realtime-media/policy'
const call = new CallController({
userId,
request: (payload, options) => room.requestServer('media', payload, options),
video: cameraVideo(() => profile),
onChange: state => render(state),
})
call.join()
// From the channel: call.connection(status), call.frame(frame). When the profile changes: call.updateVideo().The screen takes the app's classifier (NsfwClassifier): the package decides what to sample and when to blur and
clear, the app chooses the model. It is a receiver-side aid, not moderation: a modified client can skip it, and it
fails open when the model cannot run. A model that takes longer than loadTimeoutMs (60 s) to load, or
classifyTimeoutMs (15 s) to classify, counts as one that cannot run: every call waiting on it is answered, the
signal loadClassifier was given aborts so the app can stop the load, and a model delivered after that is disposed.
In tests
@codefusion-cc/realtime-media/testing is the other side of a call: fakeMediaDevices() for
navigator.mediaDevices (a request can be refused, or held like an open permission prompt until the test grants
or refuses it), FakeMediaStreamTrack (a MediaStreamTrack whose device can go away), FakeMediaStream,
fakePeerConnections() for RTCPeerConnection (offer/answer as WebRTC has them, connecting by itself once both
descriptions are applied, stats for what getStats reports; options leave connecting and ICE gathering to the test
and make a browser that keeps fewer simulcast layers than asked or refuses them, and holdNext/refuseNext/refuseAll
make a browser call wait or fail, a held operation holding those called after it), and fakeCallRoom(), whose requestServer answers a page's call commands as the
room does and records them and its answers (sent, ops(), payloads(op), answered(op)); a test changes its
answer to an op through answers, given the room's own to alter or delay, or holds an op's answers until it releases
them (hold(op)). iceAnswer() is its
answer for ICE servers. An app's call tests go through the real sessions with them, mocking nothing of this package.
The fake peer connections follow JSEP where the two sides send different kinds of media: a rollback, or an offer
taken over this side's own, returns to the last stable state; a remote offer's new section takes a transceiver
addTrack made for its kind, else one that only receives; an answer sends and receives only what both sides allow;
a received track is muted when the other side stops sending it, and unmuted when it sends again;
negotiationneeded fires in a task (under fake timers, once they advance), in stable, while something is left to
negotiate, and never once the connection is closed. The same cases run in
Chromium and Firefox on their own RTCPeerConnection (npm run test:browser), which holds the fake to them.
Chromium and Safari (libwebrtc) fail one: after offers of different kinds cross, the side that rolled back can
have its follow-up offer refused for a header extension id collision, which the fake does not reproduce.
test/meshSession.browser.test.ts connects MeshSessions in Chromium, with media, through a room in the page;
headless Firefox gathers no ICE candidates under Playwright, so it skips them.
