Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion cli/ffrwd/compiler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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=[
Expand Down
99 changes: 99 additions & 0 deletions cli/ffrwd/leaky.py
Original file line number Diff line number Diff line change
@@ -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)
64 changes: 47 additions & 17 deletions cli/ffrwd/processes.py
Original file line number Diff line number Diff line change
Expand Up @@ -4440,27 +4440,50 @@ 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):
return reader
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.
Expand All @@ -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:
Expand Down
7 changes: 2 additions & 5 deletions cli/ffrwd/timing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
104 changes: 98 additions & 6 deletions cli/tests/test_node_world.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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(
Expand All @@ -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(
Expand All @@ -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)


Expand Down
8 changes: 8 additions & 0 deletions docs/dialect.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading