diff --git a/cli/ffrwd/diagram.py b/cli/ffrwd/diagram.py index be8590c..8932b45 100644 --- a/cli/ffrwd/diagram.py +++ b/cli/ffrwd/diagram.py @@ -31,6 +31,7 @@ RowsEdge, SidecarProcess, StreamEdge, + VideoFormat, ) __all__ = ["INSTALL_HINT", "render_diagram", "render_terminal", "termaid_available"] @@ -126,7 +127,10 @@ def _rows_file_lines(process: SidecarProcess) -> list[str]: def _edge_label(edge: Edge) -> str: """What crosses this edge, as the arrow's label.""" if isinstance(edge, StreamEdge): - return f"{edge.format.container} {edge.format.codec}" + wire = edge.format + if isinstance(wire, VideoFormat) and wire.geometry is not None: + return f"{wire.container} {wire.codec}, timing, {wire.size}" + return f"{wire.container} {wire.codec}" if isinstance(edge, RowsEdge): return f"{edge.container} rows" if isinstance(edge, FeederEdge): diff --git a/cli/ffrwd/processes.py b/cli/ffrwd/processes.py index 237bff6..b513717 100644 --- a/cli/ffrwd/processes.py +++ b/cli/ffrwd/processes.py @@ -221,9 +221,10 @@ WIRE_SAMPLE_FMTS: tuple[str, ...] = ("f32", "s16") SAMPLE_FMT_CODECS: Mapping[str, str] = {"f32": "pcm_f32le", "s16": "pcm_s16le"} -# The width and height of the picture a data filter's CLOCK pad is handed: it -# reads the pts and nothing else, so each frame is made as small as a frame -# can usefully be before it crosses the pipe. +# The width and height of the picture a data filter's CLOCK pad, or a node's +# input read for its timing alone, is handed: it reads the pts and nothing +# else, so each frame is made as small as a frame can usefully be before it +# crosses the pipe. CLOCK_SIZE = 16 # The filters whose output pads are interchangeable copies of one input, and so @@ -498,6 +499,15 @@ def _read_pairs(d: Mapping[str, object], key: str) -> tuple[tuple[str, object], return tuple(_read_object(d, key).items()) +def _read_size(value: object) -> tuple[int, int] | None: + if not isinstance(value, list) or len(value) != 2: + return None + width, height = value + if not isinstance(width, int) or not isinstance(height, int): + return None + return width, height + + def _read_stream_type(value: object) -> StreamType: try: return _parse_stream_type(value) @@ -520,6 +530,10 @@ class VideoFormat: The edge into a PACKET SINK carries the encoder's output instead: `codec` is then the encoder, and `options` the validated sink options shaping it (crf, preset and kin), rendered on the producing ffmpeg's output. + + `geometry` is the size the pictures stand for where the edge carries + them smaller: a stream only timing inputs read crosses at + :data:`CLOCK_SIZE` square, and its readers are told this one. """ pix_fmt: str = DEFAULT_PIX_FMT @@ -529,6 +543,7 @@ class VideoFormat: container: str = NUT codec: str = RAWVIDEO options: tuple[tuple[str, object], ...] = () + geometry: tuple[int, int] | None = None @property def size(self) -> str | None: @@ -548,10 +563,13 @@ def to_dict(self) -> dict[str, object]: } if self.options: written["options"] = dict(self.options) + if self.geometry is not None: + written["geometry"] = {"width": self.geometry[0], "height": self.geometry[1]} return written @classmethod def from_dict(cls, d: Mapping[str, object]) -> VideoFormat: + geometry = d.get("geometry") return cls( pix_fmt=_read_text(d, "pix_fmt"), width=_read_maybe_whole(d, "width"), @@ -560,6 +578,12 @@ def from_dict(cls, d: Mapping[str, object]) -> VideoFormat: container=_read_text(d, "container", NUT), codec=_read_text(d, "codec", RAWVIDEO), options=_read_pairs(d, "options"), + geometry=None + if geometry is None + else ( + _read_whole(_read_object(d, "geometry"), "width"), + _read_whole(_read_object(d, "geometry"), "height"), + ), ) @@ -1285,6 +1309,10 @@ class SidecarProcess: # The tags the query wrote on what each ``-i`` of a node network carries, # in ``-i`` order: empty for an input it wrote none on. tags: tuple[tuple[tuple[str, str], ...], ...] = () + # The size of each picture of each ``-i`` that crosses smaller than it + # is, by its position among the input's pictures, in ``-i`` order: None + # for a picture at its own size, and empty for an input with none. + geometries: tuple[tuple[tuple[int, int] | None, ...], ...] = () @property def nodes(self) -> tuple[str, ...]: @@ -1372,6 +1400,11 @@ def to_dict(self) -> dict[str, object]: written["colors"] = [dict(one) for one in self.colors] if self.tags: written["tags"] = [dict(one) for one in self.tags] + if self.geometries: + written["geometries"] = [ + [None if size is None else list(size) for size in one] + for one in self.geometries + ] if self.network and self.graph is not None: written["graph"] = self.graph.to_dict() return written @@ -1434,6 +1467,11 @@ def from_dict(cls, d: Mapping[str, object]) -> SidecarProcess: for one in _read_list(d, "tags") if isinstance(one, dict) ), + geometries=tuple( + tuple(_read_size(size) for size in one) + for one in _read_list(d, "geometries") + if isinstance(one, list) + ), ) @@ -2071,8 +2109,12 @@ def __init__( models: Mapping[str, ModelBinding] | None = None, effects: Mapping[str, tuple[str, ...]] | None = None, anchors: Mapping[str, tuple[int, int]] | None = None, + small: Mapping[str, tuple[int, int]] | None = None, ) -> None: self.g = g + # The scales putting a picture only timing inputs read at + # CLOCK_SIZE, each with the size it stands for. + self.small = dict(small or {}) self.probes = probes or {} self.pix_fmts = pix_fmts or {} self.shapes = shapes or {} @@ -3037,7 +3079,13 @@ def _node_delay(self, name: str) -> int | None: def _node_frames(self, name: str) -> int | None: """The frames node `name` holds past its clock: its window, the waits - its interval inputs add and its outputs' latency, at its pictures' rate.""" + its interval inputs add and its outputs' latency, at its pictures' rate. + + A node reading no picture holds that time in frames of the bound: the + pictures of the input its clock comes from, or + :data:`LONGEST_FRAME_SECONDS` each where that input has none + (:meth:`_bound_frame_seconds`). + """ if self._paths is None: self._paths = paths_of(self.g, self.probes) paths = self._paths @@ -3053,7 +3101,11 @@ def _node_frames(self, name: str) -> int | None: if own <= 0: return 0 video = next((ref for ref in node.inputs if ref_type(self.g, ref) == "video"), None) - rate = paths.rate(video) if video is not None else None + if video is None: + timed = refs[0] if refs else next(iter(node.inputs), None) + seconds = self._picture_seconds(timed) if timed is not None else None + return math.ceil(own / (seconds or LONGEST_FRAME_SECONDS)) + rate = paths.rate(video) return None if rate is None else math.ceil(own * rate) def _node_delays(self, names: Sequence[str]) -> dict[str, int | None]: @@ -3667,7 +3719,7 @@ def _one_to_one(self, name: str) -> bool: return False if self.external.get(name, False): return self._shape(name).one_to_one - return self.g.nodes[name].filter in SPLIT_FILTERS + return self.g.nodes[name].filter in SPLIT_FILTERS or name in self.small def _anchor(self, ref: FrameRef) -> str: """The point `ref` is in lockstep with: back through one-to-one nodes.""" @@ -4738,6 +4790,7 @@ def _timing_wire(self, ref: FrameRef) -> StreamFormat: width=size[0] if size else None, height=size[1] if size else None, timebase=_timebase(meta.fps) if meta else None, + geometry=self.small.get(_ref_node(self._past_splits(ref)) or ""), ) def _carried_pix_fmt(self, ref: FrameRef) -> str | None: @@ -4763,6 +4816,107 @@ def _carried_pix_fmt(self, ref: FrameRef) -> str | None: meta = self._origin_meta(current) return meta.pix_fmt if meta is not None else None + def timing_reads(self) -> dict[FrameRef, tuple[list[str], tuple[int, int]]]: + """The pictures a node region reads off an ffmpeg for their timing + alone, each with the region's members and the picture's size. + + A picture qualifies where every port the region hands it to reads it + for its timing: a port reading its pixels shares the one read, which + then crosses whole. One of no known size, or one already no bigger + than :data:`CLOCK_SIZE` square, does not. + """ + found: dict[FrameRef, tuple[list[str], tuple[int, int]]] = {} + for members in self._regions(): + if not any(name in self.node_shapes for name in members): + continue + inside = set(members) + for ref, _ in self._region_reads(members): + producer = _ref_node(ref) + if ref_type(self.g, ref) != "video" or ( + producer is not None and self.external.get(producer, False) + ): + continue + size = self._picture_size(ref) + if size is None or size[0] * size[1] <= CLOCK_SIZE * CLOCK_SIZE: + continue + reads = [ + (name, position) + for name in members + if self.external[name] + for position, read in enumerate(self.g.nodes[name].inputs) + if self._served_by(read, ref, inside) + ] + if reads and all(self._timing_port(name, at) for name, at in reads): + found[ref] = (list(members), size) + return found + + def _served_by(self, read: FrameRef, ref: FrameRef, inside: Collection[str]) -> bool: + """Whether a region member's `read` is handed the region's read of + `ref`: itself, a copy bound to it, or a split inside the region.""" + seen: set[FrameRef] = set() + while read not in seen: + seen.add(read) + read = self.same_reads.get(read, read) + if read == ref: + return True + producer = _ref_node(read) + if producer is None or producer not in inside or self.external[producer]: + return False + read = self.g.nodes[producer].inputs[0] + return False + + def _timing_port(self, name: str, position: int) -> bool: + """Whether node `name` reads its `position`-th input as a picture for + its timing alone.""" + shape = self.node_shapes.get(name) + port = shape.input(self.g.nodes[name].ports[position]) if shape is not None else None + return port is not None and port.kind == "video" and port.accepts.wants == "timing" + + def shrunk( + self, reads: Mapping[FrameRef, tuple[list[str], tuple[int, int]]] + ) -> tuple[Graph, dict[str, tuple[int, int]]]: + """The graph with each of `reads` scaled to :data:`CLOCK_SIZE` square + before its region reads it, and each scale's name with the size its + picture stands for. The scale keeps the pixel format.""" + before: dict[str, list[Node]] = {} + rewired: dict[str, Node] = {} + tags = dict(self.g.stream_tags) + made: dict[str, tuple[int, int]] = {} + for ref, (members, size) in reads.items(): + readers = [ + member + for member in members + if ref in rewired.get(member, self.g.nodes[member]).inputs + ] + if not readers: + continue + name = _unique_alias(f"{ref}_small", self.g.nodes) + for member in readers: + node = rewired.get(member, self.g.nodes[member]) + rewired[member] = replace( + node, inputs=[name if read == ref else read for read in node.inputs] + ) + first = next(member for member in self.g.nodes if member in readers) + before.setdefault(first, []).append( + Node( + id=name, + filter="scale", + args={"width": CLOCK_SIZE, "height": CLOCK_SIZE}, + inputs=[ref], + outputs=["video"], + ) + ) + said = {**tags.get(self._past_splits(ref), {}), **tags.get(ref, {})} + if said: + tags[name] = said + made[name] = size + nodes: dict[str, Node] = {} + for member, node in self.g.nodes.items(): + for scale in before.get(member, []): + nodes[scale.id] = scale + nodes[member] = rewired.get(member, node) + return replace(self.g, nodes=nodes, stream_tags=tags), made + def _reads_timing(self, name: str, ref: FrameRef) -> bool: """Whether node `name` reads `ref`, or a split's copy of it, for its frames' times alone.""" @@ -5223,6 +5377,7 @@ def rewrite(ref: FrameRef) -> FrameRef: tags=self._region_tags(incoming, alias_of, read_order) if sidecar.node_network else (), + geometries=self._region_geometries(incoming, read_as, read_order), reads_rows=any(e.annotations for e in self.edges if e.target == sidecar.id), writes_rows=any(e.annotations for e in self.edges if e.source == sidecar.id), rows_modules=self._rows_modules(sidecar, members), @@ -5273,6 +5428,28 @@ def _region_colors( return () return tuple(found.get(alias, ()) for alias in order) + def _region_geometries( + self, + incoming: Sequence[StreamEdge], + read_as: Mapping[FrameRef, FrameRef], + order: Sequence[str], + ) -> tuple[tuple[tuple[int, int] | None, ...], ...]: + """The size of each picture a node network's ``-i`` carries smaller + than it is, by its position among the input's pictures, in ``-i`` + order.""" + found: dict[str, list[tuple[int, int] | None]] = {} + for edge in incoming: + wire = edge.format + if not isinstance(wire, VideoFormat) or wire.geometry is None: + continue + alias, _, index = src_parts(read_as[edge.ref]) + sizes = found.setdefault(alias, []) + sizes.extend([None] * (index + 1 - len(sizes))) + sizes[index] = wire.geometry + if not found: + return () + return tuple(tuple(found.get(alias, ())) for alias in order) + def _region_tags( self, incoming: Sequence[StreamEdge], @@ -5575,9 +5752,17 @@ def partition( # this module rather than beside it. from . import startup - plan = _Partitioner( + partitioner = _Partitioner( g, external, probes, pix_fmts, shapes, audio_wires, models, effects, anchors - ).run() + ) + reads = partitioner.timing_reads() + if reads: + shrunk, small = partitioner.shrunk(reads) + partitioner = _Partitioner( + shrunk, external, probes, pix_fmts, shapes, audio_wires, models, effects, anchors, + small, + ) + plan = partitioner.run() plan = startup.arrange(plan) startup.check(plan) return plan diff --git a/cli/ffrwd/wasm.py b/cli/ffrwd/wasm.py index 48e95ea..1d789bb 100644 --- a/cli/ffrwd/wasm.py +++ b/cli/ffrwd/wasm.py @@ -2019,7 +2019,9 @@ def _argv( ``-pad ''`` right after its own ``-i``, ``{"row": ..., "rendition": {...}}`` with absent attributes omitted -- a pad with none gets no flag. A node network's input carrying a raw picture gets its ``"color"`` there - too (:attr:`SidecarProcess.colors`), since NUT writes none. + too (:attr:`SidecarProcess.colors`), since NUT writes none, and the + ``"geometry"`` of each picture crossing smaller than it is + (:attr:`SidecarProcess.geometries`). `writes` is the mirror on the other side: one path per rows document the process writes, in document order, since a process writing several of @@ -2068,6 +2070,12 @@ def _argv( tags = process.tags[index] if index < len(process.tags) else () if tags: pad["tags"] = dict(tags) + sizes = process.geometries[index] if index < len(process.geometries) else () + if sizes: + pad["geometry"] = [ + None if size is None else {"width": size[0], "height": size[1]} + for size in sizes + ] if pad: argv += ["-pad", json.dumps(pad)] if any(grant.effect == "gpu" for grant in process.grants): diff --git a/cli/tests/test_node_world.py b/cli/tests/test_node_world.py index 0730b0d..4f93b6b 100644 --- a/cli/tests/test_node_world.py +++ b/cli/tests/test_node_world.py @@ -21,6 +21,7 @@ from ffrwd import shapes, wasm from ffrwd.compiler import compile_all +from ffrwd.diagram import render_diagram from ffrwd.errors import ErrorCode, FfrwdError from ffrwd.execute import plan_argv from ffrwd.ir import Graph @@ -1502,6 +1503,61 @@ def test_a_packets_function_over_a_node_hands_back_the_stream_still_coded( assert argv["ffmpeg1"][argv["ffmpeg1"].index("-c:0") + 1] == "copy" +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.""" + hold = {"kind": "hold", "anchor": {"kind": "first-frame"}, "lead": 0.3, "port_param": "port"} + if "v" in bound: + return _shape( + [_clock("v"), _input("a", "audio", {"kind": "lockstep"}), + _input("feed_audio", "audio", hold)], + [_output("v", "video"), _output("a", "audio")], + {"kind": "input", "port": "v"}, + ) + return _shape( + [_clock("a", "audio", window=1024), _input("feed_audio", "audio", hold)], + [_output("a", "audio")], + {"kind": "input", "port": "a"}, + ) + + +@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])"], +) +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.""" + monkeypatch.setitem(SHAPES, "switch.wasm", _switch) + monkeypatch.setitem(_PARAMS, "switch.wasm", {"port": {"type": "integer"}}) + monkeypatch.setitem( + _DECLARATIONS, + "switch", + "CREATE FUNCTION switch(v video_stream DEFAULT NULL, a audio_stream DEFAULT NULL, " + "feed_audio audio_stream DEFAULT NULL, port number DEFAULT 9000) " + "RETURNS STRUCT(v video_stream, a audio_stream) " + "AS 'switch.wasm', 'switch' LANGUAGE wasm;", + ) + probes = _probes() + monkeypatch.setattr( + "ffrwd.compiler.probe_path", lambda path, args=(), **kw: probes[path[0]] + ) + compiled = compile_all( + _declared( + f"COPY (SELECT burn(f.video[1]), ({call}).a " + "FROM input('f.mp4', realtime => true) f, input('a.mp4') ad) TO 'out.mkv'" + ), + describe=_node, + shape=_Asked(), + ) + assert compiled.plan is not None + sound = [e for e in compiled.plan.stream_edges if e.ref == "src:f:a:0"] + assert sound and all(e.live for e in sound) + + def test_a_live_query_feeding_a_node_later_than_its_bound_is_refused_at_compile( monkeypatch: pytest.MonkeyPatch, ) -> None: @@ -1639,23 +1695,31 @@ def test_a_timing_input_beside_a_reader_of_the_same_picture_binds_its_stream( ) feeder = _feeder(argv) assert feeder.count("-map") == 1 + assert "-filter_complex" not in feeder, "the one picture crosses whole" assert feeder[feeder.index("-pix_fmt:0") + 1] == "rgba" sidecar = argv["sidecar0"] assert sidecar[sidecar.index("-filter_complex") + 1] == ( "[v=0:v]spot=every=30[spots=n1];[v=0:v][boxes=n1]boxes_mask[v=out0]" ) + assert "-pad" not in sidecar or "geometry" not in sidecar[sidecar.index("-pad") + 1] -def test_a_timing_input_alone_takes_the_picture_in_the_format_it_has( +def test_a_timing_input_alone_takes_the_picture_small_in_the_format_it_has( monkeypatch: pytest.MonkeyPatch, ) -> None: + """A picture only timing inputs read crosses at 16x16, its own size said + in its ``-pad``, and the plan's diagram says so.""" monkeypatch.setitem(SHAPES, "spot.wasm", _taking(_spot, "video", _TIMING)) - feeder = _feeder(_plan_argv( - "COPY (SELECT spot(f.video[1]) FROM input('f.mp4') f) TO 'spots.ndjson'", - monkeypatch, - pix_fmt="yuv444p", - )) + query = "COPY (SELECT spot(f.video[1]) FROM input('f.mp4') f) TO 'spots.ndjson'" + argv = _plan_argv(query, monkeypatch, pix_fmt="yuv444p") + feeder = _feeder(argv) + assert feeder[feeder.index("-filter_complex") + 1] == "[0:v:0]scale=width=16:height=16[out0]" assert feeder[feeder.index("-pix_fmt:0") + 1] == "yuv444p" + (sidecar,) = [words for pid, words in argv.items() if pid.startswith("sidecar")] + pad = json.loads(sidecar[sidecar.index("-pad") + 1]) + assert pad["geometry"] == [{"width": 320, "height": 240}] + plan = _plan(query, monkeypatch, pix_fmt="yuv444p") + assert "-->|nut rawvideo, timing, 16x16|" in render_diagram([], plan) @pytest.mark.parametrize(("source", "named"), [("yuv420p", "yuv444p"), ("yuv444p", "yuv420p")]) @@ -1674,7 +1738,9 @@ def test_a_timing_input_after_a_format_takes_the_format_it_names( pix_fmt=source, ) (writer,) = [words for pid, words in argv.items() if pid.startswith("ffmpeg")] - assert writer[writer.index("-filter_complex") + 1].endswith(f"format=pix_fmts={named}[out0]") + assert writer[writer.index("-filter_complex") + 1].endswith( + f"format=pix_fmts={named},scale=width=16:height=16[out0]" + ) assert writer[writer.index("-pix_fmt:0") + 1] == named _with_boxes_mask(monkeypatch) mixed = _plan_argv( @@ -1689,6 +1755,8 @@ def test_a_timing_input_after_a_format_takes_the_format_it_names( ] assert feeder[feeder.index("-pix_fmt:0") + 1] == "rgba" assert feeder[feeder.index("-pix_fmt:1") + 1] == named + pad = json.loads(mixed["sidecar0"][mixed["sidecar0"].index("-pad") + 1]) + assert pad["geometry"] == [None, {"width": 320, "height": 240}] def test_explain_says_timing_for_an_input_read_for_its_timing( diff --git a/docs/dialect.md b/docs/dialect.md index 79d81b5..8825ac0 100644 --- a/docs/dialect.md +++ b/docs/dialect.md @@ -597,11 +597,13 @@ dest := 'path' | STDOUT | ( value-expression ) | sink(value, ...) query compiled against. A node at a COPY's TO is asked once for its ports and again with the streams the SELECT binds there. - An input the module reads for its timing alone (`wants` `timing`) - is handed the stream in the format it already has: nothing converts - or conforms it, and where another port of the region reads the same - stream, it binds that one. `boxes_mask(v, detect(v))` sends the - picture once, in the format `detect` takes. `explain` says `timing` - for it. + is handed the stream in the format it already has, scaled to 16x16 + before it leaves ffmpeg: nothing converts or conforms it, and the + node is told the picture's own size. Where another port of the + region reads the same stream, it binds that one, which crosses + whole. `boxes_mask(v, detect(v))` sends the picture once, in the + format `detect` takes. `explain` says `timing` for it, and its + diagram labels the small picture's edge `timing, 16x16`. - One region of a sidecar holds the nodes the query wires together, and everything one process hands another travels as one NUT. A signature only a node can carry (kinds mixed, a stream left out, a diff --git a/sidecar/NODE-CLI.md b/sidecar/NODE-CLI.md index b511203..3fbb2ce 100644 --- a/sidecar/NODE-CLI.md +++ b/sidecar/NODE-CLI.md @@ -36,6 +36,15 @@ after an `-i` says what its NUT does not carry: a raw picture's colour, in ffmpeg's names, and tags for its streams beside their own, which reach a node in `stream-info.tags` (a hold input anchored `tagged` reads them). + -pad '{"geometry": [null, {"width": 1280, "height": 720}]}' + +gives the size of an input's pictures where the wire carries them smaller, +by position among its video streams (the k-th is `[N:v:k]`; null keeps the +header's). A picture so given is read only by timing inputs: they are told +this size in its `video-format`, an output `like` one of them takes it, and +each frame is checked against the header's size. A picture something reads +the pixels of is refused. + ## Nodes -m = @@ -276,8 +285,9 @@ tick that held them. each frame's `pts` and `duration` and the stream's info, and no bytes: the host copies and converts nothing for it, a port feed's picture is not conformed, and `fetch` (or `same`) on its frames is a fault. Hand - it the stream in whatever raw format the source has cheapest; its - `accepts` formats are not checked. An audio timing input is re-cut by + it the stream in whatever raw format the source has cheapest, at any + size, with its own size in `-pad`'s `geometry`; its `accepts` formats + are not checked. An audio timing input is re-cut by sample count as any audio input is. - **Ordinal.** `tick.ordinal` is the tick's number in the run, from 0, counted over every instance: on every worker the same tick has the same diff --git a/sidecar/ffrwd-wasm/src/hold.rs b/sidecar/ffrwd-wasm/src/hold.rs index 6072a61..0a2b980 100644 --- a/sidecar/ffrwd-wasm/src/hold.rs +++ b/sidecar/ffrwd-wasm/src/hold.rs @@ -274,6 +274,9 @@ struct MemberQueue { bytes: usize, /// The last frame's pts, or the end of the last run of samples. last: Option, + /// The pts the last frame or run arrived with, which a run of samples + /// keeps once the ticks have taken it. + newest: Option, } impl MemberQueue { @@ -284,6 +287,7 @@ impl MemberQueue { queue: VecDeque::new(), bytes: 0, last: None, + newest: None, } } @@ -508,6 +512,7 @@ impl Group { PortKind::Audio => frame.pts.saturating_add(frame.duration.unwrap_or(0)), _ => frame.pts, }); + queue.newest = Some(frame.pts); queue.push(frame); } @@ -842,9 +847,10 @@ impl Group { return None; } } - let last = match lead.queue.back() { - Some(frame) => frame.pts, - None => feed.shown[source.lead].as_ref()?.pts, + let last = match (self.members[source.lead].kind, lead.queue.back()) { + (PortKind::Audio, _) => lead.newest?, + (_, Some(frame)) => frame.pts, + (_, None) => feed.shown[source.lead].as_ref()?.pts, }; let turn = feed .offset @@ -1484,6 +1490,39 @@ mod tests { assert_eq!(ended[0].1.ends, Some(2)); } + #[test] + fn a_sound_only_feed_keeps_its_foretold_end_through_its_last_tick() { + let grid = Grid::exact(); + for port_fed in [false, true] { + let (mut g, _) = group( + hold(Anchor::SharedClock, 0.0, None, None), + vec![audio_member(1)], + port_fed, + ); + g.arrive(0, samples(0, 1024)); + g.arrive(0, samples(1024, 1024)); + g.arrive(0, samples(2048, 512)); + g.source_close(0); + let foretold: Vec> = (0..4) + .map(|k| { + g.tick(k * 1024, Some((k + 1) * 1024), KHZ48, &grid)[0] + .feed + .as_ref() + .and_then(|f| f.ends) + }) + .collect(); + let told = foretold.iter().position(Option::is_some).expect("foretold"); + assert!( + foretold[told..3].iter().all(|ends| *ends == Some(2048)), + "from its foretelling to the tick its last sample is in ({foretold:?})" + ); + assert_eq!(foretold[3], None, "ended"); + let ended = g.take_ended(); + assert_eq!(ended.len(), 1); + assert_eq!(ended[0].1.ends, Some(2048)); + } + } + #[test] fn a_jump_in_the_source_is_a_new_feed_on_the_next_tick() { let (mut g, _) = group( diff --git a/sidecar/ffrwd-wasm/src/main.rs b/sidecar/ffrwd-wasm/src/main.rs index cb7cd6d..6652463 100644 --- a/sidecar/ffrwd-wasm/src/main.rs +++ b/sidecar/ffrwd-wasm/src/main.rs @@ -418,6 +418,17 @@ pub(crate) struct PadSpec { /// the query's `tags` column says them: `smart_timed` makes a feed timed. #[serde(default)] pub(crate) tags: std::collections::BTreeMap, + /// The size of each of the input's pictures, by its position among them, + /// where the wire carries it smaller: a picture only timing inputs read + /// crosses at a few pixels, and its readers are told this size. + #[serde(default)] + pub(crate) geometry: Vec>, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)] +pub(crate) struct PadGeometry { + pub(crate) width: u32, + pub(crate) height: u32, } #[derive(Debug, Clone, PartialEq, Deserialize)] @@ -4612,6 +4623,7 @@ mod pad_spec_tests { }, color: None, tags: Default::default(), + geometry: Vec::new(), } ); } diff --git a/sidecar/ffrwd-wasm/src/node_graph.rs b/sidecar/ffrwd-wasm/src/node_graph.rs index bb61ddd..3124f17 100644 --- a/sidecar/ffrwd-wasm/src/node_graph.rs +++ b/sidecar/ffrwd-wasm/src/node_graph.rs @@ -62,6 +62,9 @@ struct StreamDef { /// The header an input carried, kept for an output that writes the same /// stream on. header: Option, + /// The format an input's pictures cross the wire in, where `-pad` gives + /// them another size: what each frame is checked against. + wire: Option, frame_rate: Option<(u64, u64)>, /// How the command line names it, for a refusal. spelling: String, @@ -162,6 +165,7 @@ fn input_stream( decode_delay: u32::try_from(stream.decode_delay).unwrap_or(u32::MAX), latency: None, header: Some(stream.clone()), + wire: None, frame_rate: stream.frame_rate, spelling, rendition, @@ -169,6 +173,22 @@ fn input_stream( }) } +/// `def` with the size `-pad` gives its pictures, where the wire carries them +/// at another: a node reads the stream's own size, and each frame is checked +/// against the wire's. +fn sized(mut def: StreamDef, geometry: Option) -> Result { + let (Some(geometry), StreamFormat::Video(video)) = (geometry, &mut def.format) else { + return Ok(def); + }; + let wire = *video; + video.width = geometry.width; + video.height = geometry.height; + video.frame_len = crate::frame_len_for(video.pix_fmt, geometry.width, geometry.height) + .with_context(|| format!("-pad geometry of {}", def.spelling))?; + def.wire = Some(wire); + Ok(def) +} + fn class_of(format: &StreamFormat) -> StreamClass { match format { StreamFormat::Video(_) => StreamClass::Video, @@ -345,9 +365,10 @@ fn pump( running = match &def.format { StreamFormat::Video(_) | StreamFormat::Audio(_) => { let format = Format { - media: match &def.format { - StreamFormat::Video(v) => Media::Video(*v), - StreamFormat::Audio(a) => Media::Audio(*a), + media: match (&def.format, def.wire) { + (_, Some(wire)) => Media::Video(wire), + (StreamFormat::Video(v), None) => Media::Video(*v), + (StreamFormat::Audio(a), None) => Media::Audio(*a), _ => unreachable!("matched as frames"), }, time_base: def.base, @@ -436,8 +457,15 @@ fn open( let id = streams.len() as u32; match input_stream(args, input, index, stream) { Ok(def) => { - classes.entry(class_of(&def.format)).or_default().push(id); - streams.push(def); + let class = classes.entry(class_of(&def.format)).or_default(); + let geometry = args + .pads + .get(input) + .and_then(Option::as_ref) + .filter(|_| matches!(def.format, StreamFormat::Video(_))) + .and_then(|pad| pad.geometry.get(class.len()).copied().flatten()); + class.push(id); + streams.push(sized(def, geometry)?); } Err(_) if stream.class() == nut::ANNOTATION_CLASS => { streams.push(StreamDef { @@ -447,6 +475,7 @@ fn open( decode_delay: 0, latency: None, header: None, + wire: None, frame_rate: None, spelling: format!("the annotation stream of input {input}"), rendition: RenditionMeta::default(), @@ -607,6 +636,7 @@ fn open( decode_delay: 0, latency: Some(0.0), header: None, + wire: None, frame_rate: None, spelling: format!("[{label}], the rows of {}", call.module), rendition: RenditionMeta::default(), @@ -738,6 +768,21 @@ fn open( .copied() .filter(|id| !bytes_read.contains(id)) .collect(); + for id in input_ids.iter().flatten() { + let def = &streams[*id as usize]; + if let (Some(wire), StreamFormat::Video(video)) = (def.wire, &def.format) { + if !timing.contains(id) { + bail!( + "{} crosses at {}x{} and -pad gives it {}x{}, and something reads its pixels; only inputs read for their timing take a picture of another size", + def.spelling, + wire.width, + wire.height, + video.width, + video.height + ); + } + } + } let read: HashSet = consumers.keys().copied().collect(); for spec in &mut lanes { for port in spec.ports.iter_mut() { @@ -1011,6 +1056,7 @@ fn output_streams( base, latency: Some(output.latency), header, + wire: None, frame_rate, spelling: String::new(), rendition: RenditionMeta::default(), @@ -1462,6 +1508,7 @@ fn listen_for( decode_delay: 0, latency: None, header: None, + wire: None, frame_rate: None, spelling: format!( "the data on 127.0.0.1:{} for input '{}' of {name}", @@ -1549,6 +1596,7 @@ fn listen_on( decode_delay: 0, latency: None, header: None, + wire: None, frame_rate: None, spelling: format!( "the feed on 127.0.0.1:{port} for input '{}' of {name}", @@ -1900,6 +1948,7 @@ fn open_host(call: &NodeCall, pads: &[(String, u32)], defs: &[StreamDef]) -> Res decode_delay: 0, latency: Some(def.latency.unwrap_or(0.0) + latency), header: None, + wire: None, frame_rate: None, spelling: String::new(), rendition: RenditionMeta::default(), @@ -2136,4 +2185,96 @@ mod tests { ); let _ = std::fs::remove_dir_all(&dir); } + + #[test] + fn a_timing_input_is_told_the_size_its_pad_gives_and_an_output_like_it_takes_that_size() { + let Some(probe) = probe() else { + eprintln!("wasm32-wasip2 is not installed: the geometry test has no module to run"); + return; + }; + let dir = std::env::temp_dir().join(format!("ffrwd-geometry-{}", std::process::id())); + std::fs::create_dir_all(&dir).expect("a scratch directory"); + let tenths = nut::TimeBase { num: 1, den: 10 }; + let clock = nut::Stream::video("rgba", 4, 4, tenths).expect("rgba"); + let small = nut::Stream::video("yuv420p", 16, 16, tenths).expect("yuv420p"); + let mut wire = Vec::new(); + { + let mut muxer = nut::Muxer::with_streams(&mut wire, &[clock, small]).expect("headers"); + for k in 0..3i64 { + muxer.write_frame_to(0, k, &[k as u8; 64]).expect("a frame"); + muxer + .write_frame_to(1, k, &[k as u8; 16 * 16 * 3 / 2]) + .expect("a frame"); + } + } + let input = dir.join("in.nut"); + std::fs::write(&input, &wire).expect("write the input"); + let port = { + let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("a port"); + listener.local_addr().expect("its address").port() + }; + let matte = dir.join("matte.nut"); + let run = |chain: &str| { + let argv: Vec = [ + "-f", + "nut", + "-i", + &input.display().to_string(), + "-pad", + r#"{"geometry": [null, {"width": 1280, "height": 720}]}"#, + "-m", + &format!("shape_probe={probe}"), + "-filter_complex", + chain, + "-map", + "[m]", + "-f", + "nut", + &matte.display().to_string(), + ] + .iter() + .map(|a| a.to_string()) + .collect(); + let args = crate::parse_args(argv).expect("a command line"); + let crate::Modules::Network { bindings, wiring } = &args.modules else { + panic!("a network"); + }; + run_carrying(&args, bindings, wiring) + }; + + let carried = run(&format!( + "[v=0:v][size=0:v:1]shape_probe=port={port}[spots=s][matte=m]" + )) + .expect("the run"); + assert_eq!( + carried.bytes.load(Ordering::Relaxed), + 3 * 64, + "the timing input's pictures are not carried" + ); + let written = std::fs::read(&matte).expect("the matte"); + let mut demuxer = nut::Demuxer::open(written.as_slice()).expect("read the NUT headers"); + assert_eq!(demuxer.stream().pix_fmt(), Some("gray")); + assert_eq!(demuxer.stream().video_geometry(), Some((1280, 720))); + let mut buf = Vec::new(); + let mut frames = Vec::new(); + while let Some(packet) = demuxer.read_packet(&mut buf).expect("a packet") { + frames.push((packet.pts, buf.len())); + } + assert_eq!( + frames, + vec![(0, 1280 * 720), (1, 1280 * 720), (2, 1280 * 720)] + ); + + let refused = run(&format!( + "[v=0:v:1][size=0:v]shape_probe=port={port}[spots=s][matte=m]" + )) + .expect_err("a picture read for its pixels at another size is refused"); + assert!( + refused + .to_string() + .contains("crosses at 16x16 and -pad gives it 1280x720"), + "{refused:#}" + ); + let _ = std::fs::remove_dir_all(&dir); + } } diff --git a/sidecar/modules/shape-probe/src/lib.rs b/sidecar/modules/shape-probe/src/lib.rs index 71bd001..98fdff5 100644 --- a/sidecar/modules/shape-probe/src/lib.rs +++ b/sidecar/modules/shape-probe/src/lib.rs @@ -16,6 +16,8 @@ //! is handed, at its pts. //! - `mask`: `v`'s size in gray, left out when `v` is not bound. //! - `copy`: `v` itself, frame for frame. +//! - `matte`: `size`'s size in gray, black, one per frame of it; left out +//! when `size` is not bound. //! - `canvas`: a video of the size `canvas` names, only when it does. //! - `spots`: data, one message per tick, as late as twelve of `v`'s frames //! where the call says `v`'s rate and half a second where it does not. @@ -36,7 +38,7 @@ use ffrwd::av::node_types::{ Accepts, Anchor, Binding, BoundStream, Clock, Hold, InputPort, Interval, LikeInput, Message, NodeShape, OutputFormat, OutputPort, Pairing, PortKind, RowsUse, }; -use ffrwd::av::types::{Meta, Rational, VideoFormat, Wants}; +use ffrwd::av::types::{Meta, Rational, RawFrame, VideoFormat, Wants}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use serde::Deserialize; @@ -217,6 +219,17 @@ fn shape(params: &Params, bound: &[Binding]) -> NodeShape { })), )); } + if bound.iter().any(|b| b.input == "size") { + outputs.push(output( + "matte", + PortKind::Video, + Some(OutputFormat::Like(LikeInput { + port: "size".to_string(), + pixel_format: Some("gray".to_string()), + sample_format: None, + })), + )); + } if let Some(canvas) = ¶ms.canvas { outputs.push(output( "canvas", @@ -300,6 +313,10 @@ struct ShapeProbe; /// Whether the query reads `copy`, as `init` was told. static LATCHED_COPY: AtomicBool = AtomicBool::new(false); +/// The bytes of one `matte` picture, where the query reads it: `size`'s +/// width by its height, as `init` was told them. +static MATTE_LEN: AtomicU64 = AtomicU64::new(0); + /// The calls this instance has had. static CALLS: AtomicU64 = AtomicU64::new(0); @@ -344,6 +361,16 @@ impl Guest for ShapeProbe { HINT.with(|h| *h.borrow_mut() = hint); FETCH_SIZE.store(params(¶ms_text)?.fetch_size, Ordering::Relaxed); LATCHED_COPY.store(latched.iter().any(|l| l == "copy"), Ordering::Relaxed); + let matte = bound + .iter() + .find(|b| b.port == "size") + .and_then(|b| match &b.format { + Some(OutputFormat::Video(v)) => Some(u64::from(v.width) * u64::from(v.height)), + _ => None, + }) + .filter(|_| latched.iter().any(|l| l == "matte")) + .unwrap_or(0); + MATTE_LEN.store(matte, Ordering::Relaxed); Ok(()) } @@ -404,6 +431,19 @@ impl Guest for ShapeProbe { tick.fetch(id, frame.index); } } + let matte = MATTE_LEN.load(Ordering::Relaxed) as usize; + if matte > 0 { + for frame in &frames { + items.push(Emission { + port: "matte".to_string(), + payload: Payload::Frame(RawFrame { + pts: frame.pts, + duration: frame.duration, + data: vec![0; matte], + }), + }); + } + } let times: Vec = frames.iter().map(|f| f.pts.to_string()).collect(); size = format!(r#","size":[{}]"#, times.join(",")); }