pablovela5620 Claude Fable 5 commited on
Commit
4c983fd
·
1 Parent(s): 326f401

Settle generation scratch in the worker that owns it

Browse files

Four-reviewer simplify pass over the Stop change. A stopped run leaked
its scratch tree: Gradio closes the generator before the sentinel stops
the worker, so the worker_finished cleanup guard never fired. Releasing
in the worker's own finally covers return, failure, and stop alike, and
deletes the flag plus the conditional cleanup. The stop test now proves
a mid-run stop deterministically instead of an abort-at-entry, and the
session fixture takes data_dir instead of being duplicated.

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

Files changed (2) hide show
  1. fdanyone_app.py +17 -23
  2. tests/test_app_helpers.py +20 -28
fdanyone_app.py CHANGED
@@ -518,8 +518,6 @@ class Session:
518
  """Decoded source stills for the diffusion pane, keyed by canonical frame."""
519
  events: queue.Queue[str | None] = field(default_factory=queue.Queue)
520
  """Stage labels the pipeline hooks push from the worker thread."""
521
- worker_finished: bool = False
522
- """Whether the last ``_pump`` call's pipeline thread returned."""
523
 
524
 
525
  class RunCancelled(FourDAnyoneError):
@@ -557,7 +555,6 @@ def _pump(session: Session, work: Callable[[], T]) -> Iterator[str]:
557
 
558
  outcome: list[T] = []
559
  failure: list[BaseException] = []
560
- session.worker_finished = False
561
 
562
  def target() -> None:
563
  try:
@@ -575,7 +572,6 @@ def _pump(session: Session, work: Callable[[], T]) -> Iterator[str]:
575
  break
576
  yield label
577
  worker.join()
578
- session.worker_finished = True
579
  if failure:
580
  raise failure[0]
581
  return outcome[0]
@@ -803,31 +799,29 @@ def generate_phase(link: Link, session: Session) -> Iterator[str]:
803
 
804
  if session.prepared is None:
805
  raise FourDAnyoneError("The motion phase did not finish.")
 
806
  session.source_frames = _decode_source_stills(session.spec)
807
  link.recording.send_blueprint(diffusion_blueprint(), make_active=True)
808
  yield "Generating six views."
809
 
810
  def work() -> dict:
811
- return generate_run(
812
- session.prepared, seed=session.spec.seed, on_denoise_step=_denoise_hook(link, session)
813
- )
 
 
 
 
 
814
 
815
- try:
816
- steps: Iterator[str] = _pump(session, work)
817
- while True:
818
- try:
819
- label: str = next(steps)
820
- except StopIteration as done:
821
- session.summary = done.value
822
- break
823
- yield f"Denoising: {label} of {SETTINGS.num_inference_steps}."
824
- finally:
825
- # Stop closes this generator while the worker is still inside
826
- # generate_run, and its scratch directory must not be deleted out from
827
- # under it. Only a run whose pipeline call returned is safe to settle;
828
- # a cancelled one leaks its scratch into the Space's ephemeral disk.
829
- if session.worker_finished and session.prepared is not None:
830
- release_run(session.prepared)
831
 
832
  _status(link.recording, "generate: six views published")
833
  yield "Publishing the result."
 
518
  """Decoded source stills for the diffusion pane, keyed by canonical frame."""
519
  events: queue.Queue[str | None] = field(default_factory=queue.Queue)
520
  """Stage labels the pipeline hooks push from the worker thread."""
 
 
521
 
522
 
523
  class RunCancelled(FourDAnyoneError):
 
555
 
556
  outcome: list[T] = []
557
  failure: list[BaseException] = []
 
558
 
559
  def target() -> None:
560
  try:
 
572
  break
573
  yield label
574
  worker.join()
 
575
  if failure:
576
  raise failure[0]
577
  return outcome[0]
 
799
 
800
  if session.prepared is None:
801
  raise FourDAnyoneError("The motion phase did not finish.")
802
+ prepared: PreparedRun = session.prepared
803
  session.source_frames = _decode_source_stills(session.spec)
804
  link.recording.send_blueprint(diffusion_blueprint(), make_active=True)
805
  yield "Generating six views."
806
 
807
  def work() -> dict:
808
+ # The worker owns the scratch, so it must also settle it: a Stop closes
809
+ # the outer generator, which then never sees the pipeline call end.
810
+ try:
811
+ return generate_run(
812
+ prepared, seed=session.spec.seed, on_denoise_step=_denoise_hook(link, session)
813
+ )
814
+ finally:
815
+ release_run(prepared)
816
 
817
+ steps: Iterator[str] = _pump(session, work)
818
+ while True:
819
+ try:
820
+ label: str = next(steps)
821
+ except StopIteration as done:
822
+ session.summary = done.value
823
+ break
824
+ yield f"Denoising: {label} of {SETTINGS.num_inference_steps}."
 
 
 
 
 
 
 
 
825
 
826
  _status(link.recording, "generate: six views published")
827
  yield "Publishing the result."
tests/test_app_helpers.py CHANGED
@@ -4,6 +4,7 @@ from __future__ import annotations
4
 
5
  import os
6
  import pickle
 
7
  import traceback
8
  import uuid
9
  from collections.abc import Callable
@@ -75,7 +76,7 @@ def test_probe_clip_rejects_a_missing_file(tmp_path: Path) -> None:
75
  fdanyone_app.probe_clip(tmp_path / "absent.mp4", 0.0)
76
 
77
 
78
- def _pump_session() -> fdanyone_app.Session:
79
  """A Session with only the fields `_pump` touches."""
80
 
81
  return fdanyone_app.Session(
@@ -85,7 +86,7 @@ def _pump_session() -> fdanyone_app.Session:
85
  start_time=0.0,
86
  seed=0,
87
  fps=Fraction(25, 1),
88
- data_dir=Path("/nonexistent"),
89
  )
90
  )
91
 
@@ -107,45 +108,38 @@ def test_pump_yields_hook_labels_and_returns_the_work_result() -> None:
107
  assert done.value == "prepared"
108
  break
109
  assert seen == ["bboxes", "smplx"]
110
- assert session.worker_finished is True
111
 
112
 
113
- def test_stop_sentinel_unwinds_the_pump_worker(tmp_path: Path) -> None:
114
  """Stop must reach the pipeline thread itself, not just Gradio's generator.
115
 
116
  ``cancels`` closes the request generator, but the worker thread (or the
117
  forked ZeroGPU child) never hears it; the sentinel checked by the hooks is
118
- the only thing that actually frees the GPU.
 
119
  """
120
 
121
- session: fdanyone_app.Session = _pump_session()
122
- spec: fdanyone_app.RunSpec = fdanyone_app.RunSpec(
123
- token=session.spec.token,
124
- video_path=session.spec.video_path,
125
- start_time=0.0,
126
- seed=0,
127
- fps=Fraction(25, 1),
128
- data_dir=tmp_path,
129
- )
130
- session = fdanyone_app.Session(spec)
131
  assert fdanyone_app.request_stop(None) == "Nothing is running."
132
- fdanyone_app.request_stop(spec)
133
- assert (tmp_path / "stop-requested").exists()
134
 
135
- ticks: list[int] = []
 
136
 
137
  def work() -> str:
138
- # What every pipeline hook does on entry, once per stage or step.
139
- for index in range(1000):
140
- fdanyone_app._check_stop(spec)
141
- ticks.append(index)
 
 
142
  return "finished"
143
 
 
 
 
 
 
144
  with pytest.raises(fdanyone_app.RunCancelled):
145
- list(fdanyone_app._pump(session, work))
146
- assert ticks == []
147
- # The worker joined before re-raising, so the thread is truly gone.
148
- assert session.worker_finished is True
149
 
150
 
151
  def test_pump_reraises_a_worker_failure_on_the_caller_thread() -> None:
@@ -158,8 +152,6 @@ def test_pump_reraises_a_worker_failure_on_the_caller_thread() -> None:
158
 
159
  with pytest.raises(FourDAnyoneError, match="boom"):
160
  list(fdanyone_app._pump(session, work))
161
- # The worker joined before re-raising, so releasing its scratch is safe.
162
- assert session.worker_finished is True
163
 
164
 
165
  def test_blueprints_build_for_every_phase() -> None:
 
4
 
5
  import os
6
  import pickle
7
+ import threading
8
  import traceback
9
  import uuid
10
  from collections.abc import Callable
 
76
  fdanyone_app.probe_clip(tmp_path / "absent.mp4", 0.0)
77
 
78
 
79
+ def _pump_session(data_dir: Path = Path("/nonexistent")) -> fdanyone_app.Session:
80
  """A Session with only the fields `_pump` touches."""
81
 
82
  return fdanyone_app.Session(
 
86
  start_time=0.0,
87
  seed=0,
88
  fps=Fraction(25, 1),
89
+ data_dir=data_dir,
90
  )
91
  )
92
 
 
108
  assert done.value == "prepared"
109
  break
110
  assert seen == ["bboxes", "smplx"]
 
111
 
112
 
113
+ def test_stop_sentinel_unwinds_a_running_pump_worker(tmp_path: Path) -> None:
114
  """Stop must reach the pipeline thread itself, not just Gradio's generator.
115
 
116
  ``cancels`` closes the request generator, but the worker thread (or the
117
  forked ZeroGPU child) never hears it; the sentinel checked by the hooks is
118
+ the only thing that actually frees the GPU. The barrier makes the order
119
+ deterministic: the worker is provably mid-run when Stop arrives.
120
  """
121
 
 
 
 
 
 
 
 
 
 
 
122
  assert fdanyone_app.request_stop(None) == "Nothing is running."
 
 
123
 
124
+ session: fdanyone_app.Session = _pump_session(tmp_path)
125
+ stop_requested: threading.Event = threading.Event()
126
 
127
  def work() -> str:
128
+ # A pipeline hook checks the sentinel once per stage; model the stage
129
+ # boundary the run is inside when the user presses Stop.
130
+ fdanyone_app._check_stop(session.spec)
131
+ session.events.put("running")
132
+ assert stop_requested.wait(timeout=5.0)
133
+ fdanyone_app._check_stop(session.spec)
134
  return "finished"
135
 
136
+ pump = fdanyone_app._pump(session, work)
137
+ assert next(pump) == "running"
138
+ fdanyone_app.request_stop(session.spec)
139
+ assert (tmp_path / "stop-requested").exists()
140
+ stop_requested.set()
141
  with pytest.raises(fdanyone_app.RunCancelled):
142
+ list(pump)
 
 
 
143
 
144
 
145
  def test_pump_reraises_a_worker_failure_on_the_caller_thread() -> None:
 
152
 
153
  with pytest.raises(FourDAnyoneError, match="boom"):
154
  list(fdanyone_app._pump(session, work))
 
 
155
 
156
 
157
  def test_blueprints_build_for_every_phase() -> None: