diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1eaef05..25c82e2 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -37,6 +37,11 @@ jobs: with: components: rustfmt, clippy + - name: Set optional environment variables (Windows) + if: runner.os == 'Windows' + shell: pwsh + run: Add-Content -Path $env:GITHUB_ENV -Value "VCPKG_ROOT=C:\vcpkg\" + - name: Install system dependencies (Linux) if: runner.os == 'Linux' shell: bash @@ -51,6 +56,7 @@ jobs: libclang-dev \ ocl-icd-opencl-dev \ libclfft-dev \ + libopus-dev \ || true # Some runner images may not ship libclfft-dev. Fall back to building from source so the @@ -64,6 +70,10 @@ jobs: sudo ldconfig || true fi + - name: Install system dependencies (Windows) + if: runner.os == 'Windows' + run: C:\vcpkg\vcpkg install opus --vcpkg-root C:\vcpkg + - name: Format run: cargo fmt --check diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index be5a664..ee27c8f 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -58,6 +58,11 @@ jobs: - uses: dtolnay/rust-toolchain@stable # The frontend is a git submodule and ships with prebuilt artifacts in frontend/dist/. + - name: Set optional environment variables (Windows) + if: runner.os == 'Windows' + shell: pwsh + run: Add-Content -Path $env:GITHUB_ENV -Value "VCPKG_ROOT=C:\vcpkg\" + - name: Install system dependencies (Linux) if: runner.os == 'Linux' shell: bash @@ -70,6 +75,7 @@ jobs: soapysdr-tools \ ocl-icd-opencl-dev \ libclfft-dev \ + libopus-dev \ || true # If libclfft-dev isn't available on the runner image, build clFFT from source so the @@ -89,7 +95,11 @@ jobs: run: | set -euo pipefail brew update - brew install pkg-config soapysdr + brew install pkg-config soapysdr opus + + - name: Install system dependencies (Windows) + if: runner.os == 'Windows' + run: C:\vcpkg\vcpkg install opus --vcpkg-root C:\vcpkg - name: Verify frontend/dist exists shell: bash @@ -210,6 +220,7 @@ jobs: apt-get install -y --no-install-recommends \ ocl-icd-opencl-dev \ libclfft-dev \ + libopus-dev \ || true # Build SoapySDR from source for build-time linking (we do not ship/bundle it). diff --git a/Cargo.lock b/Cargo.lock index 8a83305..7c7bab7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1224,6 +1224,16 @@ dependencies = [ "unicode-width", ] +[[package]] +name = "interop" +version = "0.1.0" +dependencies = [ + "bindgen", + "cc", + "pkg-config", + "vcpkg", +] + [[package]] name = "ipnet" version = "2.11.0" @@ -1500,6 +1510,7 @@ dependencies = [ "dashmap", "futures", "inquire", + "interop", "novasdr-core", "num-complex", "rand 0.8.5", @@ -2829,6 +2840,12 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" +[[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + [[package]] name = "version_check" version = "0.9.5" diff --git a/Cargo.toml b/Cargo.toml index 66b0276..94a8df4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,6 +4,7 @@ members = [ "crates/novasdr-core", "crates/novasdr-server", "crates/ws_probe", + "crates/interop", ] [workspace.package] diff --git a/README.md b/README.md index 19b8c8c..335c332 100644 --- a/README.md +++ b/README.md @@ -63,7 +63,8 @@ sudo apt-get update && sudo apt-get install -y --no-install-recommends \ nodejs npm \ ocl-icd-opencl-dev ocl-icd-libopencl1 \ libclfft-dev \ - libusb-1.0-0-dev + libusb-1.0-0-dev \ + libopus-dev ``` @@ -79,7 +80,8 @@ sudo dnf install -y \ swig python3 python3-devel python3-numpy \ nodejs npm \ ocl-icd ocl-icd-devel \ - libusb1-devel + libusb1-devel \ + opus-devel ``` @@ -95,7 +97,8 @@ sudo pacman -Sy --noconfirm --needed \ swig python python-numpy \ nodejs npm \ ocl-icd opencl-headers \ - libusb + libusb \ + opus ``` @@ -111,7 +114,8 @@ sudo zypper --non-interactive refresh && sudo zypper --non-interactive install - swig python3 python3-devel python3-numpy \ nodejs npm \ OpenCL-Headers ocl-icd-devel \ - libusb-1_0-devel + libusb-1_0-devel \ + libopus-devel ``` @@ -125,7 +129,8 @@ brew update && brew install \ llvm \ swig python \ node \ - libusb + libusb \ + opus ``` diff --git a/crates/interop/Cargo.toml b/crates/interop/Cargo.toml new file mode 100644 index 0000000..786ee2b --- /dev/null +++ b/crates/interop/Cargo.toml @@ -0,0 +1,19 @@ +[package] +name = "interop" +version = "0.1.0" +edition = "2021" +authors = ["Alexander Sholokhov "] +keywords = ["ffi", "sdr"] +build = "build.rs" +license = "GPL-3.0-only" + +# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html + +[dependencies] + +[build-dependencies] +bindgen = { version = "0.66.1", default-features = false, features = ["runtime"] } +cc = "1.0" +pkg-config = "0.3.9" +vcpkg = "0.2.15" + diff --git a/crates/interop/build.rs b/crates/interop/build.rs new file mode 100644 index 0000000..8b9aad9 --- /dev/null +++ b/crates/interop/build.rs @@ -0,0 +1,89 @@ +use std::env; +use std::path::PathBuf; + +fn probe_opus_pkg_config() -> Option> { + match pkg_config::Config::new() + .atleast_version("1.3") + .probe("opus") + { + Err(e) => { + eprintln!("pkg_config: {}", e); + None + } + Ok(lib) => Some(lib.include_paths), + } +} + +fn probe_opus_path_patch(paths: Vec) -> Vec { + // MacOs homebrew Opus specific patch + // we get path like /opt/homebrew/Cellar/opus/1.6.1/include/opus , + // but it should be /opt/homebrew/Cellar/opus/1.6.1/include + + for path in paths.iter() { + if path.join("opus/opus.h").exists() { + return paths; + } + } + + let mut op_extra_path: Option = None; + for path in paths.iter() { + if let Some(parent) = path.parent() { + if parent.join("opus/opus.h").exists() { + op_extra_path = Some(parent.into()); + break; + } + } + } + + let Some(extra_path) = op_extra_path else { + return paths; + }; + + let mut paths = paths; + paths.push(extra_path); + paths +} + +fn do_opus() { + println!("cargo:rerun-if-changed=./interop/opus_wrapper.h"); + println!("cargo:rerun-if-changed=./interop/opus_wrapper.c"); + + let include_paths = if env::var_os("VCPKG_ROOT").is_some() { + env::set_var("VCPKGRS_DYNAMIC", "1"); + let pkg = vcpkg::Config::new() + .find_package("opus") + .expect("cant't find opus package"); + pkg.include_paths + } else { + let path = probe_opus_pkg_config().expect("Couldn't find opus"); + probe_opus_path_patch(path) + }; + + let bindgen_builder = bindgen::Builder::default() + .trust_clang_mangling(false) + .size_t_is_usize(true) + .header("opus_wrapper.h"); + + let mut cc_builder = cc::Build::new(); + cc_builder.file("opus_wrapper.c"); + + for elm in include_paths.iter() { + cc_builder.include(elm); + } + + let bindings = bindgen_builder + .allowlist_function("_opus_.*") + .generate() + .expect("Unable to generate bindings"); + + let out_path = PathBuf::from(env::var("OUT_DIR").unwrap()); + bindings + .write_to_file(out_path.join("opus_bindings.rs")) + .expect("Couldn't write opus_bindings!"); + + cc_builder.compile("opus-wrapper") +} + +fn main() { + do_opus(); +} diff --git a/crates/interop/opus_wrapper.c b/crates/interop/opus_wrapper.c new file mode 100644 index 0000000..efff5d8 --- /dev/null +++ b/crates/interop/opus_wrapper.c @@ -0,0 +1,101 @@ +#include "opus_wrapper.h" +#include + +int32_t _opus_application_audio() +{ + return OPUS_APPLICATION_AUDIO; +} + +int32_t _opus_application_voip() +{ + return OPUS_APPLICATION_VOIP; +} + +int32_t _opus_application_lowdelay() +{ + return OPUS_APPLICATION_RESTRICTED_LOWDELAY; +} + +int32_t _opus_bitrate_max() +{ + return OPUS_BITRATE_MAX; +} + +int32_t _opus_bitrate_auto() +{ + return OPUS_AUTO; +} + +void *_opus_encoder_create(int32_t fs, int32_t channels, int32_t application, int32_t *err) +{ + // allowed fs: 48000, 24000, 16000, 12000, 8000 + // allowed channels: 1, 2 + int local_err = 0; + void *res = opus_encoder_create(fs, channels, application, &local_err); + *err = local_err; // safe translation C-int to Rust's i32 + return res; +} + +void _opus_encoder_destroy(void *enc) +{ + opus_encoder_destroy(enc); +} + +int32_t _opus_set_bitrate(void *enc, int32_t bitrate) +{ + // Rates from 500 to 512000 bits per second are meaningful + return opus_encoder_ctl(enc, OPUS_SET_BITRATE(bitrate)); +} + +int32_t _opus_set_complexity(void *enc, int32_t complexity) +{ + return opus_encoder_ctl(enc, OPUS_SET_COMPLEXITY(complexity)); +} + +int32_t _opus_encode_i16(void *enc, const int16_t *pcm, size_t frame_size, uint8_t *data, size_t max_data_bytes) +{ + // To encode a frame, opus_encode() or opus_encode_float() must be called with exactly one frame (2.5, 5, 10, 20, 40 or 60 ms) of audio data: + return opus_encode(enc, pcm, frame_size, data, max_data_bytes); +} + +int32_t _opus_encode_float(void *enc, const float *pcm, size_t frame_size, uint8_t *data, size_t max_data_bytes) +{ + // To encode a frame, opus_encode() or opus_encode_float() must be called with exactly one frame (2.5, 5, 10, 20, 40 or 60 ms) of audio data: + return opus_encode_float(enc, pcm, frame_size, data, max_data_bytes); +} + +void *_opus_decoder_create(int32_t fs, int32_t channels, int32_t *error) +{ + // allowed fs: 48000, 24000, 16000, 12000, 8000 + // allowed channels: 1, 2 + int local_err = 0; + void *res = opus_decoder_create(fs, channels, &local_err); + *error = local_err; + return res; +} + +int32_t _opus_decode_i16(void *dec, const uint8_t *data, size_t len, int16_t *pcm, size_t frame_size, int32_t decode_fec) +{ + // decode_fec: Flag (0 or 1) to request that any in-band forward error correction data be + return opus_decode(dec, data, len, pcm, frame_size, decode_fec); +} + +int32_t _opus_decode_float(void *dec, const uint8_t *data, size_t len, float *pcm, size_t frame_size, int32_t decode_fec) +{ + return opus_decode_float(dec, data, len, pcm, frame_size, decode_fec); +} + +void _opus_decoder_destroy(void *dec) +{ + opus_decoder_destroy(dec); +} + +const char *_opus_get_version_string(void) +{ + return opus_get_version_string(); +} + +const char *_opus_strerror(int32_t error) +{ + return opus_strerror(error); +} diff --git a/crates/interop/opus_wrapper.h b/crates/interop/opus_wrapper.h new file mode 100644 index 0000000..b6277e9 --- /dev/null +++ b/crates/interop/opus_wrapper.h @@ -0,0 +1,35 @@ +#pragma once + +#include +#include + +#ifdef __cplusplus +extern "C" +{ +#endif + +int32_t _opus_application_audio(); +int32_t _opus_application_voip(); +int32_t _opus_application_lowdelay(); + +int32_t _opus_bitrate_max(); +int32_t _opus_bitrate_auto(); + +void *_opus_encoder_create(int32_t fs, int32_t channels, int32_t application, int32_t *err); +void _opus_encoder_destroy(void *enc); +int32_t _opus_set_bitrate(void *enc, int32_t bitrate); +int32_t _opus_set_complexity(void *enc, int32_t complexity); +int32_t _opus_encode_i16(void *enc, const int16_t *pcm, size_t frame_size, uint8_t *data, size_t max_data_bytes); +int32_t _opus_encode_float(void *enc, const float *pcm, size_t frame_size, uint8_t *data, size_t max_data_bytes); + +void *_opus_decoder_create(int32_t fs, int32_t channels, int32_t *error); +void _opus_decoder_destroy(void *dec); +int32_t _opus_decode_i16(void *dec, const uint8_t *data, size_t len, int16_t *pcm, size_t frame_size, int32_t decode_fec); +int32_t _opus_decode_float(void *dec, const uint8_t *data, size_t len, float *pcm, size_t frame_size, int32_t decode_fec); + +const char *_opus_get_version_string(void); +const char *_opus_strerror(int32_t error); + +#ifdef __cplusplus +} +#endif diff --git a/crates/interop/src/lib.rs b/crates/interop/src/lib.rs new file mode 100644 index 0000000..9365348 --- /dev/null +++ b/crates/interop/src/lib.rs @@ -0,0 +1 @@ +pub mod opus; diff --git a/crates/interop/src/opus.rs b/crates/interop/src/opus.rs new file mode 100644 index 0000000..797d01f --- /dev/null +++ b/crates/interop/src/opus.rs @@ -0,0 +1,229 @@ +use core::str; +use std::ffi; +use std::fmt; +use std::fmt::Debug; + +mod inner { + include!(concat!(env!("OUT_DIR"), "/opus_bindings.rs")); +} + +pub fn get_version_string() -> Result { + let s = unsafe { ffi::CStr::from_ptr(inner::_opus_get_version_string()) }; + s.to_str().map(|x| x.to_string()) +} + +fn opus_strerr(err: i32) -> Result { + let s = unsafe { ffi::CStr::from_ptr(inner::_opus_strerror(err)) }; + s.to_str().map(|x| x.to_string()) +} + +#[derive(Debug)] +pub enum SampleRate { + Hz8000, + Hz12000, + Hz16000, + Hz24000, + Hz48000, +} + +impl SampleRate { + pub fn as_int32(&self) -> i32 { + match self { + SampleRate::Hz8000 => 8000, + SampleRate::Hz12000 => 12000, + SampleRate::Hz16000 => 16000, + SampleRate::Hz24000 => 24000, + SampleRate::Hz48000 => 48000, + } + } +} + +#[derive(Debug)] +pub enum Channels { + Mono, + Stereo, +} + +impl Channels { + pub fn as_int32(&self) -> i32 { + match self { + Channels::Mono => 1, + Channels::Stereo => 2, + } + } +} + +#[derive(Debug)] +pub enum Application { + Voip, + Audio, + LowDelay, +} + +impl Application { + pub fn as_int32(&self) -> i32 { + unsafe { + match self { + Application::Voip => inner::_opus_application_voip(), + Application::Audio => inner::_opus_application_audio(), + Application::LowDelay => inner::_opus_application_lowdelay(), + } + } + } +} + +#[derive(Debug)] +pub enum Bitrate { + BitsPerSecond(i32), + Max, + Auto, +} + +impl Bitrate { + pub fn as_int32(&self) -> i32 { + unsafe { + match self { + Bitrate::BitsPerSecond(x) => *x, + Bitrate::Max => inner::_opus_bitrate_max(), + Bitrate::Auto => inner::_opus_bitrate_auto(), + } + } + } +} + +#[allow(non_camel_case_types)] +#[derive(Debug)] +pub enum OpusError { + OPUS_OK, + OPUS_BAD_ARG, + OPUS_BUFFER_TOO_SMALL, + OPUS_INTERNAL_ERROR, + OPUS_INVALID_PACKET, + OPUS_UNIMPLEMENTED, + OPUS_INVALID_STATE, + OPUS_ALLOC_FAIL, + UNKNOWN(i32), +} + +impl From for OpusError { + fn from(value: i32) -> Self { + match value { + 0 => OpusError::OPUS_OK, + -1 => OpusError::OPUS_BAD_ARG, + -2 => OpusError::OPUS_BUFFER_TOO_SMALL, + -3 => OpusError::OPUS_INTERNAL_ERROR, + -4 => OpusError::OPUS_INVALID_PACKET, + -5 => OpusError::OPUS_UNIMPLEMENTED, + -6 => OpusError::OPUS_INVALID_STATE, + -7 => OpusError::OPUS_ALLOC_FAIL, + x => OpusError::UNKNOWN(x), + } + } +} + +impl From<&OpusError> for i32 { + fn from(val: &OpusError) -> Self { + match val { + OpusError::OPUS_OK => 0, + OpusError::OPUS_BAD_ARG => -1, + OpusError::OPUS_BUFFER_TOO_SMALL => -2, + OpusError::OPUS_INTERNAL_ERROR => -3, + OpusError::OPUS_INVALID_PACKET => -4, + OpusError::OPUS_UNIMPLEMENTED => -5, + OpusError::OPUS_INVALID_STATE => -6, + OpusError::OPUS_ALLOC_FAIL => -7, + OpusError::UNKNOWN(x) => *x, + } + } +} + +impl fmt::Display for OpusError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + f, + "{:?}: '{}'", + self, + opus_strerr(self.into()).unwrap_or_default() + ) + } +} + +impl std::error::Error for OpusError {} + +#[derive(Debug)] +pub struct Encoder { + ptr: std::ptr::NonNull<::std::os::raw::c_void>, + channels: usize, +} + +impl Encoder { + pub fn new( + sample_rate: SampleRate, + channels: Channels, + application: Application, + ) -> Result { + let mut err = 0; + let p_enc = unsafe { + inner::_opus_encoder_create( + sample_rate.as_int32(), + channels.as_int32(), + application.as_int32(), + &mut err, + ) + }; + + if err == 0 { + Ok(Encoder { + ptr: unsafe { std::ptr::NonNull::new_unchecked(p_enc) }, + channels: channels.as_int32() as usize, + }) + } else { + Err(OpusError::from(err)) + } + } + + pub fn set_bitrate(&mut self, bitrate: Bitrate) -> Result<(), OpusError> { + let rc: i32 = unsafe { inner::_opus_set_bitrate(self.ptr.as_ptr(), bitrate.as_int32()) }; + if rc == 0 { + Ok(()) + } else { + Err(rc.into()) + } + } + + pub fn set_complexity(&mut self, complexity: i32) -> Result<(), OpusError> { + let rc: i32 = unsafe { inner::_opus_set_complexity(self.ptr.as_ptr(), complexity) }; + if rc == 0 { + Ok(()) + } else { + Err(rc.into()) + } + } + + pub fn encode(&self, input: &[i16], output: &mut [u8]) -> Result { + let rc: i32 = unsafe { + inner::_opus_encode_i16( + self.ptr.as_ptr(), + input.as_ptr(), + input.len() / self.channels, + output.as_mut_ptr(), + output.len(), + ) + }; + if rc >= 0 { + Ok(rc as usize) + } else { + Err(rc.into()) + } + } +} + +impl Drop for Encoder { + fn drop(&mut self) { + unsafe { + inner::_opus_encoder_destroy(self.ptr.as_ptr()); + } + } +} + +unsafe impl Send for Encoder {} diff --git a/crates/novasdr-core/src/config.rs b/crates/novasdr-core/src/config.rs index 03a32a2..aa3efd9 100644 --- a/crates/novasdr-core/src/config.rs +++ b/crates/novasdr-core/src/config.rs @@ -211,6 +211,7 @@ pub enum WaterfallCompression { pub enum AudioCompression { Adpcm, Flac, + Opus, } #[derive(Debug, Clone, Copy, Deserialize, PartialEq, Eq, Default)] @@ -652,6 +653,7 @@ impl Config { let audio_compression_str = match input.audio_compression { AudioCompression::Adpcm => "adpcm".to_string(), AudioCompression::Flac => "flac".to_string(), + AudioCompression::Opus => "opus".to_string(), }; Ok(Runtime { diff --git a/crates/novasdr-server/Cargo.toml b/crates/novasdr-server/Cargo.toml index 7f4d3e3..2f9cd8b 100644 --- a/crates/novasdr-server/Cargo.toml +++ b/crates/novasdr-server/Cargo.toml @@ -15,6 +15,7 @@ chrono = { version = "0.4.39", default-features = false, features = ["clock", "s dashmap = "6.1.0" futures = "0.3.31" inquire = "0.7.5" +interop = { path = "../interop" } novasdr-core = { path = "../novasdr-core" } num-complex = "0.4.6" rand = "0.8.5" diff --git a/crates/novasdr-server/src/main.rs b/crates/novasdr-server/src/main.rs index f39f1bd..441ae17 100644 --- a/crates/novasdr-server/src/main.rs +++ b/crates/novasdr-server/src/main.rs @@ -15,6 +15,7 @@ mod update_check; mod ws; use anyhow::Context; +use interop::opus; use novasdr_core::config; use std::io::IsTerminal; use std::path::Path; @@ -249,6 +250,11 @@ fn main() -> anyhow::Result<()> { } } + tracing::info!( + version = opus::get_version_string().unwrap_or_default(), + "Opus" + ); + let available_threads = std::thread::available_parallelism() .map(|n| n.get()) .unwrap_or(1); diff --git a/crates/novasdr-server/src/setup.rs b/crates/novasdr-server/src/setup.rs index f4e2f9a..80f52b9 100644 --- a/crates/novasdr-server/src/setup.rs +++ b/crates/novasdr-server/src/setup.rs @@ -1087,6 +1087,20 @@ fn edit_receiver(receiver: &mut Value) -> anyhow::Result<()> { )?; input.insert("audio_sps".to_string(), json!(audio_sps)); + { + let current = input + .get("audio_compression") + .and_then(Value::as_str) + .unwrap_or("adpcm"); + let labels = vec!["opus".to_string(), "adpcm".to_string()]; + let default_idx = labels.iter().position(|s| s == current).unwrap_or(0); + let selected = Select::new("Audio compression", labels) + .with_starting_cursor(default_idx) + .prompt() + .context("prompt audio compression")?; + input.insert("audio_compression".to_string(), json!(selected)); + } + let waterfall_size = prompt_usize( "Waterfall width (waterfall_size)", input diff --git a/crates/novasdr-server/src/ws/audio.rs b/crates/novasdr-server/src/ws/audio.rs index 88b57ca..191a4a7 100644 --- a/crates/novasdr-server/src/ws/audio.rs +++ b/crates/novasdr-server/src/ws/audio.rs @@ -6,6 +6,7 @@ use axum::{ response::IntoResponse, }; use futures::{SinkExt, StreamExt}; +use interop::opus; use novasdr_core::{ config::AudioCompression, dsp::{ @@ -22,8 +23,8 @@ use num_complex::Complex32; use realfft::{ComplexToReal, RealFftPlanner}; use rustfft::{Fft as RustFft, FftPlanner}; use serde_json::json; -use std::net::SocketAddr; use std::sync::Arc; +use std::{mem, net::SocketAddr}; fn with_audio_unique_id(basic_info: String, unique_id: &str) -> String { let Ok(mut v) = serde_json::from_str::(&basic_info) else { @@ -113,25 +114,30 @@ fn squelch_features(bins: &[Complex32]) -> SquelchFeatures { } const AUDIO_FRAME_MAGIC: [u8; 4] = *b"NSDA"; -const AUDIO_FRAME_VERSION: u8 = 1; -const AUDIO_FRAME_HEADER_LEN: usize = 36; +const AUDIO_FRAME_END_MARK: u16 = 0xaabb; +const AUDIO_FRAME_VERSION: u8 = 2; +const AUDIO_FRAME_HEADER_LEN: usize = 40; #[derive(Clone, Copy, Debug)] #[repr(u8)] enum AudioWireCodec { AdpcmIma = 1, + Opus = 2, } -fn build_audio_frame( +fn build_audio_frame_multi( codec: AudioWireCodec, frame_num: u64, l: i32, m: f64, r: i32, pwr: f32, - payload: &[u8], + payload: Vec>, ) -> Vec { - let mut out = Vec::with_capacity(AUDIO_FRAME_HEADER_LEN + payload.len()); + let expected_capacity = payload + .iter() + .fold(AUDIO_FRAME_HEADER_LEN, |acc, x| acc + 2 + x.len()); + let mut out = Vec::with_capacity(expected_capacity); out.extend_from_slice(&AUDIO_FRAME_MAGIC); out.push(AUDIO_FRAME_VERSION); out.push(codec as u8); @@ -141,7 +147,13 @@ fn build_audio_frame( out.extend_from_slice(&m.to_le_bytes()); out.extend_from_slice(&r.to_le_bytes()); out.extend_from_slice(&pwr.to_le_bytes()); - out.extend_from_slice(payload); + out.extend_from_slice(&(payload.len() as u16).to_le_bytes()); + for frame in payload { + out.extend_from_slice(&(frame.len() as u16).to_le_bytes()); + out.extend(frame); + } + out.extend_from_slice(&AUDIO_FRAME_END_MARK.to_le_bytes()); + debug_assert_eq!(expected_capacity, out.len()); out } @@ -807,13 +819,13 @@ pub struct AudioPipeline { pcm_accum_i16: Vec, pcm_accum_offset: usize, packet_samples: usize, - pwr_sum: f32, - pwr_frames: usize, dc: DcBlocker, agc: Agc, fm_prev: Complex32, last_agc: (AgcSpeed, Option, Option), squelch: SquelchState, + opus_encoder: Option, + opus_wrk_buf: Vec, } impl AudioPipeline { @@ -831,13 +843,61 @@ impl AudioPipeline { let frame_samples = audio_fft_size / 2; - // Batch ~20ms of PCM per websocket frame to reduce packet rate and browser-side scheduling - // overhead (too many tiny frames can stutter). - let target_packet_sec = 0.020_f64; - let min_packet = ((sample_rate as f64) * target_packet_sec).ceil().max(1.0) as usize; - let mut packet_samples = frame_samples.max(min_packet); - packet_samples = packet_samples.div_ceil(8) * 8; - packet_samples = packet_samples.clamp(frame_samples, 8192); + let packet_samples = match compression { + AudioCompression::Adpcm => { + // Batch ~20ms of PCM per websocket frame to reduce packet rate and browser-side scheduling + // overhead (too many tiny frames can stutter). + let target_packet_sec = 0.020_f64; + let min_packet = + ((sample_rate as f64) * target_packet_sec).ceil().max(1.0) as usize; + let mut packet_samples = frame_samples.max(min_packet); + packet_samples = packet_samples.div_ceil(8) * 8; + packet_samples.clamp(frame_samples, 8192) + } + AudioCompression::Opus => { + // number of milliseconds per chunk. opus allowed values: 5, 10, 20, 40, 60. + let ms = 20; + sample_rate * ms / 1000 + } + AudioCompression::Flac => { + return Err(anyhow::anyhow!( + "FLAC audio was removed; configure audio_compression = \"opus\" or \"adpcm\"" + )) + } + }; + + let (opus_encoder, opus_wrk_buf) = if compression == AudioCompression::Opus { + let opus_sample_rate = match sample_rate { + 8000 => opus::SampleRate::Hz8000, + 12000 => opus::SampleRate::Hz12000, + 16000 => opus::SampleRate::Hz16000, + 24000 => opus::SampleRate::Hz24000, + 48000 => opus::SampleRate::Hz48000, + x => return Err(anyhow::anyhow!("Unsupported sample rate {x} for Opus codec. Valid values are: [8000, 12000, 16000, 24000, 48000]")), + }; + + let mut opus_encoder = opus::Encoder::new( + opus_sample_rate, + opus::Channels::Mono, + opus::Application::LowDelay, + ) + .map_err(|e| anyhow::anyhow!("Opus create error: {e}"))?; + + // 40kbps Opus produces excellent quality for VoIP needs. + if let Err(e) = opus_encoder.set_bitrate(opus::Bitrate::BitsPerSecond(40000)) { + tracing::warn!(error = ?e, "opus. unsuccess set_bitrate"); + } + + if let Err(e) = opus_encoder.set_complexity(2) { + tracing::warn!(error = ?e, "opus. unsuccess set_complexity"); + } + + // 120ms with 48000sps, doubled. More than enough for Opus encoder output buffer. + let max_wrk_buf_size = 120 * 48000 * 2 / 1000; + (Some(opus_encoder), vec![0; max_wrk_buf_size]) + } else { + (None, vec![]) + }; Ok(Self { compression, @@ -858,8 +918,6 @@ impl AudioPipeline { pcm_accum_i16: Vec::with_capacity(packet_samples * 4), pcm_accum_offset: 0, packet_samples, - pwr_sum: 0.0, - pwr_frames: 0, // Keep the DC blocker cutoff low so AM has real low end; bass boost is frontend-only. dc: DcBlocker::new((sample_rate / 20).max(128)), // Match reference defaults. @@ -867,6 +925,8 @@ impl AudioPipeline { fm_prev: Complex32::new(0.0, 0.0), last_agc: (AgcSpeed::Default, None, None), squelch: SquelchState::new(), + opus_encoder, + opus_wrk_buf, }) } @@ -883,8 +943,6 @@ impl AudioPipeline { self.agc.reset(); self.pcm_accum_i16.clear(); self.pcm_accum_offset = 0; - self.pwr_sum = 0.0; - self.pwr_frames = 0; } pub fn process( @@ -1069,15 +1127,15 @@ impl AudioPipeline { float_to_i16_centered(audio_out, &mut self.pcm_frame_i16, 32768.0); self.pcm_accum_i16.extend_from_slice(&self.pcm_frame_i16); - self.pwr_sum += spectrum_slice.iter().map(|c| c.norm_sqr()).sum::(); - self.pwr_frames += 1; + let pwr = spectrum_slice.iter().map(|c| c.norm_sqr()).sum::(); - if self.compression != AudioCompression::Adpcm { - return Err(anyhow::anyhow!( - "FLAC audio was removed; configure audio_compression = \"adpcm\"" - )); - } + let audio_wire_codec = match self.compression { + AudioCompression::Adpcm => AudioWireCodec::AdpcmIma, + AudioCompression::Opus => AudioWireCodec::Opus, + AudioCompression::Flac => unreachable!(), + }; + let mut acc_frames: Vec> = Vec::new(); loop { let available = self .pcm_accum_i16 @@ -1091,27 +1149,55 @@ impl AudioPipeline { let block = &self.pcm_accum_i16[self.pcm_accum_offset..end]; self.pcm_accum_offset = end; - let frames = self.pwr_frames.max(1) as f32; - let pwr = self.pwr_sum / frames; - self.pwr_sum = 0.0; - self.pwr_frames = 0; + let payload = match self.compression { + AudioCompression::Adpcm => ima_adpcm::encode_block_i16_mono(block), + AudioCompression::Opus => { + let Some(opus_encoder) = self.opus_encoder.as_ref() else { + return Err(anyhow::anyhow!("Opus encoder is None. Impossible.")); + }; + let size = opus_encoder + .encode(block, &mut self.opus_wrk_buf) + .map_err(|e| anyhow::anyhow!("Opus encode chunk error: {e}"))?; + self.opus_wrk_buf[0..size].to_vec() + } + AudioCompression::Flac => unreachable!(), + }; + + let audio_frame_size_threshold = 700; // keep frame size less than N bytes if possible + let collected = acc_frames.iter().map(|x| x.len()).sum::(); + if collected + payload.len() > audio_frame_size_threshold { + let taken_vec = mem::replace(&mut acc_frames, vec![payload]); + out_packets.push(build_audio_frame_multi( + audio_wire_codec, + frame_num, + 0, + params.m, + spectrum_slice.len() as i32, + pwr, + taken_vec, + )); + } else { + acc_frames.push(payload); + } - let payload = ima_adpcm::encode_block_i16_mono(block); - out_packets.push(build_audio_frame( - AudioWireCodec::AdpcmIma, + if self.pcm_accum_offset >= self.packet_samples * 4 { + self.pcm_accum_i16.drain(0..self.pcm_accum_offset); + self.pcm_accum_offset = 0; + } + } + + if !acc_frames.is_empty() { + out_packets.push(build_audio_frame_multi( + audio_wire_codec, frame_num, 0, params.m, spectrum_slice.len() as i32, pwr, - &payload, + acc_frames, )); - - if self.pcm_accum_offset >= self.packet_samples * 4 { - self.pcm_accum_i16.drain(0..self.pcm_accum_offset); - self.pcm_accum_offset = 0; - } } + Ok(out_packets) }