Spaces:
Running on Zero
Running on Zero
Commit ·
df58091
1
Parent(s): f5d0a1f
Offer the finished recording as a download
Browse filesEach link's recording now tees into a footerless FileSink beside the
live binary stream (binary_stream first, then set_sinks with both — the
call replaces the sink set). publish concatenates the four part files
into one recording; footerless parts are what make that a single valid
RRD instead of one with duelling manifests. A hidden DownloadButton
appears with the merged file when the run finishes, and each new run
hides it again.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
- fdanyone/rerun_streaming.py +31 -0
- fdanyone_app.py +59 -17
- tests/test_app_helpers.py +8 -5
fdanyone/rerun_streaming.py
CHANGED
|
@@ -283,6 +283,37 @@ def check_stop(spec: RunSpec) -> None:
|
|
| 283 |
raise RunCancelled("Stopped.")
|
| 284 |
|
| 285 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 286 |
def request_stop(spec: RunSpec | None) -> str:
|
| 287 |
"""Write the run's stop sentinel; the pipeline unwinds at its next hook."""
|
| 288 |
|
|
|
|
| 283 |
raise RunCancelled("Stopped.")
|
| 284 |
|
| 285 |
|
| 286 |
+
LINK_ORDER: tuple[str, ...] = ("begin", "prepare", "gpu", "publish")
|
| 287 |
+
"""Chain links in stream order, and therefore the part-concatenation order."""
|
| 288 |
+
|
| 289 |
+
|
| 290 |
+
def rrd_part(spec: RunSpec, name: str) -> Path:
|
| 291 |
+
"""Where one link's footerless recording part lands on disk."""
|
| 292 |
+
|
| 293 |
+
part_dir: Path = spec.data_dir / "rrd"
|
| 294 |
+
part_dir.mkdir(parents=True, exist_ok=True)
|
| 295 |
+
return part_dir / f"{name}.rrd"
|
| 296 |
+
|
| 297 |
+
|
| 298 |
+
def merge_rrd_parts(spec: RunSpec) -> Path:
|
| 299 |
+
"""Concatenate the links' parts into the one recording a visitor downloads.
|
| 300 |
+
|
| 301 |
+
Parts are written with ``write_footer=False``: footered parts concatenate
|
| 302 |
+
into a file with multiple RRD manifests, which ``rrd verify`` rejects.
|
| 303 |
+
"""
|
| 304 |
+
|
| 305 |
+
import shutil
|
| 306 |
+
|
| 307 |
+
merged: Path = spec.data_dir / f"4danyone-{spec.token[:8]}.rrd"
|
| 308 |
+
with merged.open("wb") as out:
|
| 309 |
+
for name in LINK_ORDER:
|
| 310 |
+
part: Path = spec.data_dir / "rrd" / f"{name}.rrd"
|
| 311 |
+
if part.is_file():
|
| 312 |
+
with part.open("rb") as source:
|
| 313 |
+
shutil.copyfileobj(source, out, length=1 << 20)
|
| 314 |
+
return merged
|
| 315 |
+
|
| 316 |
+
|
| 317 |
def request_stop(spec: RunSpec | None) -> str:
|
| 318 |
"""Write the run's stop sentinel; the pipeline unwinds at its next hook."""
|
| 319 |
|
fdanyone_app.py
CHANGED
|
@@ -30,10 +30,12 @@ from fdanyone.rerun_streaming import (
|
|
| 30 |
Session,
|
| 31 |
begin_phase,
|
| 32 |
generate_phase,
|
|
|
|
| 33 |
motion_phase,
|
| 34 |
new_spec,
|
| 35 |
publish_phase,
|
| 36 |
request_stop,
|
|
|
|
| 37 |
smoke_generate_phase,
|
| 38 |
smoke_motion_phase,
|
| 39 |
smoke_publish_phase,
|
|
@@ -55,8 +57,17 @@ class Link:
|
|
| 55 |
|
| 56 |
return self.stream.read()
|
| 57 |
|
|
|
|
|
|
|
| 58 |
|
| 59 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 60 |
"""Open this callback's own recording, in whatever process the callback runs in.
|
| 61 |
|
| 62 |
Called first thing in every link, GPU ones included, where "this callback"
|
|
@@ -75,12 +86,18 @@ def open_link(token: str) -> Link:
|
|
| 75 |
|
| 76 |
rr.cleanup_if_forked_child()
|
| 77 |
recording: rr.RecordingStream = rr.RecordingStream(APPLICATION_ID, recording_id=token)
|
| 78 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 79 |
|
| 80 |
|
| 81 |
def begin(
|
| 82 |
video: str | None, start_time: float, seed: int
|
| 83 |
-
) -> Iterator[tuple[RunSpec, bytes | None, str, Any]]:
|
| 84 |
"""Validate the input on CPU, open a recording, and switch to the outputs."""
|
| 85 |
|
| 86 |
if video is None:
|
|
@@ -89,17 +106,29 @@ def begin(
|
|
| 89 |
spec: RunSpec = new_spec(Path(video), float(start_time), int(seed))
|
| 90 |
except FourDAnyoneError as exc:
|
| 91 |
raise gr.Error(str(exc)) from None
|
| 92 |
-
link: Link = open_link(spec.token)
|
| 93 |
-
|
| 94 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 95 |
|
| 96 |
|
| 97 |
def prepare_cpu(spec: RunSpec) -> Iterator[tuple[bytes | None, str]]:
|
| 98 |
"""Put the source clip on the frame timeline before any GPU work starts."""
|
| 99 |
|
| 100 |
-
link: Link = open_link(spec.token)
|
| 101 |
-
|
| 102 |
-
|
|
|
|
|
|
|
|
|
|
| 103 |
|
| 104 |
|
| 105 |
@spaces.GPU(duration=GPU_DURATION)
|
|
@@ -115,7 +144,7 @@ def run_gpu(spec: RunSpec) -> Iterator[tuple[dict | None, bytes | None, str]]:
|
|
| 115 |
is the only process that may build the recording it logs through.
|
| 116 |
"""
|
| 117 |
|
| 118 |
-
link: Link = open_link(spec.token)
|
| 119 |
session: Session = Session(spec)
|
| 120 |
try:
|
| 121 |
motion: Iterator[str] = (
|
|
@@ -132,6 +161,8 @@ def run_gpu(spec: RunSpec) -> Iterator[tuple[dict | None, bytes | None, str]]:
|
|
| 132 |
)
|
| 133 |
for label in generate:
|
| 134 |
yield None, link.read(), label
|
|
|
|
|
|
|
| 135 |
except FourDAnyoneError as exc:
|
| 136 |
# The banner keeps the reason on screen after the toast dismisses.
|
| 137 |
yield None, link.read(), f"Failed: {exc}"
|
|
@@ -140,26 +171,32 @@ def run_gpu(spec: RunSpec) -> Iterator[tuple[dict | None, bytes | None, str]]:
|
|
| 140 |
message: str = f"{type(exc).__name__}: {exc}"
|
| 141 |
yield None, link.read(), f"Failed: {message}"
|
| 142 |
raise gr.Error(message) from exc
|
| 143 |
-
|
|
|
|
|
|
|
| 144 |
|
| 145 |
|
| 146 |
-
def publish_cpu(spec: RunSpec, summary: dict | None) -> Iterator[tuple[bytes | None, str]]:
|
| 147 |
"""Attach the finished rig and its six videos, off the GPU allocation."""
|
| 148 |
|
| 149 |
if summary is None:
|
| 150 |
raise gr.Error("The generation phase did not publish a result.")
|
| 151 |
-
link: Link = open_link(spec.token)
|
| 152 |
try:
|
| 153 |
label: str = (
|
| 154 |
smoke_publish_phase(link.recording, spec)
|
| 155 |
if STREAM_SMOKE
|
| 156 |
else publish_phase(link.recording, spec, summary)
|
| 157 |
)
|
|
|
|
| 158 |
except Exception as exc:
|
| 159 |
message: str = f"{type(exc).__name__}: {exc}"
|
| 160 |
-
yield link.read(), f"Failed: {message}"
|
| 161 |
raise gr.Error(message) from exc
|
| 162 |
-
|
|
|
|
|
|
|
|
|
|
| 163 |
|
| 164 |
|
| 165 |
DESCRIPTION: str = """
|
|
@@ -255,6 +292,9 @@ def build_demo() -> gr.Blocks:
|
|
| 255 |
status: gr.Markdown = gr.Markdown(
|
| 256 |
"Upload a clip, or pick the example, then press Run.", elem_id="run-status"
|
| 257 |
)
|
|
|
|
|
|
|
|
|
|
| 258 |
with gr.Column(scale=3):
|
| 259 |
viewer: Rerun = Rerun(
|
| 260 |
streaming=True,
|
|
@@ -262,10 +302,12 @@ def build_demo() -> gr.Blocks:
|
|
| 262 |
panel_states={"time": "collapsed", "blueprint": "hidden", "selection": "hidden"},
|
| 263 |
)
|
| 264 |
|
| 265 |
-
started = run_button.click(
|
|
|
|
|
|
|
| 266 |
prepared = started.success(prepare_cpu, spec, [viewer, status])
|
| 267 |
generated = prepared.success(run_gpu, spec, [summary, viewer, status])
|
| 268 |
-
published = generated.success(publish_cpu, [spec, summary], [viewer, status])
|
| 269 |
# `cancels` detaches the UI at once, but the pipeline itself only halts
|
| 270 |
# when it reads the sentinel `request_stop` writes: the worker thread
|
| 271 |
# (and on ZeroGPU, the forked child) never sees a Gradio cancellation.
|
|
|
|
| 30 |
Session,
|
| 31 |
begin_phase,
|
| 32 |
generate_phase,
|
| 33 |
+
merge_rrd_parts,
|
| 34 |
motion_phase,
|
| 35 |
new_spec,
|
| 36 |
publish_phase,
|
| 37 |
request_stop,
|
| 38 |
+
rrd_part,
|
| 39 |
smoke_generate_phase,
|
| 40 |
smoke_motion_phase,
|
| 41 |
smoke_publish_phase,
|
|
|
|
| 57 |
|
| 58 |
return self.stream.read()
|
| 59 |
|
| 60 |
+
def close(self) -> None:
|
| 61 |
+
"""Release this link's file sink before the next link (or fork) runs.
|
| 62 |
|
| 63 |
+
The part is already complete — a footerless part is whole at the last
|
| 64 |
+
flush — so this is descriptor hygiene, not a correctness dependency.
|
| 65 |
+
"""
|
| 66 |
+
|
| 67 |
+
self.recording.disconnect()
|
| 68 |
+
|
| 69 |
+
|
| 70 |
+
def open_link(token: str, part: Path) -> Link:
|
| 71 |
"""Open this callback's own recording, in whatever process the callback runs in.
|
| 72 |
|
| 73 |
Called first thing in every link, GPU ones included, where "this callback"
|
|
|
|
| 86 |
|
| 87 |
rr.cleanup_if_forked_child()
|
| 88 |
recording: rr.RecordingStream = rr.RecordingStream(APPLICATION_ID, recording_id=token)
|
| 89 |
+
# Order is load-bearing: binary_stream() installs the stream sink, and
|
| 90 |
+
# set_sinks REPLACES the sink set — so the stream must be re-installed
|
| 91 |
+
# beside the file. Footerless parts are what make the final concatenation
|
| 92 |
+
# a single valid recording.
|
| 93 |
+
stream: rr.BinaryStream = recording.binary_stream()
|
| 94 |
+
recording.set_sinks(stream, rr.FileSink(str(part), write_footer=False))
|
| 95 |
+
return Link(recording=recording, stream=stream)
|
| 96 |
|
| 97 |
|
| 98 |
def begin(
|
| 99 |
video: str | None, start_time: float, seed: int
|
| 100 |
+
) -> Iterator[tuple[RunSpec, bytes | None, str, Any, Any]]:
|
| 101 |
"""Validate the input on CPU, open a recording, and switch to the outputs."""
|
| 102 |
|
| 103 |
if video is None:
|
|
|
|
| 106 |
spec: RunSpec = new_spec(Path(video), float(start_time), int(seed))
|
| 107 |
except FourDAnyoneError as exc:
|
| 108 |
raise gr.Error(str(exc)) from None
|
| 109 |
+
link: Link = open_link(spec.token, rrd_part(spec, "begin"))
|
| 110 |
+
try:
|
| 111 |
+
begin_phase(link.recording, spec)
|
| 112 |
+
yield (
|
| 113 |
+
spec,
|
| 114 |
+
link.read(),
|
| 115 |
+
"Preparing the source clip.",
|
| 116 |
+
gr.Tabs(selected="outputs"),
|
| 117 |
+
gr.DownloadButton(visible=False),
|
| 118 |
+
)
|
| 119 |
+
finally:
|
| 120 |
+
link.close()
|
| 121 |
|
| 122 |
|
| 123 |
def prepare_cpu(spec: RunSpec) -> Iterator[tuple[bytes | None, str]]:
|
| 124 |
"""Put the source clip on the frame timeline before any GPU work starts."""
|
| 125 |
|
| 126 |
+
link: Link = open_link(spec.token, rrd_part(spec, "prepare"))
|
| 127 |
+
try:
|
| 128 |
+
label: str = source_phase(link.recording, spec)
|
| 129 |
+
yield link.read(), label
|
| 130 |
+
finally:
|
| 131 |
+
link.close()
|
| 132 |
|
| 133 |
|
| 134 |
@spaces.GPU(duration=GPU_DURATION)
|
|
|
|
| 144 |
is the only process that may build the recording it logs through.
|
| 145 |
"""
|
| 146 |
|
| 147 |
+
link: Link = open_link(spec.token, rrd_part(spec, "gpu"))
|
| 148 |
session: Session = Session(spec)
|
| 149 |
try:
|
| 150 |
motion: Iterator[str] = (
|
|
|
|
| 161 |
)
|
| 162 |
for label in generate:
|
| 163 |
yield None, link.read(), label
|
| 164 |
+
# Read before close: disconnect drops the binary sink with the file.
|
| 165 |
+
closing: bytes | None = link.read()
|
| 166 |
except FourDAnyoneError as exc:
|
| 167 |
# The banner keeps the reason on screen after the toast dismisses.
|
| 168 |
yield None, link.read(), f"Failed: {exc}"
|
|
|
|
| 171 |
message: str = f"{type(exc).__name__}: {exc}"
|
| 172 |
yield None, link.read(), f"Failed: {message}"
|
| 173 |
raise gr.Error(message) from exc
|
| 174 |
+
finally:
|
| 175 |
+
link.close()
|
| 176 |
+
yield session.summary, closing, "Publishing the result."
|
| 177 |
|
| 178 |
|
| 179 |
+
def publish_cpu(spec: RunSpec, summary: dict | None) -> Iterator[tuple[bytes | None, str, Any]]:
|
| 180 |
"""Attach the finished rig and its six videos, off the GPU allocation."""
|
| 181 |
|
| 182 |
if summary is None:
|
| 183 |
raise gr.Error("The generation phase did not publish a result.")
|
| 184 |
+
link: Link = open_link(spec.token, rrd_part(spec, "publish"))
|
| 185 |
try:
|
| 186 |
label: str = (
|
| 187 |
smoke_publish_phase(link.recording, spec)
|
| 188 |
if STREAM_SMOKE
|
| 189 |
else publish_phase(link.recording, spec, summary)
|
| 190 |
)
|
| 191 |
+
closing: bytes | None = link.read()
|
| 192 |
except Exception as exc:
|
| 193 |
message: str = f"{type(exc).__name__}: {exc}"
|
| 194 |
+
yield link.read(), f"Failed: {message}", gr.skip()
|
| 195 |
raise gr.Error(message) from exc
|
| 196 |
+
finally:
|
| 197 |
+
link.close()
|
| 198 |
+
merged: Path = merge_rrd_parts(spec)
|
| 199 |
+
yield closing, label, gr.DownloadButton(value=str(merged), visible=True)
|
| 200 |
|
| 201 |
|
| 202 |
DESCRIPTION: str = """
|
|
|
|
| 292 |
status: gr.Markdown = gr.Markdown(
|
| 293 |
"Upload a clip, or pick the example, then press Run.", elem_id="run-status"
|
| 294 |
)
|
| 295 |
+
download: gr.DownloadButton = gr.DownloadButton(
|
| 296 |
+
"Download the recording (.rrd)", visible=False
|
| 297 |
+
)
|
| 298 |
with gr.Column(scale=3):
|
| 299 |
viewer: Rerun = Rerun(
|
| 300 |
streaming=True,
|
|
|
|
| 302 |
panel_states={"time": "collapsed", "blueprint": "hidden", "selection": "hidden"},
|
| 303 |
)
|
| 304 |
|
| 305 |
+
started = run_button.click(
|
| 306 |
+
begin, [video, start_time, seed], [spec, viewer, status, tabs, download]
|
| 307 |
+
)
|
| 308 |
prepared = started.success(prepare_cpu, spec, [viewer, status])
|
| 309 |
generated = prepared.success(run_gpu, spec, [summary, viewer, status])
|
| 310 |
+
published = generated.success(publish_cpu, [spec, summary], [viewer, status, download])
|
| 311 |
# `cancels` detaches the UI at once, but the pipeline itself only halts
|
| 312 |
# when it reads the sentinel `request_stop` writes: the worker thread
|
| 313 |
# (and on ZeroGPU, the forked child) never sees a Gradio cancellation.
|
tests/test_app_helpers.py
CHANGED
|
@@ -256,7 +256,7 @@ def test_a_recording_made_before_a_fork_cannot_be_flushed_after_it() -> None:
|
|
| 256 |
)
|
| 257 |
|
| 258 |
|
| 259 |
-
def test_open_link_gives_a_forked_child_a_recording_it_can_flush() -> None:
|
| 260 |
"""The fix: each link builds its own stream, in whatever process it runs in.
|
| 261 |
|
| 262 |
The token is the ``recording_id``, so the parent's rows and the child's are
|
|
@@ -264,25 +264,28 @@ def test_open_link_gives_a_forked_child_a_recording_it_can_flush() -> None:
|
|
| 264 |
"""
|
| 265 |
|
| 266 |
token: str = uuid.uuid4().hex
|
| 267 |
-
parent: fdanyone_app.Link = fdanyone_app.open_link(token)
|
| 268 |
parent.recording.log("log", rr.TextLog("parent"))
|
| 269 |
assert parent.read()
|
|
|
|
| 270 |
|
| 271 |
def body() -> int:
|
| 272 |
-
child: fdanyone_app.Link = fdanyone_app.open_link(token)
|
| 273 |
assert child.recording.get_recording_id() == token
|
| 274 |
child.recording.log("log", rr.TextLog("child"))
|
| 275 |
payload: bytes | None = child.read()
|
| 276 |
return OK if payload else WRONG
|
| 277 |
|
| 278 |
assert _in_child(body) == OK
|
|
|
|
|
|
|
| 279 |
|
| 280 |
|
| 281 |
RRD_MAGIC: bytes = b"RRF2"
|
| 282 |
"""First four bytes of every RRD document the SDK emits."""
|
| 283 |
|
| 284 |
|
| 285 |
-
def test_a_stream_per_link_adds_no_framing_the_viewer_did_not_already_get() -> None:
|
| 286 |
"""Why splitting one stream into four costs the browser viewer nothing.
|
| 287 |
|
| 288 |
Each ``read`` already returns a whole RRD document, magic bytes and manifest
|
|
@@ -294,7 +297,7 @@ def test_a_stream_per_link_adds_no_framing_the_viewer_did_not_already_get() -> N
|
|
| 294 |
where that shows up rather than in a blank viewer on the Space.
|
| 295 |
"""
|
| 296 |
|
| 297 |
-
link: fdanyone_app.Link = fdanyone_app.open_link(uuid.uuid4().hex)
|
| 298 |
for row in range(3):
|
| 299 |
link.recording.log("log", rr.TextLog(f"row {row}"))
|
| 300 |
payload: bytes | None = link.read()
|
|
|
|
| 256 |
)
|
| 257 |
|
| 258 |
|
| 259 |
+
def test_open_link_gives_a_forked_child_a_recording_it_can_flush(tmp_path: Path) -> None:
|
| 260 |
"""The fix: each link builds its own stream, in whatever process it runs in.
|
| 261 |
|
| 262 |
The token is the ``recording_id``, so the parent's rows and the child's are
|
|
|
|
| 264 |
"""
|
| 265 |
|
| 266 |
token: str = uuid.uuid4().hex
|
| 267 |
+
parent: fdanyone_app.Link = fdanyone_app.open_link(token, tmp_path / "parent.rrd")
|
| 268 |
parent.recording.log("log", rr.TextLog("parent"))
|
| 269 |
assert parent.read()
|
| 270 |
+
parent.close()
|
| 271 |
|
| 272 |
def body() -> int:
|
| 273 |
+
child: fdanyone_app.Link = fdanyone_app.open_link(token, tmp_path / "child.rrd")
|
| 274 |
assert child.recording.get_recording_id() == token
|
| 275 |
child.recording.log("log", rr.TextLog("child"))
|
| 276 |
payload: bytes | None = child.read()
|
| 277 |
return OK if payload else WRONG
|
| 278 |
|
| 279 |
assert _in_child(body) == OK
|
| 280 |
+
# The dual sink teed the parent's rows into its part file as well.
|
| 281 |
+
assert (tmp_path / "parent.rrd").stat().st_size > 0
|
| 282 |
|
| 283 |
|
| 284 |
RRD_MAGIC: bytes = b"RRF2"
|
| 285 |
"""First four bytes of every RRD document the SDK emits."""
|
| 286 |
|
| 287 |
|
| 288 |
+
def test_a_stream_per_link_adds_no_framing_the_viewer_did_not_already_get(tmp_path: Path) -> None:
|
| 289 |
"""Why splitting one stream into four costs the browser viewer nothing.
|
| 290 |
|
| 291 |
Each ``read`` already returns a whole RRD document, magic bytes and manifest
|
|
|
|
| 297 |
where that shows up rather than in a blank viewer on the Space.
|
| 298 |
"""
|
| 299 |
|
| 300 |
+
link: fdanyone_app.Link = fdanyone_app.open_link(uuid.uuid4().hex, tmp_path / "part.rrd")
|
| 301 |
for row in range(3):
|
| 302 |
link.recording.log("log", rr.TextLog(f"row {row}"))
|
| 303 |
payload: bytes | None = link.read()
|