#![cfg(all(
feature = "cli",
feature = "std",
feature = "tnc",
feature = "micE",
feature = "kiss",
feature = "fx25",
feature = "wav"
))]
#[path = "../src/bin/yodel/shared.rs"]
#[allow(dead_code, unused_imports)]
mod shared;
#[path = "../src/bin/yodel/serve.rs"]
#[allow(dead_code, unused_imports)]
mod yodel_bin;
use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::sync::mpsc::{self, RecvTimeoutError};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use yodel_bin::serve::{
MAX_CLIENTS, PcmSink, SampleSink, ServeStats, kiss_bytes, run_stream, run_tcp,
};
use yodel::SampleRate;
use yodel::ax25::{Address, UiFrame};
use yodel::kiss::{KissCommand, KissDeframer};
use yodel::tnc::{DefaultTncReceiver, TncConfig, TncReceiver, TncTransmitter};
const DEADLINE: Duration = Duration::from_secs(30);
fn addr(call: &[u8], ssid: u8) -> Address {
Address::new(call, ssid).unwrap()
}
fn config() -> TncConfig {
TncConfig::bell_202(SampleRate::new(48_000).unwrap()).unwrap()
}
fn frame_body(src_ssid: u8, info: &[u8]) -> Vec<u8> {
let frame = UiFrame::new(addr(b"APRS", 0), addr(b"N0CALL", src_ssid), info);
let mut buf = [0u8; 330];
let len = frame.build(&mut buf).unwrap();
buf[..len].to_vec()
}
fn frame_samples(body: &[u8]) -> Vec<i16> {
TncTransmitter::new(config())
.frame_samples_i16(body)
.collect()
}
fn decode_all(samples: &[i16]) -> Vec<Vec<u8>> {
let mut rx: DefaultTncReceiver = TncReceiver::new(config()).unwrap();
let mut out = Vec::new();
for &s in samples {
if let Some(frame) = rx.push_i16(s) {
let mut buf = [0u8; 330];
let len = frame.ui_frame().build(&mut buf).unwrap();
out.push(buf[..len].to_vec());
}
}
out
}
#[derive(Clone, Default)]
struct SharedSink(Arc<Mutex<Vec<i16>>>);
impl SampleSink for SharedSink {
fn write_samples(&mut self, samples: &[i16]) -> Result<(), String> {
self.0.lock().unwrap().extend_from_slice(samples);
Ok(())
}
fn finish(&mut self) -> Result<(), String> {
Ok(())
}
}
struct PacedAudio {
rx: mpsc::Receiver<Vec<i16>>,
current: std::vec::IntoIter<i16>,
}
impl Iterator for PacedAudio {
type Item = Result<i16, String>;
fn next(&mut self) -> Option<Self::Item> {
loop {
if let Some(s) = self.current.next() {
return Some(Ok(s));
}
match self.rx.recv() {
Ok(burst) => self.current = burst.into_iter(),
Err(_) => return None,
}
}
}
}
type BridgeHandles = (
u16,
mpsc::SyncSender<Vec<i16>>,
SharedSink,
mpsc::Receiver<Result<ServeStats, String>>,
);
fn start_tcp_bridge() -> BridgeHandles {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (audio_tx, audio_rx) = mpsc::sync_channel::<Vec<i16>>(16);
let sink = SharedSink::default();
let (done_tx, done_rx) = mpsc::channel();
let mut bridge_sink = sink.clone();
std::thread::spawn(move || {
let rx_audio = PacedAudio {
rx: audio_rx,
current: Vec::new().into_iter(),
};
let result = run_tcp(listener, config(), false, rx_audio, &mut bridge_sink);
let _ = done_tx.send(result);
});
(port, audio_tx, sink, done_rx)
}
fn connect(port: u16) -> TcpStream {
let stream = TcpStream::connect(("127.0.0.1", port)).unwrap();
stream.set_read_timeout(Some(DEADLINE)).unwrap();
std::thread::sleep(Duration::from_millis(200));
stream
}
fn read_kiss_frame(stream: &mut TcpStream) -> Vec<u8> {
let mut deframer = KissDeframer::<400>::new();
let mut byte = [0u8; 1];
let start = Instant::now();
loop {
assert!(start.elapsed() < DEADLINE, "timed out reading a KISS frame");
match stream.read(&mut byte) {
Ok(0) => panic!("socket closed before a KISS frame arrived"),
Ok(_) => {
if let Some(result) = deframer.push(byte[0]) {
let frame = result.expect("well-formed KISS frame");
assert_eq!(frame.command(), KissCommand::Data);
return frame.payload().to_vec();
}
}
Err(e) => panic!("reading the client socket: {e}"),
}
}
}
fn wait_for_sink(sink: &SharedSink, pred: impl Fn(&[i16]) -> bool) -> Vec<i16> {
let start = Instant::now();
loop {
{
let samples = sink.0.lock().unwrap();
if pred(&samples) {
return samples.clone();
}
}
assert!(start.elapsed() < DEADLINE, "timed out waiting for TX audio");
std::thread::sleep(Duration::from_millis(20));
}
}
#[test]
fn tcp_bridge_round_trips_both_directions() {
let (port, audio_tx, sink, done_rx) = start_tcp_bridge();
let mut client = connect(port);
let rx_body = frame_body(1, b">rx via radio");
audio_tx.send(frame_samples(&rx_body)).unwrap();
assert_eq!(read_kiss_frame(&mut client), rx_body);
let tx_body = frame_body(2, b">tx via client");
client.write_all(&kiss_bytes(&tx_body)).unwrap();
client.flush().unwrap();
let samples = wait_for_sink(&sink, |s| decode_all(s).len() == 1);
assert_eq!(decode_all(&samples), vec![tx_body]);
drop(audio_tx);
let stats = done_rx
.recv_timeout(DEADLINE)
.expect("bridge must shut down at audio EOF")
.expect("bridge must exit cleanly");
assert_eq!(
stats,
ServeStats {
rx_frames: 1,
tx_frames: 1
}
);
let mut rest = Vec::new();
assert_eq!(client.read_to_end(&mut rest).unwrap_or(0), rest.len());
}
#[test]
fn tcp_bridge_broadcasts_to_every_client() {
let (port, audio_tx, sink, done_rx) = start_tcp_bridge();
let mut first = connect(port);
let mut second = connect(port);
let rx_body = frame_body(3, b">to everyone");
audio_tx.send(frame_samples(&rx_body)).unwrap();
assert_eq!(read_kiss_frame(&mut first), rx_body);
assert_eq!(read_kiss_frame(&mut second), rx_body);
let tx_body = frame_body(4, b">second speaks");
second.write_all(&kiss_bytes(&tx_body)).unwrap();
second.flush().unwrap();
let samples = wait_for_sink(&sink, |s| decode_all(s).len() == 1);
assert_eq!(decode_all(&samples), vec![tx_body]);
drop(audio_tx);
let stats = done_rx
.recv_timeout(DEADLINE)
.expect("bridge must shut down")
.expect("bridge must exit cleanly");
assert_eq!(
stats,
ServeStats {
rx_frames: 1,
tx_frames: 1
}
);
}
#[test]
fn tcp_bridge_survives_client_disconnect() {
let (port, audio_tx, _sink, done_rx) = start_tcp_bridge();
let leaver = connect(port);
let mut stayer = connect(port);
drop(leaver);
let rx_body = frame_body(5, b">still here");
audio_tx.send(frame_samples(&rx_body)).unwrap();
assert_eq!(read_kiss_frame(&mut stayer), rx_body);
drop(audio_tx);
assert!(
done_rx
.recv_timeout(DEADLINE)
.expect("bridge must shut down")
.is_ok()
);
}
#[test]
fn tcp_bridge_shuts_down_without_clients() {
let (_port, audio_tx, _sink, done_rx) = start_tcp_bridge();
drop(audio_tx);
match done_rx.recv_timeout(DEADLINE) {
Ok(result) => assert_eq!(result.unwrap(), ServeStats::default()),
Err(RecvTimeoutError::Timeout) => panic!("idle bridge failed to shut down"),
Err(e) => panic!("bridge thread lost: {e}"),
}
}
#[test]
fn disconnected_clients_free_their_slots() {
let (port, audio_tx, _sink, done_rx) = start_tcp_bridge();
for _ in 0..MAX_CLIENTS {
drop(connect(port));
}
std::thread::sleep(Duration::from_millis(500));
let mut fresh = connect(port);
let rx_body = frame_body(9, b">after the churn");
audio_tx.send(frame_samples(&rx_body)).unwrap();
assert_eq!(
read_kiss_frame(&mut fresh),
rx_body,
"a client after {MAX_CLIENTS} connect/disconnect cycles must still be served"
);
drop(audio_tx);
assert!(
done_rx
.recv_timeout(DEADLINE)
.expect("bridge must shut down")
.is_ok()
);
}
#[test]
fn sink_failure_shuts_down_with_endless_audio() {
struct FailingSink;
impl SampleSink for FailingSink {
fn write_samples(&mut self, _samples: &[i16]) -> Result<(), String> {
Err("sink refused the burst".to_owned())
}
fn finish(&mut self) -> Result<(), String> {
Ok(())
}
}
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (done_tx, done_rx) = mpsc::channel();
std::thread::spawn(move || {
let rx_audio = std::iter::repeat(Ok(0i16));
let mut sink = FailingSink;
let _ = done_tx.send(run_tcp(listener, config(), false, rx_audio, &mut sink));
});
let mut client = connect(port);
client
.write_all(&kiss_bytes(&frame_body(10, b">into a failing sink")))
.unwrap();
client.flush().unwrap();
let result = done_rx
.recv_timeout(DEADLINE)
.expect("a failing sink must end the run even with endless audio");
assert!(
result.is_err(),
"the sink error must be reported, got {result:?}"
);
}
#[test]
fn stream_bridge_round_trips_in_memory() {
#[derive(Clone, Default)]
struct SharedWriter(Arc<Mutex<Vec<u8>>>);
impl Write for SharedWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
let rx_body = frame_body(6, b">over the air");
let tx_body = frame_body(7, b">from the host");
let kiss_in = kiss_bytes(&tx_body);
let kiss_out = SharedWriter::default();
let mut tx_audio: Vec<u8> = Vec::new();
let rx_audio = frame_samples(&rx_body).into_iter().map(Ok);
let stats = {
let mut sink = PcmSink { out: &mut tx_audio };
run_stream(
&kiss_in[..],
kiss_out.clone(),
config(),
false,
rx_audio,
&mut sink,
)
.expect("the stream bridge must run to EOF")
};
assert_eq!(
stats,
ServeStats {
rx_frames: 1,
tx_frames: 1
}
);
let mut deframer = KissDeframer::<400>::new();
let mut heard = Vec::new();
for &byte in kiss_out.0.lock().unwrap().iter() {
if let Some(Ok(frame)) = deframer.push(byte) {
assert_eq!(frame.command(), KissCommand::Data);
heard.push(frame.payload().to_vec());
}
}
assert_eq!(heard, vec![rx_body]);
let samples: Vec<i16> = tx_audio
.chunks_exact(2)
.map(|b| i16::from_le_bytes([b[0], b[1]]))
.collect();
assert_eq!(decode_all(&samples), vec![tx_body]);
}