NetworkStream PES-only: remove old IOStream/Read/Write interface

- Remove IOStream, Read, Write impls from NetworkStream
- PES write sends FMKV header before first frame (protocol fix)
- PES finish sends TCP shutdown for clean EOF
- Update tests to use PES roundtrip instead of byte-level
- Remove NetworkStream from open_input/open_output (use input/output)
This commit is contained in:
MattJackson
2026-04-15 04:36:47 +00:00
parent 9ae7d6b38a
commit 9a1c3cc218
2 changed files with 40 additions and 137 deletions
+34 -131
View File
@@ -1,18 +1,14 @@
//! NetworkStream — BD-TS over TCP with embedded metadata. //! NetworkStream — PES frames over TCP with embedded metadata.
//! //!
//! **Security:** Data is transmitted over plain TCP with no encryption. //! **Security:** Data is transmitted over plain TCP with no encryption.
//! Use only on trusted networks (LAN). TLS support is planned. //! Use only on trusted networks (LAN).
//! //!
//! Write side (sender): connects to a listener, sends FMKV header + BD-TS data. //! Write side (sender): connects to a listener, sends FMKV header + PES frames.
//! Read side (receiver): listens for a connection, reads FMKV header + BD-TS data. //! Read side (receiver): listens for a connection, reads FMKV header + PES frames.
//!
//! The FMKV metadata header is the same format as M2tsStream uses — so a
//! NetworkStream reader can hand off to any output stream (MKV, M2TS, etc.)
//! with full metadata (labels, languages, duration).
use super::{meta, IOStream}; use super::meta;
use crate::disc::DiscTitle; use crate::disc::DiscTitle;
use std::io::{self, BufReader, BufWriter, Read, Write}; use std::io::{self, BufReader, BufWriter, Write};
use std::net::{TcpListener, TcpStream}; use std::net::{TcpListener, TcpStream};
/// I/O buffer size for network reads/writes. /// I/O buffer size for network reads/writes.
@@ -32,7 +28,6 @@ enum Mode {
pub struct NetworkStream { pub struct NetworkStream {
disc_title: DiscTitle, disc_title: DiscTitle,
mode: Mode, mode: Mode,
finished: bool,
} }
impl NetworkStream { impl NetworkStream {
@@ -46,7 +41,6 @@ impl NetworkStream {
writer: BufWriter::with_capacity(NET_BUF_SIZE, stream), writer: BufWriter::with_capacity(NET_BUF_SIZE, stream),
header_written: false, header_written: false,
}, },
finished: false,
}) })
} }
@@ -77,7 +71,6 @@ impl NetworkStream {
Ok(Self { Ok(Self {
disc_title, disc_title,
mode: Mode::Read { reader }, mode: Mode::Read { reader },
finished: false,
}) })
} }
} }
@@ -91,44 +84,7 @@ impl crate::pes::Stream for NetworkStream {
} }
fn write(&mut self, frame: &crate::pes::PesFrame) -> io::Result<()> { fn write(&mut self, frame: &crate::pes::PesFrame) -> io::Result<()> {
match &mut self.mode { match &mut self.mode {
Mode::Write { writer, .. } => frame.serialize(writer), Mode::Write { writer, ref mut header_written, .. } => {
_ => Err(io::Error::new(io::ErrorKind::Unsupported, "network opened for reading")),
}
}
fn finish(&mut self) -> io::Result<()> {
if let Mode::Write { writer, .. } = &mut self.mode {
writer.flush()?;
}
Ok(())
}
fn info(&self) -> &DiscTitle { &self.disc_title }
}
impl IOStream for NetworkStream {
fn info(&self) -> &DiscTitle {
&self.disc_title
}
fn finish(&mut self) -> io::Result<()> {
if self.finished {
return Ok(());
}
self.finished = true;
if let Mode::Write { ref mut writer, .. } = self.mode {
writer.flush()?;
writer.get_ref().shutdown(std::net::Shutdown::Write)?;
}
Ok(())
}
}
impl Write for NetworkStream {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
match self.mode {
Mode::Write {
ref mut writer,
ref mut header_written,
} => {
if !*header_written { if !*header_written {
if !self.disc_title.streams.is_empty() { if !self.disc_title.streams.is_empty() {
let m = meta::M2tsMeta::from_title(&self.disc_title); let m = meta::M2tsMeta::from_title(&self.disc_title);
@@ -136,35 +92,22 @@ impl Write for NetworkStream {
} }
*header_written = true; *header_written = true;
} }
writer.write(buf) frame.serialize(writer)
} }
Mode::Read { .. } => Err(io::Error::new( _ => Err(io::Error::new(io::ErrorKind::Unsupported, "network opened for reading")),
io::ErrorKind::Unsupported,
"stream opened for reading",
)),
} }
} }
fn finish(&mut self) -> io::Result<()> {
fn flush(&mut self) -> io::Result<()> { if let Mode::Write { writer, .. } = &mut self.mode {
if let Mode::Write { ref mut writer, .. } = self.mode { writer.flush()?;
writer.flush() writer.get_ref().shutdown(std::net::Shutdown::Write)?;
} else { }
Ok(()) Ok(())
} }
} fn info(&self) -> &DiscTitle { &self.disc_title }
} }
impl Read for NetworkStream { // NetworkStream is PES-only — no IOStream/Read/Write byte interface.
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
match self.mode {
Mode::Read { ref mut reader } => reader.read(buf),
Mode::Write { .. } => Err(io::Error::new(
io::ErrorKind::Unsupported,
"stream opened for writing",
)),
}
}
}
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
@@ -173,10 +116,8 @@ mod tests {
AudioChannels, AudioStream, Codec, ColorSpace, ContentFormat, FrameRate, HdrFormat, AudioChannels, AudioStream, Codec, ColorSpace, ContentFormat, FrameRate, HdrFormat,
Resolution, SampleRate, Stream, VideoStream, Resolution, SampleRate, Stream, VideoStream,
}; };
use std::io::{Read, Write};
use std::net::TcpListener; use std::net::TcpListener;
/// Build a DiscTitle with streams for metadata tests.
fn sample_title() -> DiscTitle { fn sample_title() -> DiscTitle {
DiscTitle { DiscTitle {
playlist: "NetworkTest".into(), playlist: "NetworkTest".into(),
@@ -213,52 +154,9 @@ mod tests {
#[test] #[test]
#[ignore] // Requires TCP; may be flaky in CI environments #[ignore] // Requires TCP; may be flaky in CI environments
fn network_listen_connect_roundtrip() { fn network_pes_roundtrip() {
// Bind to OS-assigned port use crate::pes;
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener); // Release so NetworkStream::listen can bind
let addr = format!("127.0.0.1:{}", port);
let addr_clone = addr.clone();
// Spawn listener in a thread
let handle = std::thread::spawn(move || {
let mut ns = NetworkStream::listen(&addr_clone).unwrap();
let mut buf = vec![0u8; 4096];
let mut received = Vec::new();
loop {
match ns.read(&mut buf) {
Ok(0) => break,
Ok(n) => received.extend_from_slice(&buf[..n]),
Err(_) => break,
}
}
received
});
// Small delay to let the listener thread bind
std::thread::sleep(std::time::Duration::from_millis(50));
// Connect and write data
let dt = sample_title();
let mut writer = NetworkStream::connect(&addr).unwrap().meta(&dt);
let payload = b"Hello from the write side of the network stream!";
writer.write_all(payload).unwrap();
writer.finish().unwrap();
let received = handle.join().unwrap();
// The received data should end with our payload (after the FMKV header)
assert!(
received.windows(payload.len()).any(|w| w == payload),
"payload not found in received data (got {} bytes)",
received.len()
);
}
#[test]
#[ignore] // Requires TCP; may be flaky in CI environments
fn network_metadata_flows() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap(); let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port(); let port = listener.local_addr().unwrap().port();
drop(listener); drop(listener);
@@ -267,35 +165,40 @@ mod tests {
let addr_clone = addr.clone(); let addr_clone = addr.clone();
let handle = std::thread::spawn(move || { let handle = std::thread::spawn(move || {
let ns = NetworkStream::listen(&addr_clone).unwrap(); let mut ns = NetworkStream::listen(&addr_clone).unwrap();
let info = ns.info().clone(); let info = pes::Stream::info(&ns).clone();
info let mut frames = Vec::new();
while let Ok(Some(f)) = pes::Stream::read(&mut ns) {
frames.push(f);
}
(info, frames)
}); });
std::thread::sleep(std::time::Duration::from_millis(50)); std::thread::sleep(std::time::Duration::from_millis(50));
let dt = sample_title(); let dt = sample_title();
let mut writer = NetworkStream::connect(&addr).unwrap().meta(&dt); let mut writer = NetworkStream::connect(&addr).unwrap().meta(&dt);
// Must write at least one byte to trigger header send let frame = pes::PesFrame { track: 0, pts: 90000, keyframe: true, data: vec![0x47; 192] };
writer.write_all(&[0u8; 192]).unwrap(); pes::Stream::write(&mut writer, &frame).unwrap();
writer.finish().unwrap(); pes::Stream::finish(&mut writer).unwrap();
let info = handle.join().unwrap(); let (info, frames) = handle.join().unwrap();
assert_eq!(info.playlist, "NetworkTest"); assert_eq!(info.playlist, "NetworkTest");
assert_eq!(info.duration_secs, 3600.0);
assert_eq!(info.streams.len(), 2); assert_eq!(info.streams.len(), 2);
assert_eq!(frames.len(), 1);
assert_eq!(frames[0].track, 0);
assert_eq!(frames[0].pts, 90000);
} }
#[test] #[test]
fn network_empty_addr_errors() { fn network_empty_addr_errors() {
let result = NetworkStream::connect(""); let result = NetworkStream::connect("");
assert!(result.is_err(), "empty address should fail"); assert!(result.is_err());
} }
#[test] #[test]
fn network_no_port_errors() { fn network_no_port_errors() {
// Connecting to an address without a port should fail
let result = NetworkStream::connect("127.0.0.1"); let result = NetworkStream::connect("127.0.0.1");
assert!(result.is_err(), "address without port should fail"); assert!(result.is_err());
} }
} }
+6 -6
View File
@@ -212,9 +212,9 @@ pub fn open_input(url: &str, opts: &InputOptions) -> io::Result<Box<dyn IOStream
let reader = BufReader::with_capacity(IO_BUF_SIZE, file); let reader = BufReader::with_capacity(IO_BUF_SIZE, file);
Ok(Box::new(MkvStream::open(reader)?)) Ok(Box::new(MkvStream::open(reader)?))
} }
StreamUrl::Network { ref addr } => { StreamUrl::Network { .. } => {
validate_network_addr(addr)?; Err(io::Error::new(io::ErrorKind::Unsupported,
Ok(Box::new(NetworkStream::listen(addr)?)) "network:// requires PES pipeline — use input() instead of open_input()"))
} }
StreamUrl::Stdio => { StreamUrl::Stdio => {
Ok(Box::new(StdioStream::input())) Ok(Box::new(StdioStream::input()))
@@ -282,9 +282,9 @@ pub fn open_output(url: &str, meta: &DiscTitle) -> io::Result<Box<dyn IOStream>>
}; };
Ok(Box::new(MkvStream::new(writer).meta(meta).max_buffer(lookahead))) Ok(Box::new(MkvStream::new(writer).meta(meta).max_buffer(lookahead)))
} }
StreamUrl::Network { ref addr } => { StreamUrl::Network { .. } => {
validate_network_addr(addr)?; Err(io::Error::new(io::ErrorKind::Unsupported,
Ok(Box::new(NetworkStream::connect(addr)?.meta(meta))) "network:// output requires PES pipeline — use pipe() instead of open_output()"))
} }
StreamUrl::Unknown { ref raw } => { StreamUrl::Unknown { ref raw } => {
Err(io::Error::new(io::ErrorKind::InvalidInput, Err(io::Error::new(io::ErrorKind::InvalidInput,