koan-core 0.29.1

Core library for koan — bit-perfect music player. Audio engine, player, database, format strings.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
//! StreamingSource — a Read+Seek adapter over a shared, incrementally-filled byte buffer.
//!
//! The download thread writes chunks into a `StreamBuffer` while the Symphonia decoder
//! reads from a `StreamingSource` backed by the same buffer. The source blocks briefly
//! when the read position catches up to the write head, enabling true streaming decode
//! without waiting for the full download to complete.

use std::io::{self, Read, Seek, SeekFrom};
use std::sync::{Arc, Condvar, Mutex};
use std::time::Duration;

/// Longest a read may block waiting for bytes before the stream counts as dead.
/// A download that stops advancing must surface as an error, not park the decode
/// thread forever holding the ring buffer producer and the whole buffered track.
const READ_TIMEOUT: Duration = Duration::from_secs(30);

/// State shared between the download writer and the decoder reader.
struct Inner {
    data: Vec<u8>,
    /// Total expected byte length. `None` if not yet known (no Content-Length).
    total_len: Option<u64>,
    /// Set to true when the download thread has finished cleanly — the buffer
    /// holds the whole source and readers see EOF past its end.
    done: bool,
    /// Set to true when the download died before delivering everything.
    /// Reads past the buffered bytes fail rather than reporting EOF.
    failed: bool,
    /// How long a read blocks for new bytes before giving up.
    read_timeout: Duration,
}

/// A shared, growable byte buffer that the download thread writes into.
///
/// Clone it to get additional handles; all clones share the same underlying data.
#[derive(Clone)]
pub struct StreamBuffer {
    inner: Arc<(Mutex<Inner>, Condvar)>,
}

impl StreamBuffer {
    /// Create a new empty buffer. `total_len` may be provided once Content-Length is known.
    pub fn new(total_len: Option<u64>) -> Self {
        Self::with_read_timeout(total_len, READ_TIMEOUT)
    }

    fn with_read_timeout(total_len: Option<u64>, read_timeout: Duration) -> Self {
        Self {
            inner: Arc::new((
                Mutex::new(Inner {
                    data: Vec::new(),
                    total_len,
                    done: false,
                    failed: false,
                    read_timeout,
                }),
                Condvar::new(),
            )),
        }
    }

    /// Append downloaded bytes. Called by the download thread.
    pub fn push(&self, chunk: &[u8]) {
        let (lock, cvar) = &*self.inner;
        let mut inner = lock.lock().unwrap();
        inner.data.extend_from_slice(chunk);
        cvar.notify_all();
    }

    /// Signal that the download delivered everything. Readers see EOF past the
    /// buffered bytes.
    pub fn finish(&self) {
        let (lock, cvar) = &*self.inner;
        let mut inner = lock.lock().unwrap();
        inner.done = true;
        cvar.notify_all();
    }

    /// Signal that the download died. Readers past the buffered bytes get a
    /// broken-pipe error — reporting EOF here would silently truncate the track.
    pub fn fail(&self) {
        let (lock, cvar) = &*self.inner;
        let mut inner = lock.lock().unwrap();
        inner.failed = true;
        cvar.notify_all();
    }

    /// True when this is the last handle: nothing can read what is written from
    /// here on, so a writer holding it has no reason to keep going.
    pub fn is_abandoned(&self) -> bool {
        Arc::strong_count(&self.inner) == 1
    }

    /// Total bytes received so far.
    pub fn bytes_downloaded(&self) -> u64 {
        let (lock, _) = &*self.inner;
        lock.lock().unwrap().data.len() as u64
    }

    /// Total expected length (from Content-Length), if known.
    pub fn total_len(&self) -> Option<u64> {
        let (lock, _) = &*self.inner;
        lock.lock().unwrap().total_len
    }

    /// Create a `StreamingSource` that reads from this buffer starting at offset 0.
    pub fn reader(&self) -> StreamingSource {
        StreamingSource {
            inner: self.inner.clone(),
            pos: 0,
        }
    }
}

/// A `Read + Seek` view into a `StreamBuffer`.
///
/// Blocks on `read` when the read position is at or beyond the write head,
/// until more bytes arrive or the download finishes.
pub struct StreamingSource {
    inner: Arc<(Mutex<Inner>, Condvar)>,
    pos: u64,
}

impl Read for StreamingSource {
    fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
        if buf.is_empty() {
            return Ok(0);
        }

        let (lock, cvar) = &*self.inner;

        // Wait until there is data at `pos`, or the download ends one way or another.
        let guard = lock.lock().unwrap();
        let timeout = guard.read_timeout;
        let (inner, wait) = cvar
            .wait_timeout_while(guard, timeout, |s| {
                s.data.len() as u64 <= self.pos && !s.done && !s.failed
            })
            .unwrap();

        let available = inner.data.len() as u64;
        if available <= self.pos {
            if inner.failed {
                return Err(io::Error::new(
                    io::ErrorKind::BrokenPipe,
                    "stream download failed before delivering the whole track",
                ));
            }
            if wait.timed_out() {
                return Err(io::Error::new(
                    io::ErrorKind::TimedOut,
                    "stream download stalled",
                ));
            }
            // Done and no more data — EOF.
            return Ok(0);
        }

        let start = self.pos as usize;
        let end = (start + buf.len()).min(inner.data.len());
        let n = end - start;
        buf[..n].copy_from_slice(&inner.data[start..end]);
        self.pos += n as u64;
        Ok(n)
    }
}

impl Seek for StreamingSource {
    fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
        let (lock, _) = &*self.inner;
        let inner = lock.lock().unwrap();

        let new_pos: i64 = match pos {
            SeekFrom::Start(n) => n as i64,
            SeekFrom::Current(n) => self.pos as i64 + n,
            SeekFrom::End(n) => {
                // For End seeks we need total_len. If not known yet, use current data len.
                let len = inner.total_len.unwrap_or(inner.data.len() as u64) as i64;
                len + n
            }
        };

        if new_pos < 0 {
            return Err(io::Error::new(
                io::ErrorKind::InvalidInput,
                "seek before beginning of stream",
            ));
        }

        self.pos = new_pos as u64;
        Ok(self.pos)
    }
}

// Symphonia requires MediaSource: Read + Seek + Send + Any
impl symphonia::core::io::MediaSource for StreamingSource {
    fn is_seekable(&self) -> bool {
        // Seekable only if the total length is known (needed for seek-to-end math).
        // Forward seeks always work; backward seeks require buffered data already present.
        // We advertise seekable=true and handle backward seeks via the buffered Vec.
        true
    }

    fn byte_len(&self) -> Option<u64> {
        let (lock, _) = &*self.inner;
        lock.lock().unwrap().total_len
    }
}

#[cfg(test)]
mod tests {
    use std::io::{self, Read, Seek, SeekFrom};

    use super::*;

    fn filled_buffer(data: &[u8]) -> StreamBuffer {
        let buf = StreamBuffer::new(Some(data.len() as u64));
        buf.push(data);
        buf.finish();
        buf
    }

    #[test]
    fn new_buffer_starts_empty() {
        let buf = StreamBuffer::new(Some(1024));
        assert_eq!(buf.bytes_downloaded(), 0);
        assert_eq!(buf.total_len(), Some(1024));
    }

    #[test]
    fn new_buffer_unknown_total() {
        let buf = StreamBuffer::new(None);
        assert_eq!(buf.total_len(), None);
    }

    #[test]
    fn read_all_data_available() {
        let data = b"hello streaming world";
        let buf = filled_buffer(data);
        let mut src = buf.reader();

        let mut out = Vec::new();
        src.read_to_end(&mut out).unwrap();
        assert_eq!(out, data);
    }

    #[test]
    fn read_partial_then_rest() {
        let data = b"abcdefghij";
        let buf = filled_buffer(data);
        let mut src = buf.reader();

        let mut first = [0u8; 4];
        let n = src.read(&mut first).unwrap();
        assert_eq!(n, 4);
        assert_eq!(&first, b"abcd");

        let mut rest = Vec::new();
        src.read_to_end(&mut rest).unwrap();
        assert_eq!(rest, b"efghij");
    }

    #[test]
    fn seek_from_start() {
        let data = b"0123456789";
        let buf = filled_buffer(data);
        let mut src = buf.reader();

        let pos = src.seek(SeekFrom::Start(5)).unwrap();
        assert_eq!(pos, 5);

        let mut out = [0u8; 3];
        src.read_exact(&mut out).unwrap();
        assert_eq!(&out, b"567");
    }

    #[test]
    fn seek_from_current() {
        let data = b"0123456789";
        let buf = filled_buffer(data);
        let mut src = buf.reader();

        src.seek(SeekFrom::Start(2)).unwrap();
        let pos = src.seek(SeekFrom::Current(3)).unwrap();
        assert_eq!(pos, 5);

        let mut out = [0u8; 2];
        src.read_exact(&mut out).unwrap();
        assert_eq!(&out, b"56");
    }

    #[test]
    fn seek_from_end() {
        let data = b"0123456789";
        let buf = filled_buffer(data);
        let mut src = buf.reader();

        // SeekFrom::End(0) should position at total_len (EOF).
        let pos = src.seek(SeekFrom::End(0)).unwrap();
        assert_eq!(pos, 10);

        // SeekFrom::End(-3) should position at offset 7.
        let pos = src.seek(SeekFrom::End(-3)).unwrap();
        assert_eq!(pos, 7);

        let mut out = [0u8; 3];
        src.read_exact(&mut out).unwrap();
        assert_eq!(&out, b"789");
    }

    #[test]
    fn seek_before_start_errors() {
        let data = b"hello";
        let buf = filled_buffer(data);
        let mut src = buf.reader();

        let result = src.seek(SeekFrom::Current(-1));
        assert!(result.is_err());
    }

    #[test]
    fn is_complete_when_done() {
        let buf = StreamBuffer::new(Some(5));
        buf.push(b"hello");
        // Not yet finished.
        assert!(buf.bytes_downloaded() != 0); // just check bytes_downloaded works
        buf.finish();
        // After finish, a reader should see EOF immediately.
        let mut src = buf.reader();
        let mut out = Vec::new();
        src.read_to_end(&mut out).unwrap();
        assert_eq!(out, b"hello");
    }

    #[test]
    fn byte_len_returns_total() {
        use symphonia::core::io::MediaSource;
        let buf = StreamBuffer::new(Some(42));
        let src = buf.reader();
        assert_eq!(src.byte_len(), Some(42));
    }

    #[test]
    fn is_seekable_true() {
        use symphonia::core::io::MediaSource;
        let buf = StreamBuffer::new(None);
        let src = buf.reader();
        assert!(src.is_seekable());
    }

    #[test]
    fn push_increments_bytes_downloaded() {
        let buf = StreamBuffer::new(Some(10));
        buf.push(b"hello");
        assert_eq!(buf.bytes_downloaded(), 5);
        buf.push(b"world");
        assert_eq!(buf.bytes_downloaded(), 10);
    }

    #[test]
    fn multiple_readers_independent_positions() {
        let data = b"0123456789";
        let buf = filled_buffer(data);

        let mut r1 = buf.reader();
        let mut r2 = buf.reader();

        r1.seek(SeekFrom::Start(7)).unwrap();

        let mut out1 = [0u8; 3];
        r1.read_exact(&mut out1).unwrap();
        assert_eq!(&out1, b"789");

        let mut out2 = [0u8; 3];
        r2.read_exact(&mut out2).unwrap();
        assert_eq!(&out2, b"012");
    }

    #[test]
    fn failed_download_errors_instead_of_reporting_eof() {
        let buf = StreamBuffer::new(Some(1000));
        buf.push(b"partial");
        buf.fail();

        let mut src = buf.reader();
        let mut out = [0u8; 7];
        src.read_exact(&mut out).unwrap();
        assert_eq!(&out, b"partial");

        // Past the buffered bytes: an error, never a clean EOF — Ok(0) here
        // would end the track early and look like a short file.
        let err = src.read(&mut out).unwrap_err();
        assert_eq!(err.kind(), io::ErrorKind::BrokenPipe);
    }

    #[test]
    fn failed_download_wakes_a_blocked_reader() {
        let buf = StreamBuffer::new(Some(1000));
        let mut src = buf.reader();

        let writer = buf.clone();
        let waiter = std::thread::spawn(move || {
            let mut out = [0u8; 8];
            src.read(&mut out).map(|_| ()).map_err(|e| e.kind())
        });

        std::thread::sleep(std::time::Duration::from_millis(20));
        writer.fail();

        assert_eq!(waiter.join().unwrap(), Err(io::ErrorKind::BrokenPipe));
    }

    #[test]
    fn stalled_download_times_out() {
        // A download with a Content-Length that never arrives: the read must
        // give up rather than park the decode thread forever.
        let buf = StreamBuffer::with_read_timeout(Some(1000), std::time::Duration::from_millis(20));
        let mut src = buf.reader();
        let mut out = [0u8; 8];
        let err = src.read(&mut out).unwrap_err();
        assert_eq!(err.kind(), io::ErrorKind::TimedOut);
    }

    #[test]
    fn abandoned_when_no_other_handles_remain() {
        let buf = StreamBuffer::new(None);
        assert!(buf.is_abandoned());

        let reader = buf.reader();
        assert!(!buf.is_abandoned());

        drop(reader);
        assert!(buf.is_abandoned());
    }

    #[test]
    fn test_partial_availability() {
        // Push only 500 bytes without finish() — simulates in-progress download.
        let data: Vec<u8> = (0u8..=255).cycle().take(1000).collect();
        let buf = StreamBuffer::new(Some(1000));
        buf.push(&data[..500]);

        assert_eq!(buf.bytes_downloaded(), 500);
        assert_eq!(buf.total_len(), Some(1000));

        // Spawn thread to call finish() after brief delay so reader doesn't block forever.
        let buf2 = buf.clone();
        std::thread::spawn(move || {
            std::thread::sleep(std::time::Duration::from_millis(10));
            buf2.finish();
        });

        let mut src = buf.reader();
        let mut out = Vec::new();
        src.read_to_end(&mut out).unwrap();
        assert_eq!(out.len(), 500);
        assert_eq!(out, &data[..500]);
    }
}