pablovela5620 Claude Fable 5 commited on
Commit
260a292
·
1 Parent(s): 9ae2cb2

Apply the thermo-review survivors (app side)

Browse files

The chain is three links — begin absorbed the accidental prepare_cpu
boundary (same process, one recording, one part fewer). The generation
result crosses the fork as a typed RunResult. merge_rrd_parts requires
every part non-empty and lands atomically via os.replace. pump owns its
worker's lifetime: an early consumer close requests the stop sentinel
and joins with a bounded window, and run_gpu closes its phase generators
before disconnecting the recording. Vendored: space-streaming@c8893d6.
Validated end-to-end locally: 479 s, all 19 stages once, merged rrd from
exactly three parts.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

PROVENANCE.md CHANGED
@@ -6,8 +6,8 @@
6
  | --- | --- |
7
  | Source repository | <https://github.com/pablovela5620/4DAnyone-5090> |
8
  | Branch | `space-streaming` |
9
- | Commit | `d2ac77a55accd22dac2323207e5fac40ee7661e0` |
10
- | Synced | 2026-08-29T04:45:52Z |
11
 
12
  ## Excluded from the copy
13
 
 
6
  | --- | --- |
7
  | Source repository | <https://github.com/pablovela5620/4DAnyone-5090> |
8
  | Branch | `space-streaming` |
9
+ | Commit | `c8893d657278fe6bd0b3447a87a6cf29f4b4ded4` |
10
+ | Synced | 2026-08-29T05:06:46Z |
11
 
12
  ## Excluded from the copy
13
 
fdanyone/model/loader.py CHANGED
@@ -403,6 +403,11 @@ def load_pipeline(
403
  key = _cache_key(checkpoint_path, assets, device, settings)
404
  if _PIPELINE_CACHE is not None and _PIPELINE_CACHE[0] == key:
405
  return _PIPELINE_CACHE[1]
 
 
 
 
 
406
  loaded = _load_pipeline_uncached(
407
  checkpoint_path=checkpoint_path,
408
  assets=assets,
 
403
  key = _cache_key(checkpoint_path, assets, device, settings)
404
  if _PIPELINE_CACHE is not None and _PIPELINE_CACHE[0] == key:
405
  return _PIPELINE_CACHE[1]
406
+ # Evict BEFORE building: holding the old stack while the new one loads
407
+ # would put two ~11 GB stacks in memory at once — the exact situation
408
+ # the single slot exists to prevent.
409
+ _PIPELINE_CACHE = None
410
+ gc.collect()
411
  loaded = _load_pipeline_uncached(
412
  checkpoint_path=checkpoint_path,
413
  assets=assets,
fdanyone/model/turbo_lora.py CHANGED
@@ -124,12 +124,12 @@ def merge_wan_turbo_lora(
124
  up = adapter.get_tensor(up_key).to(device=compute_device, dtype=torch.float32)
125
  delta = torch.mm(up, down)
126
  parameter.add_(delta.to(device=parameter.device, dtype=parameter.dtype))
 
 
 
 
127
  finally:
128
  torch.backends.cuda.matmul.allow_tf32 = tf32
129
- for key, target_name in direct_targets.items():
130
- parameter = parameters[target_name]
131
- delta = adapter.get_tensor(key)
132
- parameter.add_(delta.to(device=parameter.device, dtype=parameter.dtype))
133
 
134
  return TurboLoraReport(
135
  lora_modules=len(down_keys),
 
124
  up = adapter.get_tensor(up_key).to(device=compute_device, dtype=torch.float32)
125
  delta = torch.mm(up, down)
126
  parameter.add_(delta.to(device=parameter.device, dtype=parameter.dtype))
127
+ for key, target_name in direct_targets.items():
128
+ parameter = parameters[target_name]
129
+ delta = adapter.get_tensor(key)
130
+ parameter.add_(delta.to(device=parameter.device, dtype=parameter.dtype))
131
  finally:
132
  torch.backends.cuda.matmul.allow_tf32 = tf32
 
 
 
 
133
 
134
  return TurboLoraReport(
135
  lora_modules=len(down_keys),
fdanyone/pipeline.py CHANGED
@@ -238,7 +238,7 @@ def _build_conditioning(
238
  stderr=subprocess.STDOUT,
239
  )
240
  completed: bool = False
241
- failure: FourDAnyoneError | None = None
242
 
243
  def _validate_target_render() -> None:
244
  return_code: int = render_process.wait()
@@ -278,7 +278,9 @@ def _build_conditioning(
278
 
279
  def wait_for_targets() -> None:
280
  # Re-entry (the pipeline-level cleanup) must repeat the original
281
- # verdict, not mask an informative failure with a generic one.
 
 
282
  nonlocal completed, failure
283
  if completed:
284
  if failure is not None:
@@ -287,7 +289,7 @@ def _build_conditioning(
287
  completed = True
288
  try:
289
  _validate_target_render()
290
- except FourDAnyoneError as exc:
291
  failure = exc
292
  raise
293
 
@@ -465,9 +467,15 @@ def prepare_run(
465
  # Re-decode the worker-produced source before it becomes a model tensor.
466
  verify_lossless_video(clip, conditioning.source_video)
467
  except BaseException:
468
- if conditioning is not None and conditioning.target_waiter is not None:
469
- conditioning.wait_for_target_skeletons()
470
- _discard_scratch(scratch)
 
 
 
 
 
 
471
  raise
472
  return PreparedRun(
473
  settings=settings,
@@ -486,6 +494,25 @@ def prepare_run(
486
  )
487
 
488
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
489
  def warm_generation_models(
490
  *,
491
  model_dir: str,
@@ -526,7 +553,7 @@ def generate_run(
526
  on_denoise_step: DenoiseStepHook | None = None,
527
  on_denoise_step_start: DenoiseStepStartHook | None = None,
528
  on_generate_stage: GenerateStageHook | None = None,
529
- ) -> dict:
530
  """Generate and publish every view, consuming the prepared run.
531
 
532
  The prompt embedding is fixed by ``prepare_run``, because whether it exists
@@ -568,9 +595,24 @@ def generate_run(
568
  )
569
  summary["result_dir"] = str(prepared.result_dir)
570
  summary["motion_dir"] = str(prepared.motion_dir)
571
- return summary
572
- finally:
573
- release_run(prepared)
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
574
 
575
 
576
  def release_run(prepared: PreparedRun) -> None:
@@ -608,7 +650,7 @@ def run_pipeline(
608
  on_denoise_step: DenoiseStepHook | None = None,
609
  on_denoise_step_start: DenoiseStepStartHook | None = None,
610
  prompt_embedding_path: Path | None = None,
611
- ) -> dict:
612
  """Execute inference and publish reusable GVHMR plus 4DAnyone results."""
613
 
614
  if seed < 0:
 
238
  stderr=subprocess.STDOUT,
239
  )
240
  completed: bool = False
241
+ failure: Exception | None = None
242
 
243
  def _validate_target_render() -> None:
244
  return_code: int = render_process.wait()
 
278
 
279
  def wait_for_targets() -> None:
280
  # Re-entry (the pipeline-level cleanup) must repeat the original
281
+ # verdict, not mask an informative failure with a generic one — and
282
+ # that holds for ANY failure, not just the domain error: a JSON or OS
283
+ # error forgotten here would turn into a silent pass on re-entry.
284
  nonlocal completed, failure
285
  if completed:
286
  if failure is not None:
 
289
  completed = True
290
  try:
291
  _validate_target_render()
292
+ except Exception as exc:
293
  failure = exc
294
  raise
295
 
 
467
  # Re-decode the worker-produced source before it becomes a model tensor.
468
  verify_lossless_video(clip, conditioning.source_video)
469
  except BaseException:
470
+ # The waiter can itself raise (its renderer failed too); the scratch
471
+ # must still be discarded, and the original failure must propagate.
472
+ try:
473
+ if conditioning is not None and conditioning.target_waiter is not None:
474
+ conditioning.wait_for_target_skeletons()
475
+ except Exception:
476
+ LOGGER.exception("Deferred target rendering also failed during unwinding")
477
+ finally:
478
+ _discard_scratch(scratch)
479
  raise
480
  return PreparedRun(
481
  settings=settings,
 
494
  )
495
 
496
 
497
+ @dataclass(frozen=True)
498
+ class RunResult:
499
+ """What one finished run publishes across the process boundary.
500
+
501
+ Frozen and picklable: on ZeroGPU it travels from the forked worker back
502
+ through a Gradio ``State``. ``metadata`` is the exported summary document,
503
+ kept for CLI and metadata consumers; the typed fields are the contract.
504
+ """
505
+
506
+ result_dir: Path
507
+ """Published result directory: six views, cameras, metadata."""
508
+ motion_dir: Path
509
+ """Published GVHMR motion directory."""
510
+ total_elapsed_seconds: float
511
+ """Wall time from pipeline start to the exported result."""
512
+ metadata: dict
513
+ """The full exported summary, exactly as written to disk."""
514
+
515
+
516
  def warm_generation_models(
517
  *,
518
  model_dir: str,
 
553
  on_denoise_step: DenoiseStepHook | None = None,
554
  on_denoise_step_start: DenoiseStepStartHook | None = None,
555
  on_generate_stage: GenerateStageHook | None = None,
556
+ ) -> RunResult:
557
  """Generate and publish every view, consuming the prepared run.
558
 
559
  The prompt embedding is fixed by ``prepare_run``, because whether it exists
 
595
  )
596
  summary["result_dir"] = str(prepared.result_dir)
597
  summary["motion_dir"] = str(prepared.motion_dir)
598
+ result = RunResult(
599
+ result_dir=prepared.result_dir,
600
+ motion_dir=prepared.motion_dir,
601
+ total_elapsed_seconds=float(summary["total_pipeline_elapsed_seconds"]),
602
+ metadata=summary,
603
+ )
604
+ except BaseException:
605
+ # Release failures during unwinding must not replace the failure that
606
+ # actually explains the run; they are logged instead.
607
+ try:
608
+ release_run(prepared)
609
+ except Exception:
610
+ LOGGER.exception("Run release failed while unwinding a generation failure")
611
+ raise
612
+ # A release failure on the success path stays loud: the renderer's own
613
+ # verdict is part of the run's validity.
614
+ release_run(prepared)
615
+ return result
616
 
617
 
618
  def release_run(prepared: PreparedRun) -> None:
 
650
  on_denoise_step: DenoiseStepHook | None = None,
651
  on_denoise_step_start: DenoiseStepStartHook | None = None,
652
  prompt_embedding_path: Path | None = None,
653
+ ) -> RunResult:
654
  """Execute inference and publish reusable GVHMR plus 4DAnyone results."""
655
 
656
  if seed < 0:
fdanyone/rerun_streaming.py CHANGED
@@ -33,7 +33,7 @@ from fdanyone.errors import FourDAnyoneError
33
  from fdanyone.model.inference import _tensor_frames
34
  from fdanyone.model.tiny_decoder import decode_tiny_target_video, load_tiny_wan_decoder
35
  from fdanyone.motion.body import BodyMotion, load_body_motion
36
- from fdanyone.pipeline import PreparedRun, generate_run, prepare_run, warm_generation_models
37
  from fdanyone.video import choose_canonical_fps
38
  from fdanyone.viz import (
39
  BOX_COLOR,
@@ -110,8 +110,10 @@ STREAM_SMOKE: bool = os.environ.get("FDANYONE_STREAM_SMOKE") == "1"
110
  SMOKE_DELAY: float = float(os.environ.get("FDANYONE_SMOKE_DELAY", "0.1"))
111
  """Seconds between smoke yields; raise it to eyeball each phase in a browser."""
112
 
113
- SMOKE_SUMMARY: dict = {"result_dir": "", "total_pipeline_elapsed_seconds": 0.0}
114
- """Stand-in for the pipeline's published metadata; the smoke run writes no files."""
 
 
115
 
116
  # Body evaluation resolves SMPL-X relative to the vendored package's repository
117
  # root, which is not where a Space keeps its models. An absolute root wins the
@@ -281,8 +283,8 @@ class Session:
281
  """The run this link was rebuilt from."""
282
  prepared: PreparedRun | None = None
283
  """Result of the motion phase, consumed by the generation phase."""
284
- summary: dict | None = None
285
- """Published metadata the generation phase returns, set once it succeeds."""
286
  source_frames: dict[int, RgbFrame] = field(default_factory=dict)
287
  """Decoded source stills for the diffusion pane, keyed by canonical frame."""
288
  events: queue.Queue[str | None] = field(default_factory=queue.Queue)
@@ -311,7 +313,7 @@ def check_stop(spec: RunSpec) -> None:
311
  raise RunCancelled("Stopped.")
312
 
313
 
314
- LINK_ORDER: tuple[str, ...] = ("begin", "prepare", "gpu", "publish")
315
  """Chain links in stream order, and therefore the part-concatenation order."""
316
 
317
 
@@ -328,17 +330,24 @@ def merge_rrd_parts(spec: RunSpec) -> Path:
328
 
329
  Parts are written with ``write_footer=False``: footered parts concatenate
330
  into a file with multiple RRD manifests, which ``rrd verify`` rejects.
 
 
 
331
  """
332
 
333
  import shutil
334
 
335
  merged: Path = spec.data_dir / f"4danyone-{spec.token[:8]}.rrd"
336
- with merged.open("wb") as out:
337
- for name in LINK_ORDER:
338
- part: Path = spec.data_dir / "rrd" / f"{name}.rrd"
339
- if part.is_file():
340
- with part.open("rb") as source:
341
- shutil.copyfileobj(source, out, length=1 << 20)
 
 
 
 
342
  return merged
343
 
344
 
@@ -376,12 +385,24 @@ def pump(session: Session, work: Callable[[], T]) -> Iterator[str]:
376
 
377
  worker: threading.Thread = threading.Thread(target=target, daemon=True)
378
  worker.start()
379
- while True:
380
- label: str | None = session.events.get()
381
- if label is None:
382
- break
383
- yield label
384
- worker.join()
 
 
 
 
 
 
 
 
 
 
 
 
385
  if failure:
386
  raise failure[0]
387
  return outcome[0]
@@ -699,7 +720,7 @@ def generate_phase(recording: rr.RecordingStream, session: Session) -> Iterator[
699
  log_status(recording, "generate: six views published")
700
 
701
 
702
- def publish_phase(recording: rr.RecordingStream, spec: RunSpec, summary: dict) -> str:
703
  """Place the finished rig and its six videos, then switch to the final layout.
704
 
705
  Everything this reads is a file on disk, so it runs outside the GPU
@@ -707,10 +728,10 @@ def publish_phase(recording: rr.RecordingStream, spec: RunSpec, summary: dict) -
707
  """
708
 
709
  recording.reset_time()
710
- log_result(recording, Path(summary["result_dir"]))
711
  log_status(recording, "done")
712
  recording.send_blueprint(result_blueprint(spec.fps), make_active=True)
713
- elapsed: float = float(summary["total_pipeline_elapsed_seconds"])
714
  return f"Done in {elapsed:.1f}s. Six views at {float(spec.fps):.3f} FPS."
715
 
716
 
 
33
  from fdanyone.model.inference import _tensor_frames
34
  from fdanyone.model.tiny_decoder import decode_tiny_target_video, load_tiny_wan_decoder
35
  from fdanyone.motion.body import BodyMotion, load_body_motion
36
+ from fdanyone.pipeline import PreparedRun, RunResult, generate_run, prepare_run, warm_generation_models
37
  from fdanyone.video import choose_canonical_fps
38
  from fdanyone.viz import (
39
  BOX_COLOR,
 
110
  SMOKE_DELAY: float = float(os.environ.get("FDANYONE_SMOKE_DELAY", "0.1"))
111
  """Seconds between smoke yields; raise it to eyeball each phase in a browser."""
112
 
113
+ SMOKE_SUMMARY: RunResult = RunResult(
114
+ result_dir=Path(""), motion_dir=Path(""), total_elapsed_seconds=0.0, metadata={}
115
+ )
116
+ """Stand-in for the pipeline's published result; the smoke run writes no files."""
117
 
118
  # Body evaluation resolves SMPL-X relative to the vendored package's repository
119
  # root, which is not where a Space keeps its models. An absolute root wins the
 
283
  """The run this link was rebuilt from."""
284
  prepared: PreparedRun | None = None
285
  """Result of the motion phase, consumed by the generation phase."""
286
+ summary: RunResult | None = None
287
+ """Published result the generation phase returns, set once it succeeds."""
288
  source_frames: dict[int, RgbFrame] = field(default_factory=dict)
289
  """Decoded source stills for the diffusion pane, keyed by canonical frame."""
290
  events: queue.Queue[str | None] = field(default_factory=queue.Queue)
 
313
  raise RunCancelled("Stopped.")
314
 
315
 
316
+ LINK_ORDER: tuple[str, ...] = ("begin", "gpu", "publish")
317
  """Chain links in stream order, and therefore the part-concatenation order."""
318
 
319
 
 
330
 
331
  Parts are written with ``write_footer=False``: footered parts concatenate
332
  into a file with multiple RRD manifests, which ``rrd verify`` rejects.
333
+ Every link's part must exist and be non-empty — silently skipping one
334
+ would offer a plausible-looking partial recording — and the merge lands
335
+ atomically so a failure never leaves a truncated file at the final path.
336
  """
337
 
338
  import shutil
339
 
340
  merged: Path = spec.data_dir / f"4danyone-{spec.token[:8]}.rrd"
341
+ parts: list[Path] = [spec.data_dir / "rrd" / f"{name}.rrd" for name in LINK_ORDER]
342
+ missing: list[str] = [part.name for part in parts if not part.is_file() or part.stat().st_size == 0]
343
+ if missing:
344
+ raise FourDAnyoneError(f"The recording is incomplete; missing parts: {missing}.")
345
+ staging: Path = merged.with_name(merged.name + ".tmp")
346
+ with staging.open("wb") as out:
347
+ for part in parts:
348
+ with part.open("rb") as source:
349
+ shutil.copyfileobj(source, out, length=1 << 20)
350
+ os.replace(staging, merged)
351
  return merged
352
 
353
 
 
385
 
386
  worker: threading.Thread = threading.Thread(target=target, daemon=True)
387
  worker.start()
388
+ try:
389
+ while True:
390
+ label: str | None = session.events.get()
391
+ if label is None:
392
+ break
393
+ yield label
394
+ worker.join()
395
+ finally:
396
+ if worker.is_alive():
397
+ # The consumer vanished mid-run (client disconnect, generator
398
+ # close). The worker only halts at its next sentinel check, so
399
+ # request the stop and give it a bounded window; a daemon thread
400
+ # cannot block interpreter exit, but an unbounded join here would
401
+ # block the closing request thread for a whole denoise step.
402
+ request_stop(session.spec)
403
+ worker.join(timeout=5.0)
404
+ if worker.is_alive():
405
+ LOGGER.warning("Pipeline worker still unwinding after its consumer closed.")
406
  if failure:
407
  raise failure[0]
408
  return outcome[0]
 
720
  log_status(recording, "generate: six views published")
721
 
722
 
723
+ def publish_phase(recording: rr.RecordingStream, spec: RunSpec, summary: RunResult) -> str:
724
  """Place the finished rig and its six videos, then switch to the final layout.
725
 
726
  Everything this reads is a file on disk, so it runs outside the GPU
 
728
  """
729
 
730
  recording.reset_time()
731
+ log_result(recording, summary.result_dir)
732
  log_status(recording, "done")
733
  recording.send_blueprint(result_blueprint(spec.fps), make_active=True)
734
+ elapsed: float = summary.total_elapsed_seconds
735
  return f"Done in {elapsed:.1f}s. Six views at {float(spec.fps):.3f} FPS."
736
 
737
 
fdanyone_app.py CHANGED
@@ -1,6 +1,7 @@
1
  """Gradio and ZeroGPU wiring for the 4DAnyone Rerun streaming runtime.
2
 
3
- The four-link event chain is ``begin -> prepare_cpu -> run_gpu -> publish_cpu``.
 
4
  ZeroGPU runs decorated callbacks in forked children, so each link creates its
5
  own recording after ``rr.cleanup_if_forked_child()`` and uses the shared token
6
  as its recording ID. The merged ``run_gpu`` remains one allocation because its
@@ -26,6 +27,7 @@ from fdanyone.rerun_streaming import (
26
  EXAMPLE_DIR,
27
  GPU_DURATION,
28
  STREAM_SMOKE,
 
29
  RunSpec,
30
  Session,
31
  begin_phase,
@@ -75,7 +77,7 @@ def open_link(token: str, part: Path) -> Link:
75
  means the forked child. A stream inherited across the fork cannot be used:
76
  the SDK compares pids on flush and raises "Fork detected during flush". A
77
  stream built here belongs to this process, and sharing ``token`` as the
78
- recording id is what makes the viewer treat all four links as one recording.
79
 
80
  ``cleanup_if_forked_child`` drops whatever streams were inherited. The SDK
81
  registers it with ``os.register_at_fork`` already, and ``rr.init`` calls it
@@ -98,8 +100,14 @@ def open_link(token: str, part: Path) -> Link:
98
 
99
  def begin(
100
  video: str | None, start_time: float, seed: int
101
- ) -> Iterator[tuple[RunSpec, bytes | None, str, Any, Any]]:
102
- """Validate the input on CPU, open a recording, and switch to the outputs."""
 
 
 
 
 
 
103
 
104
  if video is None:
105
  raise gr.Error("Upload a video, or pick the bundled example.")
@@ -117,23 +125,14 @@ def begin(
117
  gr.Tabs(selected="outputs"),
118
  gr.DownloadButton(visible=False),
119
  )
120
- finally:
121
- link.close()
122
-
123
-
124
- def prepare_cpu(spec: RunSpec) -> Iterator[tuple[bytes | None, str]]:
125
- """Put the source clip on the frame timeline before any GPU work starts."""
126
-
127
- link: Link = open_link(spec.token, rrd_part(spec, "prepare"))
128
- try:
129
  label: str = source_phase(link.recording, spec)
130
- yield link.read(), format_status(spec, label)
131
  finally:
132
  link.close()
133
 
134
 
135
  @spaces.GPU(duration=GPU_DURATION)
136
- def run_gpu(spec: RunSpec) -> Iterator[tuple[dict | None, bytes | None, str]]:
137
  """Recover the motion and generate the six views, inside one allocation.
138
 
139
  The two phases have to share a process. ``PreparedRun`` carries the decoded
@@ -147,15 +146,17 @@ def run_gpu(spec: RunSpec) -> Iterator[tuple[dict | None, bytes | None, str]]:
147
 
148
  link: Link = open_link(spec.token, rrd_part(spec, "gpu"))
149
  session: Session = Session(spec)
 
 
150
  try:
151
- motion: Iterator[str] = (
152
  smoke_motion_phase(link.recording, session)
153
  if STREAM_SMOKE
154
  else motion_phase(link.recording, session)
155
  )
156
  for label in motion:
157
  yield None, link.read(), format_status(spec, label)
158
- generate: Iterator[str] = (
159
  smoke_generate_phase(link.recording, session)
160
  if STREAM_SMOKE
161
  else generate_phase(link.recording, session)
@@ -173,11 +174,16 @@ def run_gpu(spec: RunSpec) -> Iterator[tuple[dict | None, bytes | None, str]]:
173
  yield None, link.read(), f"Failed: {message}"
174
  raise gr.Error(message) from exc
175
  finally:
 
 
 
 
 
176
  link.close()
177
  yield session.summary, closing, format_status(spec, "Publishing the result.")
178
 
179
 
180
- def publish_cpu(spec: RunSpec, summary: dict | None) -> Iterator[tuple[bytes | None, str, Any]]:
181
  """Attach the finished rig and its six videos, off the GPU allocation."""
182
 
183
  if summary is None:
@@ -306,14 +312,13 @@ def build_demo() -> gr.Blocks:
306
  started = run_button.click(
307
  begin, [video, start_time, seed], [spec, viewer, status, tabs, download]
308
  )
309
- prepared = started.success(prepare_cpu, spec, [viewer, status])
310
- generated = prepared.success(run_gpu, spec, [summary, viewer, status])
311
  published = generated.success(publish_cpu, [spec, summary], [viewer, status, download])
312
  # `cancels` detaches the UI at once, but the pipeline itself only halts
313
  # when it reads the sentinel `request_stop` writes: the worker thread
314
  # (and on ZeroGPU, the forked child) never sees a Gradio cancellation.
315
  stop_button.click(
316
- request_stop, spec, status, cancels=[started, prepared, generated, published]
317
  )
318
  return demo
319
 
 
1
  """Gradio and ZeroGPU wiring for the 4DAnyone Rerun streaming runtime.
2
 
3
+ The three-link event chain is ``begin -> run_gpu -> publish_cpu``, one link
4
+ per process boundary.
5
  ZeroGPU runs decorated callbacks in forked children, so each link creates its
6
  own recording after ``rr.cleanup_if_forked_child()`` and uses the shared token
7
  as its recording ID. The merged ``run_gpu`` remains one allocation because its
 
27
  EXAMPLE_DIR,
28
  GPU_DURATION,
29
  STREAM_SMOKE,
30
+ RunResult,
31
  RunSpec,
32
  Session,
33
  begin_phase,
 
77
  means the forked child. A stream inherited across the fork cannot be used:
78
  the SDK compares pids on flush and raises "Fork detected during flush". A
79
  stream built here belongs to this process, and sharing ``token`` as the
80
+ recording id is what makes the viewer treat all the links as one recording.
81
 
82
  ``cleanup_if_forked_child`` drops whatever streams were inherited. The SDK
83
  registers it with ``os.register_at_fork`` already, and ``rr.init`` calls it
 
100
 
101
  def begin(
102
  video: str | None, start_time: float, seed: int
103
+ ) -> Iterator[tuple[Any, bytes | None, str, Any, Any]]:
104
+ """Validate, open a recording, switch to the outputs, and log the source.
105
+
106
+ One CPU link: the earlier split between "begin" and "prepare" callbacks
107
+ was an accidental process boundary — both ran on the CPU in this process,
108
+ each with its own recording lifecycle and part file. The first yield lands
109
+ the tab switch immediately; the second follows with the source clip.
110
+ """
111
 
112
  if video is None:
113
  raise gr.Error("Upload a video, or pick the bundled example.")
 
125
  gr.Tabs(selected="outputs"),
126
  gr.DownloadButton(visible=False),
127
  )
 
 
 
 
 
 
 
 
 
128
  label: str = source_phase(link.recording, spec)
129
+ yield gr.skip(), link.read(), format_status(spec, label), gr.skip(), gr.skip()
130
  finally:
131
  link.close()
132
 
133
 
134
  @spaces.GPU(duration=GPU_DURATION)
135
+ def run_gpu(spec: RunSpec) -> Iterator[tuple[RunResult | None, bytes | None, str]]:
136
  """Recover the motion and generate the six views, inside one allocation.
137
 
138
  The two phases have to share a process. ``PreparedRun`` carries the decoded
 
146
 
147
  link: Link = open_link(spec.token, rrd_part(spec, "gpu"))
148
  session: Session = Session(spec)
149
+ motion: Iterator[str] | None = None
150
+ generate: Iterator[str] | None = None
151
  try:
152
+ motion = (
153
  smoke_motion_phase(link.recording, session)
154
  if STREAM_SMOKE
155
  else motion_phase(link.recording, session)
156
  )
157
  for label in motion:
158
  yield None, link.read(), format_status(spec, label)
159
+ generate = (
160
  smoke_generate_phase(link.recording, session)
161
  if STREAM_SMOKE
162
  else generate_phase(link.recording, session)
 
174
  yield None, link.read(), f"Failed: {message}"
175
  raise gr.Error(message) from exc
176
  finally:
177
+ # Close the phase generators FIRST: their unwinding (pump's bounded
178
+ # worker join) must finish while the recording is still connected.
179
+ for phase in (motion, generate):
180
+ if phase is not None:
181
+ phase.close()
182
  link.close()
183
  yield session.summary, closing, format_status(spec, "Publishing the result.")
184
 
185
 
186
+ def publish_cpu(spec: RunSpec, summary: RunResult | None) -> Iterator[tuple[bytes | None, str, Any]]:
187
  """Attach the finished rig and its six videos, off the GPU allocation."""
188
 
189
  if summary is None:
 
312
  started = run_button.click(
313
  begin, [video, start_time, seed], [spec, viewer, status, tabs, download]
314
  )
315
+ generated = started.success(run_gpu, spec, [summary, viewer, status])
 
316
  published = generated.success(publish_cpu, [spec, summary], [viewer, status, download])
317
  # `cancels` detaches the UI at once, but the pipeline itself only halts
318
  # when it reads the sentinel `request_stop` writes: the worker thread
319
  # (and on ZeroGPU, the forked child) never sees a Gradio cancellation.
320
  stop_button.click(
321
+ request_stop, spec, status, cancels=[started, generated, published]
322
  )
323
  return demo
324
 
tests/test_pipeline_lifecycle.py CHANGED
@@ -62,11 +62,17 @@ def test_generate_run_settles_and_discards_on_success(
62
  conditioning: _Conditioning = _Conditioning()
63
  prepared: PreparedRun = _prepared(tmp_path, conditioning)
64
  monkeypatch.setattr(fdanyone.model.inference, "generate_views", lambda **_: object())
65
- monkeypatch.setattr(fdanyone.output, "export_result", lambda **_: {"fps": Fraction(25, 1)})
 
 
 
 
66
 
67
- summary: dict = generate_run(prepared, seed=0)
68
 
69
- assert summary["result_dir"] == str(prepared.result_dir)
 
 
70
  assert conditioning.settled == 1
71
  assert not prepared.scratch.exists()
72
 
@@ -94,7 +100,11 @@ def test_generate_run_discards_even_when_settlement_raises(
94
  conditioning: _Conditioning = _Conditioning(settlement_error=RuntimeError("render worker died"))
95
  prepared: PreparedRun = _prepared(tmp_path, conditioning)
96
  monkeypatch.setattr(fdanyone.model.inference, "generate_views", lambda **_: object())
97
- monkeypatch.setattr(fdanyone.output, "export_result", lambda **_: {})
 
 
 
 
98
 
99
  with pytest.raises(RuntimeError, match="render worker died"):
100
  generate_run(prepared, seed=0)
 
62
  conditioning: _Conditioning = _Conditioning()
63
  prepared: PreparedRun = _prepared(tmp_path, conditioning)
64
  monkeypatch.setattr(fdanyone.model.inference, "generate_views", lambda **_: object())
65
+ monkeypatch.setattr(
66
+ fdanyone.output,
67
+ "export_result",
68
+ lambda **_: {"fps": Fraction(25, 1), "total_pipeline_elapsed_seconds": 1.5},
69
+ )
70
 
71
+ result = generate_run(prepared, seed=0)
72
 
73
+ assert result.result_dir == prepared.result_dir
74
+ assert result.total_elapsed_seconds == 1.5
75
+ assert result.metadata["result_dir"] == str(prepared.result_dir)
76
  assert conditioning.settled == 1
77
  assert not prepared.scratch.exists()
78
 
 
100
  conditioning: _Conditioning = _Conditioning(settlement_error=RuntimeError("render worker died"))
101
  prepared: PreparedRun = _prepared(tmp_path, conditioning)
102
  monkeypatch.setattr(fdanyone.model.inference, "generate_views", lambda **_: object())
103
+ monkeypatch.setattr(
104
+ fdanyone.output,
105
+ "export_result",
106
+ lambda **_: {"total_pipeline_elapsed_seconds": 0.0},
107
+ )
108
 
109
  with pytest.raises(RuntimeError, match="render worker died"):
110
  generate_run(prepared, seed=0)