From 75df49650d36156ad7de5782441dc494839edbe5 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 18:48:56 -0500 Subject: [PATCH 01/16] fix(engine): #82 #83 #84 - run-time warnings reach the caller, carried frame rates, peak RSS clamp #82: a warning a step raises about what it is writing (the >=6 dB level spread on an unmatched join) is emitted as a 'warning' event and folded into the job's warnings, prefixed with the step - the server's own log is not a consumer surface. #83: host_memory_peak_rss_mb is held at or above host_memory_rss_mb; the two readings come from different sources and disagree by a megabyte at idle, which makes 'peak - rss' read as un-comparable. #84: AudioVideo carries the frame rate it is meant to play at - a task's own fps, the rate of a file it read, a chain's fps - and result.fps defaults to that before falling back to 8. Declaring a rate the frames contradict now warns. Also: docs/proposals/acknowledged-cost-binding.md, the #85 assessment. --- docs/SERVER.md | 1 + docs/TASKS.md | 10 +- docs/proposals/acknowledged-cost-binding.md | 177 ++++++++++++++++++++ dw/events.py | 17 ++ dw/host_memory.py | 17 ++ dw/pipeline_processors/chain.py | 6 +- dw/result.py | 57 ++++++- dw/server/jobs.py | 16 +- dw/tasks/audio_utils.py | 12 +- dw/tasks/concat_videos.py | 8 +- dw/tasks/dissolve_videos.py | 6 +- dw/tasks/interpolate_frames.py | 9 +- dw/tasks/pair_audio.py | 4 +- dw/tasks/stabilize.py | 2 +- dw/tasks/task.py | 2 +- dw/tasks/video_utils.py | 6 +- dw/workflow_schema.json | 2 +- tests/test_concat_videos.py | 65 +++++++ tests/test_dissolve_videos.py | 12 ++ tests/test_host_memory.py | 26 +++ tests/test_job_progress.py | 43 +++++ tests/test_result.py | 83 +++++++++ tests/test_video_utils.py | 9 + 23 files changed, 564 insertions(+), 26 deletions(-) create mode 100644 docs/proposals/acknowledged-cost-binding.md diff --git a/docs/SERVER.md b/docs/SERVER.md index 5a0ef570..e31dc823 100644 --- a/docs/SERVER.md +++ b/docs/SERVER.md @@ -173,6 +173,7 @@ Every event in the stream carries a `seq` and an `event` name: | `pipeline_step` | each denoise step | `step`, `total_steps`. Emitted for a pipeline that takes a `callback_on_step_end`, and for a `ModularPipeline` (H3, LTX-2, Qwen-Image), which takes none - there the denoise block's own progress bar is what reports | | `phase` | the step changes what it is doing | `phase`, `detail` | | `pipeline_released` | a step with `release_pipeline` drops its pipeline | `step`, `index`, `gpu_memory_allocated_mb` and `gpu_memory_allocated_before_mb` (both `null` where the backend cannot say). Emitted between the step's generation and its files being written, which is where the release happens - so the ordering is readable off the event stream rather than by trying to poll memory through a sub-second write | +| `warning` | a step finds something wrong with what it is about to write | `message`, plus a `kind` and the figures behind it (`level_spread`: `spread_db`, `measure`, `command`; `fps_mismatch`: `declared_fps`, `source_fps`). Also appended to the job's `warnings`, prefixed with the step it fired in - the event keeps the moment, `warnings` keeps it where a caller polling the finished job will look, since a warning about the artifact outlives the run that noticed it | | `workflow_end` | the run finishes | `manifest` | A step spends most of its wall clock outside the denoise loop, and diff --git a/docs/TASKS.md b/docs/TASKS.md index a1339842..cef619f2 100644 --- a/docs/TASKS.md +++ b/docs/TASKS.md @@ -204,7 +204,7 @@ video generation" in the workflow guide): | `crossfade_ms` | No | Equal-power crossfade at each audio seam, drawn from the trimmed material - no effect when `trim_frames` is 0, and validation warns when one is written there (default: 75) | | `audio_bleed_ms` | No | How long the outgoing video's tail rings on over the head of the next one, at seams with nothing trimmed to crossfade (default: 0, off) | | `seam_fade_ms` | No | Fade on each side of a seam that gets neither a crossfade nor a bleed - for tonal material, not for a continuous bed (default: 3, just enough not to click) | -| `fps` | No | Frame rate of the videos - required to join audio when trimming | +| `fps` | No | Frame rate of the videos - required to join audio when trimming, and the rate the joined file is written at unless `result.fps` overrides it | | `match_levels` | No | Even the shots' loudness out before joining - `"rms"` for perceived level (the measurement `get_gallery_metadata` reports as `mean_dbfs`), `"peak"` for the loudest sample. Off by default | | `match_levels_dbfs` | No | The level `match_levels` moves every shot to (default: -1 dBFS for `peak`, -20 dBFS for `rms`) | @@ -294,8 +294,10 @@ same" means, and `"peak"` matches the loudest sample, which is the safer choice on material with big transients. A shot whose gain would clip at the target is held just below full scale and the log says so. Left off - the default, so nothing existing changes - a spread of 6 dB or more across the -tracks being joined is logged as a warning rather than passing in silence. -`dissolve_videos` takes the same pair. +tracks being joined is reported as a warning rather than passing in silence: +on the job's `warnings` and as a `warning` event in its stream, not only in +the server's log, since the caller who can act on it is the one who asked for +the run. `dissolve_videos` takes the same pair. ### dissolve_videos @@ -327,7 +329,7 @@ montage cut to a score wants: | `fade_in_frames` | No | Frames over which the first video rises out of `fade_color` (default: 0) | | `fade_out_frames` | No | Frames over which the last video sinks into it (default: 0) | | `fade_color` | No | The RGB colour the fades come from and go to (default: black) | -| `fps` | No | Frame rate of the videos - required to crossfade audio at a dissolve | +| `fps` | No | Frame rate of the videos - required to crossfade audio at a dissolve, and the rate the dissolved file is written at unless `result.fps` overrides it | | `match_levels` | No | Even the shots' loudness out before joining - `"rms"` or `"peak"`, as with [`concat_videos`](#concat_videos). Off by default | | `match_levels_dbfs` | No | The level `match_levels` moves every shot to (default: -1 dBFS for `peak`, -20 dBFS for `rms`) | diff --git a/docs/proposals/acknowledged-cost-binding.md b/docs/proposals/acknowledged-cost-binding.md new file mode 100644 index 00000000..ff0a41af --- /dev/null +++ b/docs/proposals/acknowledged-cost-binding.md @@ -0,0 +1,177 @@ +# Proposal: bind `acknowledged_cost` to an estimate, not just to a boolean + +Status: **design only** - written in answer to issue #85 (forum feedback), +after reading the gate and every path by which a run's size is decided. +No code changes yet. Written by the implementer agent (model `opus`, +provider `anthropic`). + +## The question asked + +> The acknowledgement should bind to an estimated budget or operation +> fingerprint, not just a boolean, otherwise the plan can change after +> consent while the flag remains true. + +## What the gate is today + +It is a boolean, checked in the MCP layer only: + +- `dw_mcp/diagnose.py:45` - `run_workflow` raises `COST_REFUSAL` unless + `acknowledged_cost` is truthy; `rerun_job` the same at `:248`. +- `dw_mcp/models.py` / `dw_mcp/workspaces.py` - `download_model`, + `delete_model`, `update_diffusers`, `delete_workspace` the same. +- `POST /api/jobs` does **not** require it. The gate is an agent-behaviour + gate - "say the number out loud to a human before you spend it" - not a + resource guard, and the HTTP API and the web UI queue work without it. + +Nothing is bound. The flag records that *something* was consented to, not +what. The figure the agent quoted comes from the catalog's `cost` block +(`list_workflows`), which is measured against the workflow's *stored +defaults*; the run is queued with the caller's `arguments`, which is a +different document. + +## Where the plan can grow after consent + +Each of these is reachable with `acknowledged_cost=true` and a quote that +was honest when it was made: + +1. **List fan-out.** A `for_each` step is expanded per entry + (`dw/for_each.py`, 32 entries max). `cost.minutes` is the measured cost + of the *default* list; a caller passing a 12-entry `shots` list through + `arguments` runs 12 shots. `cost.per_entry` exists precisely to price + this - and is advisory, unenforced, and set on no template yet. +2. **Cartesian product.** Several `previous_result:` references in one step + multiply (4 images x 3 masks = 12 iterations), and the multiplicands can + come from arguments. +3. **Plain numeric arguments.** `num_images_per_prompt`, `num_frames`, + `num_inference_steps` scale the run roughly linearly and are ordinary + variables. +4. **A model download mid-run.** Weights not on disk are pulled by + `from_pretrained` when the step reaches it - tens of GB and many minutes, + appearing in no `cost` block. This is the forum comment's sharpest case: + the same workflow costs 4 minutes on a warm box and 50 on a cold one. +5. **`inline_workflow`.** No catalog entry, so no `cost` at all - the quote + is whatever the agent believed. +6. **`rerun_job(new_seed=true)`.** Draws a fresh seed, which defeats the + step cache; the rerun of a "finished instantly" job is a full generation. +7. **Sub-workflows.** A `composes-workflows` step's real cost is the child's. + +So the premise in #85 holds: the flag stays `true` while the work grows, and +nothing at queue time compares what will run against what was acknowledged. + +## What the engine already has + +Almost all of the raw material: + +- `POST /api/validate` is free, takes the caller's `arguments`, and already + folds them exactly as the run will (`argument_errors`), expands `for_each` + and reports per-entry problems. +- `expand_for_each` yields the exact member list a run will execute. +- `realize_workflow` (`dw/realize.py`) already produces the canonical + realized document - arguments folded, seed pinned, prompts inlined, + `output:latest` resolved - which is the natural thing to fingerprint. +- `hub_cache.scan_models` knows what is on disk, so "which repos this run + will have to download first" is answerable before queuing. +- The manifest and `workflow.json` make the run inspectable afterwards, + which is what makes an estimate checkable rather than decorative. + +The missing pieces are an estimate in the pre-flight answer, and a check at +queue time. + +## Proposed change + +### 1. `POST /api/validate` returns a plan + +Beside `valid`/`errors`/`warnings`, on a valid answer: + +```json +"plan": { + "fingerprint": "sha256:9f13…", + "steps": 14, + "iterations": 14, + "list_entries": {"shots": 12}, + "downloads_required": [{"repo": "MiniMaxAI/MiniMax-H3", "gb": 41.2}], + "estimate": {"minutes": 38.0, "basis": "per_entry", "confidence": "measured"} +} +``` + +- `fingerprint` is a SHA-256 over the plan-shaping inputs: the realized + workflow (`realize.py`'s document, minus the seed), the expanded step + names, and the iteration count per step. It changes when the *work* + changes and not when something cosmetic does. +- `estimate.basis` is one of `per_entry` (fixed + per-entry x N), + `catalog` (the stored `cost.minutes`, defaults only), or `unknown` + (an inline workflow with no cost block) - the honesty is in the field, + not in a fabricated number. +- `downloads_required` closes case 4 on its own, and is useful with or + without the rest of this proposal. + +### 2. `acknowledged_cost` accepts what was acknowledged + +`run_workflow` / `rerun_job` keep taking `true`, and additionally take an +object: + +```json +"acknowledged_cost": {"fingerprint": "sha256:9f13…", "minutes": 38.0} +``` + +`POST /api/jobs` recomputes the plan for the arguments it was actually +given and refuses with **409** when the fingerprint differs or the recomputed +estimate exceeds the acknowledged minutes by more than a tolerance (25%, +settable), naming both figures and what changed. The agent's recovery is to +re-quote to the user - which is the behaviour the gate was for. + +### 3. Bare `true` stays legal, and says so + +An unbound `true` keeps working - the web UI, the CLI-equivalent callers and +every existing script depend on it, and a hard requirement would be a +breaking change to every MCP consumer. But the job records which kind of +acknowledgement it got (`acknowledged: "boolean" | "bound"`), so "was this +run consented to at its actual size?" is answerable after the fact, and the +skills can teach the bound form as the normal one. + +## Alternatives considered + +- **A hard budget ceiling** (`max_minutes`, run aborted when exceeded). + Needs runtime metering the engine does not have, and killing a 90%-done + video render to honour an estimate wastes exactly the resource the gate + protects. The fingerprint check is pre-flight, which is where a refusal is + free. +- **Requiring the bound form.** Breaking, and buys little over recording + which form was used. +- **Estimating cost server-side with no acknowledgement change.** Half the + value (the agent can quote better) with none of the binding - the plan can + still change between the quote and the call. +- **Doing nothing.** Defensible: the gate is a prompt-discipline device, and + the human is in the loop by construction. But the forum comment's case - + "validates cheaply, expands at runtime" - is real on this engine today + (cases 1, 4 and 6 above), and the plan-at-validate half is cheap and + useful even alone. + +## Scope + +- `dw/plan.py` (new): fingerprint + estimate from a definition and + arguments, reusing `realize_workflow`, `expand_for_each` and `hub_cache`. +- `dw/server/app.py`: `plan` on the validate answer; the 409 check on + `POST /api/jobs`. +- `dw/server/jobs.py`: record the acknowledgement form on the job and in + `jobs.sqlite`. +- `dw_mcp/diagnose.py` + `dw_mcp/server.py`: accept the object form, surface + `plan` from `validate_workflow`. +- Docs: SERVER.md, MCP.md, WORKFLOW_GUIDE.md's authoring section, and the + three plugin skills' cost step. +- Tests: fingerprint stability against cosmetic edits, change under a longer + list, the 409, the tolerance, and the boolean path unchanged. + +Roughly a two-stage piece of work: stage 1 the plan on validate (useful on +its own), stage 2 the binding and the 409. + +## Open questions for approval + +1. Is the bound form worth it at all, given that a human is already in the + loop on every acknowledgement? (Doing nothing is a legitimate answer; + stage 1 alone is another.) +2. Tolerance: 25% of the acknowledged minutes, or a fingerprint-only check + with no numeric comparison at all? A numeric one needs `per_entry` on the + templates to be meaningful, and today no template carries it. +3. Should `POST /api/jobs` grow the gate for HTTP callers too, or stay an + MCP-layer concept? (The web UI would have to send something.) diff --git a/dw/events.py b/dw/events.py index d6f8db50..77856880 100644 --- a/dw/events.py +++ b/dw/events.py @@ -102,6 +102,23 @@ def deactivate_context(token): NON_INTERRUPTIBLE_PHASES = ("loading", "task") +def emit_warning(message, **data): + """Report something the run's result carries but its status will not. + + A warning a step discovers at run time - shots being cut together 10 dB + apart, a video about to be written at a frame rate nothing chose - is + only useful where whoever asked for the run can read it. The server's + own log is not that place: a consumer over the API or MCP sees the event + stream and the job's `warnings` list and nothing else, so a diagnostic + that only reaches the log does not exist out there (#82). + + Logged as well as emitted, because the CLI and the REPL have no event + sink and the log is the whole of their surface. + """ + logger.warning(message) + get_context().emit("warning", message=message, **data) + + def emit_phase(phase, detail=None): """Report a coarse phase change on the active run. diff --git a/dw/host_memory.py b/dw/host_memory.py index 89b48340..04291270 100644 --- a/dw/host_memory.py +++ b/dw/host_memory.py @@ -58,6 +58,23 @@ def host_memory_stats(): if all(stats[key] is not None for key in ("rss_mb", "total_mb")): break + return _hold_the_high_water_mark(stats) + + +def _hold_the_high_water_mark(stats): + """Keep peak_rss_mb >= rss_mb, which is what a high-water mark means. + + The two readings come from different places - getrusage's ru_maxrss, + quantized to whole pages and taken first, against psutil's rss taken a + moment later - so a process that has never peaked meaningfully above its + current size reports them within a megabyte of each other in either + order. `peak - rss` is the whole point of the pair (what a run took and + did not give back), and a small negative there reads as "these fields + are not comparable" rather than "nothing leaked" (#83). + """ + peak, rss = stats["peak_rss_mb"], stats["rss_mb"] + if peak is not None and rss is not None and peak < rss: + stats["peak_rss_mb"] = rss return stats diff --git a/dw/pipeline_processors/chain.py b/dw/pipeline_processors/chain.py index 690c6616..10e39f1d 100644 --- a/dw/pipeline_processors/chain.py +++ b/dw/pipeline_processors/chain.py @@ -381,9 +381,11 @@ def run_chain(pipeline, chain_definition, arguments): # is muxed in so the soundtrack has no seams if spill is None: frames = frames[: config.total_frames] - return AudioVideo(frames, config.source_audio, config.source_rate) + return AudioVideo( + frames, config.source_audio, config.source_rate, fps=config.fps + ) - return AudioVideo(frames, audio, audio_rate) + return AudioVideo(frames, audio, audio_rate, fps=config.fps) class ChainConfig: diff --git a/dw/result.py b/dw/result.py index b837b9c1..b40fe1b2 100644 --- a/dw/result.py +++ b/dw/result.py @@ -12,7 +12,7 @@ is_av_available, ) from collections.abc import Mapping -from .events import emit_phase +from .events import emit_phase, emit_warning from .security import ( SecurityError, validate_file_base_name, @@ -25,6 +25,10 @@ # Result saving constants MAX_BASE_NAME_LENGTH = 200 DEFAULT_AUDIO_SAMPLE_RATE = 44100 +# The rate a video is written at when neither the workflow nor the artifact +# says - a diffusers convention old enough that changing it would restate +# every existing workflow's output +DEFAULT_VIDEO_FPS = 8 def output_file_path(output_dir, file_name): @@ -109,16 +113,24 @@ class AudioVideo: the result mux them into one file instead of dropping the audio on the floor. """ - def __init__(self, frames, audio, sample_rate): + def __init__(self, frames, audio, sample_rate, fps=None): """ Args: frames: The video, as PIL images or an array of frames audio: Waveform for this video, shaped (channels, samples) sample_rate: Sample rate of the waveform, or None if the pipeline did not report one + fps: Frame rate these frames are meant to play at, when something + knows it - a joined video's own rate, or the rate of the file + a task read. Carried for the same reason AudioTrack carries + its sample rate: `result.fps` defaults to 8, and a step that + joins 24 fps shots writing them at 8 is three times slow with + its audio still the right length (#84). A declared + `result.fps` still wins over this """ self.frames = frames self.audio = audio self.sample_rate = sample_rate + self.fps = fps class AudioTrack: @@ -402,13 +414,9 @@ def save_artifact( if isinstance(artifact, AudioVideo): self.save_audio_video(artifact, output_path, content_type) else: - export_to_video( - artifact, output_path, fps=self.result_definition.get("fps", 8) - ) + export_to_video(artifact, output_path, fps=self.video_fps(artifact)) elif content_type == "image/gif": - export_to_gif( - artifact, output_path, fps=self.result_definition.get("fps", 8) - ) + export_to_gif(artifact, output_path, fps=self.video_fps(artifact)) elif content_type.startswith("audio"): waveforms = normalize_audio(artifact) # Declared rate > the rate a generated track carries > default @@ -472,6 +480,37 @@ def save_artifact( return [output_path] + def video_fps(self, artifact): + """The frame rate this video is written at. + + Declared `result.fps` first, then the rate the artifact carries (a + join's own rate, or the rate of the files it read), then 8. + + The order matters more than it looks: `result.fps` and a task's own + `fps` argument are separate knobs, and the one an author thinks to + set is the task's. A step that told `concat_videos` its shots are 24 + fps and said nothing on `result` used to write them at 8 - the + picture three times long against an audio track still the right + length, with nothing said about it (#84). A workflow that does + declare `result.fps` still wins, so writing at a rate other than the + source's - a deliberate slow motion - stays available, and says so. + """ + declared = self.result_definition.get("fps") + carried = getattr(artifact, "fps", None) + if declared is None: + return carried or DEFAULT_VIDEO_FPS + if carried and abs(declared - carried) > 0.01: + emit_warning( + f"Writing video at {declared} fps, but the frames it was " + f"given run at {carried} fps - the file will play " + f"{carried / declared:.2g}x speed. Drop 'fps' from the step's " + f"result to keep the source rate", + kind="fps_mismatch", + declared_fps=declared, + source_fps=carried, + ) + return declared + def save_audio_video(self, artifact, output_path, content_type): """Write a video and the audio generated with it into a single file. @@ -484,7 +523,7 @@ def save_audio_video(self, artifact, output_path, content_type): output_path: Path of the file to write content_type: MIME type of the video being written """ - fps = self.result_definition.get("fps", 8) + fps = self.video_fps(artifact) # The pipeline reports the sample rate of what it generated - the result # definition can still override it sample_rate = self.result_definition.get( diff --git a/dw/server/jobs.py b/dw/server/jobs.py index 73f2d55d..98346c4e 100644 --- a/dw/server/jobs.py +++ b/dw/server/jobs.py @@ -342,7 +342,9 @@ def __init__(self, spec): self.started_at = None self.finished_at = None self.manifest = [] - self.warnings = spec.get("warnings", []) + # A copy: run-time warnings are appended to this list (see + # _note_progress) and the spec is what a rerun is built from + self.warnings = list(spec.get("warnings", [])) self.error = None self.traceback = None # Which run this job turned out to be - reported by the worker's @@ -400,6 +402,18 @@ def _note_progress(self, event): elif kind == "pipeline_step": self.denoise_step = event.get("step") self.denoise_total_steps = event.get("total_steps") + elif kind == "warning": + # Both channels, on purpose: the event log keeps the moment it + # happened, `warnings` keeps it where a caller who polled the + # finished job will actually look, since a warning about the + # artifact outlives the run that noticed it (#82). The step it + # fired in is the run's, not the warning's - the engine warns + # from inside a step without knowing which one it is + message = event.get("message") + if message: + named = f"{self.step_name}: {message}" if self.step_name else message + if named not in self.warnings: + self.warnings.append(named) elif kind == "step_start": self.step_name = event.get("step") self.step_index = event.get("index") diff --git a/dw/tasks/audio_utils.py b/dw/tasks/audio_utils.py index 81d5f6b7..b0aba909 100644 --- a/dw/tasks/audio_utils.py +++ b/dw/tasks/audio_utils.py @@ -14,6 +14,7 @@ import soundfile import torch +from ..events import emit_warning from ..security import ( validate_path, validate_url, @@ -656,10 +657,17 @@ def warn_on_level_spread(waveforms, command="concat_videos", measure="rms"): return None spread = max(levels) - min(levels) if spread >= LEVEL_SPREAD_WARN_DB: - logger.warning( + # emit_warning rather than logger.warning: this is a property of the + # file the run is about to write, and the caller reading the job is + # the one who can act on it (#82) + emit_warning( f"{command}: the tracks being joined span {spread:.1f} dB " f"({measure} {min(levels):.1f} to {max(levels):.1f} dBFS) - the cut " - f"will be audible as a level jump. Pass match_levels to even them out" + f"will be audible as a level jump. Pass match_levels to even them out", + kind="level_spread", + command=command, + spread_db=round(spread, 1), + measure=measure, ) return spread diff --git a/dw/tasks/concat_videos.py b/dw/tasks/concat_videos.py index 829c7698..a2369edf 100644 --- a/dw/tasks/concat_videos.py +++ b/dw/tasks/concat_videos.py @@ -147,4 +147,10 @@ def concat_videos( ) logger.debug(f"Concatenated {len(videos)} videos into {len(frames)} frames") - return AudioVideo(frames, audio, sample_rate) + # The rate the caller declared, else the rate the first input carries - + # either beats the result's 8 fps default (#84) + written_fps = fps or next( + (v.fps for v in videos if getattr(v, "fps", None)), + None, + ) + return AudioVideo(frames, audio, sample_rate, fps=written_fps) diff --git a/dw/tasks/dissolve_videos.py b/dw/tasks/dissolve_videos.py index 66841afd..2f281c92 100644 --- a/dw/tasks/dissolve_videos.py +++ b/dw/tasks/dissolve_videos.py @@ -108,7 +108,11 @@ def dissolve_videos( f"Dissolved {len(clips)} videos into {len(frames)} frames " f"({dissolve_frames}-frame seams)" ) - return AudioVideo(frames, audio, sample_rate) + written_fps = fps or next( + (v.fps for v in loaded if getattr(v, "fps", None)), + None, + ) + return AudioVideo(frames, audio, sample_rate, fps=written_fps) def _dissolve_join(previous, following, overlap): diff --git a/dw/tasks/interpolate_frames.py b/dw/tasks/interpolate_frames.py index 7513ab84..c0cfbe5f 100644 --- a/dw/tasks/interpolate_frames.py +++ b/dw/tasks/interpolate_frames.py @@ -49,6 +49,7 @@ def interpolate_frames(video, device="cpu", **kwargs): # An AudioVideo or a frame array unwraps to its frames; a PIL list passes # through by identity. The soundtrack does not survive - the frame count # changes, so pair_audio is how it comes back + source_fps = getattr(video, "fps", None) video = frames_as_pil_list(video) if len(video) < 2: raise ValueError(f"Need at least 2 frames to interpolate, got {len(video)}") @@ -69,7 +70,13 @@ def interpolate_frames(video, device="cpu", **kwargs): frames = _interpolate_2x(frames, model) logger.info(f"Interpolation complete: {len(video)} -> {len(frames)} frames") - return AudioVideo(frames, None, None) + # Multiplied, not carried: interpolation adds frames between the ones it + # was given, so playing them back at the source rate would run the clip + # `multiplier` times long. The rate that keeps the source's duration is + # the source's times the multiplier (#84) + return AudioVideo( + frames, None, None, fps=source_fps * multiplier if source_fps else None + ) def _interpolate_2x(frames, model): diff --git a/dw/tasks/pair_audio.py b/dw/tasks/pair_audio.py index ec21e3cd..8e19848f 100644 --- a/dw/tasks/pair_audio.py +++ b/dw/tasks/pair_audio.py @@ -68,4 +68,6 @@ def pair_audio(video, audio, sample_rate=None): # cost a copy of the whole thing for nothing frames = video.frames if isinstance(video, AudioVideo) else video logger.debug(f"Pairing frames with audio at {rate} Hz") - return AudioVideo(frames, as_channels_samples(waveform), rate) + return AudioVideo( + frames, as_channels_samples(waveform), rate, fps=getattr(video, "fps", None) + ) diff --git a/dw/tasks/stabilize.py b/dw/tasks/stabilize.py index e4b72ee8..e2fcc2e4 100644 --- a/dw/tasks/stabilize.py +++ b/dw/tasks/stabilize.py @@ -117,5 +117,5 @@ def stabilize_video(clip, smooth=0): ) if isinstance(clip, AudioVideo): - return AudioVideo(held, clip.audio, clip.sample_rate) + return AudioVideo(held, clip.audio, clip.sample_rate, fps=clip.fps) return held diff --git a/dw/tasks/task.py b/dw/tasks/task.py index eb32869b..dcdb2a92 100644 --- a/dw/tasks/task.py +++ b/dw/tasks/task.py @@ -270,7 +270,7 @@ def _per_frame(image, process): frames = [process(frame) for frame in frames_as_pil_list(image)] audio = getattr(image, "audio", None) sample_rate = getattr(image, "sample_rate", None) - return AudioVideo(frames, audio, sample_rate) + return AudioVideo(frames, audio, sample_rate, fps=getattr(image, "fps", None)) @register_command("upscale", implementation="dw.tasks.upscale.upscale_image") diff --git a/dw/tasks/video_utils.py b/dw/tasks/video_utils.py index fbbcd2ca..9626b143 100644 --- a/dw/tasks/video_utils.py +++ b/dw/tasks/video_utils.py @@ -288,7 +288,11 @@ def _decode_audio_video(handle): f"Decoded {len(frames)} frames and " f"{audio.shape[1] if audio is not None else 0} audio samples" ) - return AudioVideo(frames, audio, sample_rate if audio is not None else None) + # The file's own rate travels with it: a step that joins videos read + # from disk knows what to write them back at without being told (#84) + return AudioVideo( + frames, audio, sample_rate if audio is not None else None, fps=frame_rate + ) # How far a decoded track may be off the frames' own duration and still be diff --git a/dw/workflow_schema.json b/dw/workflow_schema.json index 3cefb295..16caebc4 100644 --- a/dw/workflow_schema.json +++ b/dw/workflow_schema.json @@ -1119,7 +1119,7 @@ "type": "string" }, "fps": { - "description": "Frames per second - only used when output is video", + "description": "Frames per second - only used when output is video. Defaults to the rate the frames themselves carry (a chain's 'fps', a join's own rate, the rate of the file a task read) and to 8 only when nothing knows better. Set it to write at a rate other than the source's - a deliberate slow motion - which the run then warns about.", "type": "integer", "default": 8 }, diff --git a/tests/test_concat_videos.py b/tests/test_concat_videos.py index 115e95e6..a4567213 100644 --- a/tests/test_concat_videos.py +++ b/tests/test_concat_videos.py @@ -398,3 +398,68 @@ def test_shots_already_at_one_level_draw_no_warning(self, caplog): concat_videos([audio_video(4, 0.5), audio_video(4, 0.45)]) assert "level jump" not in caplog.text + + +class TestWarningsReachTheCaller: + """A warning that only reaches the server's log does not exist from + outside it - see issue #82, where match_levels verified but the spread + warning was invisible over MCP.""" + + def test_the_level_spread_warning_is_emitted_as_an_event(self): + from dw.events import RunContext, activate_context, deactivate_context + + events = [] + token = activate_context(RunContext(on_event=events.append)) + try: + concat_videos([audio_video(4, 0.5), audio_video(4, 0.05)]) + finally: + deactivate_context(token) + + warnings = [e for e in events if e["event"] == "warning"] + assert len(warnings) == 1 + assert warnings[0]["kind"] == "level_spread" + assert warnings[0]["command"] == "concat_videos" + assert warnings[0]["spread_db"] == pytest.approx(20.0, abs=0.2) + assert "match_levels" in warnings[0]["message"] + + def test_matched_shots_emit_nothing(self): + from dw.events import RunContext, activate_context, deactivate_context + + events = [] + token = activate_context(RunContext(on_event=events.append)) + try: + concat_videos( + [audio_video(4, 0.5), audio_video(4, 0.05)], match_levels="rms" + ) + finally: + deactivate_context(token) + + assert [e for e in events if e["event"] == "warning"] == [] + + +class TestFrameRateTravelsWithTheJoin: + """result.fps defaults to 8, so a 24 fps cut that says nothing there + used to be written three times slow against audio of the right length - + see issue #84.""" + + def test_the_tasks_fps_is_carried_to_the_result(self): + result = concat_videos([frames(4), frames(4)], fps=24) + + assert result.fps == 24 + + def test_an_input_videos_rate_is_carried_when_the_task_is_told_nothing(self): + first = AudioVideo(frames(4), None, None, fps=30) + + result = concat_videos([first, frames(4)]) + + assert result.fps == 30 + + def test_the_tasks_own_fps_wins_over_its_inputs(self): + first = AudioVideo(frames(4), None, None, fps=30) + + result = concat_videos([first, frames(4)], fps=24) + + assert result.fps == 24 + + def test_nothing_is_carried_when_nothing_knows(self): + assert concat_videos([frames(4), frames(4)]).fps is None diff --git a/tests/test_dissolve_videos.py b/tests/test_dissolve_videos.py index d7935797..c5949007 100644 --- a/tests/test_dissolve_videos.py +++ b/tests/test_dissolve_videos.py @@ -127,3 +127,15 @@ def test_off_by_default_but_a_wide_spread_warns(self, caplog): assert abs(result.audio[0][-1]) == pytest.approx(0.05) assert "level jump" in caplog.text + + +class TestFrameRateTravelsWithTheDissolve: + """As with concat_videos - the rate the step was told is the rate the + file is written at, rather than result.fps's default of 8 (#84).""" + + def test_the_tasks_fps_is_carried_to_the_result(self): + result = dissolve_videos( + [frames(8, 0), frames(8, 255)], dissolve_frames=2, fps=24 + ) + + assert result.fps == 24 diff --git a/tests/test_host_memory.py b/tests/test_host_memory.py index 2c1433da..abd4f5b8 100644 --- a/tests/test_host_memory.py +++ b/tests/test_host_memory.py @@ -95,3 +95,29 @@ def test_the_repl_prints_host_memory_either_way(capsys, gpu_available): out = capsys.readouterr().out assert "Worker RSS: 1234.5 MB" in out assert "32000.0 MB of 64000.0 MB" in out + + +def test_the_peak_is_never_below_the_resident_figure(monkeypatch): + """peak - rss is what the pair exists to answer; a small negative there + reads as 'these fields are not comparable' - see issue #83.""" + monkeypatch.setattr(host_memory, "_peak_rss_mb", lambda: 764.1484375) + monkeypatch.setattr( + host_memory, + "_psutil_stats", + lambda: {"rss_mb": 764.79296875, "total_mb": 64000.0, "available_mb": 32000.0}, + ) + + stats = host_memory.host_memory_stats() + + assert stats["peak_rss_mb"] == stats["rss_mb"] == 764.79296875 + + +def test_a_genuine_peak_is_left_alone(monkeypatch): + monkeypatch.setattr(host_memory, "_peak_rss_mb", lambda: 33044.98) + monkeypatch.setattr( + host_memory, + "_psutil_stats", + lambda: {"rss_mb": 2561.69, "total_mb": 64000.0, "available_mb": 32000.0}, + ) + + assert host_memory.host_memory_stats()["peak_rss_mb"] == 33044.98 diff --git a/tests/test_job_progress.py b/tests/test_job_progress.py index b3f8b664..aa93aafd 100644 --- a/tests/test_job_progress.py +++ b/tests/test_job_progress.py @@ -170,3 +170,46 @@ def test_a_queued_event_is_stamped_against_creation_until_the_job_starts(): job.add_event({"event": "job_status", "status": "queued"}) assert job.events_after(-1)[0]["at"] >= 0 + + +class TestRuntimeWarnings: + """A warning a step raises about what it is writing has to reach the + caller, not just the server's log - see issue #82, where `match_levels` + verified over MCP but the spread warning it replaces did not exist out + there at all.""" + + def test_a_warning_event_lands_on_the_jobs_warnings(self): + job = running_job( + {"event": "step_start", "step": "join", "index": 0, "total_steps": 1}, + {"event": "warning", "message": "the tracks span 9.9 dB", "kind": "x"}, + ) + + assert job.warnings == ["join: the tracks span 9.9 dB"] + assert slim_job(job.detail())["warnings"] == ["join: the tracks span 9.9 dB"] + + def test_it_keeps_its_place_in_the_event_stream_too(self): + job = running_job({"event": "warning", "message": "a spread"}) + + assert [e["event"] for e in job.events_after(-1)] == ["warning"] + + def test_a_warning_before_any_step_is_carried_unnamed(self): + job = running_job({"event": "warning", "message": "a spread"}) + + assert job.warnings == ["a spread"] + + def test_the_same_warning_twice_is_carried_once(self): + job = running_job( + {"event": "step_start", "step": "join", "index": 0, "total_steps": 1}, + {"event": "warning", "message": "a spread"}, + {"event": "warning", "message": "a spread"}, + ) + + assert job.warnings == ["join: a spread"] + + def test_the_argument_warnings_a_job_was_queued_with_survive(self): + job = Job({"workflow_name": "shot", "warnings": ["unknown argument 'fsp'"]}) + job.add_event({"event": "warning", "message": "a spread"}) + + assert job.warnings == ["unknown argument 'fsp'", "a spread"] + # the spec is what a rerun is built from - appending must not edit it + assert job.spec["warnings"] == ["unknown argument 'fsp'"] diff --git a/tests/test_result.py b/tests/test_result.py index e02ca748..e122a352 100644 --- a/tests/test_result.py +++ b/tests/test_result.py @@ -1093,3 +1093,86 @@ def test_a_non_mp4_content_type_raises(self, tmp_path): if __name__ == "__main__": pytest.main([__file__, "-v"]) + + +class TestVideoFrameRate: + """`result.fps` and a task's own `fps` are separate knobs, and the one an + author sets is the task's - so the rate the frames carry decides when the + result declares none. See issue #84.""" + + def save(self, result_definition, artifact): + result = Result(result_definition) + result.add_result(artifact) + with ( + patch("dw.result.encode_video") as encode, + patch("dw.result.export_to_video") as export, + patch("dw.result.is_av_available", return_value=True), + ): + with tempfile.TemporaryDirectory() as temp_dir: + result.save(temp_dir, "test") + return encode, export + + def test_the_carried_rate_is_used_when_the_result_declares_none(self): + artifact = AudioVideo("frames", torch.zeros((2, 100)), 48000, fps=24) + + encode, _ = self.save({"content_type": "video/mp4"}, artifact) + + assert encode.call_args.kwargs["fps"] == 24 + + def test_a_declared_rate_still_wins(self): + artifact = AudioVideo("frames", torch.zeros((2, 100)), 48000, fps=24) + + encode, _ = self.save({"content_type": "video/mp4", "fps": 12}, artifact) + + assert encode.call_args.kwargs["fps"] == 12 + + def test_the_old_default_holds_when_nothing_knows_the_rate(self): + artifact = AudioVideo("frames", torch.zeros((2, 100)), 48000) + + encode, _ = self.save({"content_type": "video/mp4"}, artifact) + + assert encode.call_args.kwargs["fps"] == 8 + + def test_a_plain_frame_list_still_defaults_to_eight(self): + result = Result({"content_type": "video/mp4"}) + result.add_result([Image.new("RGB", (8, 8))]) + with ( + patch("dw.result.export_to_video") as export, + tempfile.TemporaryDirectory() as temp_dir, + ): + result.save(temp_dir, "test") + + assert export.call_args.kwargs["fps"] == 8 + + def test_declaring_a_rate_the_frames_contradict_warns(self): + from dw.events import RunContext, activate_context, deactivate_context + + events = [] + token = activate_context(RunContext(on_event=events.append)) + try: + self.save( + {"content_type": "video/mp4", "fps": 8}, + AudioVideo("frames", torch.zeros((2, 100)), 48000, fps=24), + ) + finally: + deactivate_context(token) + + warning = next(e for e in events if e["event"] == "warning") + assert warning["kind"] == "fps_mismatch" + assert warning["declared_fps"] == 8 + assert warning["source_fps"] == 24 + + def test_declaring_the_rate_the_frames_carry_warns_about_nothing(self): + from dw.events import RunContext, activate_context, deactivate_context + + events = [] + token = activate_context(RunContext(on_event=events.append)) + try: + self.save( + {"content_type": "video/mp4", "fps": 24}, + AudioVideo("frames", torch.zeros((2, 100)), 48000, fps=24), + ) + finally: + deactivate_context(token) + + assert [e for e in events if e["event"] == "warning"] == [] diff --git a/tests/test_video_utils.py b/tests/test_video_utils.py index 76430ca6..4c7850bb 100644 --- a/tests/test_video_utils.py +++ b/tests/test_video_utils.py @@ -307,6 +307,15 @@ def test_frames_and_audio_come_back_together(self, tmp_path): assert video.sample_rate == 8000 assert video.audio.shape[0] == 2 + def test_the_files_own_frame_rate_comes_back_with_it(self, tmp_path): + """A step that joins videos read from disk knows what to write them + back at without being told - see issue #84.""" + from dw.tasks.video_utils import load_audio_video + + path = self.write_video(tmp_path / "shot.mp4", fps=24, num_frames=24) + + assert load_audio_video(path).fps == 24 + def test_audio_is_fitted_to_the_frames_own_duration(self, tmp_path): """The codec pads the last block; joined shot after shot that padding would walk the sound off the picture.""" From 80baba72e73586c5928d046a637339fef929a8dc Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 19:32:51 -0500 Subject: [PATCH 02/16] fix(engine): #88 #89 #90 #92 - sub-workflow paths, one save per composed step, fps wording #88: the fps_mismatch warning reported source/declared, so 24 fps frames written at 8 were described as "3x speed" when they play at one third; and concat_videos' own `fps` docstring - which is what get_task shows - never got #84's sentence about result.fps. #90: a sub-workflow step's path is resolved the way run_workflow's workflow_path is - a catalog name, with or without .json, beside the referencing file first, then this run's workflows root, then the read-only sources dw.serve pins in DW_WORKFLOW_PATH. A stored template is composed rather than copied. A name that reaches nothing says where it looked. Progress from inside a composed run now carries the queued run's own step counter (parent_index/parent_total_steps), not the child's 1-of-1. #89: validate_workflow resolves those paths, validates the workflow each one names under steps[N].workflow.path, refuses a composition cycle, and warns about an argument the child declares no variable for. #92: a composing step that declares a saving result owns the file - the child's last step no longer writes a second copy under its own name, so a composed run stops doubling storage and the manifest stops carrying two entries under a step name the caller never wrote. A composed child's other files carry the composing step's name. --- .gitignore | 3 + docs/WORKFLOW_GUIDE.md | 39 ++++- dw/result.py | 5 +- dw/serve.py | 11 +- dw/server/app.py | 5 +- dw/server/jobs.py | 15 +- dw/tasks/concat_videos.py | 4 +- dw/workflow.py | 274 ++++++++++++++++++++++++++--- dw/workflow_schema.json | 2 +- dw/workflow_sources.py | 99 +++++++++++ dw/workspace.py | 5 + dw_mcp/server.py | 10 +- tests/test_result.py | 4 + tests/test_runs.py | 7 +- tests/test_server.py | 44 +++++ tests/test_workflow.py | 358 ++++++++++++++++++++++++++++++++++++++ 16 files changed, 846 insertions(+), 39 deletions(-) diff --git a/.gitignore b/.gitignore index 0c9a77e2..253a8407 100644 --- a/.gitignore +++ b/.gitignore @@ -165,3 +165,6 @@ workflows/*/assets/ /output/ /*.mp4 /*.wav + +# Claude Code worktrees +.claude/worktrees/ diff --git a/docs/WORKFLOW_GUIDE.md b/docs/WORKFLOW_GUIDE.md index f26d58b3..7a11624a 100644 --- a/docs/WORKFLOW_GUIDE.md +++ b/docs/WORKFLOW_GUIDE.md @@ -151,7 +151,22 @@ Invoke another workflow file: } ``` -Paths can be relative to the current file or use `builtin:` to reference built-in workflows in `dw/workflows/`. +`path` is read the way `run_workflow`'s `workflow_path` is: a catalog name as +`list_workflows` reports it (`templates/minimax/reference-to-video`), with or +without `.json`; a path relative to the file that names it (`../models/x.json`); +or `builtin:name.json` for the packaged fragments in `dw/workflows/`. A name +resolves beside the referencing file first, then against the run's own +`workflows/` directory, then against each read-only source the server lists - +so a stored template can be composed without copying it into the workspace. A +path that lands outside every source is refused, and one that resolves nowhere +is a validation error rather than a run that fails on its first step. + +When the composing step declares a `result`, that is where the composed output +is written, once: the child's own last step does not save it a second time +under its own name. A composing step that declares no `result` (or one with no +`content_type`) leaves the saving to the child, as before. The child's other +steps write into the same run directory, with the composing step's name +leading their file names. ## Cross-Step Data Flow @@ -543,6 +558,28 @@ subfolder written is one a later workflow can name: validation error at its JSON path. `file_base_name` is a name, not a path: a separator there is refused, and `subfolder` is the way to place a file. +### Composing a stored workflow + +A step with a `workflow` block runs another workflow as one step of this one, +with `arguments` handed down as that workflow's variables. Its `path` is read +the way `run_workflow`'s `workflow_path` is - a catalog name from +`list_workflows`, with or without `.json`, a path relative to the file that +names it, or `builtin:name.json` - and resolves beside the referencing file +first, then in this workspace's `workflows/`, then in each read-only source +the server lists. A stored template is composed by its catalog name; copying +it into the workspace to reach it is no longer necessary, and a copy silently +stops tracking the original. + +Declare a `result` on the composing step and the composed output is saved +there, once, under that step's name and subfolder - the composed workflow's +own last step does not write a second copy. Its other steps write into the +same run directory, prefixed with the composing step's name. + +`validate_workflow` resolves the path, so a name that reaches nothing is an +error at `steps[N].workflow.path` before anything is queued; it also validates +the workflow named, refuses a composition cycle, and warns about an argument +the composed workflow declares no variable for. + ### Being found next time The catalog derives each entry's `shape` — one of `image`, `image-set`, diff --git a/dw/result.py b/dw/result.py index b40fe1b2..5b375660 100644 --- a/dw/result.py +++ b/dw/result.py @@ -503,8 +503,9 @@ def video_fps(self, artifact): emit_warning( f"Writing video at {declared} fps, but the frames it was " f"given run at {carried} fps - the file will play " - f"{carried / declared:.2g}x speed. Drop 'fps' from the step's " - f"result to keep the source rate", + f"{declared / carried:.2g}x speed " + f"({carried / declared:.2g} times as long). Drop 'fps' from " + f"the step's result to keep the source rate", kind="fps_mismatch", declared_fps=declared, source_fps=carried, diff --git a/dw/serve.py b/dw/serve.py index 59c1d543..caed92cf 100644 --- a/dw/serve.py +++ b/dw/serve.py @@ -192,10 +192,19 @@ def main(): # search path is what lets an example run as it shipped: the workspace's # own library is still searched first and is still the only one written # to. Pinned in the environment, so the worker resolves as the API does - from .workspace import ASSETS_SUBDIR, example_libraries, set_library_fallbacks + from .workspace import ( + ASSETS_SUBDIR, + WORKFLOWS_SUBDIR, + example_libraries, + set_library_fallbacks, + ) example_dirs = example_libraries(args.examples_dirs) set_library_fallbacks(PROMPTS_SUBDIR, example_dirs[PROMPTS_SUBDIR]) + # The workflow trees themselves, so a sub-workflow step can compose a + # stored template by the name list_workflows reports rather than a copy + # of it in this workspace (#90) + set_library_fallbacks(WORKFLOWS_SUBDIR, args.examples_dirs) # The shared library goes ahead of the examples and behind the # workspace's own, which is the order 'asset:' resolves in: a workspace # name shadows a shared one, and a shared one shadows an example's. diff --git a/dw/server/app.py b/dw/server/app.py index 371c589d..12b85a10 100644 --- a/dw/server/app.py +++ b/dw/server/app.py @@ -1353,7 +1353,10 @@ def validate_workflow( "error": None, "errors": [], "warnings": workflow_argument_warnings(definition) - + entry_field_warnings(definition, request.arguments), + + entry_field_warnings(definition, request.arguments) + # An argument a sub-workflow step passes to a workflow that + # declares no variable for it - dropped in silence at run time + + candidate.sub_workflow_warnings(), } if request.arguments: # Naming what was checked is the difference between 'the stored diff --git a/dw/server/jobs.py b/dw/server/jobs.py index 98346c4e..cc97a68d 100644 --- a/dw/server/jobs.py +++ b/dw/server/jobs.py @@ -363,6 +363,7 @@ def __init__(self, spec): self.phase_detail = None self.phase_started_at = None self.step_name = None + self.parent_step = None self.step_index = None self.total_steps = None self.denoise_step = None @@ -416,8 +417,15 @@ def _note_progress(self, event): self.warnings.append(named) elif kind == "step_start": self.step_name = event.get("step") - self.step_index = event.get("index") - self.total_steps = event.get("total_steps") + # A sub-workflow counts its own steps from zero; what a caller + # watching a composed run needs is where the run it queued has + # got to, so the parent's counter wins when the event carries + # one and the step name stays the child's (#90) + self.parent_step = event.get("parent_step") + self.step_index = event.get("parent_index", event.get("index")) + self.total_steps = event.get( + "parent_total_steps", event.get("total_steps") + ) # A new step's denoise loop has not started; the previous step's # count would read as this one's progress self.denoise_step = None @@ -432,6 +440,9 @@ def progress(self): now = time.time() summary = { "step": self.step_name, + # The step of the queued workflow the one above is running + # inside, for a composed run; null when they are the same thing + "parent_step": self.parent_step, "step_index": self.step_index, "total_steps": self.total_steps, "phase": self.phase, diff --git a/dw/tasks/concat_videos.py b/dw/tasks/concat_videos.py index a2369edf..df9c60e5 100644 --- a/dw/tasks/concat_videos.py +++ b/dw/tasks/concat_videos.py @@ -61,7 +61,9 @@ def concat_videos( hard cut on tonal material, which a bleed would only stutter. It is the wrong tool for a continuous bed such as a laugh track or room tone: a fade only deepens the hole a bleed is there to cover - fps: Frame rate of the videos - required to join audio when trimming + fps: Frame rate of the videos - required to join audio when + trimming, and the rate the joined file is written at unless + the step's result.fps overrides it match_levels: Even the shots' loudness out before joining - "rms" matches perceived level (the measurement get_gallery_metadata reports as mean_dbfs), "peak" matches the diff --git a/dw/workflow.py b/dw/workflow.py index 00b898a3..3b2890c9 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -73,6 +73,11 @@ InvalidInputError, UntrustedWorkflowError, ) +from .workflow_sources import ( + builtin_root, + resolve_sub_workflow, + SubWorkflowNotFound, +) logger = logging.getLogger("dw") @@ -270,6 +275,22 @@ class Workflow: # leaves no manifest of its own, since its steps are already rolled up # into the parent's _run_dir_inherited = False + # Where the parent step that delegated to this workflow sits in the + # run the caller queued: {"step", "index", "total_steps"}. A child + # counts its own steps from zero, so without this a composed run + # reported "step 1 of 1" from inside the first of the parent's three + # (#90) - and "is this nearly finished" is the whole question progress + # answers. Handed straight down to a grandchild, so the numbers always + # describe the run that was queued + _parent_progress = None + # Whether the parent step that composed this workflow declares a + # `result` of its own. It does the saving then, and this run's last step + # does not: the two used to write the same artifact twice, once under + # the parent step's name and subfolder and once under the child's, with + # the child's entry shadowing a manifest key the caller never wrote + # (#92). A child whose parent declares nothing still saves, since + # otherwise the output would exist nowhere + _final_save_owned_by_parent = False def __init__(self, workflow_definition, output_dir, file_spec, workflow_dir=None): self.workflow_definition = workflow_definition @@ -292,11 +313,37 @@ def argument_template(self): def variables(self): return self.workflow_definition.get("variables", {}) + def step_save_name(self, workflow_id, step_name, index): + """The base name a step's files are written under. + + Inside a composed child the parent step's name leads, so two steps + composing the same workflow do not both want one name and get told + apart by a '-2' suffix that says nothing about which step made it + (#92). + """ + base = f"{workflow_id}-{step_name}.{index}" + parent = self._parent_progress + return f"{parent['step']}.{base}" if parent else base + + def _parent_progress_fields(self): + """The queued run's own step counter, on an event a sub-workflow + emits - empty for a top-level run, whose index is already that.""" + parent = self._parent_progress + if not parent: + return {} + return { + "parent_step": parent["step"], + "parent_index": parent["index"], + "parent_total_steps": parent["total_steps"], + } + def step_file_prefix(self, step_name): """Naming prefix for files a step writes on its own (chain segment spills), matching the workflow-id-step naming its results are saved under.""" - return f"{self.name}-{step_name}" + prefix = f"{self.name}-{step_name}" + parent = self._parent_progress + return f"{parent['step']}.{prefix}" if parent else prefix @property def effective_output_dir(self): @@ -385,10 +432,138 @@ def expanded_definition(self, arguments=None, source_indices=None): definition = replace_variables(definition, variables) return expand_for_each(definition, source_indices) - def validation_errors(self, arguments=None): + def resolve_sub_workflow_path(self, path): + """Where one sub-workflow step's `path` resolves to, as + (path, root) - the same resolution create_step_action does, asked + ahead of the run so validation can answer for free what used to cost + a queued job to find out (#89). + + Raises SubWorkflowNotFound, SecurityError or InvalidInputError, + each carrying the message the run would have failed with. + """ + confine_to = self.workflow_dir + if path.startswith("builtin:"): + builtin_name = path.replace("builtin:", "") + if ( + not builtin_name.endswith(".json") + or "/" in builtin_name + or "\\" in builtin_name + ): + raise InvalidInputError(f"Invalid builtin workflow name: {builtin_name}") + confine_to = builtin_root() + resolved = os.path.join(confine_to, builtin_name) + if not os.path.isfile(resolved): + raise SubWorkflowNotFound(path, [resolved]) + return validate_workflow_path(resolved, confine_to), confine_to + if confine_to is None and not os.path.isabs(path): + confine_to = catalog_root_dir(self.file_spec) + resolved, confine_to = resolve_sub_workflow( + path, os.path.dirname(self.file_spec), confine_to + ) + return validate_workflow_path(resolved, confine_to), confine_to + + def sub_workflow_errors(self, expanded, source_indices=None, composing=None): + """Every sub-workflow step whose `path` names nothing this server can + reach, composes a workflow already on the chain, or resolves to a + workflow that does not itself validate. + + `composing` is the resolved path of every workflow above this one, + which is what makes a cycle an error here rather than a recursion + the run discovers. + """ + errors = [] + composing = list(composing or []) + for index, step in enumerate(expanded.get("steps", []) or []): + reference = step.get("workflow") + if not isinstance(reference, dict) or not isinstance( + reference.get("path"), str + ): + continue + source = source_indices[index] if source_indices else index + where = f"steps[{source}].workflow.path" + path = reference["path"] + try: + resolved, root = self.resolve_sub_workflow_path(path) + except (SubWorkflowNotFound, SecurityError, InvalidInputError) as e: + errors.append({"path": where, "message": str(e)}) + continue + if resolved in composing: + errors.append( + { + "path": where, + "message": ( + f"Sub-workflow '{path}' composes a workflow that " + "is already composing it - a cycle: " + + " -> ".join(composing + [resolved]) + ), + } + ) + continue + try: + child = workflow_from_file(resolved, self.output_dir, root) + except Exception as e: + errors.append( + {"path": where, "message": f"Sub-workflow '{path}': {e}"} + ) + continue + for error in child.validation_errors(composing=composing + [resolved]): + errors.append( + { + "path": f"{where} -> {error['path']}", + "message": f"Sub-workflow '{path}': {error['message']}", + } + ) + return errors + + def sub_workflow_warnings(self, expanded=None): + """An argument a sub-workflow step passes down that the workflow it + composes declares no variable for - dropped in silence at run time, + and composition is exactly where a name drifts (#89).""" + warnings = [] + try: + expanded = ( + expanded if expanded is not None else self.expanded_definition() + ) + except Exception: + return warnings + for index, step in enumerate(expanded.get("steps", []) or []): + reference = step.get("workflow") + if not isinstance(reference, dict): + continue + passed = reference.get("arguments") + if not isinstance(passed, dict) or not isinstance( + reference.get("path"), str + ): + continue + try: + resolved, root = self.resolve_sub_workflow_path(reference["path"]) + child = workflow_from_file(resolved, self.output_dir, root) + except Exception: + # An unresolvable path is an error, reported by + # sub_workflow_errors - not a second complaint here + continue + declared = child.workflow_definition.get("variables") or {} + for name in sorted(set(passed) - set(declared)): + warnings.append( + { + "path": f"steps[{index}].workflow.arguments.{name}", + "message": ( + f"'{reference['path']}' declares no variable " + f"'{name}' - the value is dropped. Declared: " + + (", ".join(sorted(declared)) or "") + ), + } + ) + return warnings + + def validation_errors(self, arguments=None, composing=None): """Every schema violation in the definition, as [{path, message}]; empty when it validates. `arguments` are the caller's, so a - for_each over a list the caller supplies is checked as it will run.""" + for_each over a list the caller supplies is checked as it will run. + + `composing` carries the chain of sub-workflows above this one, so a + workflow that composes itself is an error rather than a recursion. + """ errors = validate_data_all(self.workflow_definition, load_schema("workflow")) # Only once the shape is known good: the passes below walk the # steps array and a definition that fails the schema may have no @@ -415,9 +590,11 @@ def validation_errors(self, arguments=None): # reported against 'variables' as a whole rather than escaping # as an unhandled exception return [{"path": "variables", "message": str(e)}] - return previous_result_reference_errors( - expanded, source_indices - ) + subfolder_errors(expanded, source_indices) + return ( + previous_result_reference_errors(expanded, source_indices) + + subfolder_errors(expanded, source_indices) + + self.sub_workflow_errors(expanded, source_indices, composing) + ) def _undeclared_variable_errors(self, arguments=None): """Every 'variable:' reference naming nothing the workflow declares. @@ -701,6 +878,7 @@ def run( step=step_data["name"], index=i, total_steps=len(steps), + **self._parent_progress_fields(), ) # Seeds resolve most-specific-first: pipeline > step > workflow @@ -730,10 +908,21 @@ def run( # A sub-workflow step is never cacheable: its files roll up # from the child's own manifest, which a hit does not rebuild. is_cacheable = "workflow" not in step_data and cache_enabled_this_run + # The last step of a composed child whose parent does the + # saving (#92) - its files are the parent step's, written + # once, under the parent's name and subfolder + parent_saves_this = ( + self._final_save_owned_by_parent and i == len(steps) - 1 + ) step_data_snapshot = None if is_cacheable: try: step_data_snapshot = copy.deepcopy(step_data) + if parent_saves_this: + # Keyed apart from the same step run standalone: + # this entry's result was never saved here, so a + # standalone hit on it would report no files + step_data_snapshot["__saved_by_parent__"] = True except Exception as ex: # A realized argument that cannot be deep-copied (an # open handle, a live model object) just means this @@ -776,6 +965,24 @@ def run( step_seed, get_device(), ) + if isinstance(step_action, Workflow): + # The child reports into this run's counter rather than + # its own, and a grandchild reports into the same one + step_action._parent_progress = self._parent_progress or { + "step": step_data["name"], + "index": i, + "total_steps": len(steps), + } + # Only when the parent's own result would write + # something: a result block that names no content_type, + # or says save: false, saves nothing, and suppressing + # the child's save for it would lose the artifact + parent_result = step_data.get("result") + step_action._final_save_owned_by_parent = bool( + isinstance(parent_result, dict) + and parent_result.get("content_type") + and parent_result.get("save", True) + ) reused = cached_result is not None if reused: logger.info(f"Step '{step.name}' unchanged - reusing cached result") @@ -830,9 +1037,13 @@ def run( ) if not reused: - saved_files = result.save( - self.step_output_dir(step_data), - f"{workflow_id}-{step.name}.{i}", + saved_files = ( + [] + if parent_saves_this + else result.save( + self.step_output_dir(step_data), + self.step_save_name(workflow_id, step.name, i), + ) ) if is_cacheable: step_cache.put( @@ -857,7 +1068,11 @@ def run( } if reused: manifest_entry["reused"] = True - self.manifest.append(manifest_entry) + # No entry at all for a step the parent saves for: the + # parent's own entry names the same files, under the step + # name the caller wrote (#92) + if not parent_saves_this: + self.manifest.append(manifest_entry) # roll the child's saves up so job history and the gallery see # every file self.manifest.extend(sub_manifest) @@ -870,6 +1085,7 @@ def run( step=step.name, index=i, total_steps=len(steps), + **self._parent_progress_fields(), **step_end_data, ) logger.debug(f"Step {step.name} completed with result: {result}") @@ -1149,24 +1365,28 @@ def create_step_action( os.path.dirname(os.path.abspath(__file__)), "workflows" ) path = os.path.join(confine_to, builtin_name) - # Handle relative paths. A template under templates/ names a - # model config as '../models/x.json'; collapsing the '..' here - # is what lets the validator judge where the path actually - # lands rather than refusing the spelling - containment is - # still checked on the resolved path below - elif not os.path.isabs(path): - base_dir = os.path.dirname(self.file_spec) - path = os.path.normpath(os.path.join(base_dir, path)) - # An unconfined run (no workflow_dir - a bare CLI - # invocation) used to rely on the '..' regex alone to - # stop a relative reference from leaving the file's own - # directory; normalising the path removes that guard, so - # here confine it to the catalog root instead - the - # referencing file's nearest ancestor literally named - # 'workflows', which still lets it climb to a sibling - # folder like models/ but not out of the catalog - if confine_to is None: + # Everything else - a relative path, or a catalog name as + # list_workflows reports it - goes through the search path. + # A template under templates/ names a model config as + # '../models/x.json', so a path relative to the referencing + # file still resolves first and the '..' is collapsed here, + # which is what lets the validator judge where the path + # actually lands rather than refusing the spelling; + # containment is still checked on the resolved path below. + # An unconfined run (no workflow_dir - a bare CLI + # invocation) used to rely on the '..' regex alone to stop a + # relative reference from leaving the file's own directory; + # normalising the path removes that guard, so confine it to + # the catalog root instead - the referencing file's nearest + # ancestor literally named 'workflows', which still lets it + # climb to a sibling folder like models/ but not out of the + # catalog + else: + if confine_to is None and not os.path.isabs(path): confine_to = catalog_root_dir(self.file_spec) + path, confine_to = resolve_sub_workflow( + path, os.path.dirname(self.file_spec), confine_to + ) # Validate the resolved path - confined when this workflow # itself is (an inline/server-submitted run), so a diff --git a/dw/workflow_schema.json b/dw/workflow_schema.json index 16caebc4..4fba0ddc 100644 --- a/dw/workflow_schema.json +++ b/dw/workflow_schema.json @@ -1090,7 +1090,7 @@ "type": "object", "properties": { "path": { - "description": "The path to the workflow file. Use 'builtin:' for built-in workflows.", + "description": "The workflow this step composes: a catalog name as list_workflows reports it (with or without '.json'), a path relative to the workflow that names it ('../models/x.json'), or 'builtin:name.json' for one of the packaged fragments. A name is resolved first beside the referencing file, then against this run's own workflows directory, then against each read-only source on the server's search path - so a stored template can be composed without copying it into the workspace. A path outside every source is refused. When the composing step declares a 'result', that is where the composed output is saved and the composed workflow's own last step does not save it again.", "type": "string" }, "arguments": { diff --git a/dw/workflow_sources.py b/dw/workflow_sources.py index 94788c4f..323d899b 100644 --- a/dw/workflow_sources.py +++ b/dw/workflow_sources.py @@ -163,3 +163,102 @@ def listing(sources): for name in source.names(): found.setdefault(name, source) return dict(sorted(found.items())) + + +def fallback_roots(primary=None): + """The read-only workflow roots a sub-workflow name is resolved against + after the directory the run is confined to. + + Pinned in DW_WORKFLOW_PATH by dw.serve, the same way the prompt and + asset libraries are, so the worker resolves a composed template exactly + as the API would. + """ + from .workspace import WORKFLOWS_SUBDIR, library_fallbacks + + return library_fallbacks(WORKFLOWS_SUBDIR, primary) + + +def _candidate_names(name): + """A name as written, and with .json appended when it has no extension - + the catalog reports names without it, and run_workflow's workflow_path + takes either (#90).""" + names = [name] + if not name.endswith(".json"): + names.append(f"{name}.json") + return names + + +class SubWorkflowNotFound(Exception): + """A sub-workflow step's path names nothing on the search path.""" + + def __init__(self, path, tried): + self.path = path + self.tried = list(tried) + detail = "\n ".join(self.tried) + super().__init__( + f"Sub-workflow '{path}' could not be resolved. It is read as a " + "catalog name from list_workflows (with or without .json), or a " + "path relative to the workflow that names it. Looked in:" + f"\n {detail}" + ) + + +def resolve_sub_workflow(path, base_dir, confine_to): + """Where a sub-workflow step's `path` resolves to, and the root the + child is confined to, as (path, root). + + Order, first hit wins: + + 1. relative to the directory of the workflow that names it - which is + how a template reaches '../models/x.json', and stays first so an + existing composition keeps meaning what it did + 2. the same, with '.json' supplied + 3. the run's own workflow root (a workspace's workflows/), by catalog + name, with or without '.json' + 4. each read-only root on the search path, the same way - which is + what lets a stored template be composed rather than copied (#90) + + An absolute path is taken as written and confined to whichever root + holds it, so the sandbox still refuses one that belongs to no source. + + Raises SubWorkflowNotFound, naming every candidate it looked at. + """ + roots = [] + if confine_to: + roots.append(os.path.abspath(os.path.expanduser(str(confine_to)))) + for root in fallback_roots(roots[0] if roots else None): + if root not in roots: + roots.append(root) + + tried = [] + if os.path.isabs(path): + candidate = os.path.normpath(path) + tried.append(candidate) + for root in roots: + source = WorkflowSource(root, EXAMPLES_ORIGIN, False) + if source.contains(candidate) and os.path.isfile(candidate): + return candidate, root + # No root holds it - hand it back confined as it was, so the + # security layer writes the refusal it always did + return candidate, confine_to + + if base_dir: + for name in _candidate_names(path): + candidate = os.path.normpath(os.path.join(base_dir, name)) + tried.append(candidate) + if os.path.isfile(candidate): + return candidate, confine_to + + for root in roots: + source = WorkflowSource(root, EXAMPLES_ORIGIN, False) + # resolve_in_source supplies '.json' itself, and refuses a name that + # would traverse out of the root + candidate = resolve_in_source(source, path) + if candidate is None: + tried.append(f"{os.path.join(root, path)} (outside the root)") + continue + tried.append(candidate) + if os.path.isfile(candidate): + return candidate, root + + raise SubWorkflowNotFound(path, tried) diff --git a/dw/workspace.py b/dw/workspace.py index 756ddd64..0b0f8c65 100644 --- a/dw/workspace.py +++ b/dw/workspace.py @@ -292,10 +292,15 @@ def set_workspace(workspace): # DW_PROMPT_DIR and DW_ASSET_DIR PROMPT_PATH_ENV_VAR = "DW_PROMPT_PATH" ASSET_PATH_ENV_VAR = "DW_ASSET_PATH" +# The same idea for workflows, which a sub-workflow step names: a stored +# template lives in an examples tree the workspace's own workflows/ cannot +# reach, so composing one used to mean copying it in (#90) +WORKFLOW_PATH_ENV_VAR = "DW_WORKFLOW_PATH" LIBRARY_PATH_ENV_VARS = { PROMPTS_SUBDIR: PROMPT_PATH_ENV_VAR, ASSETS_SUBDIR: ASSET_PATH_ENV_VAR, + WORKFLOWS_SUBDIR: WORKFLOW_PATH_ENV_VAR, } diff --git a/dw_mcp/server.py b/dw_mcp/server.py index a1c861de..b1d38605 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -595,7 +595,15 @@ def validate_workflow( A `result.subfolder` or `file_base_name` that cannot be written (a `..`, a backslash, a separator in `file_base_name`) is reported here - at its JSON path, after `for_each` expansion.""" + at its JSON path, after `for_each` expansion. + + A sub-workflow step is resolved too: a `workflow.path` that names + nothing this server can reach is an error at + `steps[N].workflow.path` (the message lists where it looked), the + workflow it names is validated in turn under that path, a + composition cycle is refused, and an argument passed down that the + composed workflow declares no variable for comes back as a + warning.""" return authoring.validate_workflow( client, workflow=workflow, diff --git a/tests/test_result.py b/tests/test_result.py index e122a352..fa197410 100644 --- a/tests/test_result.py +++ b/tests/test_result.py @@ -1161,6 +1161,10 @@ def test_declaring_a_rate_the_frames_contradict_warns(self): assert warning["kind"] == "fps_mismatch" assert warning["declared_fps"] == 8 assert warning["source_fps"] == 24 + # 24 fps frames written at 8 play in slow motion, not fast - the + # factor is declared/source, and it pointed the other way (#88) + assert "0.33x speed" in warning["message"] + assert "3 times as long" in warning["message"] def test_declaring_the_rate_the_frames_carry_warns_about_nothing(self): from dw.events import RunContext, activate_context, deactivate_context diff --git a/tests/test_runs.py b/tests/test_runs.py index fb4aece9..18713d9a 100644 --- a/tests/test_runs.py +++ b/tests/test_runs.py @@ -697,8 +697,11 @@ def test_a_parents_subfolder_does_not_move_a_childs_files( } Workflow(parent, str(tmp_path / "out"), str(tree / "Parent.json")).run({}) (run,) = (tmp_path / "out" / "Parent").iterdir() - assert (run / "runs_test-gen0.0-0.0.png").is_file() - assert not (run / "final" / "runs_test-gen0.0-0.0.png").exists() + # The parent step's name leads a composed child's file names, so two + # steps composing one workflow are told apart by the step that made + # them rather than by a '-2' suffix (#92) + assert (run / "child.runs_test-gen0.0-0.0.png").is_file() + assert not (run / "final" / "child.runs_test-gen0.0-0.0.png").exists() def test_a_chain_spill_lands_in_the_steps_subfolder(self, tmp_path, fake_pipeline): # save_segments writes through the pipeline wrapper's output_dir, diff --git a/tests/test_server.py b/tests/test_server.py index 7c4f0679..9b3fbc40 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -516,6 +516,50 @@ def test_validate_accepts_a_stored_workflow_name(server, tmp_path): assert any("guidance_scael" in w for w in result["warnings"]) +def test_validate_reports_a_sub_workflow_path_that_resolves_nowhere(server): + """The pre-flight is documented as "this will run", and a composed step + naming a workflow the server cannot reach used to come back valid and + fail 0.6 s into the job (#89).""" + with server(success_script) as client: + workflow = { + "id": "QaSubPathProbe", + "steps": [ + { + "name": "sub", + "workflow": { + "path": "templates/does-not-exist-at-all", + "arguments": {}, + }, + "result": {"content_type": "image/jpeg"}, + } + ], + } + + result = client.post("/api/validate", json={"workflow": workflow}).json() + + assert result["valid"] is False + assert [e["path"] for e in result["errors"]] == ["steps[0].workflow.path"] + assert "does-not-exist-at-all" in result["errors"][0]["message"] + + +def test_validate_accepts_a_sub_workflow_step_naming_a_stored_workflow(server): + with server(success_script) as client: + workflow = { + "id": "QaSubPathProbe", + "steps": [ + { + "name": "sub", + "workflow": {"path": "Basic", "arguments": {}}, + "result": {"content_type": "image/jpeg"}, + } + ], + } + + result = client.post("/api/validate", json={"workflow": workflow}).json() + + assert result["valid"] is True, result + + def test_validate_requires_exactly_one_workflow_source(server): with server(success_script) as client: assert client.post("/api/validate", json={}).status_code == 400 diff --git a/tests/test_workflow.py b/tests/test_workflow.py index 8f21ff18..402142f8 100644 --- a/tests/test_workflow.py +++ b/tests/test_workflow.py @@ -863,3 +863,361 @@ def test_a_variable_cycle_is_a_validation_error_at_variables(tmp_path): assert [e["path"] for e in errors] == ["variables"] assert "a -> b -> a" in errors[0]["message"] + + +class TestSubWorkflowNameResolution: + """A sub-workflow step's path reads like run_workflow's workflow_path: + a catalog name, with or without .json, resolved across the same search + path the server lists (#90).""" + + def _catalog(self, tmp_path, parent_path_value): + import json + + workflows = tmp_path / "workflows" + (workflows / "minimax").mkdir(parents=True) + child = { + "id": "child", + "steps": [ + { + "name": "noop", + "task": { + "command": "get_dict_value", + "arguments": {"dict": {"k": 1}, "key": "k"}, + }, + } + ], + } + (workflows / "minimax" / "ref2va.json").write_text(json.dumps(child)) + parent = { + "id": "parent", + "steps": [ + { + "name": "sub", + "workflow": {"path": parent_path_value, "arguments": {}}, + } + ], + } + parent_path = workflows / "parent.json" + parent_path.write_text(json.dumps(parent)) + return workflows, parent_path + + def _resolve(self, tmp_path, path_value, workflow_dir=None): + from dw.workflow import workflow_from_file + + workflows, parent_path = self._catalog(tmp_path, path_value) + workflow = workflow_from_file( + str(parent_path), + str(tmp_path / "outputs"), + str(workflow_dir) if workflow_dir else str(workflows), + ) + return workflow.create_step_action( + workflow.workflow_definition["steps"][0], + shared_components={}, + previous_pipelines={}, + default_seed=42, + device="cpu", + ) + + def test_a_catalog_name_without_the_extension_resolves(self, tmp_path): + action = self._resolve(tmp_path, "minimax/ref2va") + + assert action.name == "child" + + def test_a_name_in_a_read_only_source_resolves(self, tmp_path, monkeypatch): + """The examples tree list_workflows reports as a source - reachable + without copying the template into the workspace.""" + import json + + from dw.workspace import WORKFLOW_PATH_ENV_VAR + + examples = tmp_path / "examples" + (examples / "templates").mkdir(parents=True) + (examples / "templates" / "stored.json").write_text( + json.dumps( + { + "id": "stored", + "steps": [ + { + "name": "noop", + "task": { + "command": "get_dict_value", + "arguments": {"dict": {"k": 1}, "key": "k"}, + }, + } + ], + } + ) + ) + workspace = tmp_path / "space" / "workflows" + workspace.mkdir(parents=True) + parent = { + "id": "parent", + "steps": [ + { + "name": "sub", + "workflow": {"path": "templates/stored", "arguments": {}}, + } + ], + } + parent_path = workspace / "parent.json" + parent_path.write_text(json.dumps(parent)) + monkeypatch.setenv(WORKFLOW_PATH_ENV_VAR, str(examples)) + + from dw.workflow import workflow_from_file + + workflow = workflow_from_file( + str(parent_path), str(tmp_path / "outputs"), str(workspace) + ) + action = workflow.create_step_action( + workflow.workflow_definition["steps"][0], + shared_components={}, + previous_pipelines={}, + default_seed=42, + device="cpu", + ) + + assert action.name == "stored" + # confined to the root it was read from, not to the workspace + assert action.workflow_dir == str(examples) + + def test_a_name_that_resolves_nowhere_says_where_it_looked(self, tmp_path): + from dw.workflow_sources import SubWorkflowNotFound + + with pytest.raises(SubWorkflowNotFound) as exc_info: + self._resolve(tmp_path, "minimax/does-not-exist") + + message = str(exc_info.value) + assert "does-not-exist" in message + assert "Looked in" in message + + def test_a_relative_path_beside_the_file_still_wins(self, tmp_path): + """The '../models/x.json' form every template uses is unchanged.""" + action = self._resolve(tmp_path, "minimax/ref2va.json") + + assert action.name == "child" + + +class TestComposedStepSavesOnce: + """A sub-workflow step that declares a result owns the file: the child's + last step used to save the same artifact a second time, under its own + step name, into the run root (#92).""" + + def _compose(self, tmp_path, parent_result=True): + import json + + workflows = tmp_path / "workflows" + workflows.mkdir() + child = { + "id": "child", + "steps": [ + { + "name": "write", + "task": { + "command": "compose_text", + "arguments": {"parts": ["hello"]}, + }, + "result": {"content_type": "text/plain"}, + } + ], + } + (workflows / "child.json").write_text(json.dumps(child)) + step = {"name": "sub", "workflow": {"path": "child.json", "arguments": {}}} + if parent_result: + step["result"] = {"content_type": "text/plain", "subfolder": "final"} + parent = {"id": "parent", "steps": [step]} + parent_path = workflows / "parent.json" + parent_path.write_text(json.dumps(parent)) + + from dw.workflow import workflow_from_file + + workflow = workflow_from_file( + str(parent_path), str(tmp_path / "outputs"), str(workflows) + ) + workflow.run({}, {}) + return workflow + + def _written(self, tmp_path): + return sorted( + os.path.relpath(os.path.join(directory, name), str(tmp_path / "outputs")) + for directory, _dirs, files in os.walk(str(tmp_path / "outputs")) + for name in files + if name.endswith(".txt") + ) + + def test_the_artifact_is_written_once(self, tmp_path): + self._compose(tmp_path) + + written = self._written(tmp_path) + assert len(written) == 1, written + assert "final" in written[0] + + def test_the_manifest_names_only_the_step_the_caller_wrote(self, tmp_path): + workflow = self._compose(tmp_path) + + assert [entry["step"] for entry in workflow.manifest] == ["sub"] + + def test_a_parent_that_declares_no_result_leaves_the_child_saving(self, tmp_path): + workflow = self._compose(tmp_path, parent_result=False) + + written = self._written(tmp_path) + assert len(written) == 1, written + assert [entry["step"] for entry in workflow.manifest] == ["sub", "write"] + + def test_a_composed_file_carries_the_parent_step_name(self, tmp_path): + self._compose(tmp_path, parent_result=False) + + assert os.path.basename(self._written(tmp_path)[0]).startswith("sub.child-write") + + +class TestSubWorkflowValidation: + """A sub-workflow path that cannot resolve is a validation error, not a + run that fails 0.6 s in after the pre-flight said valid (#89).""" + + def _tree(self, tmp_path): + workflows = tmp_path / "workflows" + workflows.mkdir() + return workflows + + def _parent(self, workflows, path_value, arguments=None): + import json + + parent = { + "id": "parent", + "steps": [ + { + "name": "sub", + "workflow": {"path": path_value, "arguments": arguments or {}}, + "result": {"content_type": "image/jpeg"}, + } + ], + } + parent_path = workflows / "parent.json" + parent_path.write_text(json.dumps(parent)) + + from dw.workflow import workflow_from_file + + return workflow_from_file( + str(parent_path), str(workflows.parent / "outputs"), str(workflows) + ) + + def test_a_path_that_resolves_nowhere_is_an_error(self, tmp_path): + workflow = self._parent( + self._tree(tmp_path), "templates/does-not-exist-at-all" + ) + + errors = workflow.validation_errors() + + assert [e["path"] for e in errors] == ["steps[0].workflow.path"] + assert "does-not-exist-at-all" in errors[0]["message"] + + def test_a_path_outside_the_root_is_an_error(self, tmp_path): + import json + + workflows = self._tree(tmp_path) + outside = tmp_path / "outside.json" + outside.write_text(json.dumps({"id": "x", "steps": []})) + workflow = self._parent(workflows, str(outside)) + + errors = workflow.validation_errors() + + assert [e["path"] for e in errors] == ["steps[0].workflow.path"] + + def test_a_workflow_that_composes_itself_is_a_cycle(self, tmp_path): + import json + + workflows = self._tree(tmp_path) + definition = { + "id": "loop", + "steps": [ + { + "name": "sub", + "workflow": {"path": "loop.json", "arguments": {}}, + "result": {"content_type": "image/jpeg"}, + } + ], + } + path = workflows / "loop.json" + path.write_text(json.dumps(definition)) + + from dw.workflow import workflow_from_file + + workflow = workflow_from_file( + str(path), str(tmp_path / "outputs"), str(workflows) + ) + errors = workflow.validation_errors() + + assert any("cycle" in e["message"] for e in errors), errors + + def test_a_child_that_does_not_validate_is_reported_under_the_step( + self, tmp_path + ): + import json + + workflows = self._tree(tmp_path) + (workflows / "child.json").write_text( + json.dumps({"id": "child", "steps": [{"name": "broken"}]}) + ) + workflow = self._parent(workflows, "child") + + errors = workflow.validation_errors() + + assert errors + assert errors[0]["path"].startswith("steps[0].workflow.path -> ") + + def test_a_resolvable_child_validates_clean(self, tmp_path): + import json + + workflows = self._tree(tmp_path) + (workflows / "child.json").write_text( + json.dumps( + { + "id": "child", + "variables": {"prompt": "a cat"}, + "steps": [ + { + "name": "noop", + "task": { + "command": "compose_text", + "arguments": {"parts": ["variable:prompt"]}, + }, + "result": {"content_type": "text/plain"}, + } + ], + } + ) + ) + workflow = self._parent(workflows, "child", {"prompt": "a dog"}) + + assert workflow.validation_errors() == [] + assert workflow.sub_workflow_warnings() == [] + + def test_an_argument_the_child_does_not_declare_warns(self, tmp_path): + import json + + workflows = self._tree(tmp_path) + (workflows / "child.json").write_text( + json.dumps( + { + "id": "child", + "variables": {"prompt": "a cat"}, + "steps": [ + { + "name": "noop", + "task": { + "command": "compose_text", + "arguments": {"parts": ["variable:prompt"]}, + }, + "result": {"content_type": "text/plain"}, + } + ], + } + ) + ) + workflow = self._parent(workflows, "child", {"promt": "a dog"}) + + warnings = workflow.sub_workflow_warnings() + + assert [w["path"] for w in warnings] == [ + "steps[0].workflow.arguments.promt" + ] + assert "declares no variable" in warnings[0]["message"] From 6306f09dde092b5700f057bdc04a8f6fae99c43a Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 19:34:57 -0500 Subject: [PATCH 03/16] docs(catalog): #91 - say what a 'cost' is, and propose deriving one from job history The listing answers cost_basis: curated, and list_workflows says what that means: a cost block is a figure a maintainer measured once on the devices it names, never derived, so null means nobody wrote one down rather than 'this box has never run it'. The measured-from-history half is a proposal (docs/proposals/measured-cost-from-job-history.md) - it changes what the field is for every consumer, and the comparability rule is a decision. --- docs/SERVER.md | 6 +- .../measured-cost-from-job-history.md | 158 ++++++++++++++++++ dw/server/app.py | 7 + dw_mcp/catalog.py | 11 +- dw_mcp/server.py | 11 +- tests/test_server.py | 10 ++ 6 files changed, 198 insertions(+), 5 deletions(-) create mode 100644 docs/proposals/measured-cost-from-job-history.md diff --git a/docs/SERVER.md b/docs/SERVER.md index e31dc823..fdb3c4b1 100644 --- a/docs/SERVER.md +++ b/docs/SERVER.md @@ -265,7 +265,11 @@ The editor's forms come from these; they are just as usable from scripts: the maintainer-measured `{device, name, vram_gb, minutes}` runs, or `null` when nobody has measured it - a list-driven workflow's `cost` entry may also carry a measured `per_entry` (`{variable, minutes, entries}`), the - cost of one entry of the list it was measured against. A `models/` entry + cost of one entry of the list it was measured against. The response's + `cost_basis` says what that is - `curated`: figures a maintainer measured + once and wrote into the workflow, never derived from this server's own job + history, so `null` means nobody wrote one down rather than "this box has + never run it". A `models/` entry takes its `shape` and `traits` from the template it configures and keeps its own `cost`. A list-driven workflow (one with a `for_each` step) also carries `lists`: per list variable, the fields an entry takes, the steps diff --git a/docs/proposals/measured-cost-from-job-history.md b/docs/proposals/measured-cost-from-job-history.md new file mode 100644 index 00000000..3a4874e0 --- /dev/null +++ b/docs/proposals/measured-cost-from-job-history.md @@ -0,0 +1,158 @@ +# Proposal: derive a workflow's cost from this server's own job history + +Status: **design only** - written in answer to issue #91 (tester feedback), +which asks for three things and gets the cheapest of them shipped without a +proposal. This document covers the expensive one. Written by the implementer +agent (model `opus`, provider `anthropic`). + +## The report + +`list_workflows(shape="shot", traits="identity-referenced")` answers +`cost: null` for seven of eight entries, including +`templates/minimax/reference-to-video`, on a box that has run that template +five times at the same size. The consumer is instructed to quote a price +before spending GPU minutes (`run_workflow`'s own description, the MCP +instructions), and for the template its series actually uses, the catalog +says "unknown". The tester has carried a hand-kept cost table across six +cycles, re-deriving it from job history each time it is lost - from +`started_at`/`finished_at` on jobs this server stored. + +Five runs of that template, one RTX 3090, 124 frames at 960x544 / 20 steps: +511, 464.9, 478, ~460, ~467 seconds. ~7.8 minutes, spread under 10% - a +tighter figure than the 10.1 the one populated entry claims by hand. + +## What shipped without this proposal + +`GET /api/workflows` now answers `cost_basis: "curated"`, and +`list_workflows`' description says what that means: a `cost` block is a +figure a maintainer measured once on the devices it names and wrote into the +workflow; nothing derives one from job history; `null` means nobody wrote one +down, not that the run is cheap or that this box has never run it. That is +the tester's option 3, and it stops `null` reading as a data gap. + +It does not give anyone a number. + +## Why the rest is a proposal and not a fix + +`cost` is documented in `dw/workflow_schema.json` as *"Measured runs, one per +device the maintainer measured on. **Never derived**; absent means unknown."* +That sentence is load-bearing: a curated figure is a claim a person stands +behind, on a named card, at a named size. Deriving one changes what the field +*is*, for every consumer that reads it - which is the "new concept consumers +would have to learn" bar. It also decides questions with no obvious answer +(below). So: design, then Don, then code. + +## The shape + +A second field rather than a second meaning for the first: + +```json +"cost": [{"device": "cuda", "name": "RTX 3090", "vram_gb": 24, "minutes": 10.1}], +"observed": { + "device": "cuda", + "name": "NVIDIA GeForce RTX 3090", + "runs": 5, + "median_minutes": 7.8, + "p10_minutes": 7.6, + "p90_minutes": 8.5, + "since": "2026-09-08T14:02:11Z", + "comparable": "same-arguments" +} +``` + +`cost` keeps meaning exactly what it means today. `observed` is this +server's own history, always about *this* box's accelerator, and absent when +there is nothing to report. A consumer quoting a price prefers `observed` +when it is there (it is this machine, measured), falls back to `cost`, and +says "unknown" only when neither exists. `cost_basis` stays, and the listing +gains nothing else. + +### Where the numbers come from + +`jobs.sqlite` already holds, per job: the workflow name, `started_at`, +`finished_at`, `status`, `run_id`/`run_dir`, and the workspace. The manifest +in the run directory holds the realized workflow - the arguments folded in +and the seed pinned (`dw/realize.py`). The device is the server's, from +`get_server_info`. + +So the query is: finished jobs, `status = completed`, for one workflow +identity, on this device, most recent N. Median of +`finished_at - started_at`. + +### The four questions that make it a design + +1. **What counts as comparable?** A 141-frame run does not inform a + 124-frame estimate, and a different `weights_dtype` is a different model. + Three options, cheapest first: + - *Ignore it.* Report the median of every run of that workflow name, with + `runs` and a spread. Honest if the spread is published, useless for a + workflow whose variables move cost by 3x. + - *Bucket by the arguments that move cost.* Requires naming them, which + is per-workflow knowledge the catalog does not have. A `cost_drivers` + key in the workflow (`["num_frames", "num_inference_steps"]`) would + declare them, and the bucket key is those values. This is the honest + one and it is more schema. + - *Report the default-arguments runs only.* A run whose arguments equal + the workflow's variable defaults is comparable to the curated figure by + construction. Simple, exact, and thin - most real runs pass arguments. + + Recommendation: **bucket by declared drivers**, falling back to + default-arguments-only when a workflow declares none. It degrades to the + third option rather than to a wrong number. + +2. **Cold vs warm.** The tester's own `text-to-image` figures are 13.6 s and + 6.3 s - the same run, model on disk vs model resident. These are two + numbers and averaging them produces one that describes neither. The job + events already distinguish them (a `loading` phase that takes minutes vs + one that takes none), so the split is available: `median_minutes` warm, + `cold_minutes` when the run had to load. A consumer quoting a first run + of the session wants the cold one. + +3. **A cached run is not a run.** A seeded workflow whose every step hit the + step cache finishes in seconds and wrote nothing. Those jobs are already + flagged (`reused` on every manifest entry) and must be excluded, or the + median collapses toward zero for exactly the templates that get re-run + most. + +4. **Pruning.** Job history is prunable and a workspace is deletable. The + figures move when it happens, which is fine (`runs` says how much is + behind it), but a listing that quietly loses a number people relied on + should not be surprising - `since` and `runs` are what make it legible. + +### Cost of the feature + +A per-workflow aggregate over `jobs.sqlite`, cached like `workflow_details` +is (by the jobs table's own high-water mark rather than by mtime), computed +on listing. Reading the realized workflow of each candidate run to bucket by +drivers is the expensive part - one small JSON read per run, bounded by the +N most recent, and only for workflows the listing actually returns. Nothing +in the run path changes. No new storage; `finished_at - started_at` is +already recorded. + +Rough size: a new `dw/server/observed_cost.py` (aggregate + cache), a +`cost_drivers` key in the schema, the `observed` field in the compact and +full listings, the MCP description, docs, and tests. Half a day, most of it +the comparability rule. + +## What I would not do + +Overwrite `cost` with a derived figure, or let `observed` inherit the +`cost` shape closely enough to be mistaken for it. The distinction between +"a maintainer measured this on a 4090" and "this box averaged that last +week" is the whole value of reporting both. + +## Recommendation + +Ship it as `observed`, bucketed by declared `cost_drivers`, cold/warm split, +cached runs excluded. If that is more than the problem is worth, the +fallback that costs almost nothing is the third comparability option - +default-arguments runs only, with `runs` and `since` published - which would +have answered the tester's question today, because their five runs were all +at one size. + +## Open question for Don + +Is `cost` allowed to gain a sibling that is derived, or does "never derived" +apply to the whole of what the catalog says about price? If the latter, this +belongs in a separate tool (`get_workflow_history(name)`) rather than in the +listing, and the consumer pays a call for it. diff --git a/dw/server/app.py b/dw/server/app.py index 12b85a10..19611e35 100644 --- a/dw/server/app.py +++ b/dw/server/app.py @@ -1513,6 +1513,13 @@ def list_workflows( "sources": [source.to_dict() for source in sources], "workflows": sorted(details), "details": details, + # What a `cost` is, and so what a null one means. Curated: + # figures a maintainer measured once on the devices named and + # wrote into the workflow - nothing derives them from this + # server's own job history, so null means nobody wrote one + # down, not that the run is cheap or that this box has never + # run it (#91) + "cost_basis": "curated", } @app.put("/api/workflows/{name:path}") diff --git a/dw_mcp/catalog.py b/dw_mcp/catalog.py index 27997f13..99eb180a 100644 --- a/dw_mcp/catalog.py +++ b/dw_mcp/catalog.py @@ -13,7 +13,16 @@ def list_workflows( shape, traits, cost, output kinds and variable names per workflow - what choosing one needs and nothing that reading one needs. Templates only unless `include_models` or `configures` asks for the model configs - of one template. `get_workflow` has the full definition.""" + of one template. `get_workflow` has the full definition. + + `cost_basis` in the answer says what a `cost` is: `curated` means a + maintainer measured it once, on the devices the entry names, and wrote + it into the workflow. Nothing derives one from this server's own job + history, so `cost: null` means nobody wrote a figure down - not that + the run is cheap, and not that this box has never run it. For a + template with no figure, `list_jobs` on earlier runs of it carries + `started_at`/`finished_at`, which is the measurement this server + actually holds.""" params = {"view": "compact"} if shape: params["shape"] = shape diff --git a/dw_mcp/server.py b/dw_mcp/server.py index b1d38605..4dcf17f5 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -167,9 +167,14 @@ def list_workflows( further (comma-separated, all must match): has-audio, chained, image-conditioned, identity-referenced, needs-input-media, composes-workflows. Each entry carries a one-line `summary`, its - `shape` and `traits` (what it needs supplied), `cost` (measured - runs per device; null means unknown - call `get_memory` and say - so), output kinds and variable names. `lists`, present for a + `shape` and `traits` (what it needs supplied), `cost` (curated: + figures a maintainer measured once on the devices named and wrote + into the workflow, never derived from this server's job history - + so null means nobody wrote one down, not that the run is cheap; + the answer's `cost_basis` says as much. For a null one, earlier + runs of the same template in `list_jobs` carry + `started_at`/`finished_at`, which is the measurement this box + actually holds), output kinds and variable names. `lists`, present for a list-driven workflow, names per list variable the fields an entry takes, the steps run over it and the default's length; there `cost[].per_entry`, when present, is the measured cost of one diff --git a/tests/test_server.py b/tests/test_server.py index 9b3fbc40..74604d72 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -516,6 +516,16 @@ def test_validate_accepts_a_stored_workflow_name(server, tmp_path): assert any("guidance_scael" in w for w in result["warnings"]) +def test_the_listing_says_what_a_cost_is(server): + """`cost: null` covered both "nobody measured it" and "not measured for + your device"; the listing now says which kind of figure a cost is + (#91).""" + with server(success_script) as client: + answer = client.get("/api/workflows").json() + + assert answer["cost_basis"] == "curated" + + def test_validate_reports_a_sub_workflow_path_that_resolves_nowhere(server): """The pre-flight is documented as "this will run", and a composed step naming a workflow the server cannot reach used to come back valid and From 57d801b64aad86b65a4ee1f47d04bdac79c1367f Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 19:39:53 -0500 Subject: [PATCH 04/16] fix(realize): #90 - digest a sub-workflow the run could actually load; format --- dw/realize.py | 21 ++++++++++----------- dw/server/jobs.py | 4 +--- dw/workflow.py | 12 +++++------- tests/test_workflow.py | 16 ++++++---------- 4 files changed, 22 insertions(+), 31 deletions(-) diff --git a/dw/realize.py b/dw/realize.py index 5368f218..eb980019 100644 --- a/dw/realize.py +++ b/dw/realize.py @@ -31,6 +31,7 @@ resolve_output_reference, ) from .security import SecurityError, validate_workflow_path +from .workflow_sources import resolve_sub_workflow, SubWorkflowNotFound from .variables import set_variables logger = logging.getLogger("dw") @@ -210,20 +211,18 @@ def scan(value): def _digest(path, base_dir, workflow_dir): """The SHA-256 of a sub-workflow file, or None when it cannot be read. - Resolved the way `Workflow.create_step_action` resolves it - relative to - the referencing file's directory, then through `validate_workflow_path` - confined to `workflow_dir` - so a path this run could not have loaded is - not one realization reads either. + Resolved the way `Workflow.create_step_action` resolves it - beside the + referencing file, then across the workflow search path, then through + `validate_workflow_path` confined to the root it came from - so a path + this run could not have loaded is not one realization reads either, and + a catalog name the run composed is digested rather than recorded as + unreadable (#90). """ try: - candidate = ( - path - if os.path.isabs(path) - else os.path.normpath(os.path.join(base_dir or ".", path)) - ) - validated = validate_workflow_path(candidate, workflow_dir) + candidate, root = resolve_sub_workflow(path, base_dir or ".", workflow_dir) + validated = validate_workflow_path(candidate, root) with open(validated, "rb") as file: return hashlib.sha256(file.read()).hexdigest() - except (SecurityError, OSError, ValueError) as e: + except (SecurityError, OSError, ValueError, SubWorkflowNotFound) as e: logger.debug(f"No digest for sub-workflow {path}: {e}") return None diff --git a/dw/server/jobs.py b/dw/server/jobs.py index cc97a68d..c659c5dd 100644 --- a/dw/server/jobs.py +++ b/dw/server/jobs.py @@ -423,9 +423,7 @@ def _note_progress(self, event): # one and the step name stays the child's (#90) self.parent_step = event.get("parent_step") self.step_index = event.get("parent_index", event.get("index")) - self.total_steps = event.get( - "parent_total_steps", event.get("total_steps") - ) + self.total_steps = event.get("parent_total_steps", event.get("total_steps")) # A new step's denoise loop has not started; the previous step's # count would read as this one's progress self.denoise_step = None diff --git a/dw/workflow.py b/dw/workflow.py index 3b2890c9..3c559509 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -449,7 +449,9 @@ def resolve_sub_workflow_path(self, path): or "/" in builtin_name or "\\" in builtin_name ): - raise InvalidInputError(f"Invalid builtin workflow name: {builtin_name}") + raise InvalidInputError( + f"Invalid builtin workflow name: {builtin_name}" + ) confine_to = builtin_root() resolved = os.path.join(confine_to, builtin_name) if not os.path.isfile(resolved): @@ -502,9 +504,7 @@ def sub_workflow_errors(self, expanded, source_indices=None, composing=None): try: child = workflow_from_file(resolved, self.output_dir, root) except Exception as e: - errors.append( - {"path": where, "message": f"Sub-workflow '{path}': {e}"} - ) + errors.append({"path": where, "message": f"Sub-workflow '{path}': {e}"}) continue for error in child.validation_errors(composing=composing + [resolved]): errors.append( @@ -521,9 +521,7 @@ def sub_workflow_warnings(self, expanded=None): and composition is exactly where a name drifts (#89).""" warnings = [] try: - expanded = ( - expanded if expanded is not None else self.expanded_definition() - ) + expanded = expanded if expanded is not None else self.expanded_definition() except Exception: return warnings for index, step in enumerate(expanded.get("steps", []) or []): diff --git a/tests/test_workflow.py b/tests/test_workflow.py index 402142f8..ffa9b89f 100644 --- a/tests/test_workflow.py +++ b/tests/test_workflow.py @@ -1066,7 +1066,9 @@ def test_a_parent_that_declares_no_result_leaves_the_child_saving(self, tmp_path def test_a_composed_file_carries_the_parent_step_name(self, tmp_path): self._compose(tmp_path, parent_result=False) - assert os.path.basename(self._written(tmp_path)[0]).startswith("sub.child-write") + assert os.path.basename(self._written(tmp_path)[0]).startswith( + "sub.child-write" + ) class TestSubWorkflowValidation: @@ -1101,9 +1103,7 @@ def _parent(self, workflows, path_value, arguments=None): ) def test_a_path_that_resolves_nowhere_is_an_error(self, tmp_path): - workflow = self._parent( - self._tree(tmp_path), "templates/does-not-exist-at-all" - ) + workflow = self._parent(self._tree(tmp_path), "templates/does-not-exist-at-all") errors = workflow.validation_errors() @@ -1148,9 +1148,7 @@ def test_a_workflow_that_composes_itself_is_a_cycle(self, tmp_path): assert any("cycle" in e["message"] for e in errors), errors - def test_a_child_that_does_not_validate_is_reported_under_the_step( - self, tmp_path - ): + def test_a_child_that_does_not_validate_is_reported_under_the_step(self, tmp_path): import json workflows = self._tree(tmp_path) @@ -1217,7 +1215,5 @@ def test_an_argument_the_child_does_not_declare_warns(self, tmp_path): warnings = workflow.sub_workflow_warnings() - assert [w["path"] for w in warnings] == [ - "steps[0].workflow.arguments.promt" - ] + assert [w["path"] for w in warnings] == ["steps[0].workflow.arguments.promt"] assert "declares no variable" in warnings[0]["message"] From 2878c1e56895c8462f51e3f6f8158d1ab9c51f95 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 19:47:38 -0500 Subject: [PATCH 05/16] fix(engine): #90 - a missing candidate is reported as looked-at, not as refused --- dw/workflow_sources.py | 7 ++++--- tests/test_workflow.py | 4 ++++ 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/dw/workflow_sources.py b/dw/workflow_sources.py index 323d899b..74f47b1e 100644 --- a/dw/workflow_sources.py +++ b/dw/workflow_sources.py @@ -251,9 +251,10 @@ def resolve_sub_workflow(path, base_dir, confine_to): for root in roots: source = WorkflowSource(root, EXAMPLES_ORIGIN, False) - # resolve_in_source supplies '.json' itself, and refuses a name that - # would traverse out of the root - candidate = resolve_in_source(source, path) + # allow_create, because what is being asked is where the name would + # land rather than whether something is there - None means the name + # traverses out of the root, and the file check is the next line + candidate = resolve_in_source(source, path, allow_create=True) if candidate is None: tried.append(f"{os.path.join(root, path)} (outside the root)") continue diff --git a/tests/test_workflow.py b/tests/test_workflow.py index ffa9b89f..6ef16b0e 100644 --- a/tests/test_workflow.py +++ b/tests/test_workflow.py @@ -989,6 +989,10 @@ def test_a_name_that_resolves_nowhere_says_where_it_looked(self, tmp_path): message = str(exc_info.value) assert "does-not-exist" in message assert "Looked in" in message + # every candidate is a real path it looked at, '.json' supplied - + # not a name reported as refused when it was simply absent + assert "does-not-exist.json" in message + assert "outside the root" not in message def test_a_relative_path_beside_the_file_still_wins(self, tmp_path): """The '../models/x.json' form every template uses is unchanged.""" From 2dbd22e9edf9136cd24683cc5e3ba9ddb59f5da4 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 19:49:10 -0500 Subject: [PATCH 06/16] fix(engine): #90 - each candidate path named once in the not-found message --- dw/workflow_sources.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/dw/workflow_sources.py b/dw/workflow_sources.py index 74f47b1e..9646bbdd 100644 --- a/dw/workflow_sources.py +++ b/dw/workflow_sources.py @@ -193,7 +193,9 @@ class SubWorkflowNotFound(Exception): def __init__(self, path, tried): self.path = path - self.tried = list(tried) + # In order, each candidate once - two roots can resolve one name to + # the same file, and saying so twice reads as two failures + self.tried = list(dict.fromkeys(tried)) detail = "\n ".join(self.tried) super().__init__( f"Sub-workflow '{path}' could not be resolved. It is read as a " From e4f64b172c812c09015437d3d35cb5e6acf55959 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 21:16:48 -0500 Subject: [PATCH 07/16] fix(ui): match a run's step names to the flow graph's nodes Workflow.run expands for_each before it emits, so the event stream names the members (base@open) while the flow graph is drawn from the definition, which keeps for_each unexpanded and a composed step as one node. Every comparison in FlowView was between two names that could no longer be equal, so the active step and the finished ones stopped showing. A sub-workflow's inner step mismatched the same way, which also left the Progress list with no active row. runstate.ts reduces both sides to one name - a for_each member to its group, a child's step to the parent_step the engine already puts on its events. A group is finished only when every member is, and only a top-level step_end finishes a composed step. The job page also now says why a step wrote nothing: the Results panel lists the steps that ran and wrote no file, with the reason read off the definition (result.save is false, or no result.content_type), and appears when that list is the only thing there is to say. A workflow that keeps most of its steps in memory read as a run whose outputs had gone missing. Co-Authored-By: Claude --- ui/src/lib/pages/JobPage.svelte | 94 +++++++++++++++++---- ui/src/lib/pages/JobPage.test.ts | 138 ++++++++++++++++++++++++++++++- ui/src/lib/results.test.ts | 120 ++++++++++++++++++++++++++- ui/src/lib/results.ts | 101 ++++++++++++++++++++++ ui/src/lib/runstate.test.ts | 54 ++++++++++++ ui/src/lib/runstate.ts | 61 ++++++++++++++ 6 files changed, 552 insertions(+), 16 deletions(-) create mode 100644 ui/src/lib/runstate.test.ts create mode 100644 ui/src/lib/runstate.ts diff --git a/ui/src/lib/pages/JobPage.svelte b/ui/src/lib/pages/JobPage.svelte index af543af6..ad91dc09 100644 --- a/ui/src/lib/pages/JobPage.svelte +++ b/ui/src/lib/pages/JobPage.svelte @@ -10,7 +10,12 @@ import { ApiError, api, outputUrl, streamJobEvents } from '../api' import { confirmDialog } from '../confirm.svelte' import { go } from '../router.svelte' - import { groupResultFiles, sectionBySubfolder } from '../results' + import { + groupResultFiles, + sectionBySubfolder, + unsavedSteps, + } from '../results' + import { finishedNodes, flowNodeName } from '../runstate' import { stepProgress } from '../progress' import FlowView from '../editor/FlowView.svelte' import CopyButton from '../CopyButton.svelte' @@ -157,15 +162,27 @@ const steps = $derived( (events.find((e) => e.event === 'workflow_start')?.steps as string[]) ?? [], ) - const currentStep = $derived( - [...events].reverse().find((e) => e.event === 'step_start')?.step as - string | undefined, + // The step_start the run is on, as the engine named it: a for_each member + // keeps its '@', a sub-workflow's inner step its own name. The Progress + // list is written in these names, so it reads off this one. + const stepStart = $derived( + [...events].reverse().find((e) => e.event === 'step_start'), ) - const finishedSteps = $derived( - new Set( - events.filter((e) => e.event === 'step_end').map((e) => e.step as string), + const currentStep = $derived(stepStart?.step as string | undefined) + // The same step in the name the definition gives it, which is what the + // flow graph's nodes are called - see runstate.ts for why the two differ + const activeNode = $derived( + flowNodeName( + stepStart?.step as string | undefined, + stepStart?.parent_step as string | undefined, ), ) + // A sub-workflow's inner step has no row of its own in the Progress list, + // so what lights up while one runs is the composed step that queued it + const listStep = $derived( + steps.includes(currentStep ?? '') ? currentStep : activeNode, + ) + const finishedSteps = $derived(finishedNodes(events as JobEvent[], steps)) // Scoped to the step now running: its phase, and its own denoise counter const progress = $derived(stepProgress(events as JobEvent[])) const denoise = $derived(progress.denoise) @@ -246,6 +263,13 @@ const sectioned = $derived( sections.length > 1 || sections[0].subfolder !== '', ) + // Steps that ran and wrote no file, with the definition's reason for each. + // Without this a workflow that keeps most of its steps in memory shows a + // short Results list and no word about the rest, which reads as outputs + // that went missing rather than as a choice the workflow made + const unsaved = $derived( + unsavedSteps(job?.manifest, events as JobEvent[], definition), + ) const running = $derived(job !== null && !TERMINAL.includes(job.status)) // A cancel requested while loading a model or running a task step has no // checkpoint to catch it until that phase finishes - without this the UI @@ -352,13 +376,13 @@
- {step} - {#if step === currentStep && running} + {#if step === listStep && running} {#if denoise}
Workflow {/if} - {#if fileGroups.length} + {#if fileGroups.length || unsaved.length}

Results

{#if allReused} @@ -491,6 +515,27 @@ {/each} {/each} {/each} + {#if unsaved.length} +

Steps that wrote nothing

+
    + {#each unsaved as entry (entry.node)} +
  • + {entry.node} + {#if entry.members.length} + × {entry.members.length} + {/if} + + {#if entry.reason}{entry.reason.key} + {entry.reason.detail}{:else}no file written{/if} + +
  • + {/each} +
+ {/if}
{/if} @@ -682,4 +727,25 @@ .subhead + .stephead { margin-top: 0; } + /* The steps whose files are deliberately absent: the step name is what + the engine resolves, the reason is written for a person */ + .unsaved { + list-style: none; + margin: 0; + padding: 0; + font-size: var(--t-sm); + } + .unsaved li { + display: flex; + align-items: baseline; + flex-wrap: wrap; + gap: 0.5rem; + padding: 0.15rem 0; + } + .unsaved code { + font-size: 0.8rem; + } + .unsaved .why { + min-width: 0; + } diff --git a/ui/src/lib/pages/JobPage.test.ts b/ui/src/lib/pages/JobPage.test.ts index 84a20d16..bf4282b3 100644 --- a/ui/src/lib/pages/JobPage.test.ts +++ b/ui/src/lib/pages/JobPage.test.ts @@ -13,13 +13,17 @@ const stream = vi.hoisted(() => ({ const metadata = vi.hoisted(() => ({ byFile: {} as Record>, })) +// The definition the job ran, for the flow view and the unsaved reasons +const ran = vi.hoisted(() => ({ + definition: null as Record | null, +})) vi.mock('../api', () => ({ ApiError: class ApiError extends Error {}, api: { getJob: vi.fn(() => Promise.resolve(detail.job)), getJobWorkflow: vi.fn(() => - Promise.resolve({ definition: null, seed_variable: null }), + Promise.resolve({ definition: ran.definition, seed_variable: null }), ), galleryMetadata: vi.fn((name: string) => Promise.resolve({ @@ -62,9 +66,24 @@ afterEach(() => { cleanup() stream.onEvent = null metadata.byFile = {} + ran.definition = null vi.mocked(api.galleryMetadata).mockClear() }) +/** The flow view's box for one step of the definition. */ +function nodeFor(container: HTMLElement, name: string) { + return [...container.querySelectorAll('g.node')].find((node) => + node.getAttribute('aria-label')?.startsWith(`step ${name},`), + )! +} + +/** The page's "wrote nothing" rows, whitespace flattened. */ +function unsavedRows(container: HTMLElement) { + return [...container.querySelectorAll('.unsaved li')].map((li) => + li.textContent?.replace(/\s+/g, ' ').trim(), + ) +} + it('shows each image beside what made it, with a download link, and never probes a video', async () => { metadata.byFile['a.png'] = { model_name: 'org/model', @@ -149,3 +168,120 @@ it('places a live step_end under its subfolder before the manifest arrives', asy ).toBeTruthy(), ) }) + +it('lights the for_each step in the flow chart while one of its members runs', async () => { + // The graph is drawn from the definition, where for_each is one step; the + // engine reports the members, so the two only meet at the group name + ran.definition = { + steps: [ + { + name: 'shot', + pipeline: { configuration: { component_type: 'Fake' } }, + for_each: 'variable:shots', + }, + { name: 'episode', task: { command: 'mux' } }, + ], + } + detail.job = { ...job([]), status: 'running', finished_at: null } + const { container } = render(JobPage, { jobId: 'j1' }) + await waitFor(() => expect(stream.onEvent).not.toBeNull()) + stream.onEvent!({ + seq: 1, + event: 'workflow_start', + steps: ['shot@open', 'shot@reveal', 'episode'], + }) + stream.onEvent!({ seq: 2, event: 'step_start', step: 'shot@open' }) + await waitFor(() => + expect(nodeFor(container, 'shot').classList.contains('active')).toBe(true), + ) + + // One member of two down: the group is not behind us yet + stream.onEvent!({ seq: 3, event: 'step_end', step: 'shot@open', files: [] }) + await waitFor(() => + expect(nodeFor(container, 'shot').classList.contains('active')).toBe(true), + ) + expect(nodeFor(container, 'shot').classList.contains('done')).toBe(false) + + stream.onEvent!({ seq: 4, event: 'step_end', step: 'shot@reveal', files: [] }) + stream.onEvent!({ seq: 5, event: 'step_start', step: 'episode' }) + await waitFor(() => + expect(nodeFor(container, 'shot').classList.contains('done')).toBe(true), + ) + expect(nodeFor(container, 'episode').classList.contains('active')).toBe(true) +}) + +it('keeps the Progress list on the composed step while its child runs', async () => { + ran.definition = { + steps: [ + { name: 'shot1', workflow: { path: 'child.json' } }, + { name: 'episode', task: { command: 'mux' } }, + ], + } + detail.job = { ...job([]), status: 'running', finished_at: null } + const { container } = render(JobPage, { jobId: 'j1' }) + await waitFor(() => expect(stream.onEvent).not.toBeNull()) + stream.onEvent!({ + seq: 1, + event: 'workflow_start', + steps: ['shot1', 'episode'], + }) + // What a child emits: its own step name, with the parent's alongside + stream.onEvent!({ + seq: 2, + event: 'step_start', + step: 'reference_to_video_audio', + parent_step: 'shot1', + }) + await waitFor(() => + expect(nodeFor(container, 'shot1').classList.contains('active')).toBe(true), + ) + const dotFor = (name: string) => + [...container.querySelectorAll('.step')] + .find( + (row) => + row.querySelector('span:nth-child(2)')?.textContent?.trim() === name, + )! + .querySelector('.dot')! + expect(dotFor('shot1').classList.contains('active')).toBe(true) + expect(dotFor('episode').classList.contains('active')).toBe(false) +}) + +it('says why a step wrote nothing, so a deliberate non-output is not a missing one', async () => { + ran.definition = { + steps: [ + { name: 'base', pipeline: {}, result: { save: false } }, + { name: 'edit', task: { command: 'mux' } }, + { name: 'film', pipeline: {}, result: { content_type: 'video/mp4' } }, + ], + } + detail.job = job([ + { step: 'base', files: [] }, + { step: 'edit', files: [] }, + { step: 'film', files: ['final/film.mp4'], subfolder: 'final' }, + ]) + const { container } = render(JobPage, { jobId: 'j1' }) + await waitFor(() => + expect(screen.getByText('Steps that wrote nothing')).toBeTruthy(), + ) + expect(unsavedRows(container)).toEqual([ + 'base result.save is false, so the step is kept in memory and never written', + 'edit result.content_type is not declared, so there is no file type to write', + ]) + // The step that did write is in the results above, not in this list + expect(screen.getByRole('heading', { level: 3, name: 'final/' })).toBeTruthy() +}) + +it('explains a run that wrote nothing at all rather than showing no results', async () => { + ran.definition = { steps: [{ name: 'base', result: { save: false } }] } + detail.job = job([{ step: 'base', files: [] }]) + const { container } = render(JobPage, { jobId: 'j1' }) + await waitFor(() => + expect(screen.getByText('Steps that wrote nothing')).toBeTruthy(), + ) + expect( + screen.getByRole('heading', { level: 2, name: 'Results' }), + ).toBeTruthy() + expect(unsavedRows(container)).toEqual([ + 'base result.save is false, so the step is kept in memory and never written', + ]) +}) diff --git a/ui/src/lib/results.test.ts b/ui/src/lib/results.test.ts index 681eddd1..4681dbe4 100644 --- a/ui/src/lib/results.test.ts +++ b/ui/src/lib/results.test.ts @@ -1,5 +1,10 @@ import { describe, expect, it } from 'vitest' -import { groupResultFiles, sectionBySubfolder } from './results' +import { + groupResultFiles, + sectionBySubfolder, + unsavedReason, + unsavedSteps, +} from './results' import type { JobEvent } from './types' const stepEnd = ( @@ -160,3 +165,116 @@ describe('sectionBySubfolder', () => { expect(sections.map((s) => s.subfolder)).toEqual(['final', 'shots/act-1']) }) }) + +describe('unsavedReason', () => { + it('names result.save when the workflow asked for no file', () => { + expect( + unsavedReason({ result: { save: false, content_type: 'image/png' } }), + ).toEqual({ + key: 'result.save', + detail: 'is false, so the step is kept in memory and never written', + }) + }) + + it('names result.content_type when there is no type to write', () => { + expect(unsavedReason({ task: { command: 'mux' } })).toEqual({ + key: 'result.content_type', + detail: 'is not declared, so there is no file type to write', + }) + }) + + it('has no reason for a step that declares a file', () => { + expect( + unsavedReason({ result: { content_type: 'video/mp4', save: true } }), + ).toBeNull() + }) +}) + +describe('unsavedSteps', () => { + const definition = (steps: Array>) => ({ steps }) + const steps = definition([ + { name: 'base', result: { save: false } }, + { name: 'edit', task: { command: 'mux' } }, + { name: 'film', result: { content_type: 'video/mp4' } }, + ]) + + it('names every step that ran and wrote nothing, with the reason', () => { + const unsaved = unsavedSteps( + [ + { step: 'base', files: [] }, + { step: 'edit', files: [] }, + { step: 'film', files: ['final/film.mp4'], subfolder: 'final' }, + ], + [], + steps, + ) + expect(unsaved).toEqual([ + { + node: 'base', + members: [], + reason: { + key: 'result.save', + detail: 'is false, so the step is kept in memory and never written', + }, + }, + { + node: 'edit', + members: [], + reason: { + key: 'result.content_type', + detail: 'is not declared, so there is no file type to write', + }, + }, + ]) + }) + + it('collapses a for_each group into the one step the workflow declares', () => { + const unsaved = unsavedSteps( + [ + { step: 'base@open', files: [] }, + { step: 'base@reveal', files: [] }, + ], + [], + steps, + ) + expect(unsaved).toEqual([ + { + node: 'base', + members: ['base@open', 'base@reveal'], + reason: { + key: 'result.save', + detail: 'is false, so the step is kept in memory and never written', + }, + }, + ]) + }) + + it('lists only steps that have finished, so a running job is not accused', () => { + const events = [ + { seq: 1, event: 'workflow_start', steps: ['base', 'edit', 'film'] }, + { seq: 2, event: 'step_end', step: 'base', files: [] }, + ] as unknown as JobEvent[] + expect(unsavedSteps(undefined, events, steps).map((s) => s.node)).toEqual([ + 'base', + ]) + }) + + it("ignores a sub-workflow's inner steps, which are no node of this graph", () => { + const composed = definition([ + { name: 'shot1_amnesty', workflow: { path: 'child.json' } }, + ]) + const unsaved = unsavedSteps( + [ + { step: 'reference_to_video_audio', files: [] }, + { step: 'shot1_amnesty', files: ['intermediate/a.mp4'] }, + ], + [], + composed, + ) + expect(unsaved).toEqual([]) + }) + + it('has nothing to say without a definition to read the reasons off', () => { + expect(unsavedSteps([{ step: 'base', files: [] }], [], null)).toEqual([]) + }) +}) diff --git a/ui/src/lib/results.ts b/ui/src/lib/results.ts index c84e946c..0f3db546 100644 --- a/ui/src/lib/results.ts +++ b/ui/src/lib/results.ts @@ -1,3 +1,4 @@ +import { flowNodeName } from './runstate' import type { JobEvent, ManifestEntry, StepEndEvent } from './types' /** One step's output files, with where in the run directory they landed. */ @@ -86,3 +87,103 @@ export function sectionBySubfolder(groups: StepGroup[]): SubfolderSection[] { groups: bySubfolder.get(subfolder)!, })) } + +/** One step of a run that finished without writing a file. */ +export interface UnsavedStep { + /** The step's name in the definition. One entry covers a whole for_each + * group, since that is the one step the workflow declares. */ + node: string + /** The member names the engine ran it under, one per for_each entry - + * empty for a step that has no `for_each`, which is its own name. */ + members: string[] + /** Why nothing was written, read off the definition - null when the + * definition does not explain it. */ + reason: UnsavedReason | null +} + +/** The JSON key a step writes nothing because of, and what that means - + * kept apart so the page can set the key in the type the engine resolves + * literally. */ +export interface UnsavedReason { + key: string + detail: string +} + +/** Why a step's result block saves no file: the two ways the engine skips + * the save, both in `Result.save` (dw/result.py). */ +export function unsavedReason(step: Record): UnsavedReason | null { + const result = step.result + if (result && result.save === false) { + return { + key: 'result.save', + detail: 'is false, so the step is kept in memory and never written', + } + } + if (!result || typeof result.content_type !== 'string') { + return { + key: 'result.content_type', + detail: 'is not declared, so there is no file type to write', + } + } + return null +} + +/** The steps that ran and wrote nothing, with the reason the definition + * gives for each. + * + * The page groups a run's files by producing step and drops the empty + * groups, which made a workflow whose steps deliberately write nothing + * indistinguishable from a run whose outputs had gone missing: the steps + * were there, their files were not, and nothing said why. A step counts as + * having run when the manifest carries its entry - every step that finished + * gets one - or a top-level step_end named it, so a job still in flight + * never lists a step it has not reached yet, and a historical job, which + * has no events at all, reads correctly off the manifest alone. */ +export function unsavedSteps( + manifest: ManifestEntry[] | undefined, + events: JobEvent[], + definition: Record | null, +): UnsavedStep[] { + const byNode = new Map>() + for (const step of (definition?.steps ?? []) as Array>) { + if (typeof step?.name === 'string') byNode.set(step.name, step) + } + if (!byNode.size) return [] + + // An inner step of a sub-workflow is no node of the parent's graph, so + // the entry it rolls up into the manifest is no evidence about anything + // here - only a definition step, or a member of one, names a node + const nodeOf = (step: string): string | undefined => { + const node = flowNodeName(step) + return node && byNode.has(node) ? node : undefined + } + + const wrote = new Set() + for (const group of groupResultFiles(manifest, events)) { + const node = nodeOf(group.step) + if (node) wrote.add(node) + } + + const ran = new Map() + const note = (step: string) => { + const node = nodeOf(step) + if (!node) return + const members = ran.get(node) ?? [] + if (!members.includes(step)) members.push(step) + ran.set(node, members) + } + for (const entry of manifest ?? []) note(entry.step) + for (const event of events) { + if (event.event === 'step_end' && !event.parent_step && event.step) { + note(event.step as string) + } + } + + return [...ran] + .filter(([node]) => !wrote.has(node)) + .map(([node, members]) => ({ + node, + members: members.filter((member) => member !== node), + reason: unsavedReason(byNode.get(node)!), + })) +} diff --git a/ui/src/lib/runstate.test.ts b/ui/src/lib/runstate.test.ts new file mode 100644 index 00000000..43990901 --- /dev/null +++ b/ui/src/lib/runstate.test.ts @@ -0,0 +1,54 @@ +import { describe, expect, it } from 'vitest' +import { finishedNodes, flowNodeName } from './runstate' +import type { JobEvent } from './types' + +const ended = (step: string, parent_step?: string): JobEvent => + ({ seq: 0, event: 'step_end', step, parent_step }) as JobEvent + +describe('flowNodeName', () => { + it('reduces a for_each member to the group the definition declares', () => { + expect(flowNodeName('base@open')).toBe('base') + }) + + it("reduces a sub-workflow's inner step to the composed step that queued it", () => { + expect(flowNodeName('reference_to_video_audio', 'shot1_amnesty')).toBe( + 'shot1_amnesty', + ) + }) + + it('takes the parent even when the child step is a for_each group itself', () => { + expect(flowNodeName('base@open', 'shot@reveal')).toBe('shot') + }) + + it('leaves an ordinary step alone, and has nothing for a missing name', () => { + expect(flowNodeName('film')).toBe('film') + expect(flowNodeName(undefined)).toBeUndefined() + }) +}) + +describe('finishedNodes', () => { + const members = ['base@open', 'base@reveal', 'film'] + + it('finishes a for_each group only once every member has ended', () => { + expect(finishedNodes([ended('base@open')], members)).toEqual([]) + expect( + finishedNodes([ended('base@open'), ended('base@reveal')], members), + ).toEqual(['base']) + }) + + it('finishes a plain step on its own end', () => { + expect(finishedNodes([ended('film')], members)).toEqual(['film']) + }) + + it("does not finish a composed step for one of its child's inner steps", () => { + const names = ['shot1_amnesty', 'episode'] + const events = [ended('inner_a', 'shot1_amnesty')] + expect(finishedNodes(events, names)).toEqual([]) + events.push(ended('shot1_amnesty')) + expect(finishedNodes(events, names)).toEqual(['shot1_amnesty']) + }) + + it('has nothing finished before the run has reported anything', () => { + expect(finishedNodes([], members)).toEqual([]) + }) +}) diff --git a/ui/src/lib/runstate.ts b/ui/src/lib/runstate.ts new file mode 100644 index 00000000..8bd0fc9c --- /dev/null +++ b/ui/src/lib/runstate.ts @@ -0,0 +1,61 @@ +/** Where a run's event stream meets the definition it is running. + * + * The flow graph is drawn from the definition, where `for_each` is left + * unexpanded and a sub-workflow is one step, but the engine *runs* the + * definition expanded: it reports `base@open` where the definition says + * `base`, and a sub-workflow's inner step under its own name where the + * definition says the composed step. Nothing the page highlights matches + * until both sides are reduced to the same name - which is what broke the + * flow view's run state when list-driven steps arrived. */ + +import type { JobEvent } from './types' + +/** The for_each step/member separator - `dw/for_each.py`'s + * MEMBER_SEPARATOR, reserved in every step name, so the first one in a + * name always splits a member from its group. */ +const MEMBER_SEPARATOR = '@' + +/** The name the definition gives the step the engine ran. + * + * A for_each member reduces to its group (`base@open` -> `base`); anything + * a sub-workflow emitted reduces to the step that queued it, which the + * engine puts on every event a child emits as `parent_step` + * (`_parent_progress_fields`, dw/workflow.py). A plain step is its own + * name, returned as it is. */ +export function flowNodeName( + step?: string, + parentStep?: string, +): string | undefined { + const name = parentStep || step + if (!name) return undefined + const at = name.indexOf(MEMBER_SEPARATOR) + return at === -1 ? name : name.slice(0, at) +} + +/** The nodes a run has finished, given the step names its `workflow_start` + * listed. + * + * A for_each group is one node, so it is done only once every member is - + * greening it on the first would say the group is behind us while its + * remaining entries are still queued. A sub-workflow is done when the + * composed step itself ends, not when one of the child's inner steps does, + * so only the top-level step_end events count: a child's carry the parent's + * name in `parent_step`, which is exactly the marker to leave out. */ +export function finishedNodes( + events: JobEvent[], + stepNames: string[], +): string[] { + const ended = new Set( + events + .filter((event) => event.event === 'step_end' && !event.parent_step) + .map((event) => event.step as string), + ) + const members = new Map() + for (const name of stepNames) { + const node = flowNodeName(name) + if (node) members.set(node, [...(members.get(node) ?? []), name]) + } + return [...members] + .filter(([, names]) => names.every((name) => ended.has(name))) + .map(([node]) => node) +} From 0053814388c645ca680f345b59f928fed909bd68 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 21:51:29 -0500 Subject: [PATCH 08/16] chore(todos): remove outdated todos document --- todos.md | 108 ------------------------------------------------------- 1 file changed, 108 deletions(-) delete mode 100644 todos.md diff --git a/todos.md b/todos.md deleted file mode 100644 index ae81cad9..00000000 --- a/todos.md +++ /dev/null @@ -1,108 +0,0 @@ -# todos - -## chaining - -- Keyframe pre-planning: since fl2va takes first and last keyframes, generate a storyboard of keyframes first, then fill each segment between consecutive pairs. Segments become independent — no cumulative drift, and parallelizable. This is arguably a better long-video strategy than sequential chaining -- Anti-drift correction: histogram/color matching each segment's frames back to segment 0 — cheap, addresses the best-known failure mode of autoregressive video chaining -- Per-segment prompt scheduling: "prompt": ["intro shot...", "then the camera...", ...] — one prompt per segment, for narrative arcs (falls out of the chain loop almost for free) -- Chained prompt-embed reuse — LTX2I2VChained only. chain.py:255 re-runs the full prompt through the 14 GB Gemma once per segment; at 3 segments that's two redundant encodes per run. Engine change, not config. - -## performance - -- Save a pre-quantized checkpoint. The 45 s SDNQ pass re-quantizes identical weights on every cold start. Save once locally, point model_name at it, and cold starts drop to plain weight loading. Also speeds the REPL's first load. This is the one real remaining structural win for non-REPL use. -- save compiled checkpoint like above -- torch.compile with repeated_blocks. Attacks the 25 s denoise across 48 repeated blocks. It only became viable when the transformer went resident — compile and group-offload hooks fight each other, and that's gone now. But first-run compilation costs more than it saves, so it only pays off paired with #1, where the graph survives between runs. - -## ltx 2.5 round-out - -Priority ranked. Assessment: we shipped t2v + i2v + chaining out of ~11 native -capabilities. Most of the gap was two generic engine limits, not per-model work. - -1. [x] **Cache blocks for LTX-2.5.** `LTX2VideoTransformerBlock` registered in - [cache_blocks.json](dw/cache_blocks.json). Its second returned stream is audio, not - `encoder_hidden_states`, so the registry grew an `encoder_hidden_states_argument_name` - remap - without it a skipped block feeds the text embeddings back as the audio - stream (pinned by a test that reproduces exactly that). Of little use on the - distilled model's 8-step schedule; it is there for anything longer. -2. [x] **Frames and audio across a step boundary.** `Result.get_artifact_properties` - reads object artifacts by attribute, so `previous_result:step.frames` works on the - AudioVideo a task or a chain produces. Two tasks carry the pieces: `video_frames` - (frames as one 0-255 array, the shape conditions want) and `pair_audio` (puts a - soundtrack back beside frames a frames-only step returned). Ships LTX2TwoStage.json. -3. [x] **Argument objects built from named fields.** `from_arguments` in - [arguments.py](dw/arguments.py) constructs a type from the arguments it names, for - types with no `from_file()` and no media `kind` - LTX-2's conditions and IC-LoRA - references. Defers to `build_objects` when one of those arguments names a step. - Ships LTX2Keyframes.json, LTX2Extend.json, LTX2ICLora.json. -4. [-] **Diffusion decoder.** Built, run, dropped. `LTX2VideoDiffusionDecodePipeline` - works, but two things make it a bad deal on 24GB: its first three stages run on the - full volume by design (only stage 4 and the diffusion blocks tile) and the attention - mask they build is quadratic in the output grid - 70GiB at 1536x896x121, unaffected - by tile size - and a step run with `output_type: "{latent}"` also returns audio - latents that nothing outside a pipeline call can vocode, so the flow is silent. - Base resolution would fit, but that is a silent clip at the same size as LTX2.json. - Documented in RECIPES_24GB. The per-component `enable_tiling` this exposed is kept: - the engine could only tile a component literally named 'vae' before. -5. [x] **Prompt enhancement.** LTX2I2VEnhancePrompt.json declares `google/gemma-4-E2B-it` - as the `prompt_enhancer` plus its `processor`, image-conditioned off the reference - frame. Quantized `uint4` because the pipeline moves the enhancer onto the accelerator - and never moves it back. -6. [-] **Non-distilled example.** Dropped: `transformer_full` is ~38GB in bf16 and its - guidance knobs cost three transformer passes per step, so it is not a 24GB - configuration. The knobs are documented in RECIPES_24GB; no example ships them. -7. [x] **Auto duration.** LTX2I2VEnhancePrompt.json omits `num_frames` and gives the - duration head `min_seconds` / `max_seconds` instead. -8. [x] **Docs.** LTX-2.5 section in RECIPES_24GB, `from_arguments` and the - frames/audio hand-off in WORKFLOW_GUIDE, both tasks in TASKS, the dual-stream cache - block in ACCELERATION. - -### left over - -- **Stage-2 refine.** The full LTX flow re-denoises the upsampled latents at - `STAGE_2_DISTILLED_SIGMA_VALUES` with `noise_scale: 0.909375` (the standard pipeline - defaults it to 0.0, so it must be passed). Blocked on the batch dimension: a step run - with `output_type: "{latent}"` pairs its latents with the audio latents per batch item, - which drops the leading axis every pipeline's `latents` argument expects. The - `previous_result:step.frames` route keeps it, but then the audio latents have no - decoder of their own - `audio_vae` + `vocoder` are only reachable inside a pipeline - call. Needs either a standalone audio decode step or batch-preserving latent artifacts. -- **HDR** (`LTX2HDRPipeline`). Still deferred - needs pre-computed connector embeddings - from a safetensors file we have no path to produce. -- [x] **Released pipelines did not give their VRAM back.** Found and fixed. - `populate_from_pretrained_arguments` loaded each declared sub-component *into the - step's own definition dict* (`from_pretrained_arguments[name] = component`), and the - definition belongs to the workflow, which outlives every step - so `release_pipeline` - freed the wrapper while the weights stayed reachable. LTX2ICLora.json OOMed with - 22.6GiB allocated: two 22B transformers alive at once. It only showed up on workflows - that declare sub-components, which is why an sd15 pipeline (model_name only) released - cleanly and hid it. Both mutation sites now copy; a second load also keeps its - 'model_name', which load_component used to consume out of the definition. - Confirmed end to end: LTX2TwoStage.json now drops from 13.3GB to 2.8GB at its - release boundary, where it used to hold 12.2GB. LTX2ICLora.json keeps sharing its - transformer between the two steps - that is now a preference (it skips a second - 38GB load) rather than the workaround it was. - -- **Gated repo.** `Lightricks/LTX-2.5-22b-IC-LoRA-Pixel-Spatial-Upscaler` is behind a - license click-through, accepted on this machine as of 2026-08-18. - `google/gemma-4-E2B-it` needs nothing. - -## introspection - -- [x] **Task argument discovery.** Resolved 2026-08-30 with a third option - neither of the two on the table: the handlers were already thin shims - forwarding `**arguments` into implementation functions with real - signatures, so `@register_command` now records each implementation's - dotted path and the introspection layer reads that function's signature - - the same function the dispatch calls, so the schema cannot drift (and a - registry-integrity test fails if a path rots). `describe_task` serves - `/api/tasks/{command}` in describe_class's shape; the editor's task forms - consume it, and `/api/validate` flags task-argument typos like pipeline - ones. `provided=` hides dispatch-supplied parameters; `device` is always - offered. gather_inputs and the image processors stay declared free-form. - -## deferred - -- **Multi-GPU workers.** Deferred 2026-08-30: no multi-GPU hardware to test - against at the moment. The seams are already in place when it returns - - device overrides accept `cuda:N` throughout, and the server's JobManager - is the natural place to grow a worker pool. From 45fb90882f2946cd7bd28cd7528e4ee229e2afb1 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sat, 12 Sep 2026 23:41:17 -0500 Subject: [PATCH 09/16] fix(ui): keep a flow node's labels inside its box SVG text neither wraps nor takes text-overflow, so a sub-workflow step's path (`templates/minimax/composable-reference-shot.json`) or a long step name ran out of the node's right edge. Each label is now cut to a character budget for its font with an ellipsis where it was cut - a name from the end, a path from the start, since a path is told apart by how it ends - and the node carries a tooltip with whatever did not fit. A clipPath on the node catches what a wider glyph set still pushes past. Co-Authored-By: Claude Opus 5 (1M context) --- ui/src/lib/editor/FlowView.svelte | 53 ++++++++++++++++++++++++++++-- ui/src/lib/editor/FlowView.test.ts | 52 +++++++++++++++++++++++++++++ 2 files changed, 103 insertions(+), 2 deletions(-) diff --git a/ui/src/lib/editor/FlowView.svelte b/ui/src/lib/editor/FlowView.svelte index 7299f932..0c7d18cb 100644 --- a/ui/src/lib/editor/FlowView.svelte +++ b/ui/src/lib/editor/FlowView.svelte @@ -23,6 +23,40 @@ const graph = $derived(dataFlowGraph(workflow)) + // SVG text neither wraps nor takes text-overflow, so a label longer than + // the box ran out of its right edge. Budgets are characters at the box's + // inner width (BOX_W less the 10px inset each side) for each line's font + // - bold 12px sans for the name, 10px mono for the detail - and the + // clipPath below catches what a wider glyph set still pushes past. + const NAME_CHARS = 20 + const NAME_CHARS_WITH_TAG = 15 // the entry tag sits in the top-right corner + const DETAIL_CHARS = 26 + + /** The text cut to `max` characters with an ellipsis where it was cut: a + * name is told apart by how it starts, a path by how it ends. */ + function fit(text: string, max: number, keep: 'head' | 'tail'): string { + if (text.length <= max) return text + return keep === 'head' + ? text.slice(0, max - 1) + '…' + : '…' + text.slice(text.length - max + 1) + } + /** The parts of a node's labels that did not fit, in full, for its tooltip. */ + function overflowTitle(node: FlowNode): string { + const nameShown = fit( + node.name, + node.isEntryPoint ? NAME_CHARS_WITH_TAG : NAME_CHARS, + 'head', + ) + return [ + nameShown === node.name ? null : node.name, + fit(node.detail, DETAIL_CHARS, 'tail') === node.detail + ? null + : node.detail, + ] + .filter(Boolean) + .join('\n') + } + // Layered left-to-right layout: a node's layer is one past the deepest // producer that feeds it directly, so entry points (no previous_result // input) sit in the first column and depth reads as real dependency @@ -157,6 +191,9 @@ > + + + {#each layout.edgeLines ?? [] as edge, i (i)} @@ -180,14 +217,26 @@ class:active={node.name === activeStep} class:done={stateOf(node.name) === 'done'} transform={`translate(${pos.x}, ${pos.y})`} + clip-path="url(#flow-nodebox)" aria-label={`step ${node.name}, ${kindLabel(node.kind)}${stateOf(node.name) ? ', ' + stateOf(node.name) : ''}${node.isEntryPoint ? ', entry point' : ''}${fanIn ? ', fan-in: ' + fanIn.label : ''}`} {...nodeAttributes(node.name)} > + {#if overflowTitle(node)} + {overflowTitle(node)} + {/if} - {node.name} + {fit( + node.name, + node.isEntryPoint ? NAME_CHARS_WITH_TAG : NAME_CHARS, + 'head', + )} {kindLabel(node.kind)} {#if node.detail} - {node.detail} + {fit(node.detail, DETAIL_CHARS, 'tail')} {/if} {#if node.isEntryPoint}