Skip to main content

koan_core/upnp/
stream.rs

1//! A stream koan processed itself, served to a renderer.
2//!
3//! When the output's profile does anything, the renderer cannot be handed the
4//! original file: what the profile does to the samples would never reach it.
5//! The session is decoded and processed here as it would be for this device,
6//! into the same ring, and the encoder drains the ring instead of an audio
7//! callback. Its output is one endless FLAC (or WAV) stream per session,
8//! served on a token like a file. Track changes inside the session are the
9//! timeline's boundaries; the renderer is never asked to switch.
10//!
11//! The renderer paces the stream. Encoded chunks go into a `Pipe`, a log
12//! that every open connection reads at its own place. The encoder runs at
13//! most `AHEAD` bytes past the connection furthest along, so a renderer that
14//! stops reading (paused, or its buffer full) holds the encoder, which holds
15//! the ring, which holds the decoder.
16//!
17//! Renderers open a URL more than once, and not one after another: Kodi opens
18//! one connection to play from and another that reads a few hundred kilobytes
19//! and hangs up, at the same moment. So no connection replaces another, and
20//! until one has read `EARLY` bytes a new connection starts at the top of the
21//! stream. After that a new one is a reconnect, and joins where the stream
22//! is; `Pipe::origin_ms` says how far into the stream that was, which is
23//! what the renderer's position then counts from. A connection that falls
24//! more than `AHEAD` behind the furthest skips forward.
25
26use std::collections::VecDeque;
27use std::io::Write;
28use std::net::TcpStream;
29use std::sync::Arc;
30use std::sync::atomic::{AtomicU64, Ordering};
31use std::thread;
32use std::time::Duration;
33
34use flacenc::component::{BitRepr, Stream, StreamInfo};
35use flacenc::error::Verify;
36use flacenc::source::{Fill, FrameBuf};
37use parking_lot::{Condvar, Mutex};
38
39use super::serve::Request;
40use crate::audio::buffer::PlaybackTimeline;
41use crate::player::state::QueueItemId;
42
43/// Frames per FLAC block, and per chunk of the pipe.
44const BLOCK: usize = 4096;
45
46/// How far the encoder runs ahead of the connection furthest along, and how
47/// far behind it a connection may fall, in bytes.
48const AHEAD: usize = 512 * 1024;
49
50/// How far a connection reads before a new one joins the stream where it is
51/// rather than at its start.
52const EARLY: usize = 1024 * 1024;
53
54/// Bytes of audio between two ICY metadata blocks.
55const ICY_INTERVAL: usize = 16 * 1024;
56
57#[derive(Debug, Clone, Copy, PartialEq, Eq)]
58pub enum Encoding {
59    Flac,
60    Wav,
61}
62
63impl Encoding {
64    pub fn extension(self) -> &'static str {
65        match self {
66            Self::Flac => "flac",
67            Self::Wav => "wav",
68        }
69    }
70}
71
72/// What the stream carries: the session's output rate and channels, and the
73/// integer depth the samples are dithered to.
74#[derive(Debug, Clone, Copy, PartialEq, Eq)]
75pub struct Format {
76    pub encoding: Encoding,
77    pub rate: u32,
78    pub channels: u16,
79    pub bits: u8,
80}
81
82impl Format {
83    /// The depth a source's samples are sent at: its own, within what both
84    /// encodings carry, and 24 for a source that has none (a lossy one).
85    pub fn bits_for(source: Option<u16>) -> u8 {
86        source.map_or(24, |b| b.clamp(16, 24) as u8)
87    }
88}
89
90struct Chunk {
91    bytes: Vec<u8>,
92    /// Bytes of the stream before this chunk, the header aside.
93    offset: usize,
94    /// Frames of the stream before this chunk's first.
95    start: u64,
96    /// The track playing at this chunk's first frame.
97    track: Option<QueueItemId>,
98}
99
100impl Chunk {
101    fn end(&self) -> usize {
102        self.offset + self.bytes.len()
103    }
104}
105
106#[derive(Default)]
107struct State {
108    header: Vec<u8>,
109    /// The chunks some connection may still read, oldest first.
110    log: VecDeque<Arc<Chunk>>,
111    /// The index of `log`'s first chunk in the stream.
112    base: usize,
113    /// Each open connection, oldest first: its id, the index of the next
114    /// chunk it reads, and the frame it joined at.
115    readers: Vec<(u64, usize, u64)>,
116    /// The furthest any connection has read, as the index of the next chunk.
117    read_to: usize,
118    /// Bytes encoded so far, and the frame the next chunk starts at.
119    written: usize,
120    next_frame: u64,
121    connections: u64,
122    finished: bool,
123    closed: bool,
124}
125
126impl State {
127    /// The next chunk after the furthest any connection has read, and how
128    /// far into the stream that is.
129    fn furthest(&self) -> (usize, usize) {
130        let index = self.read_to.max(self.base);
131        let offset = self
132            .log
133            .get(index - self.base)
134            .map_or(self.written, |c| c.offset);
135        (index, offset)
136    }
137
138    /// Whether a new connection starts at the top: nobody has read far yet.
139    fn early(&self) -> bool {
140        self.base == 0 && self.furthest().1 <= EARLY
141    }
142
143    /// Drop what no connection can be served any more: everything before
144    /// `AHEAD` behind the furthest, once past the early part.
145    fn trim(&mut self) {
146        if self.early() {
147            return;
148        }
149        let (_, furthest) = self.furthest();
150        while self.log.front().is_some_and(|c| c.end() + AHEAD < furthest) {
151            self.log.pop_front();
152            self.base += 1;
153        }
154    }
155}
156
157/// Between the encoder and the connections the renderer has open.
158pub struct Pipe {
159    state: Mutex<State>,
160    changed: Condvar,
161    rate: u32,
162    origin: AtomicU64,
163    mime: String,
164    title: Box<dyn Fn(QueueItemId) -> String + Send + Sync>,
165}
166
167impl Pipe {
168    pub fn new(
169        format: Format,
170        mime: &str,
171        title: impl Fn(QueueItemId) -> String + Send + Sync + 'static,
172    ) -> Arc<Self> {
173        Arc::new(Self {
174            state: Mutex::default(),
175            changed: Condvar::new(),
176            rate: format.rate,
177            origin: AtomicU64::new(0),
178            mime: mime.to_string(),
179            title: Box::new(title),
180        })
181    }
182
183    /// Where in the stream the renderer's connection began, which its
184    /// position counts from.
185    pub fn origin_ms(&self) -> u64 {
186        self.origin.load(Ordering::Acquire) * 1000 / self.rate.max(1) as u64
187    }
188
189    /// End the stream now: the encoder and every connection let go.
190    pub fn close(&self) {
191        self.state.lock().closed = true;
192        self.changed.notify_all();
193    }
194
195    fn set_header(&self, header: Vec<u8>) {
196        self.state.lock().header = header;
197    }
198
199    /// Add a chunk, waiting while the encoder is `AHEAD` of every connection.
200    /// False once the pipe is closed.
201    fn push(&self, bytes: Vec<u8>, start: u64, frames: u64, track: Option<QueueItemId>) -> bool {
202        let mut state = self.state.lock();
203        while state.written >= state.furthest().1 + AHEAD && !state.closed {
204            self.changed.wait(&mut state);
205        }
206        if state.closed {
207            return false;
208        }
209        let offset = state.written;
210        state.written += bytes.len();
211        state.next_frame = start + frames;
212        state.log.push_back(Arc::new(Chunk {
213            bytes,
214            offset,
215            start,
216            track,
217        }));
218        state.trim();
219        drop(state);
220        self.changed.notify_all();
221        true
222    }
223
224    /// Everything has been written: a connection that reaches the end ends.
225    fn finish(&self) {
226        self.state.lock().finished = true;
227        self.changed.notify_all();
228    }
229
230    fn connect(&self) -> Reader {
231        let mut state = self.state.lock();
232        state.connections += 1;
233        let id = state.connections;
234        // Early, from the top; otherwise a reconnect, carrying on from the
235        // furthest anything has been read.
236        let index = if state.early() { 0 } else { state.furthest().0 };
237        let origin = match index {
238            0 => 0,
239            _ => state
240                .log
241                .get(index - state.base)
242                .map_or(state.next_frame, |c| c.start),
243        };
244        state.readers.push((id, index, origin));
245        self.follow_oldest(&state);
246        Reader { id }
247    }
248
249    /// The renderer's position counts from where the connection it plays
250    /// from joined, taken to be the oldest still open: a connection opened
251    /// alongside it is a probe, and one opened after it has gone is its
252    /// reconnect. With none open, the last stands.
253    fn follow_oldest(&self, state: &State) {
254        if let Some(oldest) = state.readers.first() {
255            self.origin.store(oldest.2, Ordering::Release);
256        }
257    }
258
259    fn disconnect(&self, reader: &Reader) {
260        let mut state = self.state.lock();
261        state.readers.retain(|r| r.0 != reader.id);
262        self.follow_oldest(&state);
263        state.trim();
264        drop(state);
265        self.changed.notify_all();
266    }
267
268    /// The next chunk for `reader`, waiting for the encoder. `None` when the
269    /// stream is over: finished, or closed.
270    fn next(&self, reader: &Reader) -> Option<Arc<Chunk>> {
271        let mut state = self.state.lock();
272        loop {
273            if state.closed {
274                return None;
275            }
276            let base = state.base;
277            let at = state.readers.iter().position(|r| r.0 == reader.id)?;
278            // Fallen behind what is kept: on from the oldest kept.
279            let index = state.readers[at].1.max(base);
280            if let Some(chunk) = state.log.get(index - base).cloned() {
281                state.readers[at].1 = index + 1;
282                state.read_to = state.read_to.max(index + 1);
283                state.trim();
284                drop(state);
285                self.changed.notify_all();
286                return Some(chunk);
287            }
288            if state.finished {
289                return None;
290            }
291            self.changed.wait(&mut state);
292        }
293    }
294}
295
296struct Reader {
297    id: u64,
298}
299
300/// Serve the stream to one connection, for as long as it reads.
301pub(crate) fn serve(stream: &mut TcpStream, req: &Request, pipe: &Pipe) -> std::io::Result<()> {
302    let icy = req.header("Icy-MetaData").is_some_and(|v| v.trim() == "1");
303    // No length and no ranges: the stream is made as it is sent. CI=1 says
304    // it is converted from the original.
305    write!(
306        stream,
307        "HTTP/1.1 200 OK\r\nContent-Type: {}\r\nAccept-Ranges: none\r\n{}transferMode.dlna.org: Streaming\r\ncontentFeatures.dlna.org: DLNA.ORG_OP=00;DLNA.ORG_CI=1;DLNA.ORG_FLAGS=01700000000000000000000000000000\r\nConnection: close\r\n\r\n",
308        pipe.mime,
309        if icy {
310            format!("icy-metaint: {ICY_INTERVAL}\r\n")
311        } else {
312            String::new()
313        }
314    )?;
315    if req.method == "HEAD" {
316        return Ok(());
317    }
318    let header = pipe.state.lock().header.clone();
319    let mut out = Icy {
320        stream,
321        every: icy.then_some(ICY_INTERVAL),
322        left: ICY_INTERVAL,
323        title: String::new(),
324        sent: String::new(),
325    };
326    let reader = pipe.connect();
327    log::info!(
328        "upnp: stream connection {} ({}), from {}ms",
329        reader.id,
330        req.header("User-Agent").unwrap_or("no agent"),
331        pipe.origin_ms()
332    );
333    let mut sent = 0usize;
334    let result = (|| {
335        out.write(&header)?;
336        while let Some(chunk) = pipe.next(&reader) {
337            if let Some(track) = chunk.track {
338                out.title = (pipe.title)(track);
339            }
340            out.write(&chunk.bytes)?;
341            sent += chunk.bytes.len();
342        }
343        Ok::<_, std::io::Error>(())
344    })();
345    pipe.disconnect(&reader);
346    log::info!(
347        "upnp: stream connection {} ended after {sent} bytes{}",
348        reader.id,
349        match &result {
350            Err(e) => format!(": {e}"),
351            Ok(()) => String::new(),
352        }
353    );
354    result
355}
356
357/// Audio bytes with ICY metadata blocks between them, when the renderer asked
358/// for them: the title of the track at that point in the stream, sent when it
359/// changes.
360struct Icy<'a> {
361    stream: &'a mut TcpStream,
362    every: Option<usize>,
363    left: usize,
364    title: String,
365    sent: String,
366}
367
368impl Icy<'_> {
369    fn write(&mut self, mut bytes: &[u8]) -> std::io::Result<()> {
370        let Some(every) = self.every else {
371            return self.stream.write_all(bytes);
372        };
373        while !bytes.is_empty() {
374            let n = self.left.min(bytes.len());
375            self.stream.write_all(&bytes[..n])?;
376            bytes = &bytes[n..];
377            self.left -= n;
378            if self.left == 0 {
379                let block = if self.title == self.sent {
380                    vec![0]
381                } else {
382                    self.sent = self.title.clone();
383                    icy_block(&self.title)
384                };
385                self.stream.write_all(&block)?;
386                self.left = every;
387            }
388        }
389        Ok(())
390    }
391}
392
393/// One ICY metadata block: its length in sixteens, then `StreamTitle`, padded.
394fn icy_block(title: &str) -> Vec<u8> {
395    let title = title.replace('\'', "’");
396    let mut text = format!("StreamTitle='{title}';").into_bytes();
397    text.truncate(255 * 16);
398    let sixteens = text.len().div_ceil(16);
399    text.resize(sixteens * 16, 0);
400    let mut block = vec![sixteens as u8];
401    block.extend(text);
402    block
403}
404
405/// The encoder thread. Dropping it waits for the thread, so close the pipe
406/// first.
407pub struct Encoder {
408    thread: Option<thread::JoinHandle<()>>,
409}
410
411impl Drop for Encoder {
412    fn drop(&mut self) {
413        if let Some(thread) = self.thread.take() {
414            let _ = thread.join();
415        }
416    }
417}
418
419/// Drain `consumer` into `pipe`, encoded as `format` says, until the decoder
420/// lets go of the ring or the pipe is closed. `decoder` is woken whenever the
421/// ring has room: a decoder parked on a full ring waits for as long as half
422/// of it takes to play, and the renderer reads faster than that at first.
423pub fn start(
424    consumer: rtrb::Consumer<f32>,
425    format: Format,
426    pipe: Arc<Pipe>,
427    timeline: Arc<PlaybackTimeline>,
428    decoder: Option<thread::Thread>,
429) -> std::io::Result<Encoder> {
430    let mut codec = Codec::new(format).map_err(std::io::Error::other)?;
431    pipe.set_header(codec.header());
432    let thread = thread::Builder::new()
433        .name("koan-upnp-encode".into())
434        .spawn(move || {
435            encode(
436                consumer,
437                format,
438                &mut codec,
439                &pipe,
440                &timeline,
441                decoder.as_ref(),
442            )
443        })?;
444    Ok(Encoder {
445        thread: Some(thread),
446    })
447}
448
449fn encode(
450    mut consumer: rtrb::Consumer<f32>,
451    format: Format,
452    codec: &mut Codec,
453    pipe: &Pipe,
454    timeline: &PlaybackTimeline,
455    decoder: Option<&thread::Thread>,
456) {
457    let channels = format.channels as usize;
458    let block = BLOCK * channels;
459    let mut samples: Vec<f32> = Vec::with_capacity(block);
460    let mut ints: Vec<i32> = Vec::with_capacity(block);
461    let mut dither = Dither::new(format.bits);
462    let mut frames: u64 = 0;
463    loop {
464        let wanted = block - samples.len();
465        let ready = consumer.slots().min(wanted);
466        if ready > 0
467            && let Ok(chunk) = consumer.read_chunk(ready)
468        {
469            samples.extend(chunk);
470            if let Some(decoder) = decoder {
471                decoder.unpark();
472            }
473        }
474        let ended = consumer.is_abandoned() && consumer.is_empty();
475        if samples.len() == block || ended && !samples.is_empty() {
476            ints.clear();
477            ints.extend(samples.iter().map(|&s| dither.quantise(s)));
478            let bytes = codec.encode(&ints, samples.len() / channels);
479            let track = timeline.track_at(frames * channels as u64);
480            let count = (samples.len() / channels) as u64;
481            if !pipe.push(bytes, frames, count, track) {
482                return;
483            }
484            frames += count;
485            samples.clear();
486            continue;
487        }
488        if ended {
489            pipe.finish();
490            return;
491        }
492        if pipe.state.lock().closed {
493            return;
494        }
495        // An empty ring is the start of a session or its end; while it plays,
496        // the renderer's reading is what this thread waits on, in `push`.
497        if ready == 0 {
498            thread::park_timeout(Duration::from_millis(10));
499        }
500    }
501}
502
503/// f32 to integers at `bits`, with triangular dither of one LSB either way:
504/// truncating without it would add distortion the chain did not ask for.
505struct Dither {
506    scale: f64,
507    rng: u64,
508}
509
510impl Dither {
511    fn new(bits: u8) -> Self {
512        Self {
513            scale: (1u64 << (bits - 1)) as f64,
514            rng: 0x9E37_79B9_7F4A_7C15,
515        }
516    }
517
518    fn uniform(&mut self) -> f64 {
519        // xorshift64*
520        self.rng ^= self.rng >> 12;
521        self.rng ^= self.rng << 25;
522        self.rng ^= self.rng >> 27;
523        (self.rng.wrapping_mul(0x2545_F491_4F6C_DD1D) >> 11) as f64 / (1u64 << 53) as f64
524    }
525
526    fn quantise(&mut self, sample: f32) -> i32 {
527        let noise = self.uniform() - self.uniform();
528        (sample as f64 * self.scale + noise)
529            .round()
530            .clamp(-self.scale, self.scale - 1.0) as i32
531    }
532}
533
534enum Codec {
535    Flac {
536        config: Box<flacenc::error::Verified<flacenc::config::Encoder>>,
537        info: StreamInfo,
538        buffer: FrameBuf,
539        frame: usize,
540    },
541    Wav {
542        format: Format,
543    },
544}
545
546impl Codec {
547    fn new(format: Format) -> Result<Self, String> {
548        Ok(match format.encoding {
549            Encoding::Flac => {
550                let mut info = StreamInfo::new(
551                    format.rate as usize,
552                    format.channels as usize,
553                    format.bits as usize,
554                )
555                .map_err(|e| e.to_string())?;
556                info.set_block_sizes(BLOCK, BLOCK)
557                    .map_err(|e| e.to_string())?;
558                Self::Flac {
559                    config: Box::new(
560                        flacenc::config::Encoder::default()
561                            .into_verified()
562                            .map_err(|(_, e)| e.to_string())?,
563                    ),
564                    info,
565                    buffer: FrameBuf::with_size(format.channels as usize, BLOCK)
566                        .map_err(|e| e.to_string())?,
567                    frame: 0,
568                }
569            }
570            Encoding::Wav => Self::Wav { format },
571        })
572    }
573
574    /// The start of the stream. No length is known, so a FLAC stream says
575    /// zero samples, as the format allows, and a WAV the largest size.
576    fn header(&self) -> Vec<u8> {
577        match self {
578            Self::Flac { info, .. } => {
579                let mut sink = flacenc::bitsink::MemSink::<u8>::new();
580                Stream::with_stream_info(info.clone())
581                    .write(&mut sink)
582                    .expect("writing to memory");
583                sink.as_slice().to_vec()
584            }
585            Self::Wav { format } => {
586                let bytes = format.bits as u32 / 8;
587                let align = bytes * format.channels as u32;
588                let mut h = Vec::with_capacity(44);
589                h.extend_from_slice(b"RIFF");
590                h.extend_from_slice(&u32::MAX.to_le_bytes());
591                h.extend_from_slice(b"WAVEfmt ");
592                h.extend_from_slice(&16u32.to_le_bytes());
593                h.extend_from_slice(&1u16.to_le_bytes());
594                h.extend_from_slice(&format.channels.to_le_bytes());
595                h.extend_from_slice(&format.rate.to_le_bytes());
596                h.extend_from_slice(&(format.rate * align).to_le_bytes());
597                h.extend_from_slice(&(align as u16).to_le_bytes());
598                h.extend_from_slice(&(format.bits as u16).to_le_bytes());
599                h.extend_from_slice(b"data");
600                h.extend_from_slice(&u32::MAX.to_le_bytes());
601                h
602            }
603        }
604    }
605
606    fn encode(&mut self, ints: &[i32], frames: usize) -> Vec<u8> {
607        match self {
608            Self::Flac {
609                config,
610                info,
611                buffer,
612                frame,
613            } => {
614                if buffer.size() != frames {
615                    buffer.resize(frames);
616                }
617                buffer.fill_interleaved(ints).expect("samples in range");
618                let encoded = flacenc::encode_fixed_size_frame(config, buffer, *frame, info)
619                    .expect("frame number in range");
620                *frame += 1;
621                let mut sink = flacenc::bitsink::MemSink::<u8>::new();
622                encoded.write(&mut sink).expect("writing to memory");
623                sink.as_slice().to_vec()
624            }
625            Self::Wav { format } => {
626                let bytes = format.bits as usize / 8;
627                let mut out = Vec::with_capacity(ints.len() * bytes);
628                for s in ints {
629                    out.extend_from_slice(&s.to_le_bytes()[..bytes]);
630                }
631                out
632            }
633        }
634    }
635}
636
637/// Decode a whole stream, for tests that check what a renderer was sent.
638#[cfg(test)]
639pub(crate) fn decode(bytes: &[u8], extension: &str) -> (u32, Vec<f32>) {
640    use symphonia::core::codecs::audio::AudioDecoderOptions;
641    use symphonia::core::formats::probe::Hint;
642    use symphonia::core::formats::{FormatOptions, TrackType};
643    use symphonia::core::io::MediaSourceStream;
644    use symphonia::core::meta::MetadataOptions;
645    let mss = MediaSourceStream::new(
646        Box::new(std::io::Cursor::new(bytes.to_vec())),
647        Default::default(),
648    );
649    let mut hint = Hint::new();
650    hint.with_extension(extension);
651    let mut reader = symphonia::default::get_probe()
652        .probe(
653            &hint,
654            mss,
655            FormatOptions::default(),
656            MetadataOptions::default(),
657        )
658        .unwrap();
659    let track = reader.default_track(TrackType::Audio).unwrap();
660    let id = track.id;
661    let params = track.codec_params.as_ref().unwrap().audio().unwrap();
662    let rate = params.sample_rate.unwrap();
663    let mut decoder = symphonia::default::get_codecs()
664        .make_audio_decoder(params, &AudioDecoderOptions::default())
665        .unwrap();
666    let mut out = Vec::new();
667    while let Ok(Some(packet)) = reader.next_packet() {
668        if packet.track_id != id {
669            continue;
670        }
671        let decoded = decoder.decode(&packet).unwrap();
672        let mut samples = vec![0f32; decoded.samples_interleaved()];
673        decoded.copy_to_slice_interleaved(&mut samples);
674        out.extend(samples);
675    }
676    (rate, out)
677}
678
679#[cfg(test)]
680mod tests {
681    use super::*;
682
683    /// A stereo sine through the encoder, as the renderer would read it.
684    fn round_trip(encoding: Encoding, bits: u8) -> (Vec<f32>, Vec<f32>, u32) {
685        let format = Format {
686            encoding,
687            rate: 48_000,
688            channels: 2,
689            bits,
690        };
691        let input: Vec<f32> = (0..30_000)
692            .flat_map(|i| {
693                let s = (i as f32 * 0.031).sin() * 0.5;
694                [s, -s]
695            })
696            .collect();
697        let (mut producer, consumer) = rtrb::RingBuffer::new(input.len());
698        for s in &input {
699            producer.push(*s).unwrap();
700        }
701        drop(producer);
702        let pipe = Pipe::new(format, "audio/flac", |_| String::new());
703        let encoder = start(
704            consumer,
705            format,
706            pipe.clone(),
707            PlaybackTimeline::new(),
708            None,
709        )
710        .unwrap();
711        let mut bytes = pipe.state.lock().header.clone();
712        let reader = pipe.connect();
713        while let Some(chunk) = pipe.next(&reader) {
714            bytes.extend_from_slice(&chunk.bytes);
715        }
716        drop(encoder);
717        let (rate, output) = decode(&bytes, encoding.extension());
718        (input, output, rate)
719    }
720
721    #[test]
722    fn a_flac_stream_decodes_to_what_was_sent_within_the_dither() {
723        let (input, output, rate) = round_trip(Encoding::Flac, 24);
724        assert_eq!(rate, 48_000);
725        assert_eq!(output.len(), input.len());
726        let lsb = 1.0 / (1u32 << 23) as f32;
727        let worst = input
728            .iter()
729            .zip(&output)
730            .map(|(a, b)| (a - b).abs())
731            .fold(0f32, f32::max);
732        assert!(worst <= 2.0 * lsb, "worst error {worst}, lsb {lsb}");
733    }
734
735    #[test]
736    fn a_wav_stream_at_sixteen_bits_decodes_within_the_dither() {
737        let (input, output, _) = round_trip(Encoding::Wav, 16);
738        assert_eq!(output.len(), input.len());
739        let lsb = 1.0 / 32768.0;
740        for (a, b) in input.iter().zip(&output) {
741            assert!((a - b).abs() <= 2.0 * lsb, "{a} vs {b}");
742        }
743    }
744
745    #[test]
746    fn dither_is_unbiased_and_one_lsb_wide() {
747        let mut dither = Dither::new(16);
748        let target = 0.25 / 32768.0;
749        let n = 200_000;
750        let values: Vec<i32> = (0..n).map(|_| dither.quantise(target)).collect();
751        assert!(values.iter().all(|v| (-1..=1).contains(v)));
752        let mean = values.iter().map(|&v| v as f64).sum::<f64>() / n as f64;
753        assert!((mean - 0.25).abs() < 0.01, "mean {mean}");
754    }
755
756    #[test]
757    fn icy_blocks_are_padded_sixteens() {
758        let block = icy_block("Polar Bear – Peepers");
759        assert_eq!(block.len(), 1 + block[0] as usize * 16);
760        assert!(block[1..].starts_with(b"StreamTitle='Polar Bear"));
761        assert_eq!(icy_block("it's").len() % 16, 1);
762    }
763
764    /// Kodi's pattern: a probe that hangs up, a second connection, and a
765    /// third opened while the second is still reading, which is the one it
766    /// plays. The third must go on getting the stream after the replay.
767    #[test]
768    fn the_connection_that_replaces_an_open_one_is_fed_past_the_replay() {
769        use std::io::Read;
770        let format = Format {
771            encoding: Encoding::Wav,
772            rate: 44_100,
773            channels: 2,
774            bits: 16,
775        };
776        let (mut producer, consumer) = rtrb::RingBuffer::new(1 << 16);
777        let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
778        let feeding = stop.clone();
779        let feeder = thread::spawn(move || {
780            while !feeding.load(Ordering::Relaxed) {
781                if producer.push(0.25).is_err() {
782                    thread::sleep(Duration::from_millis(1));
783                }
784            }
785        });
786        let pipe = Pipe::new(format, "audio/wav", |_| String::new());
787        let _encoder = start(
788            consumer,
789            format,
790            pipe.clone(),
791            PlaybackTimeline::new(),
792            None,
793        )
794        .unwrap();
795        let listener = super::super::serve::Listener::start(Box::new(|_, _| {})).unwrap();
796        let token = listener.add(super::super::serve::Served::Stream {
797            pipe: pipe.clone(),
798            art: Default::default(),
799        });
800        let open = || {
801            let mut s = TcpStream::connect(("127.0.0.1", listener.port())).unwrap();
802            write!(s, "GET /t/{token}.wav HTTP/1.1\r\nRange: bytes=0-\r\n\r\n").unwrap();
803            s
804        };
805        let read = |s: &mut TcpStream, n: usize| {
806            let mut buf = vec![0u8; n];
807            s.read_exact(&mut buf).unwrap();
808        };
809        let mut probe = open();
810        read(&mut probe, 160_000);
811        drop(probe);
812        let mut second = open();
813        read(&mut second, 340_000);
814        let mut third = open();
815        // The replay, and well past it.
816        read(&mut third, 2_000_000);
817        drop(second);
818        stop.store(true, Ordering::Relaxed);
819        pipe.close();
820        feeder.join().unwrap();
821    }
822
823    /// Connections opened at the start each get the stream from its top,
824    /// and none ends another. Once one has read past `EARLY`, a new one is a
825    /// reconnect and carries on from the furthest read.
826    #[test]
827    fn early_connections_start_at_the_top_and_a_later_one_carries_on() {
828        let format = Format {
829            encoding: Encoding::Wav,
830            rate: 1000,
831            channels: 1,
832            bits: 16,
833        };
834        let pipe = Pipe::new(format, "audio/wav", |_| String::new());
835        let push = |start| assert!(pipe.push(vec![0; 300 * 1024], start, 1000, None));
836        push(0);
837        push(1000);
838        let probe = pipe.connect();
839        assert_eq!(pipe.next(&probe).unwrap().start, 0);
840        let player = pipe.connect();
841        assert_eq!(pipe.next(&player).unwrap().start, 0);
842        assert_eq!(pipe.next(&probe).unwrap().start, 1000, "the probe reads on");
843        pipe.disconnect(&probe);
844        assert_eq!(pipe.next(&player).unwrap().start, 1000);
845        assert_eq!(pipe.origin_ms(), 0);
846
847        push(2000);
848        push(3000);
849        assert_eq!(pipe.next(&player).unwrap().start, 2000);
850        assert_eq!(pipe.next(&player).unwrap().start, 3000);
851        push(4000);
852        pipe.disconnect(&player);
853        let again = pipe.connect();
854        assert_eq!(pipe.origin_ms(), 4000);
855        assert_eq!(pipe.next(&again).unwrap().start, 4000);
856    }
857
858    /// A connection opened late while the renderer still plays from another
859    /// is a probe: the position the renderer reports still counts from where
860    /// its own connection joined. Once that one goes, the newer is its
861    /// reconnect.
862    #[test]
863    fn a_late_connection_alongside_the_playing_one_moves_nothing() {
864        let format = Format {
865            encoding: Encoding::Wav,
866            rate: 1000,
867            channels: 1,
868            bits: 16,
869        };
870        let pipe = Pipe::new(format, "audio/wav", |_| String::new());
871        let push = |start| assert!(pipe.push(vec![0; 300 * 1024], start, 1000, None));
872        let player = pipe.connect();
873        for start in [0, 1000, 2000, 3000] {
874            push(start);
875            pipe.next(&player).unwrap();
876        }
877        push(4000);
878        assert_eq!(pipe.origin_ms(), 0);
879
880        let probe = pipe.connect();
881        assert_eq!(pipe.origin_ms(), 0, "the probe moves nothing");
882        assert_eq!(pipe.next(&probe).unwrap().start, 4000);
883        pipe.disconnect(&probe);
884        assert_eq!(pipe.origin_ms(), 0);
885
886        let reconnect = pipe.connect();
887        pipe.disconnect(&player);
888        assert_eq!(pipe.origin_ms(), 5000, "the reconnect carries on");
889        pipe.disconnect(&reconnect);
890    }
891}