diff --git a/implementations/main.cpp b/implementations/main.cpp index 0ee64c4..38c6623 100644 --- a/implementations/main.cpp +++ b/implementations/main.cpp @@ -1257,6 +1257,7 @@ private: setenv("AMR_MODE", EnvOr("AMR_MODE", "2").c_str(), 1); setenv("DTX", EnvOr("DTX", "0").c_str(), 1); setenv("OCTET_ALIGN", leg.octetAlign ? "1" : "0", 1); + setenv("CODEC", leg.codec.c_str(), 1); std::vector c; for (auto& s : argv) c.push_back(const_cast(s.c_str())); c.push_back(nullptr); diff --git a/implementations/media.cpp b/implementations/media.cpp index 14bf2e3..1dffd0a 100644 --- a/implementations/media.cpp +++ b/implementations/media.cpp @@ -3,28 +3,42 @@ // lint-disable-file fixed-width-types no-char-pointer /* -imsd-media — the RTP/AMR-WB media leg for a userspace VoLTE call, spawned as -its own process by the daemon (imsd) exactly as the Python prototype spawned +imsd-media — the RTP media leg for a userspace VoLTE call, spawned as its +own process by the daemon (imsd) exactly as the Python prototype spawned rtpcap.py: one media leg per call, argv-configured, torn down on SIGTERM or when the downlink dries up. Keeping it a separate process preserves the far-end-hangup contract (exit code 3 = downlink RTP stopped, which on carriers whose network BYE never reaches our SAs is the reliable teardown trigger) and isolates a media crash from the control-plane daemon. -Binds the advertised local RTP port, sends uplink octet-aligned AMR-WB frames -toward the media gateway, captures the downlink, and — with PLAY=1 — -reconstructs the media clock from RTP timestamps so pw-play stays real-time- -paced through far-end DTX silence (every missing 20 ms slot is decoded as an -FT-15 NO_DATA frame for CNG/PLC). MIC=1 feeds live pw-record audio through -libvo-amrwbenc; downlink decode is libopencore-amrwb. Both codecs are dlopen'd -so the binary has no link-time dependency on them. +Codecs (CODEC env, set by the daemon from the negotiated SDP): AMR-WB +(16 kHz; the mobile-to-mobile VoLTE codec, the default), AMR narrowband and +G.711 PCMA/PCMU (8 kHz; what a PSTN gateway offers when a landline calls). +AMR frames ride RFC 4867 payloads, octet-aligned or bandwidth-efficient +(OCTET_ALIGN, mirrored from the SDP); G.711 is raw samples. The AMR codecs +are dlopen'd (libvo-amrwbenc + libopencore-amrwb for WB, libopencore-amrnb +for NB) so the binary has no link-time dependency on them; G.711 is a table. + +Binds the advertised local RTP port, sends uplink frames toward the media +gateway, captures the downlink, and — with PLAY=1 — reconstructs the media +clock from RTP timestamps so pw-play stays real-time-paced through far-end +DTX silence (every missing 20 ms slot is decoded as a NO_DATA frame for +CNG/PLC; zeros for G.711). MIC=1 feeds live pw-record audio through the +encoder. Usage: imsd-media - imsd-media --selftest (payload pack/depay roundtrip, both formats) -Env: MIC PLAY GAIN PLAY_GAIN AMR_MODE DTX MEDIA_TIMEOUT RTP_DUMP OCTET_ALIGN -AUDIO_USER + imsd-media --selftest (payload pack/depay roundtrip, both formats, + both AMR codecs; G.711 table roundtrip) +Env: CODEC MIC PLAY GAIN PLAY_GAIN AMR_MODE DTX MEDIA_TIMEOUT RTP_DUMP +OCTET_ALIGN AUDIO_USER MIC_SRC PCM_DUMP (OCTET_ALIGN=0 selects RFC 4867 bandwidth-efficient payloads both ways; the -daemon sets it from the negotiated SDP — 1 is the default/legacy behavior.) +daemon sets it from the negotiated SDP — 1 is the default/legacy behavior. +AMR_MODE is the encoder mode of whichever AMR codec is active; default 2 +for AMR-WB (12.65 kbit/s), 7 for AMR (12.2 kbit/s). MIC_SRC= feeds +raw s16 PCM at the codec rate through the encoder instead of pw-record, +paced in real time; PCM_DUMP=1 writes the decoded downlink to .pcm +— both are test seams for driving the leg against a synthetic RTP peer +with no PipeWire and no network.) */ #include @@ -46,28 +60,125 @@ double Now() { return std::chrono::duration(Clock::now().time_since_epoch()).count(); } -// octet-aligned AMR-WB speech-frame byte sizes by frame type (mode). -constexpr int AmrwbBytes(int ft) { - switch (ft) { - case 0: return 17; case 1: return 23; case 2: return 32; case 3: return 36; - case 4: return 40; case 5: return 46; case 6: return 50; case 7: return 58; - case 8: return 60; case 9: return 5; default: return 0; - } -} +enum class Codec { AmrWb, AmrNb, Pcma, Pcmu }; -// AMR-WB speech bits per frame type (RFC 4867 table 2 / TS 26.201) — the -// exact payload bit counts of the bandwidth-efficient format. FT 14/15 -// (SPEECH_LOST/NO_DATA) carry 0 bits; unknown FTs return -1 so a corrupt -// ToC aborts the packet instead of shifting every later bit. -constexpr int AmrwbBits(int ft) { +Codec ParseCodec(std::string_view name) { + if (name == "AMR") return Codec::AmrNb; + if (name == "PCMA") return Codec::Pcma; + if (name == "PCMU") return Codec::Pcmu; + return Codec::AmrWb; +} +std::string_view CodecName(Codec c) { + switch (c) { + case Codec::AmrWb: return "AMR-WB"; + case Codec::AmrNb: return "AMR"; + case Codec::Pcma: return "PCMA"; + case Codec::Pcmu: return "PCMU"; + } + return "AMR-WB"; +} +bool IsAmr(Codec c) { return c == Codec::AmrWb || c == Codec::AmrNb; } +// Sample rate = RTP clock rate; one 20 ms frame is 320 or 160 samples. +int CodecRate(Codec c) { return c == Codec::AmrWb ? 16000 : 8000; } +int FrameSamples(Codec c) { return CodecRate(c) / 50; } + +// AMR speech bits per frame type — the exact payload bit counts of the +// bandwidth-efficient format. AMR-WB: RFC 4867 table 2 / TS 26.201; AMR: +// RFC 4867 table 1 / TS 26.101 (FT 8 = SID, 9-11 = the other systems' SID +// frames a gateway may forward). FT 14/15 (SPEECH_LOST/NO_DATA) carry 0 +// bits; unknown FTs return -1 so a corrupt ToC aborts the packet instead of +// shifting every later bit. +constexpr int AmrBits(Codec c, int ft) { + if (c == Codec::AmrWb) { + switch (ft) { + case 0: return 132; case 1: return 177; case 2: return 253; case 3: return 285; + case 4: return 317; case 5: return 365; case 6: return 397; case 7: return 461; + case 8: return 477; case 9: return 40; case 14: return 0; case 15: return 0; + default: return -1; + } + } switch (ft) { - case 0: return 132; case 1: return 177; case 2: return 253; case 3: return 285; - case 4: return 317; case 5: return 365; case 6: return 397; case 7: return 461; - case 8: return 477; case 9: return 40; case 14: return 0; case 15: return 0; + case 0: return 95; case 1: return 103; case 2: return 118; case 3: return 134; + case 4: return 148; case 5: return 159; case 6: return 204; case 7: return 244; + case 8: return 39; case 9: return 43; case 10: return 38; case 11: return 37; + case 15: return 0; default: return -1; } } +// octet-aligned speech-frame byte size by frame type: the bits rounded up +// (AMR-WB: 17 23 32 36 40 46 50 58 60, SID 5; AMR: 12 13 15 17 19 20 26 31, +// SID 5). 0 for frame types that carry no speech. +constexpr int AmrBytes(Codec c, int ft) { + int bits = AmrBits(c, ft); + return bits <= 0 ? 0 : (bits + 7) / 8; +} + +// ---- G.711 (ITU-T, the classic Sun g711.c formulation) -------------------- +// 16-bit linear <-> 8-bit companded. Byte 0xD5 (A-law) / 0xFF (mu-law) is +// digital zero. Idempotent: encoding a decoded sample gives the same byte. + +constexpr int UlawBias = 0x84; +constexpr int UlawClip = 8159; + +int SegmentOf(int val, std::span ends) { + for (std::size_t i = 0; i < ends.size(); i++) + if (val <= ends[i]) return static_cast(i); + return static_cast(ends.size()); +} + +std::uint8_t LinearToAlaw(std::int16_t pcm) { + static constexpr std::array ends = {0x1F, 0x3F, 0x7F, 0xFF, 0x1FF, 0x3FF, 0x7FF, 0xFFF}; + int val = pcm >> 3; // 13-bit magnitude space + int mask = 0xD5; + if (val < 0) { + mask = 0x55; + val = -val - 1; + } + int seg = SegmentOf(val, ends); + if (seg >= 8) return static_cast(0x7F ^ mask); + int aval = seg << 4; + aval |= seg < 2 ? (val >> 1) & 0x0F : (val >> seg) & 0x0F; + return static_cast(aval ^ mask); +} + +std::int16_t AlawToLinear(std::uint8_t a) { + a ^= 0x55; + int t = (a & 0x0F) << 4; + int seg = (a & 0x70) >> 4; + if (seg == 0) t += 8; + else if (seg == 1) t += 0x108; + else t = (t + 0x108) << (seg - 1); + return static_cast((a & 0x80) ? t : -t); +} + +std::uint8_t LinearToUlaw(std::int16_t pcm) { + static constexpr std::array ends = {0x3F, 0x7F, 0xFF, 0x1FF, 0x3FF, 0x7FF, 0xFFF, 0x1FFF}; + int val = pcm >> 2; // 14-bit magnitude space + int mask = 0xFF; + if (val < 0) { + val = -val; + mask = 0x7F; + } + if (val > UlawClip) val = UlawClip; + val += UlawBias >> 2; + int seg = SegmentOf(val, ends); + if (seg >= 8) return static_cast(0x7F ^ mask); + int uval = (seg << 4) | ((val >> (seg + 1)) & 0x0F); + return static_cast(uval ^ mask); +} + +std::int16_t UlawToLinear(std::uint8_t u) { + u = static_cast(~u); + int t = ((u & 0x0F) << 3) + UlawBias; + t <<= (u & 0x70) >> 4; + return static_cast((u & 0x80) ? (UlawBias - t) : (t - UlawBias)); +} + +std::uint8_t G711Encode(Codec c, std::int16_t pcm) { return c == Codec::Pcma ? LinearToAlaw(pcm) : LinearToUlaw(pcm); } +std::int16_t G711Decode(Codec c, std::uint8_t b) { return c == Codec::Pcma ? AlawToLinear(b) : UlawToLinear(b); } +std::uint8_t G711Zero(Codec c) { return c == Codec::Pcma ? 0xD5 : 0xFF; } + std::string EnvOr(const char* k, const char* d) { const char* v = std::getenv(k); return v ? std::string(v) : std::string(d); @@ -97,62 +208,99 @@ socklen_t MakeAddr(std::string_view ip, int port, sockaddr_storage& ss) { return sizeof(sockaddr_in); } -// ---- AMR-WB codecs via dlopen (vo-amrwbenc E_IF_*, opencore-amrwb D_IF_*) -- +// ---- AMR codecs via dlopen ------------------------------------------------ +// AMR-WB: vo-amrwbenc E_IF_* (encode), opencore-amrwb D_IF_* (decode). +// AMR: opencore-amrnb Encoder_Interface_* / Decoder_Interface_*. +// Both families speak the RFC 4867 §5.3 storage format: [header byte][speech], +// header = (FT << 3) | (Q << 2) — exactly the octet-aligned ToC byte with F=0. class Encoder { public: - bool Open() { - lib_ = dlopen("libvo-amrwbenc.so.0", RTLD_NOW); + bool Open(Codec c, int dtx) { + codec_ = c; + if (c == Codec::AmrWb) { + lib_ = dlopen("libvo-amrwbenc.so.0", RTLD_NOW); + if (!lib_) return false; + wbInit_ = reinterpret_cast(dlsym(lib_, "E_IF_init")); + wbEnc_ = reinterpret_cast(dlsym(lib_, "E_IF_encode")); + exit_ = reinterpret_cast(dlsym(lib_, "E_IF_exit")); + if (!wbInit_ || !wbEnc_ || !exit_) return false; + st_ = wbInit_(); + return st_ != nullptr; + } + lib_ = dlopen("libopencore-amrnb.so.0", RTLD_NOW); if (!lib_) return false; - init_ = reinterpret_cast(dlsym(lib_, "E_IF_init")); - enc_ = reinterpret_cast(dlsym(lib_, "E_IF_encode")); - exit_ = reinterpret_cast(dlsym(lib_, "E_IF_exit")); - if (!init_ || !enc_ || !exit_) return false; - st_ = init_(); + nbInit_ = reinterpret_cast(dlsym(lib_, "Encoder_Interface_init")); + nbEnc_ = reinterpret_cast(dlsym(lib_, "Encoder_Interface_Encode")); + exit_ = reinterpret_cast(dlsym(lib_, "Encoder_Interface_exit")); + if (!nbInit_ || !nbEnc_ || !exit_) return false; + st_ = nbInit_(dtx); // NB: DTX is an init-time choice return st_ != nullptr; } - // 320 int16 samples -> one RFC 3267 storage frame (header byte + speech). + // one 20 ms frame of s16 samples -> one storage frame (header byte + speech). std::vector Encode(std::int16_t* samples, int mode, int dtx) { std::uint8_t out[128]; - int n = enc_(st_, static_cast(mode), samples, out, static_cast(dtx)); + int n = codec_ == Codec::AmrWb + ? wbEnc_(st_, static_cast(mode), samples, out, static_cast(dtx)) + : nbEnc_(st_, mode, samples, out, 0); if (n <= 0) return {}; return std::vector(out, out + n); } - ~Encoder() { if (st_ && exit_) exit_(st_); if (lib_) dlclose(lib_); } + ~Encoder() { + if (st_ && exit_) exit_(st_); + if (lib_) dlclose(lib_); + } private: - using InitFn = void* (*)(); - using EncFn = int (*)(void*, std::int16_t, std::int16_t*, std::uint8_t*, std::int16_t); + using WbInitFn = void* (*)(); + using WbEncFn = int (*)(void*, std::int16_t, std::int16_t*, std::uint8_t*, std::int16_t); + using NbInitFn = void* (*)(int); + using NbEncFn = int (*)(void*, int, const std::int16_t*, std::uint8_t*, int); using ExitFn = void (*)(void*); - void* lib_ = nullptr; void* st_ = nullptr; - InitFn init_ = nullptr; EncFn enc_ = nullptr; ExitFn exit_ = nullptr; + Codec codec_ = Codec::AmrWb; + void* lib_ = nullptr; + void* st_ = nullptr; + WbInitFn wbInit_ = nullptr; + WbEncFn wbEnc_ = nullptr; + NbInitFn nbInit_ = nullptr; + NbEncFn nbEnc_ = nullptr; + ExitFn exit_ = nullptr; }; class Decoder { public: - bool Open() { - lib_ = dlopen("libopencore-amrwb.so.0", RTLD_NOW); + bool Open(Codec c) { + codec_ = c; + bool wb = c == Codec::AmrWb; + lib_ = dlopen(wb ? "libopencore-amrwb.so.0" : "libopencore-amrnb.so.0", RTLD_NOW); if (!lib_) return false; - init_ = reinterpret_cast(dlsym(lib_, "D_IF_init")); - dec_ = reinterpret_cast(dlsym(lib_, "D_IF_decode")); - exit_ = reinterpret_cast(dlsym(lib_, "D_IF_exit")); + init_ = reinterpret_cast(dlsym(lib_, wb ? "D_IF_init" : "Decoder_Interface_init")); + dec_ = reinterpret_cast(dlsym(lib_, wb ? "D_IF_decode" : "Decoder_Interface_Decode")); + exit_ = reinterpret_cast(dlsym(lib_, wb ? "D_IF_exit" : "Decoder_Interface_exit")); if (!init_ || !dec_ || !exit_) return false; st_ = init_(); return st_ != nullptr; } - // one storage frame ([header][speech]) -> 640 bytes of 16 kHz s16 PCM. - std::array Decode(const std::uint8_t* frame, int len) { - std::array out{}; + // one storage frame ([header][speech]) -> one 20 ms frame of s16 PCM. + std::vector Decode(const std::uint8_t* frame, int len) { + std::vector out(static_cast(FrameSamples(codec_)), 0); std::vector in(frame, frame + len); dec_(st_, in.data(), out.data(), 0); return out; } - ~Decoder() { if (st_ && exit_) exit_(st_); if (lib_) dlclose(lib_); } + ~Decoder() { + if (st_ && exit_) exit_(st_); + if (lib_) dlclose(lib_); + } private: using InitFn = void* (*)(); - using DecFn = void (*)(void*, std::uint8_t*, std::int16_t*, int); + using DecFn = void (*)(void*, const std::uint8_t*, std::int16_t*, int); using ExitFn = void (*)(void*); - void* lib_ = nullptr; void* st_ = nullptr; - InitFn init_ = nullptr; DecFn dec_ = nullptr; ExitFn exit_ = nullptr; + Codec codec_ = Codec::AmrWb; + void* lib_ = nullptr; + void* st_ = nullptr; + InitFn init_ = nullptr; + DecFn dec_ = nullptr; + ExitFn exit_ = nullptr; }; // MSB-first bit cursor over an RTP payload (bandwidth-efficient AMR-WB is @@ -170,10 +318,10 @@ struct BitReader { } }; -// bandwidth-efficient AMR-WB de-payload. Same contract as DepayOctet: +// bandwidth-efficient AMR de-payload. Same contract as DepayOctet: // [(storage-header-byte, speech-bytes)...], speech re-aligned to octets. std::vector>> -DepayBe(std::span pl) { +DepayBe(Codec c, std::span pl) { std::vector>> out; BitReader br{pl}; if (!br.Ok(4)) return out; @@ -189,9 +337,9 @@ DepayBe(std::span pl) { if (!f) break; } for (auto [ft, q] : tocs) { - int bits = AmrwbBits(ft); + int bits = AmrBits(c, ft); if (bits < 0 || !br.Ok(static_cast(bits))) break; - std::vector speech(AmrwbBytes(ft), 0); + std::vector speech(AmrBytes(c, ft), 0); for (int i = 0; i < bits; i++) if (br.Take(1)) speech[i >> 3] |= 0x80 >> (i & 7); out.emplace_back(static_cast((ft << 3) | (q ? 0x04 : 0)), std::move(speech)); @@ -202,8 +350,8 @@ DepayBe(std::span pl) { // storage-format frame (header byte + octet-aligned speech) -> RTP payload. // Octet-aligned: CMR byte + the frame verbatim (the storage header doubles // as a ToC byte with F=0). Bandwidth-efficient: 10 header bits + exactly -// AmrwbBits(ft) speech bits, final octet zero-padded. -std::vector PayloadFromFrame(std::span frame, bool octetAlign) { +// AmrBits(ft) speech bits, final octet zero-padded. +std::vector PayloadFromFrame(Codec c, std::span frame, bool octetAlign) { if (octetAlign) { std::vector pl = {0xF0}; pl.insert(pl.end(), frame.begin(), frame.end()); @@ -211,7 +359,7 @@ std::vector PayloadFromFrame(std::span frame, } int ft = (frame[0] >> 3) & 0x0F; int q = (frame[0] >> 2) & 1; - int bits = AmrwbBits(ft); + int bits = AmrBits(c, ft); if (bits < 0) bits = 0; std::vector pl((10 + bits + 7) / 8, 0); auto put = [&](int pos, int n, std::uint32_t v) { @@ -227,10 +375,10 @@ std::vector PayloadFromFrame(std::span frame, return pl; } -// octet-aligned AMR-WB de-payload: skip CMR, read ToC bytes until F=0, then +// octet-aligned AMR de-payload: skip CMR, read ToC bytes until F=0, then // the speech runs. Returns [(storage-header-byte, speech-bytes)...]. std::vector>> -DepayOctet(std::span pl) { +DepayOctet(Codec c, std::span pl) { std::vector>> out; if (pl.empty()) return out; std::size_t i = 1; // skip CMR @@ -242,7 +390,7 @@ DepayOctet(std::span pl) { } for (std::uint8_t toc : tocs) { int ft = (toc >> 3) & 0x0F; - int n = AmrwbBytes(ft); + int n = AmrBytes(c, ft); std::vector speech; if (n > 0 && i + n <= pl.size()) speech.assign(pl.begin() + i, pl.begin() + i + n); out.emplace_back(static_cast(toc & 0x7C), std::move(speech)); @@ -267,7 +415,7 @@ std::size_t RtpPayloadOffset(std::span pkt) { // is postmarketOS's standard account. `toChild` true = we write the child's // stdin (pw-play); false = we read its stdout (pw-record). Returns {pid, fd}. struct Child { pid_t pid = -1; int fd = -1; }; -Child SpawnPw(bool play, bool toChild) { +Child SpawnPw(bool play, bool toChild, int rate) { int pipefd[2]; if (pipe(pipefd) != 0) return {}; std::vector argv; @@ -280,7 +428,8 @@ Child SpawnPw(bool play, bool toChild) { } const char* tool = play ? "pw-play" : "pw-record"; const char* lat = play ? "40ms" : "20ms"; - for (const char* a : {tool, "--raw", "--rate", "16000", "--channels", "1", "--format", "s16", "--latency", lat, "-"}) + std::string rateStr = std::to_string(rate); + for (const char* a : {tool, "--raw", "--rate", rateStr.c_str(), "--channels", "1", "--format", "s16", "--latency", lat, "-"}) argv.emplace_back(a); pid_t pid = fork(); if (pid == 0) { @@ -315,6 +464,7 @@ bool ReadExact(int fd, std::uint8_t* buf, std::size_t n) { struct TxState { std::mutex lock; int pt = 0; + std::uint32_t tsStep = 320; // samples per 20 ms frame at the codec's clock std::uint32_t ssrc = 0x5EED1234; std::uint32_t seq = 1000; std::uint32_t ts = 160000; @@ -338,7 +488,7 @@ void RtpSend(int sock, TxState& st, std::span payload) { hdr[8] = (st.ssrc >> 24) & 0xFF; hdr[9] = (st.ssrc >> 16) & 0xFF; hdr[10] = (st.ssrc >> 8) & 0xFF; hdr[11] = st.ssrc & 0xFF; st.seq = (st.seq + 1) & 0xFFFF; - st.ts = (st.ts + 320) & 0xFFFFFFFF; + st.ts = (st.ts + st.tsStep) & 0xFFFFFFFF; a1 = st.dst; l1 = st.dstLen; a2 = st.latched; l2 = st.latchedLen; } std::vector pkt(hdr, hdr + 12); @@ -355,15 +505,34 @@ void RtpSend(int sock, TxState& st, std::span payload) { } std::vector>> -Depay(std::span pl, bool octetAlign) { - return octetAlign ? DepayOctet(pl) : DepayBe(pl); +Depay(Codec c, std::span pl, bool octetAlign) { + return octetAlign ? DepayOctet(c, pl) : DepayBe(c, pl); } -std::vector SilenceFrame(int ft, bool octetAlign) { - // storage frame: ToC(F=0,FT,Q=1) + zeroed speech, payloaded per mode. +// What goes out every 20 ms while the mic is not feeding frames: for AMR a +// storage frame ToC(F=0,FT,Q=1) + zeroed speech, payloaded per mode; for +// G.711 one frame of digital zero. +std::vector SilenceFrame(Codec c, int ft, bool octetAlign) { + if (!IsAmr(c)) return std::vector(static_cast(FrameSamples(c)), G711Zero(c)); std::vector f = {static_cast((ft << 3) | 0x04)}; - f.resize(1 + AmrwbBytes(ft), 0); - return PayloadFromFrame(f, octetAlign); + f.resize(1 + static_cast(AmrBytes(c, ft)), 0); + return PayloadFromFrame(c, f, octetAlign); +} + +// G.711 uplink: one 20 ms frame of samples -> one payload of companded bytes. +std::vector G711Payload(Codec c, std::span pcm) { + std::vector pl(pcm.size()); + for (std::size_t i = 0; i < pcm.size(); i++) + pl[i] = G711Encode(c, pcm[i]); + return pl; +} + +// G.711 downlink: one frame's worth of companded bytes -> PCM. +std::vector G711DecodeFrame(Codec c, std::span bytes) { + std::vector pcm(bytes.size()); + for (std::size_t i = 0; i < bytes.size(); i++) + pcm[i] = G711Decode(c, bytes[i]); + return pcm; } std::atomic Quit{false}; @@ -372,31 +541,64 @@ void OnSig(int) { Quit.store(true); } // playout queue item struct PktItem { std::uint32_t ts; std::vector payload; }; -} // namespace - -int main(int argc, char** argv) { - if (argc == 2 && std::string_view(argv[1]) == "--selftest") { - // pack->depay roundtrip of every frame type, both payload formats. - // BE carries exactly AmrwbBits(ft) bits, so the pattern's padding - // bits in the last speech byte must be zero for equality to hold. - for (int ft : {0, 1, 2, 3, 4, 5, 6, 7, 8, 9}) { - std::vector frame = { - static_cast((ft << 3) | 0x04)}; - int bits = AmrwbBits(ft); - for (int i = 0; i < AmrwbBytes(ft); i++) +// --selftest: pack->depay roundtrip of every frame type, both payload +// formats, both AMR codecs; G.711 digital zero, idempotence over the whole +// 16-bit range, and a 1 kHz sine surviving with the codec's nominal SNR. +int SelfTest() { + for (Codec c : {Codec::AmrWb, Codec::AmrNb}) { + for (int ft = 0; ft < 16; ft++) { + int bits = AmrBits(c, ft); + if (bits <= 0) continue; + // BE carries exactly AmrBits(ft) bits, so the pattern's padding + // bits in the last speech byte must be zero for equality to hold. + std::vector frame = {static_cast((ft << 3) | 0x04)}; + for (int i = 0; i < AmrBytes(c, ft); i++) frame.push_back(static_cast(0xA5 + i * 31)); if (bits % 8) frame.back() &= static_cast(0xFF << (8 - bits % 8)); for (bool oa : {true, false}) { - auto got = Depay(PayloadFromFrame(frame, oa), oa); + auto got = Depay(c, PayloadFromFrame(c, frame, oa), oa); if (got.size() != 1 || got[0].first != frame[0] || !std::equal(got[0].second.begin(), got[0].second.end(), frame.begin() + 1, frame.end())) { - std::println(std::cerr, "selftest FAIL ft={} oa={}", ft, oa); + std::println(std::cerr, "selftest FAIL {} ft={} oa={}", CodecName(c), ft, oa); return 1; } } } - std::println("selftest OK"); - return 0; } + if (LinearToAlaw(0) != 0xD5 || LinearToUlaw(0) != 0xFF) { + std::println(std::cerr, "selftest FAIL G.711 digital zero"); + return 1; + } + for (Codec c : {Codec::Pcma, Codec::Pcmu}) { + for (int v = -32768; v <= 32767; v++) { + std::uint8_t b = G711Encode(c, static_cast(v)); + if (G711Encode(c, G711Decode(c, b)) != b) { + std::println(std::cerr, "selftest FAIL {} not idempotent at {}", CodecName(c), v); + return 1; + } + } + double sig = 0; + double err = 0; + for (int i = 0; i < 8000; i++) { + double x = 10000.0 * std::sin(2 * std::numbers::pi * 1000.0 * i / 8000.0); + auto sample = static_cast(std::lround(x)); + std::int16_t d = G711Decode(c, G711Encode(c, sample)); + sig += x * x; + err += (x - d) * (x - d); + } + double snr = 10 * std::log10(sig / err); + if (snr < 30) { + std::println(std::cerr, "selftest FAIL {} sine SNR {:.1f} dB", CodecName(c), snr); + return 1; + } + } + std::println("selftest OK"); + return 0; +} + +} // namespace + +int main(int argc, char** argv) { + if (argc == 2 && std::string_view(argv[1]) == "--selftest") return SelfTest(); if (argc < 8) { std::println(std::cerr, "usage: imsd-media local rtp_port remote_ip remote_port pt secs out"); return 2; @@ -409,25 +611,35 @@ int main(int argc, char** argv) { double secs = std::atof(argv[6]); std::string out = argv[7]; + Codec codec = ParseCodec(EnvOr("CODEC", "AMR-WB")); + bool amr = IsAmr(codec); + int rate = CodecRate(codec); + auto frameSamples = static_cast(FrameSamples(codec)); bool mic = EnvBool("MIC", false); bool play = EnvBool("PLAY", false); - int amrMode = std::atoi(EnvOr("AMR_MODE", "2").c_str()); + int amrMode = std::atoi(EnvOr("AMR_MODE", codec == Codec::AmrWb ? "2" : "7").c_str()); double gain = std::atof(EnvOr("GAIN", "1.0").c_str()); double playGain = std::atof(EnvOr("PLAY_GAIN", "1.0").c_str()); int dtx = EnvBool("DTX", false) ? 1 : 0; bool octetAlign = EnvBool("OCTET_ALIGN", true); double mediaTimeout = std::atof(EnvOr("MEDIA_TIMEOUT", "6.0").c_str()); bool rtpDump = EnvBool("RTP_DUMP", false); + std::string micSrc = EnvOr("MIC_SRC", ""); + bool pcmDump = EnvBool("PCM_DUMP", false); constexpr int ExitMediaTimeout = 3; constexpr int PrimeFrames = 8; constexpr int MaxFill = 25; constexpr std::size_t PlayqMax = 256; int sock = socket(Is6(local) ? AF_INET6 : AF_INET, SOCK_DGRAM, 0); - if (sock < 0) { std::println(std::cerr, "imsd-media: socket failed"); return 1; } + if (sock < 0) { + std::println(std::cerr, "imsd-media: socket failed"); + return 1; + } int one = 1; setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, &one, sizeof one); - sockaddr_storage bindA; socklen_t bindL = MakeAddr(local, rtpPort, bindA); + sockaddr_storage bindA; + socklen_t bindL = MakeAddr(local, rtpPort, bindA); if (bind(sock, reinterpret_cast(&bindA), bindL) != 0) { std::println(std::cerr, "imsd-media: bind [{}]:{} failed", local, rtpPort); return 1; @@ -437,63 +649,98 @@ int main(int argc, char** argv) { TxState st; st.pt = pt; + st.tsStep = static_cast(frameSamples); st.dstLen = MakeAddr(rIp, rPort, st.dst); - st.latched = st.dst; st.latchedLen = st.dstLen; + st.latched = st.dst; + st.latchedLen = st.dstLen; signal(SIGTERM, OnSig); signal(SIGINT, OnSig); signal(SIGPIPE, SIG_IGN); - // ---- mic uplink thread + // ---- mic uplink thread: pw-record (or the MIC_SRC file, real-time paced) + // through the codec's encoder, one RTP packet per 20 ms frame std::atomic stop{false}; std::jthread micThread; if (mic) { micThread = std::jthread([&] { Encoder enc; - if (!enc.Open()) { - std::println(std::cerr, "imsd-media: mic: encoder unavailable; silence fallback"); + if (amr && !enc.Open(codec, dtx)) { + std::println(std::cerr, "imsd-media: mic: {} encoder unavailable; silence fallback", CodecName(codec)); return; } - Child rec = SpawnPw(false, /*toChild=*/false); - if (rec.pid < 0) return; - std::uint8_t raw[640]; + int fd = -1; + pid_t pid = -1; + if (!micSrc.empty()) { + fd = open(micSrc.c_str(), O_RDONLY); + if (fd < 0) { + std::println(std::cerr, "imsd-media: mic: cannot open MIC_SRC {}", micSrc); + return; + } + } else { + Child rec = SpawnPw(false, /*toChild=*/false, rate); + if (rec.pid < 0) return; + fd = rec.fd; + pid = rec.pid; + } + std::vector raw(frameSamples * 2); + std::vector samples(frameSamples); + auto next = Clock::now(); while (!stop.load()) { - if (!ReadExact(rec.fd, raw, 640)) { - std::println(std::cerr, "imsd-media: mic: pw-record EOF; silence fallback"); + if (!ReadExact(fd, raw.data(), raw.size())) { + std::println(std::cerr, "imsd-media: mic: {} EOF; silence fallback", pid > 0 ? "pw-record" : "MIC_SRC"); break; } - std::int16_t samples[320]; - std::memcpy(samples, raw, 640); + if (pid < 0) { + // a file delivers instantly; pace it like a microphone + next += std::chrono::milliseconds(20); + std::this_thread::sleep_until(next); + } + std::memcpy(samples.data(), raw.data(), raw.size()); if (gain != 1.0) - for (int i = 0; i < 320; i++) { - int v = static_cast(samples[i] * gain); - samples[i] = static_cast(v < -32768 ? -32768 : (v > 32767 ? 32767 : v)); + for (auto& sample : samples) { + int v = static_cast(sample * gain); + sample = static_cast(v < -32768 ? -32768 : (v > 32767 ? 32767 : v)); } - std::vector frame = enc.Encode(samples, amrMode, dtx); - if (frame.empty()) continue; + std::vector payload; + if (amr) { + std::vector frame = enc.Encode(samples.data(), amrMode, dtx); + if (frame.empty()) continue; + payload = PayloadFromFrame(codec, frame, octetAlign); + } else { + payload = G711Payload(codec, samples); + } st.micOn.store(true); { std::scoped_lock g(st.lock); st.txMic++; } - RtpSend(sock, st, PayloadFromFrame(frame, octetAlign)); + RtpSend(sock, st, payload); } st.micOn.store(false); - close(rec.fd); - kill(rec.pid, SIGKILL); - waitpid(rec.pid, nullptr, 0); + close(fd); + if (pid > 0) { + kill(pid, SIGKILL); + waitpid(pid, nullptr, 0); + } }); } - // ---- downlink playout (decode + RTP-timestamp clock reconstruction) + // ---- downlink playout (decode + RTP-timestamp clock reconstruction). + // Runs for pw-play (PLAY=1) and/or the PCM dump (PCM_DUMP=1). Decoder dec; Child playCh; bool playOn = false; - if (play) { - if (dec.Open()) { - playCh = SpawnPw(true, /*toChild=*/true); - playOn = playCh.pid >= 0; + bool decodeOn = false; + if (play || pcmDump) { + if (!amr || dec.Open(codec)) { + if (play) { + playCh = SpawnPw(true, /*toChild=*/true, rate); + playOn = playCh.pid >= 0; + } + decodeOn = playOn || pcmDump; } else { - std::println(std::cerr, "imsd-media: play: decoder unavailable; capture only"); + std::println(std::cerr, "imsd-media: play: {} decoder unavailable; capture only", CodecName(codec)); } } + std::FILE* pcmFile = decodeOn && pcmDump ? std::fopen((std::format("{}.pcm", out)).c_str(), "wb") : nullptr; std::mutex qlock; std::condition_variable qcv; std::deque playq; @@ -502,9 +749,11 @@ int main(int argc, char** argv) { std::uint64_t late = 0; std::uint64_t qdrop = 0; std::jthread playThread; - if (playOn) { + if (decodeOn) { playThread = std::jthread([&] { auto writePcm = [&](std::span pcm) { + if (pcmFile) std::fwrite(pcm.data(), 2, pcm.size(), pcmFile); + if (!playOn) return true; std::size_t bytes = pcm.size() * 2; const char* p = reinterpret_cast(pcm.data()); std::size_t off = 0; @@ -515,14 +764,20 @@ int main(int argc, char** argv) { } return true; }; - auto applyGain = [&](std::array& pcm) { + auto applyGain = [&](std::vector& pcm) { if (playGain == 1.0) return; - for (auto& s : pcm) { - int v = static_cast(s * playGain); - s = static_cast(v < -32768 ? -32768 : (v > 32767 ? 32767 : v)); + for (auto& sample : pcm) { + int v = static_cast(sample * playGain); + sample = static_cast(v < -32768 ? -32768 : (v > 32767 ? 32767 : v)); } }; - std::array zero{}; + // Decoder state is time-ordered: the CNG fill for a gap must be + // decoded BEFORE the frames that follow the gap, so depay first, + // decode in playout order. + auto decodeFrame = [&](std::span f) { + return amr ? dec.Decode(f.data(), static_cast(f.size())) : G711DecodeFrame(codec, f); + }; + std::vector zero(frameSamples, 0); std::optional expect; for (;;) { PktItem item; @@ -533,7 +788,18 @@ int main(int argc, char** argv) { item = std::move(playq.front()); playq.pop_front(); } - auto frames = Depay(item.payload, octetAlign); + std::vector> frames; + if (amr) { + for (auto& [hdr, speech] : Depay(codec, item.payload, octetAlign)) { + std::vector f = {hdr}; + f.insert(f.end(), speech.begin(), speech.end()); + frames.push_back(std::move(f)); + } + } else { + // whole frames back to back (a gateway may pack 2 at ptime 40) + for (std::size_t off = 0; off + frameSamples <= item.payload.size(); off += frameSamples) + frames.emplace_back(item.payload.begin() + static_cast(off), item.payload.begin() + static_cast(off + frameSamples)); + } if (frames.empty()) continue; int fill = 0; if (!expect) { @@ -541,36 +807,39 @@ int main(int argc, char** argv) { if (!writePcm(zero)) return; } else { std::uint32_t diff = (item.ts - *expect) & 0xFFFFFFFF; - if (diff >= 0x80000000u) { late++; continue; } - fill = static_cast(diff / 320); + if (diff >= 0x80000000u) { + late++; + continue; + } + fill = static_cast(diff / frameSamples); if (fill > MaxFill) fill = 0; } for (int i = 0; i < fill; i++) { cng++; std::uint8_t nodata = 0x7C; - auto pcm = dec.Decode(&nodata, 1); + std::vector pcm = amr ? dec.Decode(&nodata, 1) : zero; applyGain(pcm); if (!writePcm(pcm)) return; } - for (auto& [hdr, speech] : frames) { + for (auto& f : frames) { rxPlayed++; - std::vector f = {hdr}; - f.insert(f.end(), speech.begin(), speech.end()); - auto pcm = dec.Decode(f.data(), static_cast(f.size())); + std::vector pcm = decodeFrame(f); applyGain(pcm); if (!writePcm(pcm)) return; } - expect = (item.ts + 320 * static_cast(frames.size())) & 0xFFFFFFFF; + expect = (item.ts + static_cast(frameSamples * frames.size())) & 0xFFFFFFFF; } }); } // ---- main recv loop - std::vector silence = SilenceFrame(0, octetAlign); + std::vector silence = SilenceFrame(codec, 0, octetAlign); for (int i = 0; i < 5; i++) RtpSend(sock, st, silence); // latch burst std::FILE* dump = rtpDump ? std::fopen((std::format("{}.rtp", out)).c_str(), "wb") : nullptr; - double t0 = Now(), lastTx = 0, lastRx = Now(); + double t0 = Now(); + double lastTx = 0; + double lastRx = Now(); std::uint64_t rx = 0; std::uint64_t rxBytes = 0; bool gotMedia = false; @@ -588,7 +857,8 @@ int main(int argc, char** argv) { lastTx = now; } std::uint8_t buf[65535]; - sockaddr_storage src; socklen_t srcLen = sizeof src; + sockaddr_storage src; + socklen_t srcLen = sizeof src; ssize_t n = recvfrom(sock, buf, sizeof buf, 0, reinterpret_cast(&src), &srcLen); if (n <= 0) continue; lastRx = Now(); @@ -626,7 +896,10 @@ int main(int argc, char** argv) { std::uint32_t pktTs = (static_cast(buf[4]) << 24) | (buf[5] << 16) | (buf[6] << 8) | buf[7]; std::scoped_lock g(qlock); - if (playq.size() >= PlayqMax) { playq.pop_front(); qdrop++; } + if (playq.size() >= PlayqMax) { + playq.pop_front(); + qdrop++; + } playq.push_back({pktTs, std::vector(buf + off, buf + n)}); qcv.notify_one(); } @@ -637,6 +910,7 @@ int main(int argc, char** argv) { if (micThread.joinable()) micThread.join(); if (playThread.joinable()) playThread.join(); if (dump) std::fclose(dump); + if (pcmFile) std::fclose(pcmFile); if (playOn) { close(playCh.fd); int status; @@ -652,14 +926,14 @@ int main(int argc, char** argv) { // .stats sidecar (tiny; always written) if (std::FILE* sf = std::fopen((std::format("{}.stats", out)).c_str(), "w")) { std::print(sf, - "{{\"tx\": {}, \"tx_mic\": {}, \"rx\": {}, \"rx_bytes\": {}, " + "{{\"codec\": \"{}\", \"tx\": {}, \"tx_mic\": {}, \"rx\": {}, \"rx_bytes\": {}, " "\"rx_played\": {}, \"cng\": {}, \"late\": {}, \"qdrop\": {}, " "\"first_src\": \"{}\", \"dst\": \"{}:{}\", \"pt\": {}, \"mic\": {}, " "\"play\": {}, \"amr_mode\": {}, \"media_ended\": {}}}", - st.tx, st.txMic, rx, rxBytes, rxPlayed, cng, late, qdrop, firstSrc, + CodecName(codec), st.tx, st.txMic, rx, rxBytes, rxPlayed, cng, late, qdrop, firstSrc, rIp, rPort, pt, mic, play, amrMode, mediaEnded); std::fclose(sf); } - std::println("imsd-media: tx={} tx_mic={} rx={} rx_played={} cng={} late={} " "qdrop={} rx_bytes={} first_src={} media_ended={}", st.tx, st.txMic, rx, rxPlayed, cng, late, qdrop, rxBytes, firstSrc, mediaEnded); + std::println("imsd-media: codec={} tx={} tx_mic={} rx={} rx_played={} cng={} late={} " "qdrop={} rx_bytes={} first_src={} media_ended={}", CodecName(codec), st.tx, st.txMic, rx, rxPlayed, cng, late, qdrop, rxBytes, firstSrc, mediaEnded); return mediaEnded ? ExitMediaTimeout : 0; }