From 66cbee9e8df0c5d5a79d080870dc32b581d62bd5 Mon Sep 17 00:00:00 2001 From: Jon-Carlos Rivera Date: Sat, 3 Oct 2026 01:52:57 -0700 Subject: [PATCH 1/4] feat(sidecar): a timing input is told its picture's size from -pad A picture only timing inputs read may cross smaller than it is. -pad's "geometry" gives each of an input's pictures its own size, by position among them: the node is told that size in video-format, an output like the input takes it, and each frame is checked against the wire's header. A picture something reads the pixels of is refused at another size. shape-probe gains `matte`, a gray picture like `size`, to show it. Co-Authored-By: Claude Opus 5.5 --- sidecar/NODE-CLI.md | 14 ++- sidecar/ffrwd-wasm/src/main.rs | 12 ++ sidecar/ffrwd-wasm/src/node_graph.rs | 151 ++++++++++++++++++++++++- sidecar/modules/shape-probe/src/lib.rs | 42 ++++++- 4 files changed, 211 insertions(+), 8 deletions(-) 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/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(",")); } From 32686fc9acd8f7ff97fe40941ecba3049d820917 Mon Sep 17 00:00:00 2001 From: Jon-Carlos Rivera Date: Sat, 3 Oct 2026 01:53:24 -0700 Subject: [PATCH 2/4] feat(cli): a picture only timing inputs read crosses at 16x16 A node input that wants timing reads a frame's times and the stream's size, never its pixels, yet the picture it was handed crossed whole: decoded, split, relayed and dropped by the host. Where every port a node region hands a picture to reads it for its timing, the picture is now scaled to 16x16 in the ffmpeg that writes it, in the pixel format it already has, and its own size rides the input's -pad as "geometry". A picture a port of the same region reads the pixels of still crosses whole, once, for both. The plan's diagram labels such an edge "timing, 16x16". On the demo's head this is the sell clock: 103 CPU s a minute of output to 93, the split, the relay and sell each cheaper, deals and pictures as before. Co-Authored-By: Claude Opus 5.5 --- cli/ffrwd/diagram.py | 6 +- cli/ffrwd/processes.py | 187 +++++++++++++++++++++++++++++++++-- cli/ffrwd/wasm.py | 10 +- cli/tests/test_node_world.py | 27 +++-- docs/dialect.md | 12 ++- 5 files changed, 222 insertions(+), 20 deletions(-) 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..0f011e2 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 {} @@ -3667,7 +3709,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 +4780,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 +4806,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 +5367,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 +5418,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 +5742,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..0a72a90 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 @@ -1639,23 +1640,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 +1683,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 +1700,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 From 87795a7a940890d8e0f9dc1dd6d2e46c422d13b1 Mon Sep 17 00:00:00 2001 From: Jon-Carlos Rivera Date: Sat, 3 Oct 2026 01:53:33 -0700 Subject: [PATCH 3/4] fix(cli): a node clocked by its sound counts what it holds in time A node's delay was counted at its pictures' rate, so a node reading no picture (the switch with only `a` bound, 1024 samples a tick) had no size, and over one live input whose picture goes elsewhere the plan was refused UNBOUNDED_LIVE_INPUT naming it. What it holds past its clock is now counted in frames of the bound: the pictures of the input its clock comes from, or the longest frame where that input has none, the unit the edges' bounds are sized in. Co-Authored-By: Claude Opus 5.5 --- cli/ffrwd/processes.py | 14 +++++++-- cli/tests/test_node_world.py | 55 ++++++++++++++++++++++++++++++++++++ 2 files changed, 67 insertions(+), 2 deletions(-) diff --git a/cli/ffrwd/processes.py b/cli/ffrwd/processes.py index 0f011e2..b513717 100644 --- a/cli/ffrwd/processes.py +++ b/cli/ffrwd/processes.py @@ -3079,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 @@ -3095,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]: diff --git a/cli/tests/test_node_world.py b/cli/tests/test_node_world.py index 0a72a90..4f93b6b 100644 --- a/cli/tests/test_node_world.py +++ b/cli/tests/test_node_world.py @@ -1503,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: From a68722b08a68cebd6b7c39851f685fbe71e1d25d Mon Sep 17 00:00:00 2001 From: Jon-Carlos Rivera Date: Sat, 3 Oct 2026 01:53:34 -0700 Subject: [PATCH 4/4] fix(sidecar): a sound-only feed keeps its foretold end to its last tick A hold group's end was foretold from the lead member's last queued frame, or else from the picture it last showed. A sound shows nothing, so once the ticks had taken its last run the foretold end was gone, and the feed's last tick reported none where a picture's reports it. A sound's end is now foretold from the pts its last run arrived with. Co-Authored-By: Claude Opus 5.5 --- sidecar/ffrwd-wasm/src/hold.rs | 45 +++++++++++++++++++++++++++++++--- 1 file changed, 42 insertions(+), 3 deletions(-) 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(