1use 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
43const BLOCK: usize = 4096;
45
46const AHEAD: usize = 512 * 1024;
49
50const EARLY: usize = 1024 * 1024;
53
54const 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#[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 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 offset: usize,
94 start: u64,
96 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 log: VecDeque<Arc<Chunk>>,
111 base: usize,
113 readers: Vec<(u64, usize, u64)>,
116 read_to: usize,
118 written: usize,
120 next_frame: u64,
121 connections: u64,
122 finished: bool,
123 closed: bool,
124}
125
126impl State {
127 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 fn early(&self) -> bool {
140 self.base == 0 && self.furthest().1 <= EARLY
141 }
142
143 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
157pub 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 pub fn origin_ms(&self) -> u64 {
186 self.origin.load(Ordering::Acquire) * 1000 / self.rate.max(1) as u64
187 }
188
189 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 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 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 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 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 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 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
300pub(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 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
357struct 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
393fn 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
405pub 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
419pub 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 if ready == 0 {
498 thread::park_timeout(Duration::from_millis(10));
499 }
500 }
501}
502
503struct 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 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 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#[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 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 #[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 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 #[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 #[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}