diff --git a/cli/ffrwd/compiler.py b/cli/ffrwd/compiler.py index 99b1eb6..8399b8d 100644 --- a/cli/ffrwd/compiler.py +++ b/cli/ffrwd/compiler.py @@ -52,6 +52,7 @@ from .functions import WasmFunction, package_modules, script_definitions from .inputs import declared_probe, forces_demuxer, probe_options, render_options from .ir import LEAKY, Graph, Lateral +from .leaky import decode_ahead_of_leaky from .lower import ProbePath, input_option_values, lower_commands, lower_table from .parser import Resolved, parse, resolve from .probe import ProbeFailure, ProbeResult @@ -698,7 +699,9 @@ def compile_all( probe_failures=probe_failures, shapes=shapes_module.ShapeCache(shape), ) - ready = [insert_splits(insert_pts_resets(graph)) for graph in graphs] + ready = [ + insert_splits(insert_pts_resets(decode_ahead_of_leaky(graph))) for graph in graphs + ] ready[0] = replace( ready[0], laterals=[ diff --git a/cli/ffrwd/leaky.py b/cli/ffrwd/leaky.py new file mode 100644 index 0000000..0e173aa --- /dev/null +++ b/cli/ffrwd/leaky.py @@ -0,0 +1,99 @@ +"""The form a leaky takes: over coded packets, or over decoded pictures. + +``ffrwd.leaky`` over a coded picture drops whole groups, from a late packet +to the next keyframe, which is right only where everything after it hands +the packets on as they are (a publish, a file the stream is copied into). +Where anything decodes the picture after it, dropping single pictures is +cheaper on the viewer: a stall costs a picture or two rather than the wait +for a keyframe. So a leaky reading packets keeps them only when every reader +of its output takes packets; otherwise a decode goes ahead of it, and it +leaks over the pictures. + +Runs before :func:`~ffrwd.split.insert_splits` (``compiler.py``), so a +leaky's readers are its consumers themselves. + +Pure: returns a new Graph, never mutates `g`. +""" + +from __future__ import annotations + +from dataclasses import replace + +from .ir import LEAKY, FrameRef, Graph, Node +from .shapes import node_shape + +# The filter that stands for the decode: it hands on every picture as it is, +# so the ffmpeg running it decodes the packets and nothing else. +DECODE_FILTER = "null" + + +def _coded(g: Graph, ref: FrameRef) -> bool: + """Whether `ref` is a coded picture: what a node writes as packets, or + what a packet filter hands back.""" + name, _, pad = ref.partition(":") + if name in g.packet_filters: + return True + raw = g.node_shapes.get(name) + node = g.nodes.get(name) + if raw is None or node is None: + return False + outputs = node_shape(node.filter, raw).outputs + index = int(pad) if pad.isdigit() else 0 + return index < len(outputs) and outputs[index].kind == "packets" + + +def _takes_packets(g: Graph, reader: Node, ref: FrameRef) -> bool: + """Whether `reader` takes `ref` as packets: a node port reading packets, + or a packet sink or filter.""" + if reader.id in g.packet_sinks or reader.id in g.packet_filters: + return True + raw = g.node_shapes.get(reader.id) + if raw is None or ref not in reader.inputs: + return False + position = reader.inputs.index(ref) + if position >= len(reader.ports): + return False + port = node_shape(reader.filter, raw).input(reader.ports[position]) + return port is not None and port.kind == "packets" + + +def _readers_take_packets(g: Graph, ref: FrameRef) -> bool: + """Whether everything reading `ref` takes packets: nodes, and files that + copy the stream rather than encode it.""" + for node in g.nodes.values(): + if ref in node.inputs and not _takes_packets(g, node, ref): + return False + for sink in g.sinks: + codec = sink.options.get("video_codec") + if any(output.ref == ref for output in sink.outputs) and codec not in (None, "copy"): + return False + return True + + +def decode_ahead_of_leaky(g: Graph) -> Graph: + """A decode in front of each leaky reading packets whose output something + decodes anyway.""" + decodes: dict[str, Node] = {} + for node in g.nodes.values(): + if node.filter != LEAKY or not node.inputs: + continue + read = node.inputs[0] + if not _coded(g, read) or _readers_take_packets(g, node.id): + continue + decodes[node.id] = Node( + id=f"{node.id}_decode", + filter=DECODE_FILTER, + args={}, + inputs=[read], + outputs=["video"], + ) + if not decodes: + return g + nodes: dict[str, Node] = {} + for name, node in g.nodes.items(): + decode = decodes.get(name) + if decode is not None: + nodes[decode.id] = decode + node = replace(node, inputs=[decode.id, *node.inputs[1:]]) + nodes[name] = node + return replace(g, nodes=nodes) diff --git a/cli/ffrwd/processes.py b/cli/ffrwd/processes.py index b513717..2949459 100644 --- a/cli/ffrwd/processes.py +++ b/cli/ffrwd/processes.py @@ -4440,19 +4440,41 @@ def _behind_split(self, target: str, reader: str | None) -> list[str]: if any(_ref_node(read) == reader for read in self.g.nodes[name].inputs) ] + def _behind_passing(self, target: str, reader: str | None) -> list[str]: + """The members of `target` reading what `reader` hands on, past the + region's splits and leakies, which hand on the pictures they read; + empty where `reader` is neither.""" + if reader is None or self.g.nodes[reader].filter not in (*SPLIT_FILTERS, LEAKY): + return [] + inside = self.members.get(target, []) + found: list[str] = [] + ahead = [reader] + while ahead: + current = ahead.pop(0) + for name in inside: + if name in found or not any( + _ref_node(read) == current for read in self.g.nodes[name].inputs + ): + continue + if self.g.nodes[name].filter in (*SPLIT_FILTERS, LEAKY): + ahead.append(name) + else: + found.append(name) + return found + def _wire_reader(self, target: str, reader: str | None) -> str | None: """The node whose format an edge ending at `reader` carries. - Behind a region's split, the first module that names a wire format, - where every module naming one names the same: the network opens each - of them on the stream the edge brings. Where they differ, the split - itself, and the edge carries the default, as it did before a split's - modules were looked past at all. A data filter's clock pad names - none and takes what arrives. + Behind a region's splits and leakies, the first module that names a + wire format, where every module naming one names the same: the + network opens each of them on the stream the edge brings, which + reaches them as it came in. Where they differ, `reader` itself, and + the edge carries the default. A data filter's clock pad names none + and takes what arrives. """ named = [ (name, wire) - for name in self._behind_split(target, reader) + for name in self._behind_passing(target, reader) if (wire := self._named_wire(name, reader)) is not None ] if not named or any(wire != named[0][1] for _, wire in named): @@ -4460,7 +4482,8 @@ def _wire_reader(self, target: str, reader: str | None) -> str | None: return named[0][0] def _named_wire(self, name: str, split: str | None) -> object: - """The wire module `name` names for what it reads off `split`, or None. + """The wire module `name` names for what it reads off `split`, a + split or a leaky, or None. A node's is what its port accepts; an older module's, the pixel format and pcm it declared. @@ -4476,16 +4499,23 @@ def _named_wire(self, name: str, split: str | None) -> object: return None if wire == (None, None) else wire def _read_position(self, node: Node, ref: FrameRef) -> int | None: - """Where `node` reads `ref`: itself, or a pad of a split of it.""" + """Where `node` reads `ref`: itself, or what splits and leakies in + front of it hand on of it.""" for position, read in enumerate(node.inputs): - producer = _ref_node(read) - if read == ref or ( - producer is not None - and producer in self.g.nodes - and self.g.nodes[producer].filter in SPLIT_FILTERS - and self.g.nodes[producer].inputs[0] == ref - ): - return position + seen: set[str] = set() + while True: + if read == ref: + return position + producer = _ref_node(read) + if ( + producer is None + or producer in seen + or producer not in self.g.nodes + or self.g.nodes[producer].filter not in (*SPLIT_FILTERS, LEAKY) + ): + break + seen.add(producer) + read = self.g.nodes[producer].inputs[0] return None def _carries_annotations(self, ref: FrameRef, consumer: str | None) -> bool: diff --git a/cli/ffrwd/timing.py b/cli/ffrwd/timing.py index 6538b7d..cdd88ae 100644 --- a/cli/ffrwd/timing.py +++ b/cli/ffrwd/timing.py @@ -323,13 +323,10 @@ def rate(self, ref: FrameRef) -> Fraction | None: return stream_rate(self.graph, self.probes, self.shapes.get, ref) -# Filters whose pictures or sound leave at a rate their args do not say. +# Filters whose pictures leave at a rate their args do not say. A sound's +# rate is its sample rate, which retiming, selecting or stretching keeps. _RATE_LOST = frozenset( { - "ainterleave", - "aselect", - "asetpts", - "atempo", "decimate", "framestep", "interleave", diff --git a/cli/tests/test_node_world.py b/cli/tests/test_node_world.py index 4f93b6b..86a60b7 100644 --- a/cli/tests/test_node_world.py +++ b/cli/tests/test_node_world.py @@ -972,6 +972,31 @@ def test_nodes_taking_one_picture_in_different_formats_get_a_stream_each( ) +@pytest.mark.parametrize( + ("select", "node"), [("ring(p.v, spot(p.v))", "ring"), ("burn(p.v)", "burn")], + ids=["split", "alone"], +) +def test_a_picture_reaching_a_node_through_a_leaky_crosses_in_the_format_it_takes( + monkeypatch: pytest.MonkeyPatch, select: str, node: str +) -> None: + """A leaky hands on the pictures it reads, so a region holding it and a + node reading its output in rgba is handed rgba on the way in.""" + monkeypatch.setitem(SHAPES, "spot.wasm", _taking(_spot, "video", _TIMING)) + monkeypatch.setitem(SHAPES, "ring.wasm", _taking(_reader("spots", _ROWS), "video", _RGBA)) + monkeypatch.setitem(SHAPES, "burn.wasm", _taking(_burn, "video", _RGBA)) + argv = _plan_argv( + "COPY (WITH p AS (SELECT ffrwd.leaky(f.video[1], max_lateness => 0.5) AS v " + f"FROM input('f.mp4', realtime => true) f) SELECT {select} FROM p) TO 'out.mkv'", + monkeypatch, + pix_fmt="yuv420p", + ) + (sidecar,) = [words for pid, words in argv.items() if pid.startswith("sidecar")] + chain = sidecar[sidecar.index("-filter_complex") + 1] + assert chain.startswith("[0:v]leaky=") and f"{node}[v=out0]" in chain + feeder = _feeder(argv) + assert feeder[feeder.index("-pix_fmt:0") + 1] == "rgba" + + def test_each_stream_of_one_nut_is_conformed_to_the_port_it_feeds( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -1375,6 +1400,46 @@ def test_a_node_at_a_copys_to_reads_the_select_as_a_packet_sink_did( assert encoder[encoder.index("-c:0") + 1] == "libx264" +_SUB = "CREATE FUNCTION sub(relay text) RETURNS source AS 'sub.wasm', 'sub' LANGUAGE wasm;\n" + + +def test_a_leaky_over_packets_something_decodes_after_leaks_the_decoded_pictures( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """burn reads pictures, so the decode goes ahead of the leaky, which then + drops single pictures rather than whole groups.""" + monkeypatch.setitem(SHAPES, "sub.wasm", _subscribe) + argv = _plan_argv( + _SUB + "COPY (SELECT burn(ffrwd.leaky(v.video[1])) FROM sub('r') v " + "WHERE v.height = 720) TO 'out.mkv'", + monkeypatch, + ) + (decode,) = [words for words in argv.values() if "[0:v:0]null[out0]" in words] + assert decode[decode.index("-c:0") + 1] == "rawvideo" + (leaking,) = [words for words in argv.values() if any("leaky=" in w for w in words)] + assert leaking[leaking.index("-filter_complex") + 1].startswith("[0:v]leaky=") + + +def test_a_leaky_over_packets_only_publish_reads_keeps_the_packets( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setitem(SHAPES, "sub.wasm", _subscribe) + monkeypatch.setitem(SHAPES, "publish.wasm", _publish) + monkeypatch.setitem( + _PARAMS, "publish.wasm", {"relay": {"type": "string"}, "broadcast": {"type": "string"}} + ) + argv = _plan_argv( + _SUB + _PUBLISH + "COPY (SELECT ffrwd.leaky(v.video[1]), v.audio[1] FROM sub('r') v " + "WHERE v.height = 720) TO publish('https://relay', 'b')", + monkeypatch, + describe=_reporting("publish.wasm"), + ) + assert not any("[0:v:0]null[out0]" in words for words in argv.values()) + (leaking,) = [words for words in argv.values() if any("leaky=" in w for w in words)] + assert "sub=relay=r[hd=n10]" in leaking[leaking.index("-filter_complex") + 1] + assert "[n10]leaky=" in leaking[leaking.index("-filter_complex") + 1] + + def test_a_nodes_picture_into_a_node_sink_is_encoded_and_its_data_is_not_copied( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -1505,7 +1570,7 @@ def test_a_packets_function_over_a_node_hands_back_the_stream_still_coded( def _switch(params: Mapping[str, object], bound: Sequence[str]) -> dict[str, object]: """The switch: clocked by its picture, or with no `v` bound by its sound, - 1024 samples a tick.""" + as ffrwd/switch 0.5.2's ``--shape --bound a,feed_audio`` prints it.""" hold = {"kind": "hold", "anchor": {"kind": "first-frame"}, "lead": 0.3, "port_param": "port"} if "v" in bound: return _shape( @@ -1514,23 +1579,48 @@ def _switch(params: Mapping[str, object], bound: Sequence[str]) -> dict[str, obj [_output("v", "video"), _output("a", "audio")], {"kind": "input", "port": "v"}, ) + f32 = {"sample_formats": ["f32"]} + rows = {"kind": "data", "latency": 0.0, "format": {"kind": "data", "codec": "json"}} return _shape( - [_clock("a", "audio", window=1024), _input("feed_audio", "audio", hold)], - [_output("a", "audio")], + [ + {**_clock("a", "audio", window=1024), "accepts": f32}, + { + **_input( + "feed_audio", + "audio", + {**hold, "anchor": {"kind": "tagged", "tag": "smart_timed"}, + "group": "switch", "timeout": 1.0}, + ), + "accepts": {**f32, "like": "a"}, + }, + ], + [ + {"name": "a", "kind": "audio", "latency": 0.0, + "format": {"kind": "like", "port": "a"}}, + {"name": "clock", **rows}, + {"name": "feeds", **rows}, + {"name": "feed_rows", **rows}, + ], {"kind": "input", "port": "a"}, + pure=False, + one_to_one=True, ) @pytest.mark.parametrize( "call", ["switch(f.video[1], f.audio[1], feed_audio => ad.audio[1])", - "switch(a => f.audio[1], feed_audio => ad.audio[1])"], + "switch(a => f.audio[1], feed_audio => ad.audio[1])", + "switch(a => setpts(aresample(f.audio[1], 48000), 'PTS+10/TB'), " + "feed_audio => ad.audio[1])"], + ids=["picture", "sound", "sound-retimed"], ) def test_a_node_clocked_by_its_sound_beside_the_picture_of_one_live_input_compiles( monkeypatch: pytest.MonkeyPatch, call: str ) -> None: """Its window is counted in time, as a picture's is, so the edge its - sound leaves the live reader on is bounded rather than refused.""" + sound leaves the live reader on is bounded rather than refused. A + retimed sound keeps its sample rate, the rate the window counts in.""" monkeypatch.setitem(SHAPES, "switch.wasm", _switch) monkeypatch.setitem(_PARAMS, "switch.wasm", {"port": {"type": "integer"}}) monkeypatch.setitem( @@ -1554,7 +1644,9 @@ def test_a_node_clocked_by_its_sound_beside_the_picture_of_one_live_input_compil shape=_Asked(), ) assert compiled.plan is not None - sound = [e for e in compiled.plan.stream_edges if e.ref == "src:f:a:0"] + edges = compiled.plan.stream_edges + reader = next(e.source for e in edges if e.ref == "src:f:v:0") + sound = [e for e in edges if e.source == reader and e.ref != "src:f:v:0"] assert sound and all(e.live for e in sound) diff --git a/docs/dialect.md b/docs/dialect.md index 8825ac0..ad94150 100644 --- a/docs/dialect.md +++ b/docs/dialect.md @@ -1643,6 +1643,14 @@ queue, and a sink's max-lateness. `max_lateness => COALESCE(:max_lateness, 0.5)`. A sound stream, a limit out of range or not a number, and any other argument are refused by name. +- **Packets or pictures.** Over a coded picture (a subscribe's, say) it + reads packets in decode order and drops whole groups, from a late + packet to the next keyframe that is in time. It keeps that form only + where everything reading its output takes packets: a publish, or a + file the stream is copied into. Where anything decodes the picture + after it (a node reading pictures, an ffmpeg filter, an encoder), the + plan decodes ahead of it instead, and it drops single pictures: a + stall costs a picture or two, not the wait for the next keyframe. - **Where.** Right after the live input at a head, right after the subscribe at a leaf, before anything splits. It is legal anywhere on a video lane; over a file it drops nothing unless something upstream diff --git a/sidecar/ffrwd-wasm/src/feeds.rs b/sidecar/ffrwd-wasm/src/feeds.rs index 83b98bd..a99635e 100644 --- a/sidecar/ffrwd-wasm/src/feeds.rs +++ b/sidecar/ffrwd-wasm/src/feeds.rs @@ -209,6 +209,10 @@ struct Conn { impl Conn { fn new(serial: u64, socket: TcpStream, streams: usize) -> Result { + // Windows hands back a socket accepted off a nonblocking listener + // nonblocking too, and a read on it would return at once rather than + // wait out its timeout. + socket.set_nonblocking(false)?; socket.set_read_timeout(Some(POLL))?; Ok(Conn { serial, @@ -723,3 +727,32 @@ fn drain( } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_read_on_an_idle_connection_waits_out_its_timeout() { + let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).expect("bind"); + listener.set_nonblocking(true).expect("nonblocking"); + let _feeder = TcpStream::connect(listener.local_addr().expect("address")).expect("connect"); + let until = Instant::now() + Duration::from_secs(5); + let socket = loop { + match listener.accept() { + Ok((socket, _)) => break socket, + Err(e) if is_timeout(&e) && Instant::now() < until => thread::sleep(POLL / 10), + Err(e) => panic!("accept: {e}"), + } + }; + let mut conn = Conn::new(1, socket, 1).expect("conn"); + let started = Instant::now(); + let read = conn.socket.read(&mut [0u8; 16]); + assert!(read.as_ref().is_err_and(is_timeout), "{read:?}"); + assert!( + started.elapsed() >= POLL / 2, + "returned after {:?}", + started.elapsed() + ); + } +}