MSS (Media Streaming Service) — an open-source, Rust-built centralized media plane for SIP infrastructures whose media anchors in rtpengine.
If your calls flow through rtpengine (OpenSIPS/Kamailio in front, FreeSWITCH
or anything else behind), MSS can tap any call's audio without touching the
call: it asks rtpengine for a copy of each participant's media (NG protocol
subscribe request/answer — the SIPREC mechanism) and fans it out to as many
consumers as you attach — real-time transcription, ASR, recording, voice-AI
bridges. No media bugs, no dummy legs, no conference tricks, and taps are
re-creatable, so a lost pod re-subscribes instead of losing the call.
The longer mission (see docs/roadmap.md): remove every media workload from FreeSWITCH phase by phase — passive fan-out, then recording, then inline interactive media (voice-AI speech in), then mixing — until the softswitch does only call control, or nothing at all.
- Tap ingest from rtpengine ≥ mr14 (
subscribe/unsubscribe), with per-leg jitter buffering, G.711 decode (µ-law and A-law, whatever arrives), RFC 4733 DTMF, and speaker attribution by SSRC correlation. - Fan-out hub: N consumers per call, attach/detach mid-call, bounded per-consumer queues with counted drop-oldest — a slow consumer only ever hurts itself.
- Consumer transports: a WebSocket adapter speaking the Twilio Media
Streams dialect (drop-in for tooling that already consumes it), and a
native gRPC bidirectional stream (
proto/mediastream.proto) with binary frames, capability-checked audio injection, and mark/clear barge-in semantics. - Control plane: the
MediaControlgRPC API (proto/mediacontrol.proto) over sessions / attachments / playbacks, typed events onto Kafka (mss.events, gapless per-session sequence), a Redis session registry with ownership leases and automatic re-subscribe after pod loss, optional bearer-token auth, and Prometheus metrics with shipped alert rules (deploy/prometheus-alerts.yaml). - Audio injection into tapped calls via rtpengine
play media(utterance-shaped; streaming TTS arrives with Phase 3 inline legs). - A legacy façade (
TelCompat): serves a FreeSWITCH-controller-style verb API (StartStream/StartRecording/…) byte-compatibly on the same port, so an existing controller can be pointed at MSS by config flag and rolled back the same way. Optional — skip it if you have no such controller.
cargo run -p mediaserverd # needs Rust plus cmake/make/g++ (libopus is vendored)Environment (all optional except the listen address):
| Variable | What it does |
|---|---|
MSS_CONTROL_LISTEN |
ip:port for the gRPC control+data plane |
MSS_RTPENGINE_NODE |
default rtpengine NG address (ip:port) |
MSS_TAP_LOCAL_IP |
address rtpengine sends tap media to (must be routable from rtpengine) |
MSS_TAP_TRANSCODE |
on (default) asks rtpengine to transcode the tap to MSS_TAP_FORMAT; off accepts the call's own codec, which keeps the tap eligible for rtpengine's in-kernel path but refuses a call whose codec MSS cannot decode |
MSS_TAP_FORMAT |
what the tap decodes: pcmu (default), pcma or opus |
MSS_OPUS_DECODE_RATE_HZ |
Opus decode rate, default 16000; one of 8000/12000/16000/24000/48000 |
MSS_KAFKA_BROKERS |
comma list; unset = events stay in-process |
MSS_REDIS_URL |
session registry; unset = sessions die with the pod |
MSS_METRICS_LISTEN |
ip:port for Prometheus /metrics |
MSS_AUTH_TOKEN |
shared bearer secret; unset = open (lab mode) |
MSS_POD_NAME |
this pod's identity in the registry |
Then drive it:
cargo run -p control-api --example mss_ctl -- http://127.0.0.1:50551 \
create req-1 <sip-call-id> <caller-from-tag> <rtpengine-ip:port>
cargo run -p control-api --example mss_ctl -- http://127.0.0.1:50551 \
attach req-1 ws://your-consumer/wsA full docker-compose lab (OpenSIPS + FreeSWITCH + rtpengine + Redpanda +
synthetic callers and mock consumers) lives in lab/.
| Crate | What it is |
|---|---|
crates/media-core |
Sans-IO media pipeline core: RTP parse/serialize, G.711, Opus decode, RFC 4733 DTMF, jitter buffer with G.711 Appendix I concealment, per-consumer encode/resample, packet-replay harness. No sockets, no clocks, no async. |
crates/opus-ffi |
Safe wrapper over libopus (vendored, statically linked). The only crate here that contains FFI unsafe; every other crate forbids it. |
crates/rtpengine-ng |
Sans-IO rtpengine NG protocol client: bencode, subscribe request/answer, unsubscribe, play media/stop media, query, subscription SDP. |
crates/protocol |
Frozen consumer wire dialects: Twilio Media Streams JSON and audio_fork send_text control events. The serialization tests are the spec. |
crates/session-core |
Sans-IO control-plane state machine: sessions, attachments, playbacks, events — capability authorization, one authoritative attachment per session, idempotent retries. No sockets and no async; unlike the protocol/DSP cores it does read the wall clock, stamping each session's opened_at with SystemTime::now(). |
crates/control-api |
The network surface: MediaControl + MediaStream + TelCompat gRPC services over session-core. Pure-Rust protobuf build (no protoc). |
crates/sip-uas |
Sans-IO SIP UAS: transactions, dialogs, session timers, and one client transaction — the in-dialog BYE. It refuses to open an INVITE client transaction by name, so "MSS answers, it never dials" is a property of the type. No sockets. |
crates/call-events |
Stopgap Redis-stream publisher for the SIP front door's call-control events — invited, answered, end_of_interaction, ended on mss:call-events by default, paired with the AnswerSession/HangupSession RPCs (deploy.md). Idle if the stream env is emptied or Redis is unset. Delete when Kafka mss.events is the bus. |
crates/mediaserverd |
The daemon: Tokio control plane + dedicated real-time media threads (the two-world architecture), fan-out hub, consumer bridges, Kafka event pump, Redis registry keeper, metrics. |
- Two worlds. Tokio owns everything latency-tolerant. Dedicated OS threads own the packet path (recv → jitter → decode → fan-out → paced send). Bounded lock-free queues between them; the media world never blocks on the control world.
- Sans-IO cores. Protocol/DSP logic takes packets and instants as
parameters and returns values, so every media bug is reproducible by
packet replay in a unit test. (The control-plane
session-coreis sans-IO but not clockless: it stampsopened_atfrom the wall clock.) - No per-packet allocation on the hot path after session setup.
- No panics on network input. Malformed RTP/bencode/JSON is an error value.
- Supervised sessions. Every tap leg has an audio-flow watchdog; a dead task must never be a silently dead call.
- Zero
unsafeoutside dedicated FFI wrapper crates; codec/DSP math is adopted from proven libraries, never reimplemented.
The binding version of these rules is CONSTITUTION.md;
coding style (including the strictly-no-comments policy — intent lives in
names, types, tests and docs/, enforced by CI) is
docs/rust-guidelines.md.
- docs/architecture.md — the accepted design and its decision records.
- docs/roadmap.md — the phased plan with exit criteria and live status; docs/tasks.md is the ordered next-up list with a definition of done per item.
- docs/implementation-notes.md — per-module status and known limitations (sources carry no comments; this is where that context lives).
- docs/testing.md and docs/lab.md — the three test altitudes and the lab that exists today.
- CLAUDE.md — orientation for AI-assisted sessions and new engineers.
cargo test --workspace
cargo clippy --all-targets -- -D warnings
cargo fmt --all --check
grep -rn '//' crates --include='*.rs' # must output nothing
cargo deny check allToolchain is pinned in rust-toolchain.toml. CI runs exactly the gate above.
Dual-licensed under MIT or Apache-2.0, at your option. Contributions are accepted under the same terms.