4danyone-rerun / fdanyone_app.py
pablovela5620's picture
Claude Fable 5
Apply the thermo-review survivors (app side)
260a292
Raw History Blame Contribute Delete
13.5 kB
"""Gradio and ZeroGPU wiring for the 4DAnyone Rerun streaming runtime.
The three-link event chain is ``begin -> run_gpu -> publish_cpu``, one link
per process boundary.
ZeroGPU runs decorated callbacks in forked children, so each link creates its
own recording after ``rr.cleanup_if_forked_child()`` and uses the shared token
as its recording ID. The merged ``run_gpu`` remains one allocation because its
motion result contains live state that cannot cross a process boundary.
"""
from __future__ import annotations
import logging
from collections.abc import Iterator
from dataclasses import dataclass
from pathlib import Path
from typing import Any
import gradio as gr
import rerun as rr
import spaces
from gradio_rerun import Rerun
from fdanyone.errors import FourDAnyoneError
from fdanyone.rerun_streaming import (
APPLICATION_ID,
EXAMPLE_DIR,
GPU_DURATION,
STREAM_SMOKE,
RunResult,
RunSpec,
Session,
begin_phase,
format_status,
generate_phase,
merge_rrd_parts,
motion_phase,
new_spec,
publish_phase,
request_stop,
rrd_part,
smoke_generate_phase,
smoke_motion_phase,
smoke_publish_phase,
source_phase,
)
@dataclass(frozen=True)
class Link:
"""One chain link's own recording, and the sink its bytes are drained from."""
recording: rr.RecordingStream
"""Explicit stream this link and its pipeline hooks log through."""
stream: rr.BinaryStream
"""Binary sink feeding the Gradio viewer."""
def read(self) -> bytes | None:
"""Return bytes buffered since the last read."""
return self.stream.read()
def close(self) -> None:
"""Release this link's file sink before the next link (or fork) runs.
The part is already complete — a footerless part is whole at the last
flush — so this is descriptor hygiene, not a correctness dependency.
"""
self.recording.disconnect()
def open_link(token: str, part: Path) -> Link:
"""Open this callback's own recording, in whatever process the callback runs in.
Called first thing in every link, GPU ones included, where "this callback"
means the forked child. A stream inherited across the fork cannot be used:
the SDK compares pids on flush and raises "Fork detected during flush". A
stream built here belongs to this process, and sharing ``token`` as the
recording id is what makes the viewer treat all the links as one recording.
``cleanup_if_forked_child`` drops whatever streams were inherited. The SDK
registers it with ``os.register_at_fork`` already, and ``rr.init`` calls it
outright for the same reason, so this is belt and braces and a no-op outside
a child. It is called unconditionally: guarding it on the name existing
would turn a rename in the SDK into a silent return of the very bug this
function exists to prevent.
"""
rr.cleanup_if_forked_child()
recording: rr.RecordingStream = rr.RecordingStream(APPLICATION_ID, recording_id=token)
# Order is load-bearing: binary_stream() installs the stream sink, and
# set_sinks REPLACES the sink set — so the stream must be re-installed
# beside the file. Footerless parts are what make the final concatenation
# a single valid recording.
stream: rr.BinaryStream = recording.binary_stream()
recording.set_sinks(stream, rr.FileSink(str(part), write_footer=False))
return Link(recording=recording, stream=stream)
def begin(
video: str | None, start_time: float, seed: int
) -> Iterator[tuple[Any, bytes | None, str, Any, Any]]:
"""Validate, open a recording, switch to the outputs, and log the source.
One CPU link: the earlier split between "begin" and "prepare" callbacks
was an accidental process boundary — both ran on the CPU in this process,
each with its own recording lifecycle and part file. The first yield lands
the tab switch immediately; the second follows with the source clip.
"""
if video is None:
raise gr.Error("Upload a video, or pick the bundled example.")
try:
spec: RunSpec = new_spec(Path(video), float(start_time), int(seed))
except FourDAnyoneError as exc:
raise gr.Error(str(exc)) from None
link: Link = open_link(spec.token, rrd_part(spec, "begin"))
try:
begin_phase(link.recording, spec)
yield (
spec,
link.read(),
format_status(spec, "Preparing the source clip."),
gr.Tabs(selected="outputs"),
gr.DownloadButton(visible=False),
)
label: str = source_phase(link.recording, spec)
yield gr.skip(), link.read(), format_status(spec, label), gr.skip(), gr.skip()
finally:
link.close()
@spaces.GPU(duration=GPU_DURATION)
def run_gpu(spec: RunSpec) -> Iterator[tuple[RunResult | None, bytes | None, str]]:
"""Recover the motion and generate the six views, inside one allocation.
The two phases have to share a process. ``PreparedRun`` carries the decoded
clip, the open conditioning artifacts, and a completion barrier for the
skeleton renderer, so it cannot be pickled from one ZeroGPU worker into
another — and every ``@spaces.GPU`` function gets a worker of its own.
Everything in here runs in the forked child, including ``open_link``: this
is the only process that may build the recording it logs through.
"""
link: Link = open_link(spec.token, rrd_part(spec, "gpu"))
session: Session = Session(spec)
motion: Iterator[str] | None = None
generate: Iterator[str] | None = None
try:
motion = (
smoke_motion_phase(link.recording, session)
if STREAM_SMOKE
else motion_phase(link.recording, session)
)
for label in motion:
yield None, link.read(), format_status(spec, label)
generate = (
smoke_generate_phase(link.recording, session)
if STREAM_SMOKE
else generate_phase(link.recording, session)
)
for label in generate:
yield None, link.read(), format_status(spec, label)
# Read before close: disconnect drops the binary sink with the file.
closing: bytes | None = link.read()
except FourDAnyoneError as exc:
# The banner keeps the reason on screen after the toast dismisses.
yield None, link.read(), f"Failed: {exc}"
raise gr.Error(str(exc)) from None
except Exception as exc:
message: str = f"{type(exc).__name__}: {exc}"
yield None, link.read(), f"Failed: {message}"
raise gr.Error(message) from exc
finally:
# Close the phase generators FIRST: their unwinding (pump's bounded
# worker join) must finish while the recording is still connected.
for phase in (motion, generate):
if phase is not None:
phase.close()
link.close()
yield session.summary, closing, format_status(spec, "Publishing the result.")
def publish_cpu(spec: RunSpec, summary: RunResult | None) -> Iterator[tuple[bytes | None, str, Any]]:
"""Attach the finished rig and its six videos, off the GPU allocation."""
if summary is None:
raise gr.Error("The generation phase did not publish a result.")
link: Link = open_link(spec.token, rrd_part(spec, "publish"))
try:
label: str = (
smoke_publish_phase(link.recording, spec)
if STREAM_SMOKE
else publish_phase(link.recording, spec, summary)
)
closing: bytes | None = link.read()
except Exception as exc:
message: str = f"{type(exc).__name__}: {exc}"
yield link.read(), f"Failed: {message}", gr.skip()
raise gr.Error(message) from exc
finally:
link.close()
merged: Path = merge_rrd_parts(spec)
yield closing, label, gr.DownloadButton(value=str(merged), visible=True)
DESCRIPTION: str = """
# 4DAnyone × [Rerun](https://rerun.io)
One monocular clip of a person becomes six synchronized novel views. GVHMR
recovers the SMPL-X motion, a Wan 2.2 diffusion transformer generates every
view, and each phase streams into the Rerun viewer as it happens.
"""
SOURCE_CLIP_HEIGHT: int = 360
"""Display height of the source preview. A portrait clip scaled to the column
width is taller than the fold, which pushes every control below it out of sight."""
STATUS_CSS: str = """
#run-status {
display: flex;
flex-direction: column;
justify-content: center;
/* Tall enough for the progress overlay Gradio draws over a running event. */
min-height: 5rem;
/* Gradio writes `overflow: auto` inline on every block. */
overflow: visible !important;
padding: 0.75rem 1rem;
border-radius: var(--radius-lg);
background: var(--background-fill-secondary);
}
#run-status p {
font-size: 1.15rem;
line-height: 1.5;
margin: 0;
}
"""
"""Banner styling for the run status, the only readout of a multi-minute run.
Gradio otherwise gives a Markdown block the body font and a height its own text
overflows, which turns the line into a scrolling sliver. Gradio 6 takes ``css``
on ``launch``, not on the ``Blocks`` constructor."""
def build_demo() -> gr.Blocks:
"""Assemble the persistent viewer, the inputs, and the run chain.
Controls sit in a narrow column and the viewer fills a wide one beside it,
the same split every other Rerun Space here uses. Stacking them instead
pushes the viewer off the fold on a laptop, and the viewer is the demo. It
stays outside the tabs for the same reason it is created once: a tab switch
would tear down the component and drop the stream mid-run.
The chain is joined by ``success`` rather than ``then``: each link's inputs
are the previous link's outputs, so a link that failed leaves the next one
nothing to work with, and running it anyway would bury the real error under
a second, meaningless one.
"""
# Serve the example clips in place; without this Gradio copies all 209 MB
# of them into its temp cache at construction.
gr.set_static_paths([EXAMPLE_DIR])
with gr.Blocks(title="4DAnyone × Rerun") as demo:
gr.Markdown(DESCRIPTION)
spec: gr.State = gr.State(None)
summary: gr.State = gr.State(None)
with gr.Row():
with gr.Column(scale=1):
with gr.Tabs() as tabs:
with gr.Tab("Input", id="input"):
video: gr.Video = gr.Video(
label="Source clip", sources=["upload"], height=SOURCE_CLIP_HEIGHT
)
with gr.Accordion("Advanced settings", open=False):
start_time: gr.Number = gr.Number(
value=0.0, label="Start time (seconds)", minimum=0.0
)
seed: gr.Number = gr.Number(
value=0, label="Seed", precision=0, minimum=0
)
# A single input keeps the compact thumbnail gallery; a
# second column would turn it into a labeled table.
gr.Examples(
examples=[[str(clip)] for clip in sorted(EXAMPLE_DIR.glob("*.mp4"))],
inputs=[video],
cache_examples=False,
examples_per_page=4,
)
with gr.Tab("Outputs", id="outputs"):
gr.Markdown(
"The viewer beside this switches layout with the run: detections on the "
"source clip, then a grid of per-step previews, then the finished rig "
"playing on a loop."
)
with gr.Row():
run_button: gr.Button = gr.Button("Run", variant="primary")
stop_button: gr.Button = gr.Button("Stop", variant="stop")
status: gr.Markdown = gr.Markdown(
"Upload a clip, or pick the example, then press Run.", elem_id="run-status"
)
download: gr.DownloadButton = gr.DownloadButton(
"Download the recording (.rrd)", visible=False
)
with gr.Column(scale=3):
viewer: Rerun = Rerun(
streaming=True,
height=760,
panel_states={"time": "collapsed", "blueprint": "hidden", "selection": "hidden"},
)
started = run_button.click(
begin, [video, start_time, seed], [spec, viewer, status, tabs, download]
)
generated = started.success(run_gpu, spec, [summary, viewer, status])
published = generated.success(publish_cpu, [spec, summary], [viewer, status, download])
# `cancels` detaches the UI at once, but the pipeline itself only halts
# when it reads the sentinel `request_stop` writes: the worker thread
# (and on ZeroGPU, the forked child) never sees a Gradio cancellation.
stop_button.click(
request_stop, spec, status, cancels=[started, generated, published]
)
return demo
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO, format="%(asctime)s | %(levelname)s | %(message)s")
build_demo().launch(ssr_mode=False, css=STATUS_CSS, show_error=True)