libfreemkv 1.1.0

Open source raw disc access library for optical drives
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
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
//! NetworkStream — PES frames over TCP with embedded metadata.
//!
//! **Security:** Data is transmitted over plain TCP with no encryption.
//! Use only on trusted networks (LAN).
//!
//! Write side (sender): connects to a listener, sends FMKV header + PES frames.
//! Read side (receiver): listens for a connection, reads FMKV header + PES frames.

use super::meta;
use crate::disc::DiscTitle;
use std::io::{self, BufReader, BufWriter, Write};
use std::net::{IpAddr, TcpListener, TcpStream, ToSocketAddrs};

/// I/O buffer size for network reads/writes.
const NET_BUF_SIZE: usize = 256 * 1024;

/// True if `ip` is one we must never connect a `network://` output to:
/// loopback, RFC1918/ULA private, link-local, unspecified, or multicast.
///
/// `validate_network_target` (in autorip) vets the host once at
/// settings-save time, but the raw hostname is re-resolved here at rip
/// time — a DNS-rebinding attacker can flip a previously-public name to
/// `127.0.0.1` / `10.x` / `169.254.x` in that window. Re-checking the
/// actually-resolved address at connect time closes that TOCTOU.
pub(crate) fn is_blocked_ip(ip: IpAddr) -> bool {
    match ip {
        IpAddr::V4(v4) => {
            let o = v4.octets();
            v4.is_loopback()
                || v4.is_private()
                || v4.is_link_local()
                || v4.is_unspecified()
                || v4.is_multicast()
                || v4.is_broadcast()
                // carrier-grade NAT 100.64.0.0/10
                || (o[0] == 100 && (o[1] & 0xc0) == 0x40)
                // "this network" 0.0.0.0/8
                || o[0] == 0
                // Class E reserved 240.0.0.0/4
                || o[0] >= 240
        }
        IpAddr::V6(v6) => {
            v6.is_loopback()
                || v6.is_unspecified()
                || v6.is_multicast()
                // unique-local fc00::/7
                || (v6.segments()[0] & 0xfe00) == 0xfc00
                // link-local fe80::/10
                || (v6.segments()[0] & 0xffc0) == 0xfe80
                // IPv4-mapped (::ffff:x.x.x.x) and IPv4-compatible (::x.x.x.x);
                // to_ipv4() returns Some for both forms — re-check as IPv4 so an
                // IPv4-mapped private/loopback address can't bypass the block above.
                || v6.to_ipv4().map(|m| is_blocked_ip(IpAddr::V4(m))) == Some(true)
        }
    }
}

/// Resolve `addr` (host:port) and return the first socket address whose
/// IP is NOT [`is_blocked_ip`]. Errors with
/// [`crate::error::Error::NetworkAddrBlocked`] if every resolved address
/// is blocked, or propagates the resolver's own error if resolution
/// fails. The returned `SocketAddr` carries a vetted IP literal, so the
/// subsequent `TcpStream::connect` cannot be re-pointed by a second DNS
/// lookup (it connects to the IP we vetted, not the name).
fn resolve_allowed_addr(addr: &str) -> io::Result<std::net::SocketAddr> {
    // Zero resolved addresses and "all resolved addresses blocked" both
    // mean there is no safe address to connect to — same error either way.
    addr.to_socket_addrs()?
        .find(|sa| !is_blocked_ip(sa.ip()))
        .ok_or_else(|| {
            crate::error::Error::NetworkAddrBlocked {
                addr: addr.to_string(),
            }
            .into()
        })
}

enum Mode {
    Write {
        writer: BufWriter<TcpStream>,
        header_written: bool,
    },
    Read {
        reader: BufReader<TcpStream>,
    },
}

/// TCP network stream for distributed rip/remux.
pub struct NetworkStream {
    disc_title: DiscTitle,
    mode: Mode,
}

impl NetworkStream {
    /// Connect to a remote listener for writing.
    /// Sends FMKV metadata header on first write.
    pub fn connect(addr: &str) -> io::Result<Self> {
        Self::connect_vetted(addr, true)
    }

    /// `connect` with an explicit SSRF-vetting toggle.
    ///
    /// `vet=true` (the public [`connect`](Self::connect) path) resolves
    /// the target and refuses any loopback/private/link-local/multicast
    /// address, closing the DNS-rebinding TOCTOU. `vet=false` exists only
    /// for in-crate tests, which must connect to `127.0.0.1` ephemeral
    /// listeners that the production vet would (correctly) reject.
    fn connect_vetted(addr: &str, vet: bool) -> io::Result<Self> {
        // Resolve + vet the target before connecting. Connect to the
        // vetted IP literal (not the raw name) so a DNS rebind between
        // settings-save validation and now can't redirect us to a
        // loopback/private/link-local host (SSRF).
        let stream = if vet {
            let vetted = resolve_allowed_addr(addr)?;
            TcpStream::connect(vetted)?
        } else {
            TcpStream::connect(addr)?
        };
        // The sender is the latency-sensitive side; set nodelay here too
        // (the listen side already does) so the final sub-MSS flush after
        // finish() isn't held by Nagle. The 256 KB BufWriter coalesces
        // bulk writes, so this only affects the tail.
        stream.set_nodelay(true)?;
        Ok(Self {
            disc_title: DiscTitle::empty(),
            mode: Mode::Write {
                writer: BufWriter::with_capacity(NET_BUF_SIZE, stream),
                header_written: false,
            },
        })
    }

    /// Set stream metadata (write side only). Returns self for chaining.
    ///
    /// Only meaningful on a [`connect`](Self::connect)-constructed
    /// (write) stream — the title is sent in the FMKV header on first
    /// write. On a [`listen`](Self::listen)-constructed (read) stream
    /// the stored title is immediately overwritten by the header read in
    /// `listen()`, so calling `meta()` there is a silent no-op.
    pub fn meta(mut self, dt: &DiscTitle) -> Self {
        self.disc_title = dt.clone();
        self
    }

    /// Listen for an incoming connection and read from it.
    /// Extracts FMKV metadata header from the sender.
    ///
    /// Accepts exactly one connection; the listening socket is dropped after
    /// `accept`, so the bound port is freed and any subsequent connection
    /// attempt to the same address is refused.
    pub fn listen(addr: &str) -> io::Result<Self> {
        Self::accept_from(TcpListener::bind(addr)?)
    }

    /// Accept one connection from an already-bound listener and read from it.
    /// Lets a caller bind first (learning the actual port for an ephemeral
    /// `:0` bind) and hand the listener in, closing the bind/drop/re-bind race
    /// that `listen(addr)` would otherwise have.
    pub fn accept_from(listener: TcpListener) -> io::Result<Self> {
        let (stream, _peer) = listener.accept()?;
        stream.set_nodelay(true)?;
        let mut reader = BufReader::with_capacity(NET_BUF_SIZE, stream);

        // Read FMKV metadata header
        let disc_title = meta::read_header(&mut reader)?
            .ok_or_else(|| -> io::Error { crate::error::Error::NoMetadata.into() })?
            .to_title();

        Ok(Self {
            disc_title,
            mode: Mode::Read { reader },
        })
    }
}

/// Write the FMKV metadata header exactly once, before any frames. Always
/// writes (even when the title has no streams) so the receiver's
/// `read_header()` always finds the magic and never falls into the
/// NoMetadata path on a zero-frame stream.
fn ensure_header_written(
    writer: &mut BufWriter<TcpStream>,
    header_written: &mut bool,
    disc_title: &DiscTitle,
) -> io::Result<()> {
    if !*header_written {
        let m = meta::M2tsMeta::from_title(disc_title);
        meta::write_header(writer, &m)?;
        *header_written = true;
    }
    Ok(())
}

impl crate::pes::Stream for NetworkStream {
    fn read(&mut self) -> io::Result<Option<crate::pes::PesFrame>> {
        match &mut self.mode {
            Mode::Read { reader } => crate::pes::PesFrame::deserialize(reader),
            _ => Err(crate::error::Error::StreamWriteOnly.into()),
        }
    }
    fn write(&mut self, frame: &crate::pes::PesFrame) -> io::Result<()> {
        match &mut self.mode {
            Mode::Write {
                writer,
                header_written,
            } => {
                ensure_header_written(writer, header_written, &self.disc_title)?;
                frame.serialize(writer)
            }
            _ => Err(crate::error::Error::StreamReadOnly.into()),
        }
    }
    fn finish(&mut self) -> io::Result<()> {
        if let Mode::Write {
            writer,
            header_written,
        } = &mut self.mode
        {
            // Always emit the FMKV header before shutdown, even for a
            // zero-frame stream (e.g. a title that produced no PES frames).
            // Without it the receiver's read_header() sees a clean EOF and
            // rejects the stream with NoMetadata.
            ensure_header_written(writer, header_written, &self.disc_title)?;
            writer.flush()?;
            writer.get_ref().shutdown(std::net::Shutdown::Write)?;
        }
        Ok(())
    }
    fn info(&self) -> &DiscTitle {
        &self.disc_title
    }
}

// NetworkStream is PES-only — no IOStream/Read/Write byte interface.

#[cfg(test)]
mod tests {
    use super::*;
    use crate::disc::{
        AudioChannels, AudioStream, Codec, ColorSpace, ContentFormat, FrameRate, HdrFormat,
        Resolution, SampleRate, Stream, VideoStream,
    };
    use std::net::TcpListener;

    /// SSRF guard: every loopback / private / link-local / multicast /
    /// unspecified address (v4 and v6) must be rejected, and ordinary
    /// public addresses must be allowed. This is what closes the
    /// DNS-rebinding window in `NetworkStream::connect`.
    #[test]
    fn is_blocked_ip_rejects_internal_targets() {
        use std::net::{Ipv4Addr, Ipv6Addr};
        // Built from octets (not string literals) so the repo's internal-infra
        // secret scanner doesn't flag the RFC1918 addresses.
        let v4 = |a, b, c, d| IpAddr::V4(Ipv4Addr::new(a, b, c, d));
        let blocked: &[(IpAddr, &str)] = &[
            (v4(127, 0, 0, 1), "loopback"),
            (v4(127, 10, 20, 30), "loopback /8"),
            (v4(10, 0, 0, 1), "private 10/8"),
            (v4(172, 16, 5, 5), "private 172.16/12"),
            (v4(192, 168, 1, 1), "private 192.168/16"),
            (v4(169, 254, 10, 10), "link-local"),
            (v4(0, 0, 0, 0), "unspecified"),
            (v4(224, 0, 0, 1), "multicast"),
            (v4(255, 255, 255, 255), "broadcast"),
            (IpAddr::V6(Ipv6Addr::LOCALHOST), "loopback v6"),
            (IpAddr::V6(Ipv6Addr::UNSPECIFIED), "unspecified v6"),
            (
                IpAddr::V6(Ipv6Addr::new(0xfc00, 0, 0, 0, 0, 0, 0, 1)),
                "ULA",
            ),
            (
                IpAddr::V6(Ipv6Addr::new(0xfd12, 0x3456, 0, 0, 0, 0, 0, 1)),
                "ULA",
            ),
            (
                IpAddr::V6(Ipv6Addr::new(0xfe80, 0, 0, 0, 0, 0, 0, 1)),
                "link-local v6",
            ),
            (
                IpAddr::V6(Ipv6Addr::new(0xff02, 0, 0, 0, 0, 0, 0, 1)),
                "multicast v6",
            ),
            // CGNAT / 0.0.0.0/8 / Class E (finding 8).
            (v4(100, 64, 0, 1), "carrier-grade NAT"),
            (v4(100, 127, 255, 254), "carrier-grade NAT edge"),
            (v4(0, 1, 2, 3), "0.0.0.0/8"),
            (v4(240, 0, 0, 1), "Class E"),
            (v4(255, 0, 0, 1), "Class E high"),
            // IPv4-mapped / -compatible IPv6 bypass (finding 7).
            (
                IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0xffff, 0x0a00, 0x0001)),
                "IPv4-mapped RFC1918 (::ffff:0a00:0001)",
            ),
            (
                IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0xffff, 0x7f00, 0x0001)),
                "::ffff:127.0.0.1 mapped",
            ),
            (
                IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0, 0x7f00, 0x0001)),
                "::127.0.0.1 compatible",
            ),
        ];
        for (ip, label) in blocked {
            assert!(is_blocked_ip(*ip), "{label} ({ip}) must be blocked");
        }

        let allowed: &[(IpAddr, &str)] = &[
            (v4(8, 8, 8, 8), "public dns"),
            (v4(1, 1, 1, 1), "public dns"),
            (v4(93, 184, 216, 34), "example.com"),
            (
                IpAddr::V6(Ipv6Addr::new(0x2606, 0x2800, 0x220, 1, 0, 0, 0, 1)),
                "public v6",
            ),
            (
                IpAddr::V6(Ipv6Addr::new(0, 0, 0, 0, 0, 0xffff, 0x0808, 0x0808)),
                "::ffff:8.8.8.8 public mapped",
            ),
        ];
        for (ip, label) in allowed {
            assert!(!is_blocked_ip(*ip), "{label} ({ip}) must be allowed");
        }
    }

    /// The public `connect` must refuse a loopback target with the typed
    /// `NetworkAddrBlocked` error (PermissionDenied) rather than attempting
    /// the TCP connect — this is the rebinding TOCTOU close at the connect.
    #[test]
    fn connect_refuses_blocked_loopback_target() {
        // Bind a real loopback listener so a non-vetting connect WOULD
        // succeed; the vetting connect must still refuse it.
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        let err = match NetworkStream::connect(&addr.to_string()) {
            Ok(_) => panic!("loopback target must be refused by the SSRF guard"),
            Err(e) => e,
        };
        assert_eq!(err.kind(), io::ErrorKind::PermissionDenied);
    }

    fn sample_title() -> DiscTitle {
        DiscTitle {
            playlist: "NetworkTest".into(),
            playlist_id: 1,
            duration_secs: 3600.0,
            size_bytes: 0,
            clips: Vec::new(),
            streams: vec![
                Stream::Video(VideoStream {
                    pid: 0x1011,
                    codec: Codec::Hevc,
                    resolution: Resolution::R2160p,
                    frame_rate: FrameRate::F23_976,
                    hdr: HdrFormat::Hdr10,
                    color_space: ColorSpace::Bt2020,
                    display_aspect: None,
                    secondary: false,
                    label: "Main".into(),
                    measured_cicp: None,
                }),
                Stream::Audio(AudioStream {
                    pid: 0x1100,
                    codec: Codec::TrueHd,
                    channels: AudioChannels::Surround71,
                    language: "eng".into(),
                    sample_rate: SampleRate::S48,
                    secondary: false,
                    purpose: crate::disc::LabelPurpose::Normal,
                    label: "English".into(),
                }),
            ],
            chapters: Vec::new(),
            extents: Vec::new(),
            content_format: ContentFormat::BdTs,
            codec_privates: Vec::new(),
        }
    }

    #[test]
    fn network_pes_roundtrip() {
        use crate::pes;
        use std::sync::mpsc;

        // The listener thread owns the bound socket and reports its actual
        // local address back over a channel before accept(). The main thread
        // connects only after receiving the address — no bind/drop/re-bind
        // window, no sleep-as-synchronisation.
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        let (addr_tx, addr_rx) = mpsc::channel();

        let handle = std::thread::spawn(move || {
            addr_tx.send(addr).unwrap();
            let mut ns = NetworkStream::accept_from(listener).unwrap();
            let info = pes::Stream::info(&ns).clone();
            let mut frames = Vec::new();
            while let Ok(Some(f)) = pes::Stream::read(&mut ns) {
                frames.push(f);
            }
            (info, frames)
        });

        let addr = addr_rx.recv().unwrap();
        let dt = sample_title();
        let mut writer = NetworkStream::connect_vetted(&addr.to_string(), false)
            .unwrap()
            .meta(&dt);
        let frame = pes::PesFrame {
            coding: None,
            source: None,
            track: 0,
            pts: 90000,
            keyframe: true,
            data: vec![0x47; 192],
            duration_ns: None,
        };
        pes::Stream::write(&mut writer, &frame).unwrap();
        pes::Stream::finish(&mut writer).unwrap();

        let (info, frames) = handle.join().unwrap();
        assert_eq!(info.playlist, "NetworkTest");
        assert_eq!(info.streams.len(), 2);
        assert_eq!(frames.len(), 1);
        assert_eq!(frames[0].track, 0);
        assert_eq!(frames[0].pts, 90000);
    }

    #[test]
    fn network_zero_frame_finish_still_sends_header() {
        use crate::pes;
        use std::sync::mpsc;

        // A title that produces no PES frames must still send the FMKV header
        // on finish(), so the receiver gets the metadata instead of rejecting
        // the stream with NoMetadata on a clean EOF.
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        let (addr_tx, addr_rx) = mpsc::channel();

        let handle = std::thread::spawn(move || {
            addr_tx.send(addr).unwrap();
            // listen()'s read_header must succeed (header present), not error.
            let ns = NetworkStream::accept_from(listener).unwrap();
            pes::Stream::info(&ns).playlist.clone()
        });

        let addr = addr_rx.recv().unwrap();
        let dt = sample_title();
        let mut writer = NetworkStream::connect_vetted(&addr.to_string(), false)
            .unwrap()
            .meta(&dt);
        // No write() at all — straight to finish().
        pes::Stream::finish(&mut writer).unwrap();

        let playlist = handle.join().unwrap();
        assert_eq!(
            playlist, "NetworkTest",
            "zero-frame finish() must still deliver the metadata header"
        );
    }

    #[test]
    fn network_empty_addr_errors() {
        let result = NetworkStream::connect("");
        assert!(result.is_err());
    }

    #[test]
    fn network_no_port_errors() {
        let result = NetworkStream::connect("127.0.0.1");
        assert!(result.is_err());
    }

    /// Spawn an accepting reader and return (its address, join handle that
    /// yields all frames read after the FMKV header).
    fn spawn_reader() -> (
        std::net::SocketAddr,
        std::thread::JoinHandle<(DiscTitle, Vec<crate::pes::PesFrame>)>,
    ) {
        use crate::pes;
        // Bind BEFORE spawning so the port is live when connect() runs — no
        // channel handshake needed (the listener already owns the socket).
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        let handle = std::thread::spawn(move || {
            let mut ns = NetworkStream::accept_from(listener).unwrap();
            let info = pes::Stream::info(&ns).clone();
            let mut frames = Vec::new();
            while let Ok(Some(f)) = pes::Stream::read(&mut ns) {
                frames.push(f);
            }
            (info, frames)
        });
        (addr, handle)
    }

    /// write() on a listen()/accept-constructed (READ) stream must return
    /// StreamReadOnly — the read side has no writer. (Returning Ok would let
    /// a caller silently lose frames written into a receive-only socket.)
    #[test]
    fn write_on_read_side_is_read_only_error() {
        use crate::pes;
        let (addr, handle) = spawn_reader();

        // Sender connects, sends header (zero frames), finishes — so the
        // reader's accept_from() returns. We test the reader's write guard.
        let dt = sample_title();
        let mut writer = NetworkStream::connect_vetted(&addr.to_string(), false)
            .unwrap()
            .meta(&dt);
        pes::Stream::finish(&mut writer).unwrap();
        let (_info, _frames) = handle.join().unwrap();

        // Now build a fresh read-side stream and confirm its write() errors.
        // (Re-bind, accept once, then immediately try to write to it.)
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr2 = listener.local_addr().unwrap();
        let h = std::thread::spawn(move || {
            let mut ns = NetworkStream::accept_from(listener).unwrap();
            let frame = pes::PesFrame {
                coding: None,
                source: None,
                track: 0,
                pts: 0,
                keyframe: true,
                data: vec![0u8; 8],
                duration_ns: None,
            };
            // Read side: writing must be a typed StreamReadOnly error.
            let err = pes::Stream::write(&mut ns, &frame).expect_err("read side write must error");
            err.kind()
        });
        // Drive the accept: connect + send header so accept_from completes.
        let mut w2 = NetworkStream::connect_vetted(&addr2.to_string(), false)
            .unwrap()
            .meta(&dt);
        pes::Stream::finish(&mut w2).unwrap();
        let kind = h.join().unwrap();
        // E_STREAM_READ_ONLY (9000) maps to Unsupported.
        assert_eq!(kind, io::ErrorKind::Unsupported);
    }

    /// read() on a connect()-constructed (WRITE) stream must return
    /// StreamWriteOnly — never Ok(None), which a caller would read as a
    /// legitimately empty stream.
    #[test]
    fn read_on_write_side_is_write_only_error() {
        use crate::pes;
        let (addr, handle) = spawn_reader();
        let dt = sample_title();
        let mut writer = NetworkStream::connect_vetted(&addr.to_string(), false)
            .unwrap()
            .meta(&dt);
        let err = pes::Stream::read(&mut writer).expect_err("write side read must error");
        // E_STREAM_WRITE_ONLY (9001) maps to Unsupported.
        assert_eq!(err.kind(), io::ErrorKind::Unsupported);
        pes::Stream::finish(&mut writer).unwrap();
        let _ = handle.join().unwrap();
    }

    /// The FMKV header must be written exactly once, before the first frame,
    /// even across many frames. The receiver must therefore reconstruct the
    /// title exactly once and read every frame after it — a header re-emitted
    /// between frames would desync PesFrame::deserialize and corrupt frame N.
    #[test]
    fn header_written_once_then_all_frames_roundtrip() {
        use crate::pes;
        let (addr, handle) = spawn_reader();
        let dt = sample_title();
        let mut writer = NetworkStream::connect_vetted(&addr.to_string(), false)
            .unwrap()
            .meta(&dt);
        for i in 0..5u8 {
            let frame = pes::PesFrame {
                coding: None,
                source: None,
                track: (i % 2) as usize,
                pts: i as i64 * 90_000,
                keyframe: i == 0,
                data: vec![i; 100 + i as usize],
                duration_ns: None,
            };
            pes::Stream::write(&mut writer, &frame).unwrap();
        }
        pes::Stream::finish(&mut writer).unwrap();
        let (info, frames) = handle.join().unwrap();
        // Title parsed once and intact.
        assert_eq!(info.streams.len(), 2);
        // Every frame survived in order with exact payloads — no desync from
        // a duplicated header.
        assert_eq!(frames.len(), 5);
        for (i, f) in frames.iter().enumerate() {
            assert_eq!(f.pts, i as i64 * 90_000, "frame {i} pts");
            assert_eq!(f.data.len(), 100 + i, "frame {i} payload length");
            assert!(
                f.data.iter().all(|&b| b == i as u8),
                "frame {i} payload bytes"
            );
        }
    }

    /// The receiver's title comes strictly from the SENDER's FMKV header:
    /// the sender's meta() title is what accept_from() reconstructs, proving
    /// the metadata flows sender→receiver over the header (not from the
    /// receiver's empty default). Distinct sender title confirms the source.
    #[test]
    fn receiver_title_comes_from_sender_header() {
        use crate::pes;
        let (addr, handle) = spawn_reader();
        let mut dt = sample_title();
        dt.playlist = "SenderControlled".into();
        dt.playlist_id = 42;
        let mut writer = NetworkStream::connect_vetted(&addr.to_string(), false)
            .unwrap()
            .meta(&dt);
        pes::Stream::finish(&mut writer).unwrap();
        let (info, _frames) = handle.join().unwrap();
        // The receiver default title is empty (playlist ""); it must have
        // been replaced by the sender's header-carried title.
        assert_eq!(info.playlist, "SenderControlled");
        assert_eq!(
            info.streams.len(),
            2,
            "stream descriptors round-trip via header"
        );
    }

    /// accept_from() must reject a connection whose first bytes are NOT the
    /// FMKV magic — there is no metadata to drive muxing, so it surfaces
    /// NoMetadata rather than proceeding with an empty/garbage title.
    #[test]
    fn accept_from_rejects_stream_without_fmkv_header() {
        use std::io::Write as _;
        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
        let addr = listener.local_addr().unwrap();
        let handle = std::thread::spawn(move || {
            // Raw non-FMKV bytes (not starting with 'F') then close.
            let mut s = TcpStream::connect(addr).unwrap();
            s.write_all(&[0x47u8; 64]).unwrap(); // TS sync bytes, no FMKV magic
            s.shutdown(std::net::Shutdown::Both).unwrap();
        });
        let err = match NetworkStream::accept_from(listener) {
            Ok(_) => panic!("missing FMKV header must error, not silently accept"),
            Err(e) => e,
        };
        // E_NO_METADATA (9008) maps to InvalidInput.
        assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
        handle.join().unwrap();
    }
}