From 9b00eae8ac342aa873e345e144682bf56bc4e967 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sun, 13 Sep 2026 22:01:18 -0500 Subject: [PATCH 01/62] fix: #141 #142 - a task step's signature, and a soundtrack that follows the cut #141: a task step that left a *required* argument unset validated as `valid: true` and then failed the job on Python's own "resample_audio() missing 1 required positional argument: 'audio'" - the one class of mistake a free pre-flight most obviously exists for. An unknown argument sat beside it as a warning, so a step with every argument it was given rejected and every argument it needs missing still validated. `task_signature_errors` (dw/introspection.py) reports both, at the JSON path each sits at, from the same introspection `get_task` already answers `required: true` from; the command itself refuses the same thing at run time, in the validator's wording rather than Python's, for a value that arrived from a variable or an earlier step. The unknown-argument warning is dropped - reported once, by the pass whose verdict it changes. No catalog workflow trips either check. #142: music-video sliced a hardcoded 496-frame soundtrack (4 shots x 124) while everything else in it followed the `shots` list, so a two-shot run laid 20.7 s of song over 10.3 s of picture and reported `succeeded` with `warnings: []`. `pair_audio` takes `fit: "video"` now: the track is cut to the frames it is laid over, or padded with silence and warned about when it is shorter. The template hands it the whole song and the slice step is gone, so the soundtrack follows the list whatever length it is; without `fit`, a length that disagrees with the frames' is warned about rather than passing in silence. Implementer agent, model opus via provider anthropic. --- docs/WORKFLOW_GUIDE.md | 10 ++ dw/introspection.py | 124 +++++++++++++++- dw/tasks/pair_audio.py | 106 +++++++++++++- dw/tasks/task.py | 25 ++++ dw/workflow.py | 6 + tests/test_for_each.py | 10 +- tests/test_pair_audio_fit.py | 145 +++++++++++++++++++ tests/test_server.py | 7 +- tests/test_task_discovery.py | 17 ++- tests/test_task_signature_errors.py | 141 ++++++++++++++++++ tests/test_workflow.py | 8 +- workflows/templates/minimax/README.md | 2 +- workflows/templates/minimax/music-video.json | 27 ++-- 13 files changed, 590 insertions(+), 38 deletions(-) create mode 100644 tests/test_pair_audio_fit.py create mode 100644 tests/test_task_signature_errors.py diff --git a/docs/WORKFLOW_GUIDE.md b/docs/WORKFLOW_GUIDE.md index 0a11e852..d7859b9f 100644 --- a/docs/WORKFLOW_GUIDE.md +++ b/docs/WORKFLOW_GUIDE.md @@ -1677,6 +1677,16 @@ rate of their own. A mono track needs no preparation: an mp4 audio stream takes stereo and nothing else, so saving duplicates the single channel into two and emits a warning saying it did. +The track and the frames are two lengths a workflow used to have to keep equal by +hand. `"fit": "video"` derives one from the other instead: the track is cut to +exactly the frames it is laid over, or padded with silence and warned about when it +is shorter than they are. That is what a soundtrack over a cut whose length is an +argument needs - nothing in a workflow can multiply a list's length by a frame +count, so `music-video.json` sliced a fixed 496 frames of song while its cut +followed a `shots` list, and a two-shot run wrote 10.3 s of picture into a 20.7 s +container and reported `succeeded` with no warnings (#142). Left unset the track is +used as it is and a disagreement is warned about rather than passing in silence. + Which shape a pipeline argument wants is the pipeline's business, and the two LTX-2 paths differ: a keyframe condition is mapped from 0-255, so it takes the `video_frames` array, while an IC-LoRA reference goes through the video processor, which expects the diff --git a/dw/introspection.py b/dw/introspection.py index 93174a7f..1d3001b5 100644 --- a/dw/introspection.py +++ b/dw/introspection.py @@ -492,6 +492,119 @@ def unknown_task_arguments(command, argument_names): return sorted(set(argument_names) - known) +def unknown_task_argument_message(command, name): + """The wording for one argument a task command does not take.""" + return ( + f"task '{command}' does not accept argument '{name}' - its " + f"implementation's signature is the whole of what it takes, so the " + f"argument would reach Python as an unexpected keyword" + ) + + +def missing_task_arguments(command, argument_names): + """The arguments a task command requires that the given names do not supply. + + A task step's `arguments` dict is the whole of what reaches the + implementation - nothing is injected around it, so a required parameter + absent from the dict is a run that cannot start. Unlike + unknown_task_arguments this does not stop at `accepts_kwargs`: **kwargs + says more names are allowed, never that a required one may be left out + (an image processor takes any keys and still needs its `image`). + + Same never-wrong contract otherwise: empty for a command that cannot be + described at all, and 'device' is never required. + """ + try: + description = describe_task(command) + except Exception: + return [] + supplied = set(argument_names) + return sorted( + p["name"] + for p in description["parameters"] + if p.get("required") and p["name"] != "device" and p["name"] not in supplied + ) + + +def missing_task_argument_message(command, missing): + """The one wording both the static pass and the run-time guard use for a + task step that leaves a required argument unset.""" + named = ", ".join(f"'{name}'" for name in missing) + return ( + f"task '{command}' requires {named}, which the step does not supply. " + f"A task's 'arguments' are the whole of what reaches the command, so " + f"a required argument left out is a run that cannot start" + ) + + +def task_signature_errors(workflow_definition, source_indices=None): + """Every task step whose arguments its command's signature refuses, as + [{path, message}] - a required argument left unset, and an argument the + command does not take. + + The one class of mistake a free pre-flight is most obviously for, and the + one it used to let through: `validate_workflow` answered `valid: true` + and the job then failed with Python's own + "resample_audio() missing 1 required positional argument: 'audio'" + (#141). An unknown argument was a warning beside it, so a step with every + argument it was given rejected and every argument it needs missing still + validated - both are a guaranteed TypeError at the same call, so both are + errors now. + + The definition handed here has already been substituted and expanded, so + a for_each member is checked as it will run; `source_indices` maps each + expanded step back to the step the author wrote. + """ + from .for_each import MEMBER_SEPARATOR, render_path + + steps = workflow_definition.get("steps") + if not isinstance(steps, list): + return [] + + errors = [] + for index, step in enumerate(steps): + if not isinstance(step, dict): + continue + task = step.get("task") + if not isinstance(task, dict): + continue + command = task.get("command") + # 'inputs' is a list template rather than a named-argument dict - + # the command consumes it whole, so there is no name to miss + arguments = task.get("arguments") + if not isinstance(command, str) or not isinstance(arguments, dict): + continue + missing = missing_task_arguments(command, arguments.keys()) + unknown = unknown_task_arguments(command, arguments.keys()) + if not missing and not unknown: + continue + source = ( + source_indices[index] + if source_indices is not None and index < len(source_indices) + else index + ) + name = step.get("name") + where = ( + f" in member '{name}'" + if isinstance(name, str) and MEMBER_SEPARATOR in name + else "" + ) + + def report(key, message): + errors.append( + { + "path": render_path(("steps", source, "task", "arguments", key)), + "message": f"{message}{where}.", + } + ) + + if missing: + report(missing[0], missing_task_argument_message(command, missing)) + for key in unknown: + report(key, unknown_task_argument_message(command, key)) + return errors + + def _inert_crossfade_warnings(step, command, arguments): """concat_videos draws its crossfade from the trimmed-off material, so with nothing trimmed a `crossfade_ms` the author wrote does nothing. A @@ -535,14 +648,9 @@ def workflow_argument_warnings(workflow_definition): task = step.get("task") if task and isinstance(task.get("arguments"), dict): command = task.get("command") - if isinstance(command, str): - for argument_name in unknown_task_arguments( - command, task["arguments"].keys() - ): - warnings.append( - f"Step '{step.get('name')}': task '{command}' does not " - f"accept argument '{argument_name}'" - ) + # An unknown or missing task argument is an error rather than a + # warning now (task_signature_errors, #141) - reported once, by + # the pass whose verdict it changes warnings.extend(_inert_crossfade_warnings(step, command, task["arguments"])) pipeline = step.get("pipeline") if not pipeline: diff --git a/dw/tasks/pair_audio.py b/dw/tasks/pair_audio.py index 0b41cb2a..b7f9ccd4 100644 --- a/dw/tasks/pair_audio.py +++ b/dw/tasks/pair_audio.py @@ -9,11 +9,17 @@ import logging +from ..events import emit_warning from ..result import AudioVideo from .audio_utils import as_channels_samples logger = logging.getLogger("dw") +# A track and a cut are frame-aligned by construction here, so a difference +# smaller than this is rounding rather than a decision anyone can act on - +# one video frame at 24 fps is 41 ms +LENGTH_WARN_MS = 100.0 + class _Loaded: """A waveform read from a file, shaped like the artifact pair_audio expects.""" @@ -23,7 +29,87 @@ def __init__(self, audio, sample_rate): self.sample_rate = sample_rate -def pair_audio(video, audio, sample_rate=None): +def _frame_count(frames): + """How many frames the video is, or None when that cannot be told cheaply + (a lazily-decoded reader, an object with no length).""" + try: + return len(frames) + except TypeError: + shape = getattr(frames, "shape", None) + return int(shape[0]) if shape else None + + +def _fit_to_video(waveform, rate, frames, fps, fit): + """Answer the track that goes with these frames, and say when the two do + not agree. + + The lengths of a video and the track laid over it are two numbers a + workflow used to have to keep equal by hand, and nothing checked: the + music-video template sliced a soundtrack of a fixed 496 frames while its + cut followed a `shots` list, so a two-shot run wrote 10.3 s of picture + into a 20.7 s container and reported `succeeded` with no warnings (#142). + `fit: "video"` derives the length from the frames instead, and with no + `fit` the mismatch is at least said out loud. + """ + from .audio_utils import frames_to_samples, slice_samples + + if fit not in (None, "video"): + # Refused rather than ignored: a misspelled 'fit' that quietly did + # nothing is the silence this argument exists to end + raise ValueError( + f"pair_audio: 'fit' takes 'video' or nothing, got {fit!r}. " + f"'video' cuts or pads the track to the length of the frames" + ) + + count = _frame_count(frames) + if not count or not fps or not rate: + # Nothing to compare against - a frame count or a rate this layer + # cannot know is not a mismatch + return waveform + + wanted = frames_to_samples(count, fps, rate) + have = waveform.shape[1] + if abs(have - wanted) / float(rate) * 1000.0 < LENGTH_WARN_MS: + return waveform + + video_seconds = count / float(fps) + audio_seconds = have / float(rate) + if fit != "video": + emit_warning( + f"pair_audio: the track is {audio_seconds:.2f} s and the video it " + f"is laid over is {video_seconds:.2f} s ({count} frames at " + f"{fps:g} fps), so the saved file's duration and its frame count " + f"disagree. Pass 'fit': 'video' to cut or pad the track to the " + f"frames, or make the track the length of the cut.", + kind="audio_video_length_mismatch", + command="pair_audio", + audio_seconds=audio_seconds, + video_seconds=video_seconds, + ) + return waveform + + fitted = slice_samples(waveform, 0, wanted) + if wanted > have: + emit_warning( + f"pair_audio: 'fit' padded the {audio_seconds:.2f} s track with " + f"{(wanted - have) / float(rate):.2f} s of silence to reach the " + f"{video_seconds:.2f} s of video it is laid over - the last part " + f"of the cut has no soundtrack. A longer track, or fewer frames, " + f"is what covers it.", + kind="audio_padded_to_video", + command="pair_audio", + audio_seconds=audio_seconds, + video_seconds=video_seconds, + ) + else: + logger.info( + f"pair_audio: trimmed the track from {audio_seconds:.2f} s to the " + f"{video_seconds:.2f} s of video it is laid over" + ) + return fitted + + +def pair_audio(video, audio, sample_rate=None, fps=None, fit=None): """Pair a video's frames with an audio track. Args: @@ -40,6 +126,18 @@ def pair_audio(video, audio, sample_rate=None): sample_rate: Sample rate of the waveform. Required unless `audio` carries one; given here it wins, for a track whose rate was reported wrong + fps: The rate the frames play at, when the frames do not carry one - + only used to work out how long the video is, never written + fit: "video" cuts the track to the length of the frames, or pads it + with silence and warns when it is shorter than they are. This is + how a soundtrack follows a cut whose length is an argument + rather than a constant: nothing in a workflow can multiply a + list's length by a frame count, so a slice written to fit four + shots stayed 496 frames long when the list held two, and the + deliverable's audio ran twice as long as its picture with + `succeeded` and no warnings (#142). Left unset the track is used + as it is, and a length that disagrees with the frames' is + warned about rather than passing in silence Returns: One AudioVideo holding the frames and the track, at the rate the @@ -75,7 +173,9 @@ def pair_audio(video, audio, sample_rate=None): # list, an array and a tensor alike, and converting a long video here would # cost a copy of the whole thing for nothing frames = video.frames if isinstance(video, AudioVideo) else video + frame_rate = fps if fps is not None else getattr(video, "fps", None) logger.debug(f"Pairing frames with audio at {rate} Hz") - return AudioVideo( - frames, as_channels_samples(waveform), rate, fps=getattr(video, "fps", None) + waveform = _fit_to_video( + as_channels_samples(waveform), rate, frames, frame_rate, fit ) + return AudioVideo(frames, waveform, rate, fps=getattr(video, "fps", None)) diff --git a/dw/tasks/task.py b/dw/tasks/task.py index dcdb2a92..5e58b036 100644 --- a/dw/tasks/task.py +++ b/dw/tasks/task.py @@ -520,6 +520,22 @@ def command(self): """Get command name or 'unknown' if not specified""" return self.task_definition.get("command", "unknown") + def _check_required_arguments(self, arguments): + """Refuse a task whose required arguments are not all present, in the + validator's wording rather than Python's.""" + if not isinstance(arguments, dict): + # An 'inputs' list template - consumed whole, no names to miss + return + from ..introspection import ( + missing_task_argument_message, + missing_task_arguments, + ) + + missing = missing_task_arguments(self.command, arguments.keys()) + if missing: + message = missing_task_argument_message(self.command, missing) + raise ValueError(message[0].upper() + message[1:]) + def run(self, arguments, previous_pipelines={}): """ Execute the task with given arguments using the command registry. @@ -548,6 +564,15 @@ def run(self, arguments, previous_pipelines={}): # and decoding is otherwise indistinguishable from a hang emit_phase("task", detail=self.command) + # A required argument that never arrived - because it was left + # out, or because a variable or an earlier step resolved to + # nothing - used to reach Python and come back as + # "resample_audio() missing 1 required positional argument: + # 'audio'", which names the calling convention rather than the + # workflow. The static pass in validation_errors refuses the + # literal case first; this is the backstop it cannot see (#141) + self._check_required_arguments(arguments) + # Look up command in registry if self.command in _COMMAND_REGISTRY: handler = _COMMAND_REGISTRY[self.command] diff --git a/dw/workflow.py b/dw/workflow.py index 81e97d64..b25f8d05 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -29,6 +29,7 @@ ) from .locations import location_errors from .reference_limits import reference_limit_errors +from .introspection import task_signature_errors from .task_domains import task_argument_errors from .subfolders import step_subfolder, subfolder_errors from .step import Step @@ -604,6 +605,11 @@ def validation_errors(self, arguments=None, composing=None): # sample rate a silent fallback to 44100 (dw/task_domains.py, # #139, #140) + task_argument_errors(expanded, source_indices) + # A required task argument left unset validated as `valid: true` + # and then failed the job on Python's own signature error, which + # is the one mistake a free pre-flight most obviously exists for + # (dw/introspection.py, #141) + + task_signature_errors(expanded, source_indices) + self.sub_workflow_errors(expanded, source_indices, composing) ) diff --git a/tests/test_for_each.py b/tests/test_for_each.py index e44ef5b7..1c33e73a 100644 --- a/tests/test_for_each.py +++ b/tests/test_for_each.py @@ -632,11 +632,19 @@ def test_one_slice_and_one_shot_per_entry_in_list_order(self): assert names == ( ["draw_singer", "write_song"] + [f"slice@{k}" for k in self.KEYS] - + ["soundtrack"] + [f"shot@{k}" for k in self.KEYS] + ["edit", "music_video"] ) + def test_the_soundtrack_is_cut_to_the_edit_rather_than_to_a_constant(self): + """There is no 'soundtrack' step: a slice of a fixed 496 frames was + right only for a four-entry list, so a two-shot run laid 20.7 s of + song over 10.3 s of picture (#142). pair_audio derives it now.""" + got = steps_by_name(self.expanded()) + arguments = got["music_video"]["task"]["arguments"] + assert arguments["audio"] == "previous_result:write_song" + assert arguments["fit"] == "video" + def test_each_slice_starts_where_its_entry_says(self): got = steps_by_name(self.expanded()) starts = [ diff --git a/tests/test_pair_audio_fit.py b/tests/test_pair_audio_fit.py new file mode 100644 index 00000000..78154980 --- /dev/null +++ b/tests/test_pair_audio_fit.py @@ -0,0 +1,145 @@ +"""A soundtrack whose length follows the cut it is laid over. + +`music-video.json` sliced a hardcoded 496 frames of song (4 shots x 124) while +everything else in it followed the `shots` list, so a two-shot run produced a +deliverable holding 248 frames of picture in a 20.7 s container - audio twice +as long as video, `status: succeeded`, `warnings: []` (#142). `fit: "video"` +derives the length from the frames instead, and without it a disagreement is +at least said out loud. +""" + +import json +import pathlib + +import numpy +import pytest + +from dw.result import AudioVideo +from dw.tasks.pair_audio import pair_audio + +SAMPLE_RATE = 44100 +FPS = 24 + + +def song(seconds, sample_rate=SAMPLE_RATE): + return numpy.zeros((2, int(seconds * sample_rate)), dtype=numpy.float32) + + +def cut(frames, fps=FPS): + return AudioVideo([object()] * frames, None, None, fps=fps) + + +def samples(result): + return result.audio.shape[1] + + +@pytest.fixture +def warnings_emitted(): + """The messages of the warning events a call emits - what a consumer over + the API or MCP actually reads, rather than the server's log (#82).""" + from dw.events import RunContext, activate_context, deactivate_context + + messages = [] + + def record(event): + if event["event"] == "warning": + messages.append(event["message"]) + + token = activate_context(RunContext(on_event=record)) + try: + yield messages + finally: + deactivate_context(token) + + +class TestFitToTheVideo: + def test_the_reported_case_is_cut_to_the_two_shot_edit(self): + """248 frames at 24 fps is 10.33 s, from a 30 s song.""" + result = pair_audio(cut(248), song(30), sample_rate=SAMPLE_RATE, fit="video") + assert samples(result) == round(248 / FPS * SAMPLE_RATE) + + def test_the_default_four_shot_length_is_unchanged(self): + """496 frames - what the hardcoded slice used to produce.""" + result = pair_audio(cut(496), song(30), sample_rate=SAMPLE_RATE, fit="video") + assert samples(result) == round(496 / FPS * SAMPLE_RATE) + + def test_the_other_direction_pads_and_warns(self, warnings_emitted): + """Six shots outrun a 30 s song - the silent-padding direction the + report could not afford to run.""" + result = pair_audio(cut(744), song(30), sample_rate=SAMPLE_RATE, fit="video") + assert samples(result) == round(744 / FPS * SAMPLE_RATE) + assert any("padded" in w for w in warnings_emitted) + + def test_a_track_that_already_matches_is_left_alone(self, warnings_emitted): + exact = song(248 / FPS) + result = pair_audio(cut(248), exact, sample_rate=SAMPLE_RATE, fit="video") + assert samples(result) == exact.shape[1] + assert warnings_emitted == [] + + def test_fps_may_be_given_when_the_frames_carry_none(self): + result = pair_audio( + [object()] * 248, song(30), sample_rate=SAMPLE_RATE, fps=FPS, fit="video" + ) + assert samples(result) == round(248 / FPS * SAMPLE_RATE) + + def test_an_unknown_fit_is_refused_rather_than_ignored(self): + with pytest.raises(ValueError, match="'fit'"): + pair_audio(cut(248), song(30), sample_rate=SAMPLE_RATE, fit="audio") + + +class TestWithoutFit: + def test_a_mismatch_is_warned_about(self, warnings_emitted): + result = pair_audio(cut(248), song(30), sample_rate=SAMPLE_RATE) + assert samples(result) == song(30).shape[1], "the track is left as it is" + assert any("disagree" in w for w in warnings_emitted) + + def test_an_agreeing_pair_says_nothing(self, warnings_emitted): + pair_audio(cut(248), song(248 / FPS), sample_rate=SAMPLE_RATE) + assert warnings_emitted == [] + + def test_unknown_fps_is_not_a_mismatch(self, warnings_emitted): + """A frame rate this layer cannot know is not something to warn about.""" + pair_audio(cut(248, fps=None), song(30), sample_rate=SAMPLE_RATE) + assert warnings_emitted == [] + + +class TestTheTemplateItself: + def test_music_video_derives_its_soundtrack(self): + definition = json.loads( + pathlib.Path("workflows/templates/minimax/music-video.json").read_text() + ) + steps = {s["name"]: s for s in definition["steps"]} + assert "soundtrack" not in steps, "the hardcoded 496-frame slice is gone" + arguments = steps["music_video"]["task"]["arguments"] + assert arguments["audio"] == "previous_result:write_song" + assert arguments["fit"] == "video" + + def test_no_template_hardcodes_a_soundtrack_length(self): + """Every `slice_audio` in the catalog whose count is a literal is one + a list cannot resize out from under - so the literal must not be the + length of a whole cut.""" + for path in pathlib.Path("workflows").rglob("*.json"): + definition = json.loads(path.read_text()) + if not isinstance(definition, dict): + continue + for step in definition.get("steps", []): + task = step.get("task") or {} + if task.get("command") != "pair_audio": + continue + audio = task.get("arguments", {}).get("audio", "") + if not isinstance(audio, str) or not audio.startswith( + "previous_result:" + ): + continue + source = {s["name"]: s for s in definition["steps"]}.get( + audio.split(":", 1)[1] + ) + source_task = (source or {}).get("task") or {} + if source_task.get("command") != "slice_audio": + continue + count = source_task.get("arguments", {}).get("num_frames") + assert not isinstance(count, int), ( + f"{path}: '{source['name']}' cuts a literal {count}-frame " + f"soundtrack for a cut whose length may be an argument - " + f"use pair_audio's 'fit': 'video' instead (#142)" + ) diff --git a/tests/test_server.py b/tests/test_server.py index 5403ee91..fa906c15 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -2579,8 +2579,11 @@ def test_task_typos_surface_in_validation(self, server): } with server(success_script) as client: result = client.post("/api/validate", json={"workflow": workflow}).json() - assert result["valid"] - assert any("trim_framse" in warning for warning in result["warnings"]) + # An argument the command does not take reaches Python as an + # unexpected keyword, so it refuses the workflow rather than + # warning beside a `valid: true` (#141) + assert not result["valid"] + assert any("trim_framse" in error["message"] for error in result["errors"]) class TestDiffusersUpdate: diff --git a/tests/test_task_discovery.py b/tests/test_task_discovery.py index 42bab516..37c62e9d 100644 --- a/tests/test_task_discovery.py +++ b/tests/test_task_discovery.py @@ -95,7 +95,13 @@ def test_kwargs_and_unknown_commands_stay_silent(self): assert unknown_task_arguments("upscale", ["anything_at_all"]) == [] assert unknown_task_arguments("not_a_task", ["x"]) == [] - def test_workflow_warnings_cover_task_steps(self): + def test_a_task_typo_is_an_error_rather_than_a_warning(self): + """It used to be a warning beside `valid: true`, so a step with every + argument it was given rejected still validated - and then failed the + job on Python's own TypeError. task_signature_errors owns it now + (#141), and it is reported once.""" + from dw.introspection import task_signature_errors + definition = { "steps": [ { @@ -107,10 +113,11 @@ def test_workflow_warnings_cover_task_steps(self): } ] } - warnings = workflow_argument_warnings(definition) - assert len(warnings) == 1 - assert "trim_framse" in warnings[0] - assert "concat_videos" in warnings[0] + assert workflow_argument_warnings(definition) == [] + errors = task_signature_errors(definition) + assert len(errors) == 1 + assert errors[0]["path"] == "steps[0].task.arguments.trim_framse" + assert "concat_videos" in errors[0]["message"] def test_inputs_style_task_steps_are_left_alone(self): definition = { diff --git a/tests/test_task_signature_errors.py b/tests/test_task_signature_errors.py new file mode 100644 index 00000000..3f237608 --- /dev/null +++ b/tests/test_task_signature_errors.py @@ -0,0 +1,141 @@ +"""A task step whose arguments its command's signature refuses. + +`validate_workflow` used to answer `valid: true` for a task step that left a +*required* argument unset, and then the job failed on Python's own +"resample_audio() missing 1 required positional argument: 'audio'" - the one +class of mistake a free pre-flight most obviously exists for (#141). An +unknown argument sat beside it as a warning, so a step with every argument it +was given rejected and every argument it needs missing still validated. Both +are the same guaranteed TypeError at the same call, so both are errors, and +the command itself refuses the same thing at run time for a value that +arrived from a variable or an earlier step. +""" + +import json +import pathlib + +import pytest + +from dw.introspection import ( + missing_task_arguments, + task_signature_errors, + unknown_task_arguments, + workflow_argument_warnings, +) +from dw.tasks.task import Task + + +def task_step(command, arguments, name="a"): + return { + "name": name, + "task": {"command": command, "arguments": arguments}, + "result": {"content_type": "audio/wav"}, + } + + +def errors_for(command, arguments): + return task_signature_errors({"id": "sig", "steps": [task_step(command, arguments)]}) + + +class TestARequiredArgumentLeftUnset: + def test_the_reported_repro_is_refused(self): + """#141's first repro: resample_audio without its audio.""" + errors = errors_for("resample_audio", {"target_sample_rate": 16000}) + assert len(errors) == 1 + assert errors[0]["path"] == "steps[0].task.arguments.audio" + assert "resample_audio" in errors[0]["message"] + assert "'audio'" in errors[0]["message"] + + def test_a_supplied_argument_passes(self): + assert errors_for("resample_audio", {"audio": "x", "target_sample_rate": 16000}) == [] + + def test_a_reference_counts_as_supplied(self): + """The key's presence is what the check is - the value may still be a + reference that only resolves during the run.""" + assert ( + errors_for( + "resample_audio", + {"audio": "previous_result:earlier", "target_sample_rate": 16000}, + ) + == [] + ) + + def test_device_is_never_required(self): + assert "device" not in missing_task_arguments("resample_audio", ["audio"]) + + def test_kwargs_does_not_excuse_a_required_argument(self): + """An image processor takes any keys and still needs its image.""" + assert missing_task_arguments("canny", []) == ["image"] + assert missing_task_arguments("canny", ["image"]) == [] + + def test_an_unknown_command_reports_nothing(self): + assert missing_task_arguments("not_a_command", []) == [] + assert errors_for("not_a_command", {}) == [] + + +class TestAnArgumentTheCommandDoesNotTake: + def test_the_near_miss_from_the_report_is_refused(self): + """#141's second repro: crossfade_audio's parameter is `audios`, a + list, so `audio`/`other` are rejected and `audios` is missing - every + argument wrong, and it used to validate.""" + errors = errors_for( + "crossfade_audio", {"audio": "x", "other": "y", "crossfade_ms": 0} + ) + paths = [e["path"] for e in errors] + assert paths == [ + "steps[0].task.arguments.audios", + "steps[0].task.arguments.audio", + "steps[0].task.arguments.other", + ] + + def test_it_is_no_longer_also_a_warning(self): + """Reported once, by the pass whose verdict it changes.""" + definition = { + "id": "sig", + "steps": [task_step("crossfade_audio", {"audios": ["x"], "nope": 1})], + } + assert task_signature_errors(definition) + assert not [w for w in workflow_argument_warnings(definition) if "nope" in w] + + def test_a_free_form_command_accepts_anything(self): + assert unknown_task_arguments("gather_inputs", ["whatever"]) == [] + assert errors_for("gather_inputs", {"whatever": 1}) == [] + + +class TestTheRunTimeBackstop: + """The static pass sees literals. A required argument that arrived from a + variable or an earlier step and resolved to nothing reaches the command, + which refuses it in the validator's wording rather than Python's.""" + + def test_the_command_refuses_and_says_which_argument(self): + task = Task({"command": "resample_audio", "arguments": {}}, "cpu") + with pytest.raises(ValueError) as caught: + task.run({"target_sample_rate": 16000}) + message = str(caught.value) + assert "resample_audio" in message + assert "'audio'" in message + assert "positional argument" not in message + + def test_an_inputs_list_template_is_left_alone(self): + """`inputs` is consumed whole - there are no names to miss.""" + task = Task({"command": "gather_inputs", "inputs": [1, 2]}, "cpu") + task._check_required_arguments([1, 2]) + + +class TestTheCatalogItself: + """Every workflow shipped in the repo passes the new check - a template + that did not would be a run nothing could start.""" + + @pytest.mark.parametrize( + "path", + sorted( + str(p) + for p in list(pathlib.Path("workflows").rglob("*.json")) + + list(pathlib.Path("dw/workflows").glob("*.json")) + ), + ) + def test_workflow_has_no_task_signature_error(self, path): + definition = json.loads(pathlib.Path(path).read_text()) + if not isinstance(definition, dict) or "steps" not in definition: + pytest.skip("not a workflow") + assert task_signature_errors(definition) == [] diff --git a/tests/test_workflow.py b/tests/test_workflow.py index 6ef16b0e..dc668748 100644 --- a/tests/test_workflow.py +++ b/tests/test_workflow.py @@ -254,7 +254,13 @@ def child_steps(self): return [ { "name": "noop", - "task": {"command": "get_image_size", "arguments": {}}, + # A supplied 'image' - a task step that leaves a required + # argument unset no longer validates (#141), and this stub + # is about seed inheritance rather than about the task + "task": { + "command": "get_image_size", + "arguments": {"image": "asset:nothing.png"}, + }, } ] diff --git a/workflows/templates/minimax/README.md b/workflows/templates/minimax/README.md index c06e04be..619e7d66 100644 --- a/workflows/templates/minimax/README.md +++ b/workflows/templates/minimax/README.md @@ -133,4 +133,4 @@ same clip as an audio reference in each shot, as | Example | What it introduces | | ------- | ------------------ | | [dialogue-short.json](dialogue-short.json) | A five-shot sitcom scene: Z-Image draws the cast, one `for_each` step over a `shots` list generates a shot per entry - its prompt, its references, its length - on one loaded model, and `concat_videos` gathers the episode | -| [music-video.json](music-video.json) | A music video cut to a generated song, one `shots` list driving both `for_each` groups: `slice_audio` deals each entry its frame-exact piece, the shot lip-syncs to it, and `pair_audio` lays the unbroken track over the finished edit | +| [music-video.json](music-video.json) | A music video cut to a generated song, one `shots` list driving both `for_each` groups: `slice_audio` deals each entry its frame-exact piece, the shot lip-syncs to it, and `pair_audio` lays the unbroken track over the finished edit, cut to the length of the edit by `fit: "video"` so the soundtrack follows the list rather than a constant | diff --git a/workflows/templates/minimax/music-video.json b/workflows/templates/minimax/music-video.json index 71c64838..ea81b8f1 100644 --- a/workflows/templates/minimax/music-video.json +++ b/workflows/templates/minimax/music-video.json @@ -1,8 +1,13 @@ { "id": "MiniMaxH3MusicVideo", - "description": "A music video built from cuts, sung to a soundtrack that never touches a chain. The long-take way to film a song - one chained generation lip-synced end to end - degrades with every carried segment and can let the sync slip. This builds the video the way music television does instead: MiniMax-Music3 writes the song, one 'slice_audio' step per entry of 'shots' cuts it into frame-exact pieces (124 frames at 24 fps each, from each entry's 'start_frame'), and one shot per entry is generated fresh from the same Z-Image portrait plus its own slice, lip-synced to just those five seconds. 'shots' is a list, so the cut is an argument: add an entry and there is one more slice and one more shot, named for it ('slice@closeup', 'shot@closeup'), and the loaded model is reused across every shot - identical pipeline definitions share one model. No shot conditions on another shot's output, so the last cut is as clean as the first. 'concat_videos' gathers the shots in list order, and because each one covered exactly its slice's frames, the edit is sample-accurate by construction: 'pair_audio' drops the original, unbroken song over the whole cut and the mouths line up in every shot. The generation models pass through one at a time - each is released before the next loads - so the workflow peaks no higher than its largest single model. Shots generated independently also drift in loudness - 10 dB between two shots of one scene is ordinary - and no seam control can hide a level jump, because it is either side of the cut rather than at it; 'match_levels' ('rms' for perceived level, 'peak' for the loudest sample) evens the shots out before they are joined, and left null, as it is by default, a wide spread is warned about in the log rather than passing in silence.", + "description": "A music video built from cuts, sung to a soundtrack that never touches a chain. The long-take way to film a song - one chained generation lip-synced end to end - degrades with every carried segment and can let the sync slip. This builds the video the way music television does instead: MiniMax-Music3 writes the song, one 'slice_audio' step per entry of 'shots' cuts it into frame-exact pieces (124 frames at 24 fps each, from each entry's 'start_frame'), and one shot per entry is generated fresh from the same Z-Image portrait plus its own slice, lip-synced to just those five seconds. 'shots' is a list, so the cut is an argument: add an entry and there is one more slice and one more shot, named for it ('slice@closeup', 'shot@closeup'), and the loaded model is reused across every shot - identical pipeline definitions share one model. No shot conditions on another shot's output, so the last cut is as clean as the first. 'concat_videos' gathers the shots in list order, and because each one covered exactly its slice's frames, the edit is sample-accurate by construction: 'pair_audio' drops the original, unbroken song over the whole cut and the mouths line up in every shot. The soundtrack follows the list too - 'pair_audio' is given the whole song with 'fit': 'video', so it is cut to exactly the frames the edit came to, whatever length 'shots' is; a list long enough to outrun 'audio_duration' seconds of song is padded with silence and warned about rather than passing quietly. The generation models pass through one at a time - each is released before the next loads - so the workflow peaks no higher than its largest single model. Shots generated independently also drift in loudness - 10 dB between two shots of one scene is ordinary - and no seam control can hide a level jump, because it is either side of the cut rather than at it; 'match_levels' ('rms' for perceived level, 'peak' for the loudest sample) evens the shots out before they are joined, and left null, as it is by default, a wide spread is warned about in the log rather than passing in silence.", "cost": [ - {"device": "cuda", "name": "RTX 3090", "vram_gb": 24, "minutes": 35} + { + "device": "cuda", + "name": "RTX 3090", + "vram_gb": 24, + "minutes": 35 + } ], "variables": { "singer_portrait_prompt": "prompt:zimage/otter_singer_portrait", @@ -120,19 +125,6 @@ } } }, - { - "name": "soundtrack", - "task": { - "command": "slice_audio", - "arguments": { - "audio": "previous_result:write_song", - "sample_rate": "variable:sample_rate", - "start_frame": 0, - "num_frames": 496, - "fps": "variable:fps" - } - } - }, { "name": "shot", "for_each": "variable:shots", @@ -293,8 +285,9 @@ "command": "pair_audio", "arguments": { "video": "previous_result:edit", - "audio": "previous_result:soundtrack", - "sample_rate": "variable:sample_rate" + "audio": "previous_result:write_song", + "sample_rate": "variable:sample_rate", + "fit": "video" } }, "result": { From df80db427555f1dc0f8268618c6ba28f736cda27 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sun, 13 Sep 2026 22:13:43 -0500 Subject: [PATCH 02/62] feat: #122 - a step nothing references does not run Approved proposal docs/proposals/unreferenced-step-elision.md. Before the first step executes, dw/elision.py drops any step whose result no later step reads and which writes no file. `dialogue-short` cast from portraits that already exist ran its two Z-Image steps anyway and threw the pictures away - about a minute of GPU and two model loads per episode (#109). Four guardrails, as approved: a step that saves is kept (a `result` with a content type and `save` not false - exactly what Result.save asks); the last step is kept; a release moves onto the last surviving step before it (`release_pipeline` only when that step loaded the same pipeline, since moving it elsewhere would unload something the workflow never asked to unload; `release_models` always); and every elision is an emit_warning, so a misspelled reference shows up as "draw_character_a did not run" rather than as a silently different picture. Elision is transitive and runs to a fixed point. `build_plan` elides too, so `steps`, `downloads_required`, the estimate and the fingerprint are the work that will actually happen, with the dropped steps listed under `elided_steps`; the run manifest records the same list; the editor's plan lines name them. With the engine change, per Don's condition: `dialogue-short`'s two draw steps carry `save: false`, so a `from_file` cast actually skips them - and `draw_character_b`'s `release_pipeline` is dropped rather than moved when `draw_character_a` goes too, because nothing was loaded to release. The default (uncast) run is unchanged, and a test sweeps every catalog workflow to keep it that way. Test fixtures that used an unreferenced non-saving step to test something else (the plan, the step cache, pipeline release) now declare a result or a reference - noted at each. Implementer agent, model opus via provider anthropic. --- docs/MCP.md | 2 +- docs/WORKFLOW_GUIDE.md | 45 +++ dw/elision.py | 215 ++++++++++++ dw/plan.py | 7 + dw/workflow.py | 21 ++ tests/test_elision.py | 305 ++++++++++++++++++ tests/test_pipeline_caching.py | 10 +- tests/test_plan.py | 6 + tests/test_server.py | 1 + tests/test_workflow_step_cache.py | 5 + ui/src/lib/plan.ts | 8 + ui/src/lib/types.ts | 3 + .../templates/minimax/dialogue-short.json | 8 +- 13 files changed, 631 insertions(+), 5 deletions(-) create mode 100644 dw/elision.py create mode 100644 tests/test_elision.py diff --git a/docs/MCP.md b/docs/MCP.md index 6fbf1d87..ea3985ab 100644 --- a/docs/MCP.md +++ b/docs/MCP.md @@ -249,7 +249,7 @@ The session starts in `default` and stays there unless it is told otherwise. | Tool | Arguments | Purpose | | --- | --- | --- | -| `validate_workflow(workflow=None, name=None, workspace=None, arguments=None)` | exactly one of `workflow` (inline definition) or `name` (a stored workflow, as `list_workflows` reports it), optional `workspace`, optional `arguments` | Check a workflow against the schema and against real pipeline signatures. Free and instant. Validating by name uses the workflow file's own directory as the base directory, so it sees what a run would. Returns every schema violation in `errors`, each with the JSON path it sits at, so a draft is fixed in one pass, and a `previous_result:` that names no earlier step is one of them. `warnings` covers what still runs but is probably wrong - a signature mismatch, and, for a list-driven variable, an entry key no step reads, at the entry's path. `workspace` names the workspace for this one call without switching the session to it - use it to pin a job whose `output:` or `asset:` references live in a workspace other than the session's. Pass the same `arguments` you will pass to `run_workflow` and they are checked too - an undeclared or renamed variable name, a value that will not coerce to the declared type, and an `asset:`, `prompt:` or `output:` reference that names nothing this workspace can reach, each reported at `arguments.`. `checked_arguments` lists what was covered, so a `valid: true` about the stored defaults cannot be mistaken for one about your values. A reference set the model would refuse - too many images, videos or audio clips, or, for MiniMax-H3, audio as the only reference - is an error here too, rather than a failure minutes into a run you acknowledged. `run_workflow` makes the same check and refuses a bad argument rather than queuing a job that fails on its first step. A valid answer carries `plan` - the fingerprint, step count, list lengths, `downloads_required` and `estimate` (with `basis`) for the arguments given; quote from it | +| `validate_workflow(workflow=None, name=None, workspace=None, arguments=None)` | exactly one of `workflow` (inline definition) or `name` (a stored workflow, as `list_workflows` reports it), optional `workspace`, optional `arguments` | Check a workflow against the schema and against real pipeline signatures. Free and instant. Validating by name uses the workflow file's own directory as the base directory, so it sees what a run would. Returns every schema violation in `errors`, each with the JSON path it sits at, so a draft is fixed in one pass, and a `previous_result:` that names no earlier step is one of them. `warnings` covers what still runs but is probably wrong - a signature mismatch, and, for a list-driven variable, an entry key no step reads, at the entry's path. `workspace` names the workspace for this one call without switching the session to it - use it to pin a job whose `output:` or `asset:` references live in a workspace other than the session's. Pass the same `arguments` you will pass to `run_workflow` and they are checked too - an undeclared or renamed variable name, a value that will not coerce to the declared type, and an `asset:`, `prompt:` or `output:` reference that names nothing this workspace can reach, each reported at `arguments.`. `checked_arguments` lists what was covered, so a `valid: true` about the stored defaults cannot be mistaken for one about your values. A reference set the model would refuse - too many images, videos or audio clips, or, for MiniMax-H3, audio as the only reference - is an error here too, rather than a failure minutes into a run you acknowledged. `run_workflow` makes the same check and refuses a bad argument rather than queuing a job that fails on its first step. A valid answer carries `plan` - the fingerprint, step count, list lengths, `downloads_required`, `estimate` (with `basis`) and `elided_steps` for the arguments given; quote from it. `steps` counts what will run: a step nothing reads and which saves no file does not run, and is named in `elided_steps` instead | | `list_workspaces()` | — | The server's workspaces and which one this session is using. Each has its own workflows, assets and outputs; the prompt library is shared by all of them | | `use_workspace(name)` | `name` | Work in that workspace for the rest of the session - every later call reads and writes there. This is how to keep your work out of another agent's namespace rather than sharing the default one. Checked against the server, so a typo fails here rather than scoping every later call to nothing | | `create_workspace(name, use=False)` | `name`, `use` | Create a workspace. Pass use=true to switch this session to it as well; otherwise the session stays where it was and the result says so | diff --git a/docs/WORKFLOW_GUIDE.md b/docs/WORKFLOW_GUIDE.md index d7859b9f..ec184b8b 100644 --- a/docs/WORKFLOW_GUIDE.md +++ b/docs/WORKFLOW_GUIDE.md @@ -616,6 +616,20 @@ than prefixing it, so `"file_base_name": "episode"` in a `final` subfolder writes `final/episode-0.0.mp4` - name each step that sets one differently, or the second collides and picks up a `-2`. +A step that saves nothing and which no later step reads does not run at +all: the engine drops it before the first step executes and warns once per +dropped step. That is how a template whose portraits can be supplied as +`asset:` files stops paying for the steps that would have drawn them. It +follows from what the definition says, never from a value produced during +the run, so it is decided at validate time too - the `plan` a validate call +answers with counts only the steps that will run and lists the rest under +`elided_steps`. Four things keep a step: a `result` with a `content_type` +and `save` not `false`, being the last step, being read by a later step +(`previous_result:`, `gather:`, a `pipeline_reference`, a shared component), +or being read by a step that is itself kept - elision is transitive. If a +step you meant to run is named in the warnings, a reference to it is +misspelled somewhere later or it needs a `result`. + ### Composing a stored workflow A step with a `workflow` block runs another workflow as one step of this one, @@ -912,6 +926,37 @@ later step needing one of them reloads it. **Example:** [enhance-prompt.json](../workflows/templates/minimax/enhance-prompt.json) +#### A step nothing reads does not run + +Before the first step executes, the engine drops any step whose result no later step +reads and which writes no file, and warns once per dropped step saying which and why. +`dialogue-short` cast from portraits that already exist used to run its two Z-Image +steps anyway and throw the pictures away - about a minute of GPU per episode on +something nothing looked at (#122). + +Four things keep a step: + +- **it saves** - a `result` with a `content_type`, and `save` not `false`. A workflow + whose whole point is writing three images references nothing, so this is the rule that + keeps elision from being destructive. `"save": false` is how a step says it is + scaffolding. +- **it is the last step** - it is the run's answer, whatever it declares. +- **something reads it** - `previous_result:`/`from_previous_result` (including + `previous_result:step.property`), a `gather:` (which is a list of those by the time + this runs), a `pipeline_reference` naming it, or a `reused_components` entry naming a + component it shares. +- Elision is transitive, so dropping a step can drop the step it read in turn. + +`release_pipeline` on an elided step moves onto the last surviving step before it when +that step loaded the same pipeline, and `release_models` moves unconditionally - a +release that vanished with its step would leak the memory it existed to free. The plan a +validate call answers with is computed after elision, so `steps`, `downloads_required` +and the cost it quotes are the work that will actually happen, and it lists what was +dropped under `elided_steps`; the run manifest records the same list. + +If a step you expected to run is named in the warnings, the usual cause is a reference +to it spelled wrong somewhere later, or a step that was meant to declare a `result`. + ### VAE Options ```json diff --git a/dw/elision.py b/dw/elision.py new file mode 100644 index 00000000..095d90cd --- /dev/null +++ b/dw/elision.py @@ -0,0 +1,215 @@ +"""A step nothing references, and which saves nothing, does not run. + +`templates/minimax/dialogue-short` draws its two characters with Z-Image and +references those portraits from every shot. An episode can just as well be +cast from portraits that already exist - a shot entry's subject reference +takes `from_file: "asset:cast/priya.jpg"` exactly as its voice references do - +and then the two draw steps still ran and their output was discarded: roughly +55 seconds and two model loads on portraits nothing in the run looked at +(#109, #122). The recurring cast is the headline use of that template, so +paying for it every episode was the wrong default, and no argument the caller +could pass avoided it. + +So: before the first step executes, drop any step whose result no later step +reads and which writes no file. It is static - it depends only on what the +realized workflow references, never on a value produced during the run, which +is what keeps it a different thing from conditional execution (#118 closed the +step object specifically so an invented `when` is a hard error). + +Four guardrails, all of them load-bearing: + +- **A step that saves is kept.** `save` defaults to true, so a step declaring + a `result` is a deliverable unless it says otherwise - a workflow whose + whole point is writing three images references nothing. +- **The last step is kept**, whatever it declares: it is the run's answer. +- **A release moves rather than disappearing.** `release_pipeline` frees the + memory the step it sits on took; eliding it silently would leak that for + the rest of the run. It moves onto the last surviving step before it when + that step loaded the same pipeline, and is dropped when nothing was loaded + to release. +- **Every elision is a warning.** Without one a misspelled reference makes + the step feeding it vanish and the failure moves from "previous result not + found" to "the picture is wrong". The static reference check refuses an + unresolvable literal reference, which contains most of it - the warning is + what covers the rest. + +Elision is transitive: dropping a step can leave the step it read +unreferenced in turn, so it runs to a fixed point. +""" + +import logging + +from .step_cache import reference_resolves_to, referenced_result_names + +logger = logging.getLogger("dw") + + +def _saves(step): + """Whether the step writes a file. + + Exactly what `Result.save` asks: a `content_type` to write, and `save` + not turned off. `save` defaults to true, so declaring a `result` with a + content type is declaring a deliverable - and a step with no `result` at + all writes nothing, whatever it generates, which is what makes the + portrait steps droppable in the first place. + """ + result = step.get("result") + if not isinstance(result, dict): + return False + if result.get("content_type") is None: + return False + return result.get("save", True) is not False + + +def _reused_component_names(steps): + """Every component name a later step asks an earlier one to have shared. + + Component sharing is keyed on the component's name rather than on the + step's, so the step that shares is referenced without ever being named. + """ + names = set() + for step in steps: + pipeline = step.get("pipeline") + if isinstance(pipeline, dict): + for name in pipeline.get("reused_components") or []: + names.add(name) + return names + + +def _referenced_pipeline_steps(steps): + """Every step name a later `pipeline_reference` addresses.""" + names = set() + for step in steps: + reference = step.get("pipeline_reference") + if isinstance(reference, dict): + name = reference.get("reference_name") + if isinstance(name, str): + names.add(name) + return names + + +def _needed_by(step, later): + """Whether any of `later` reads this step - by result, by pipeline, or by + a component it shares.""" + name = step.get("name") + if not isinstance(name, str): + return True + + if any(reference_resolves_to(ref, name) for ref in referenced_result_names(later)): + return True + if name in _referenced_pipeline_steps(later): + return True + + pipeline = step.get("pipeline") + shared = set((pipeline or {}).get("shared_components") or []) + return bool(shared & _reused_component_names(later)) + + +def _carry_release(elided, kept): + """Move an elided step's release onto the last surviving step before it. + + `release_pipeline` names the pipeline the elided step itself loaded, so + it is only meaningful on an earlier step that loaded the same one - + moving it anywhere else would unload something the workflow did not ask + to unload. `release_models` frees the process-wide task-model cache, so + the last step that ran before this one is exactly where it belongs. With + nothing before it, nothing was loaded and the flag is dropped. + """ + from .workflow import pipeline_cache_key + + if not kept: + return False + predecessor = kept[-1] + carried = False + if elided.get("release_models"): + predecessor["release_models"] = True + carried = True + if not elided.get("release_pipeline"): + return carried + elided_pipeline = elided.get("pipeline") + kept_pipeline = predecessor.get("pipeline") + if not elided_pipeline or not kept_pipeline: + return carried + if pipeline_cache_key(elided_pipeline) == pipeline_cache_key(kept_pipeline): + predecessor["release_pipeline"] = True + carried = True + return carried + + +def elide_unreferenced_steps(steps): + """The steps that will run, and what was dropped. + + Returns (kept, elided) where `elided` is [{'step': name, 'reason': str}], + in the order the steps were written. `steps` is the expanded, substituted + list - a `for_each` member is a step like any other by then, and `gather:` + has already become the `previous_result:` list it stands for. + + The list is not copied: a kept step is the same object that went in, and + a release carried onto a survivor is written into it. + """ + if not isinstance(steps, list) or len(steps) < 2: + return steps, [] + + kept = list(steps) + elided = [] + changed = True + while changed: + changed = False + for index, step in enumerate(kept): + if index == len(kept) - 1: + # The run's answer, whatever it declares + continue + if not isinstance(step, dict) or _saves(step): + continue + if _needed_by(step, kept[index + 1 :]): + continue + reason = "nothing after it reads its result and it saves no file" + if _carry_release(step, kept[:index]): + reason += "; its release was carried onto the step before it" + elided.append({"step": step.get("name"), "reason": reason}) + del kept[index] + changed = True + break + + # Written order, not the order the fixed point happened to reach them in + order = { + step.get("name"): index + for index, step in enumerate(steps) + if isinstance(step, dict) + } + elided.sort(key=lambda entry: order.get(entry["step"], 0)) + return kept, elided + + +def elide_definition(workflow_def): + """elide_unreferenced_steps over a whole definition, in place. + + Returns the elided-step records; the definition's `steps` is replaced + when anything was dropped, so a caller that wants the untouched list must + keep its own copy. + """ + steps = workflow_def.get("steps") + kept, elided = elide_unreferenced_steps(steps) + if elided: + workflow_def["steps"] = kept + return elided + + +def warn_elided(elided): + """Say what did not run, where whoever asked for the run can read it. + + A warning rather than a log line for the reason every run-time warning is + one: a consumer over the API or MCP sees the job's `warnings` list and + nothing else (#82), and "the portrait step was skipped" is precisely the + thing that explains an otherwise inexplicable result. + """ + from .events import emit_warning + + for entry in elided: + emit_warning( + f"Step '{entry['step']}' did not run: {entry['reason']}. If it " + f"was meant to, either a later step's reference to it is " + f"misspelled or it needs a 'result' to save.", + kind="step_elided", + step=entry["step"], + ) diff --git a/dw/plan.py b/dw/plan.py index 54b7777d..5d048373 100644 --- a/dw/plan.py +++ b/dw/plan.py @@ -21,6 +21,7 @@ from huggingface_hub import model_info from huggingface_hub.utils import HFValidationError, validate_repo_id +from .elision import elide_definition from .hub_cache import scan_models from .realize import ( BUILTIN_PREFIX, @@ -88,11 +89,17 @@ def build_plan( expanded = Workflow( realized, candidate.output_dir, candidate.file_spec, candidate.workflow_dir ).expanded_definition() + # The plan is what the run does, and a run does not execute a step + # nothing reads (dw/elision.py, #122) - so the step count, the downloads + # and the fingerprint are all taken after elision, and the acknowledged + # cost is the cost of the work that happens + elided = elide_definition(expanded) entries = list_entries(definition, realized) measured_entries = list_entries(definition, definition) return { "fingerprint": fingerprint(expanded, definition, annotations), "steps": len(expanded.get("steps") or []), + "elided_steps": elided, "list_entries": entries, "cached_steps": cached_steps(definition, realized, arguments, cache_probe), "downloads_required": downloads_required( diff --git a/dw/workflow.py b/dw/workflow.py index b25f8d05..e845bc61 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -29,6 +29,7 @@ ) from .locations import location_errors from .reference_limits import reference_limit_errors +from .elision import elide_definition, warn_elided from .introspection import task_signature_errors from .task_domains import task_argument_errors from .subfolders import step_subfolder, subfolder_errors @@ -679,6 +680,9 @@ def _prepare_definition(self, workflow_def, arguments, base_dir): the run prepares - the step cache keys on the realized step, and a probe that prepared it differently would answer for a run that never happens. + + Records the steps elision dropped on `self._elided_steps` (#122) - + the same list run() warns about and writes into the manifest. """ workflow_id = workflow_def["id"] variables = workflow_def.get("variables", None) @@ -707,6 +711,13 @@ def _prepare_definition(self, workflow_def, arguments, base_dir): # run before anything loads workflow_def = expand_for_each(workflow_def) + # A step nothing after it reads, and which saves no file, does not + # run - after expansion, so a for_each member is judged like any + # other step, and before the seed and the run id, so everything + # downstream counts the steps that will actually execute + # (dw/elision.py, #122) + self._elided_steps = elide_definition(workflow_def) + # Set up random seed for reproducibility. Resolved lazily - as a # dict.get default, torch.seed() would run on every call and reseed # the global RNG even when the workflow names an explicit seed @@ -874,6 +885,9 @@ def run( # BEFORE its replacement loads, or the transition holds both at once self._prior_step_keys = prior_step_keys or {} self.manifest = [] + # What elision dropped this run, filled by _prepare_definition and + # read by the warning pass and the manifest (#122) + self._elided_steps = [] # Overwritten on the way out of the try below - a run that leaves # this alone died on an exception the manifest should say so about status = "failed" @@ -908,6 +922,10 @@ def run( workflow_def, default_seed = self._prepare_definition( workflow_def, arguments, base_dir ) + # Said out loud before anything loads: a step that vanishes + # because a reference to it is misspelled would otherwise show + # up only as a different picture (#122) + warn_elided(self._elided_steps) # A workflow that names no seed gets a fresh one every run, so no # step's cache entry can ever match again - skip the cache # wholesale rather than deep-copying every step's realized images @@ -1312,6 +1330,9 @@ def _write_run_manifest( }, "seed": seed, "arguments": arguments or {}, + # What did not run, and why - a run says what it did not do + # as well as what it did (#122) + "elided_steps": self._elided_steps, "steps": [ { **entry, diff --git a/tests/test_elision.py b/tests/test_elision.py new file mode 100644 index 00000000..b4e11eaa --- /dev/null +++ b/tests/test_elision.py @@ -0,0 +1,305 @@ +"""A step nothing references, and which saves nothing, does not run. + +`dialogue-short` cast from portraits that already exist still ran its two +Z-Image steps and threw their output away - roughly 55 s and two model loads +per episode on pictures nothing looked at (#109, #122). Elision drops them, +with four guardrails: a step that saves is kept, the last step is kept, a +release moves onto the step that runs in its place, and every elision is a +warning, because a misspelled reference would otherwise make the step feeding +it vanish silently. +""" + +import copy +import json +import pathlib + +import pytest + +from dw.elision import elide_definition, elide_unreferenced_steps +from dw.workflow import Workflow + + +def step(name, **extra): + return {"name": name, **extra} + + +def task(name, reads=None, **extra): + arguments = {"audio": f"previous_result:{reads}"} if reads else {"audio": "x"} + return step(name, task={"command": "fade_audio", "arguments": arguments}, **extra) + + +def names(steps): + return [s["name"] for s in steps] + + +class TestWhatIsDropped: + def test_a_step_nothing_reads_and_which_saves_nothing(self): + kept, elided = elide_unreferenced_steps( + [task("orphan"), task("deliverable", result={"content_type": "audio/wav"})] + ) + assert names(kept) == ["deliverable"] + assert elided[0]["step"] == "orphan" + assert "nothing after it reads" in elided[0]["reason"] + + def test_elision_is_transitive(self): + """Dropping a step can leave the one it read unreferenced in turn.""" + kept, elided = elide_unreferenced_steps( + [ + task("first"), + task("second", reads="first"), + task("kept", result={"content_type": "audio/wav"}), + ] + ) + assert names(kept) == ["kept"] + assert [e["step"] for e in elided] == ["first", "second"] + + def test_the_records_are_in_written_order(self): + _, elided = elide_unreferenced_steps( + [ + task("a"), + task("b", reads="a"), + task("last", result={"content_type": "audio/wav"}), + ] + ) + assert [e["step"] for e in elided] == ["a", "b"] + + +class TestTheGuardrails: + def test_a_step_that_saves_is_kept(self): + """`save` defaults to true - a workflow whose whole point is writing + three images references nothing.""" + steps = [ + task("picture", result={"content_type": "image/png"}), + task("last", result={"content_type": "audio/wav"}), + ] + kept, elided = elide_unreferenced_steps(steps) + assert names(kept) == ["picture", "last"] + assert elided == [] + + def test_save_false_is_what_makes_it_droppable(self): + steps = [ + task("picture", result={"content_type": "image/png", "save": False}), + task("last", result={"content_type": "audio/wav"}), + ] + kept, _ = elide_unreferenced_steps(steps) + assert names(kept) == ["last"] + + def test_the_last_step_is_always_kept(self): + kept, elided = elide_unreferenced_steps([task("a", result={"save": False})]) + assert names(kept) == ["a"] + assert elided == [] + + def test_a_referenced_step_is_kept(self): + steps = [task("source"), task("sink", reads="source")] + kept, _ = elide_unreferenced_steps(steps) + assert names(kept) == ["source", "sink"] + + def test_a_property_reference_counts(self): + """`previous_result:segment.mask` reads `segment`.""" + steps = [ + task("segment"), + step( + "use", + task={ + "command": "fade_audio", + "arguments": {"audio": "previous_result:segment.mask"}, + }, + ), + ] + kept, _ = elide_unreferenced_steps(steps) + assert names(kept) == ["segment", "use"] + + def test_a_from_previous_result_object_counts(self): + steps = [ + task("draw"), + step( + "shot", + pipeline={ + "arguments": { + "references": [ + {"reference_type": "X", "from_previous_result": "draw"} + ] + } + }, + ), + ] + kept, _ = elide_unreferenced_steps(steps) + assert names(kept) == ["draw", "shot"] + + def test_a_pipeline_reference_counts(self): + """The step is read for its loaded pipeline rather than its result.""" + steps = [ + step("load", pipeline={"configuration": {"component_type": "X"}}), + step("reuse", pipeline_reference={"reference_name": "load"}), + ] + kept, elided = elide_unreferenced_steps(steps) + assert names(kept) == ["load", "reuse"] + assert elided == [] + + def test_a_shared_component_counts(self): + """Sharing is keyed on the component's name, so the step that shares + is referenced without ever being named.""" + steps = [ + step("loader", pipeline={"shared_components": ["text_encoder"]}), + step("user", pipeline={"reused_components": ["text_encoder"]}), + ] + kept, elided = elide_unreferenced_steps(steps) + assert names(kept) == ["loader", "user"] + assert elided == [] + + +class TestReleasesMoveRatherThanDisappearing: + PIPELINE = { + "configuration": {"component_type": "ZImagePipeline"}, + "from_pretrained_arguments": {"model_name": "Tongyi-MAI/Z-Image-Turbo"}, + } + + def draw(self, name, prompt, **extra): + return step( + name, + pipeline={**copy.deepcopy(self.PIPELINE), "arguments": {"prompt": prompt}}, + **extra, + ) + + def test_a_release_carries_onto_the_step_that_loaded_the_same_pipeline(self): + """Two identical pipeline definitions share one loaded model, so the + release belongs on whichever of them still runs.""" + steps = [ + self.draw("a", "one"), + self.draw("b", "two", release_pipeline=True), + task("uses_a", reads="a", result={"content_type": "audio/wav"}), + ] + kept, elided = elide_unreferenced_steps(steps) + assert names(kept) == ["a", "uses_a"] + assert kept[0]["release_pipeline"] is True + assert "release was carried" in elided[0]["reason"] + + def test_a_release_with_nothing_before_it_is_dropped(self): + """Nothing loaded, so nothing leaks.""" + steps = [ + self.draw("a", "one", release_pipeline=True), + task("last", result={"content_type": "audio/wav"}), + ] + kept, elided = elide_unreferenced_steps(steps) + assert names(kept) == ["last"] + assert "release was carried" not in elided[0]["reason"] + + def test_a_release_is_not_moved_onto_a_different_pipeline(self): + """Moving it would unload something the workflow never asked to + unload.""" + other = step( + "other", + pipeline={"configuration": {"component_type": "FluxPipeline"}}, + result={"content_type": "image/png"}, + ) + steps = [ + other, + self.draw("a", "one", release_pipeline=True), + task("last", result={"content_type": "audio/wav"}), + ] + kept, _ = elide_unreferenced_steps(steps) + assert names(kept) == ["other", "last"] + assert "release_pipeline" not in kept[0] + + def test_release_models_always_carries(self): + """It frees the process-wide task-model cache rather than one + pipeline, so the step that ran before is exactly where it belongs.""" + steps = [ + task("earlier", result={"content_type": "audio/wav"}), + task("orphan", release_models=True), + task("last", result={"content_type": "audio/wav"}), + ] + kept, _ = elide_unreferenced_steps(steps) + assert names(kept) == ["earlier", "last"] + assert kept[0]["release_models"] is True + + +class TestTheWarning: + """Without it a misspelled reference makes the step feeding it vanish and + the failure moves from "previous result not found" to "the picture is + wrong". A log line is not enough - a consumer over the API or MCP reads + the job's warnings and nothing else (#82).""" + + def test_every_elided_step_is_named(self): + from dw.elision import warn_elided + from dw.events import RunContext, activate_context, deactivate_context + + events = [] + token = activate_context(RunContext(on_event=events.append)) + try: + warn_elided([{"step": "draw_character_a", "reason": "nothing reads it"}]) + finally: + deactivate_context(token) + + warnings = [e for e in events if e["event"] == "warning"] + assert len(warnings) == 1 + assert warnings[0]["kind"] == "step_elided" + assert warnings[0]["step"] == "draw_character_a" + assert "did not run" in warnings[0]["message"] + + +class TestDialogueShort: + """The case that raised it.""" + + PATH = "workflows/templates/minimax/dialogue-short.json" + + def definition(self): + return json.loads(pathlib.Path(self.PATH).read_text()) + + def expanded(self, definition): + return Workflow(definition, "outputs", self.PATH).expanded_definition() + + def cast_from_files(self, definition): + for entry in definition["variables"]["shots"]: + for reference in entry.get("references", []): + if "from_previous_result" in reference: + name = reference.pop("from_previous_result") + reference["from_file"] = f"asset:cast/{name}.jpg" + return definition + + def test_the_draw_steps_save_nothing(self): + """Don's condition on the engine change: without `save: false` the + deliverable guardrail keeps them and the motivating case saves + nothing.""" + steps = {s["name"]: s for s in self.definition()["steps"]} + for name in ("draw_character_a", "draw_character_b"): + assert steps[name]["result"]["save"] is False + + def test_the_default_run_is_unchanged(self): + expanded = self.expanded(self.definition()) + elided = elide_definition(expanded) + assert elided == [] + assert "draw_character_a" in names(expanded["steps"]) + + def test_a_cast_episode_draws_nothing(self): + expanded = self.expanded(self.cast_from_files(self.definition())) + elided = elide_definition(expanded) + assert [e["step"] for e in elided] == [ + "draw_character_a", + "draw_character_b", + ] + assert not [n for n in names(expanded["steps"]) if n.startswith("draw_")] + + def test_a_half_cast_episode_still_draws(self): + """Something still reads them, so they still run.""" + definition = self.definition() + entry = definition["variables"]["shots"][0] + for reference in entry.get("references", []): + reference.pop("from_previous_result", None) + expanded = self.expanded(definition) + assert elide_definition(expanded) == [] + + +class TestTheCatalogIsUnchanged: + """Elision is a property every workflow inherits - no template may lose a + step to it on its stored defaults.""" + + @pytest.mark.parametrize( + "path", sorted(str(p) for p in pathlib.Path("workflows").rglob("*.json")) + ) + def test_no_step_is_elided_on_the_defaults(self, path): + definition = json.loads(pathlib.Path(path).read_text()) + if not isinstance(definition, dict) or "steps" not in definition: + pytest.skip("not a workflow") + expanded = Workflow(definition, "outputs", path).expanded_definition() + assert elide_definition(expanded) == [] diff --git a/tests/test_pipeline_caching.py b/tests/test_pipeline_caching.py index 5a53dacb..d5158727 100644 --- a/tests/test_pipeline_caching.py +++ b/tests/test_pipeline_caching.py @@ -235,6 +235,10 @@ def step(name, **extra): "from_pretrained_arguments": {"model_name": f"model-{name}"}, "arguments": {"prompt": "test"}, }, + # Both steps save, so both run: a step that saves nothing and + # which nothing reads is elided before the run (#122), and these + # are release tests rather than elision ones + "result": {"content_type": "image/png"}, } return { @@ -321,7 +325,11 @@ def _release_models_workflow_def(release): "pipeline": { "configuration": {"component_type": "{MockPipeline}"}, "from_pretrained_arguments": {"model_name": "model-generate"}, - "arguments": {"prompt": "test"}, + # Reads the expanded prompt, which is the shape this + # exists for - and keeps the task step in the run, since + # a step nothing reads and which saves nothing is elided + # before the first step executes (#122) + "arguments": {"prompt": "previous_result:expand_prompt"}, }, }, ], diff --git a/tests/test_plan.py b/tests/test_plan.py index d9cfaae5..85b5b9e0 100644 --- a/tests/test_plan.py +++ b/tests/test_plan.py @@ -38,11 +38,16 @@ def definition(): "num_frames": "variable:frames", }, }, + # Both steps declare a result: a step that saves nothing and + # which nothing reads does not run at all now (#122), and + # these fixtures are about the plan rather than about elision + "result": {"content_type": "image/png"}, }, { "name": "shot", "for_each": "variable:shots", "task": {"command": "x", "arguments": {"prompt": "item:prompt"}}, + "result": {"content_type": "video/mp4"}, }, ], } @@ -96,6 +101,7 @@ def make(spec=None, arguments=None, **overrides): "steps", "list_entries", "cached_steps", + "elided_steps", "downloads_required", "estimate", } diff --git a/tests/test_server.py b/tests/test_server.py index fa906c15..11389e67 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -3632,6 +3632,7 @@ def test_an_old_history_database_gains_the_column(tmp_path): "steps": 0, "list_entries": {}, "cached_steps": None, + "elided_steps": [], "downloads_required": [], "estimate": None, } diff --git a/tests/test_workflow_step_cache.py b/tests/test_workflow_step_cache.py index 35862e0a..b4a90f27 100644 --- a/tests/test_workflow_step_cache.py +++ b/tests/test_workflow_step_cache.py @@ -504,6 +504,10 @@ def _two_step_def(second_reads_first, workflow_id="test_step_cache_two"): "from_pretrained_arguments": {"model_name": "model-a"}, "arguments": {"prompt": "variable:a_prompt"}, }, + # A saves, so it runs whether or not B reads it - a step + # that saves nothing and which nothing reads is elided now + # (#122), and these are cache tests rather than elision ones + "result": {"content_type": "image/png"}, }, { "name": "B", @@ -512,6 +516,7 @@ def _two_step_def(second_reads_first, workflow_id="test_step_cache_two"): "from_pretrained_arguments": {"model_name": "model-b"}, "arguments": b_arguments, }, + "result": {"content_type": "image/png"}, }, ], } diff --git a/ui/src/lib/plan.ts b/ui/src/lib/plan.ts index b9e619dd..a8a3584b 100644 --- a/ui/src/lib/plan.ts +++ b/ui/src/lib/plan.ts @@ -39,6 +39,14 @@ export function describePlan(plan: Plan): PlanLine[] { }) } + const elided = plan.elided_steps ?? [] + if (elided.length) { + lines.push({ + text: `skipped: ${elided.map((e) => e.step).join(', ')} - nothing reads them`, + tone: 'warn', + }) + } + const downloads = plan.downloads_required.map((entry) => { const name = entry.repo ?? entry.url ?? '?' return entry.gb === null ? name : `${name} (${entry.gb} GB)` diff --git a/ui/src/lib/types.ts b/ui/src/lib/types.ts index f962141e..fab97e30 100644 --- a/ui/src/lib/types.ts +++ b/ui/src/lib/types.ts @@ -239,6 +239,9 @@ export interface Plan { /** How many steps the worker's step cache would serve; null when the * worker was busy or did not answer. */ cached_steps: number | null + /** The steps that will not run because nothing reads their result and + * they save no file - already excluded from `steps` (#122). */ + elided_steps: { step: string; reason: string }[] downloads_required: { repo: string | null; url?: string; gb: number | null }[] estimate: { minutes: number | null diff --git a/workflows/templates/minimax/dialogue-short.json b/workflows/templates/minimax/dialogue-short.json index 507e1b7a..69932ac0 100644 --- a/workflows/templates/minimax/dialogue-short.json +++ b/workflows/templates/minimax/dialogue-short.json @@ -1,6 +1,6 @@ { "id": "MiniMaxH3SitcomShort", - "description": "A digital short built the way television is built: from cuts, not from one long take. Chained generation degrades with length - every segment conditions on the previous segment's output, so artifacts compound and identity drifts. A scene cut resets that completely: each shot here is generated fresh from the same two character portraits, so shot five is exactly as clean as shot one and the scene can run as long as the script does. Two Z-Image steps draw the cast (the second reuses the first's loaded pipeline - identical configurations share one model - and 'release_pipeline' frees it before the video model loads). The shots are one 'for_each' step over the 'shots' list: one entry per shot, carrying its 'name', its 'prompt', its 'references' and its 'num_frames' - the tag entry runs 141 frames where the others run 124, since length is per-shot. The list is an argument, so a six-shot scene is one more entry, not another file, and the members are named for their entries ('shot@react'); the loaded MiniMax-H3 is reused across all of them. Character consistency across cuts comes from referencing the same portraits in every shot; voice consistency comes from repeating each character's voice description verbatim in every prompt, and, when 'character_a_voice' / 'character_b_voice' name a clip ('asset:cast/priya.wav'), from the audio reference each entry lists for whoever speaks in it - an entry's reference says 'variable:character_a_voice', so one variable sets the voice in every shot that character has. Both default to null, and a reference whose file is null is left out of the list, so a run that names no voice generates exactly what it generated before the variables existed. A shot where both speak lists both; if that reads worse than one, name one voice and leave the other null. The variables are named for roles rather than for the cast of this example - 'character_a', the 'react' entry - because the beats are the reusable part and the sketch is not. The soundscape writes the laugh track. A final 'concat_videos' task is the editor, gathering the shots in list order into one episode - hard cuts, no trims, no seams to hide, because nothing was carried between them. Only the picture cuts hard: 'audio_bleed_ms' rings each shot's laugh track on over the silent opening of the next, the way a live audience carries across a cut. It and 'seam_fade_ms' are variables, so a seam is re-tuned with an argument rather than a copy of the workflow; 1800 ms is where a five-shot cut measured best, since a generated shot opens on more silence than it looks. Shots generated independently also drift in loudness - 10 dB between two shots of one scene is ordinary - and no seam control can hide a level jump, because it is either side of the cut rather than at it; 'match_levels' ('rms' for perceived level, 'peak' for the loudest sample) evens the shots out before they are joined, and left null, as it is by default, a wide spread is warned about in the log rather than passing in silence. An episode can be cast from portraits that already exist rather than drawn here: a shot entry's subject reference takes 'from_file' the way its voice references do ('asset:cast/priya.jpg'), so a recurring cast carries across episodes by name instead of being redrawn each time - which is the whole point of a cast. Note what that does not do yet: the two Z-Image steps still run and their portraits are discarded, because the engine has no way to skip a step nothing references (#109). Until it does, an episode cast entirely from files pays about a minute for two portraits it throws away.", + "description": "A digital short built the way television is built: from cuts, not from one long take. Chained generation degrades with length - every segment conditions on the previous segment's output, so artifacts compound and identity drifts. A scene cut resets that completely: each shot here is generated fresh from the same two character portraits, so shot five is exactly as clean as shot one and the scene can run as long as the script does. Two Z-Image steps draw the cast (the second reuses the first's loaded pipeline - identical configurations share one model - and 'release_pipeline' frees it before the video model loads). The shots are one 'for_each' step over the 'shots' list: one entry per shot, carrying its 'name', its 'prompt', its 'references' and its 'num_frames' - the tag entry runs 141 frames where the others run 124, since length is per-shot. The list is an argument, so a six-shot scene is one more entry, not another file, and the members are named for their entries ('shot@react'); the loaded MiniMax-H3 is reused across all of them. Character consistency across cuts comes from referencing the same portraits in every shot; voice consistency comes from repeating each character's voice description verbatim in every prompt, and, when 'character_a_voice' / 'character_b_voice' name a clip ('asset:cast/priya.wav'), from the audio reference each entry lists for whoever speaks in it - an entry's reference says 'variable:character_a_voice', so one variable sets the voice in every shot that character has. Both default to null, and a reference whose file is null is left out of the list, so a run that names no voice generates exactly what it generated before the variables existed. A shot where both speak lists both; if that reads worse than one, name one voice and leave the other null. The variables are named for roles rather than for the cast of this example - 'character_a', the 'react' entry - because the beats are the reusable part and the sketch is not. The soundscape writes the laugh track. A final 'concat_videos' task is the editor, gathering the shots in list order into one episode - hard cuts, no trims, no seams to hide, because nothing was carried between them. Only the picture cuts hard: 'audio_bleed_ms' rings each shot's laugh track on over the silent opening of the next, the way a live audience carries across a cut. It and 'seam_fade_ms' are variables, so a seam is re-tuned with an argument rather than a copy of the workflow; 1800 ms is where a five-shot cut measured best, since a generated shot opens on more silence than it looks. Shots generated independently also drift in loudness - 10 dB between two shots of one scene is ordinary - and no seam control can hide a level jump, because it is either side of the cut rather than at it; 'match_levels' ('rms' for perceived level, 'peak' for the loudest sample) evens the shots out before they are joined, and left null, as it is by default, a wide spread is warned about in the log rather than passing in silence. An episode can be cast from portraits that already exist rather than drawn here: a shot entry's subject reference takes 'from_file' the way its voice references do ('asset:cast/priya.jpg'), so a recurring cast carries across episodes by name instead of being redrawn each time - which is the whole point of a cast. An episode cast entirely from files does not pay for them: the two Z-Image steps save nothing ('save': false - the shots are the only thing that reads them), so once no shot references them the engine drops them from the run and warns that it did, and the minute and two model loads they cost are not spent (#122). Cast some shots from files and not others and they still run, because something still reads them.", "summary": "A multi-shot dialogue short: Z-Image draws the cast, each shot is generated fresh from the same portraits, then cut.", "cost": [ { @@ -147,7 +147,8 @@ "result": { "content_type": "image/jpeg", "embed_metadata": true, - "subfolder": "intermediate" + "subfolder": "intermediate", + "save": false } }, { @@ -174,7 +175,8 @@ "result": { "content_type": "image/jpeg", "embed_metadata": true, - "subfolder": "intermediate" + "subfolder": "intermediate", + "save": false } }, { From e537520e205d32c68c2c835cde0b47b1a7adf19d Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Sun, 13 Sep 2026 23:00:42 -0500 Subject: [PATCH 03/62] feat: #96 - a variable's bound is declared, checked before anything loads `validate_workflow(num_frames=61)` on an H3 template answered `valid: true`, then the run loaded the weights and the turbo LoRA and failed 138.7 s later on a check against two integers. Every term in that check is a property of the model, so the rule is now declared on the workflow and checked for free. `variable_constraints` (dw/variable_constraints.py) takes a chain step's `frame_snap` field names - one shape, not two - and a chain writes `"frame_snap": "constraint:num_frames"` rather than repeating the numbers. Checked in `validation_errors` (so POST /api/validate, validate_workflow and the pre-queue check all refuse at `arguments.` / `variables.`), at run time in `apply_constraints` before anything loads, and reported beside the default by `list_workflows` (terse) / `get_workflow(variables_only=true)` - the half that stops the next consumer picking 61. The bounds hold for the value the run will *use*, matching diffusers' `align_num_frames`, which snaps before it range-checks: 108 is accepted (it becomes 124), 346 refused (it would become 362). H3's templates declare `snap: "up"` and warn through `emit_warning` when a count is rounded - the silent half of #96. LTX-2.5's declare the `8 * n + 1` grid with *no* snap, because those pipelines floor an off-grid count: rounding up would be a second silent change to the length. tests/test_variable_constraints.py sweeps the whole catalog and pins every declared number to the diffusers symbol it derives from, rather than asserting per template. Co-Authored-By: Claude Opus 5 --- CLAUDE.md | 23 ++ docs/WORKFLOW_GUIDE.md | 53 ++- docs/proposals/preflight-argument-bounds.md | 35 ++ dw/schema.py | 3 +- dw/server/app.py | 23 +- dw/server/catalog_shape.py | 36 +- dw/variable_constraints.py | 360 ++++++++++++++++ dw/workflow.py | 23 ++ dw/workflow_schema.json | 102 +++-- dw_mcp/server.py | 22 +- tests/test_catalog_shape.py | 1 + tests/test_catalog_structure.py | 7 +- tests/test_server.py | 1 + tests/test_variable_constraints.py | 384 ++++++++++++++++++ .../templates/ltx2/chained-segments.json | 8 + workflows/templates/ltx2/extend-clip.json | 8 + .../templates/ltx2/generative-upscale.json | 8 + workflows/templates/ltx2/image-to-video.json | 8 + workflows/templates/ltx2/keyframes.json | 8 + workflows/templates/ltx2/text-to-video.json | 8 + workflows/templates/ltx2/two-stage.json | 8 + .../minimax/chain-matched-and-aligned.json | 17 +- .../minimax/chain-matched-to-audio.json | 17 +- .../minimax/chain-video-continuity.json | 17 +- .../templates/minimax/chained-segments.json | 17 +- .../minimax/composable-references.json | 12 +- .../minimax/enhance-prompt-with-image.json | 12 +- .../templates/minimax/enhance-prompt.json | 12 +- .../minimax/first-and-last-frame.json | 10 + .../minimax/generated-subject-reference.json | 10 + .../templates/minimax/image-to-video.json | 10 + .../templates/minimax/last-frame-only.json | 10 + workflows/templates/minimax/music-video.json | 10 + .../templates/minimax/reference-to-video.json | 10 + workflows/templates/minimax/storyboard.json | 10 + .../templates/minimax/video-with-audio.json | 10 + .../minimax/voice-timbre-reference.json | 10 + 37 files changed, 1264 insertions(+), 59 deletions(-) create mode 100644 dw/variable_constraints.py create mode 100644 tests/test_variable_constraints.py diff --git a/CLAUDE.md b/CLAUDE.md index 26c7053c..91034aa4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -368,6 +368,29 @@ same reason - default setup cannot load a pack. one entry in the table; `tests/test_task_domains.py` pins every entry to a real parameter of a real command so a rename cannot leave one checking nothing +- **A variable's bound is declared by the author, checked three times** — a + model's own rule about a value (H3's `num_frames` is `17 * n + 5` from 124 + to 345) is a property of the model, so it lives in the workflow rather than + in engine code, as a `variable_constraints` entry (`dw/variable_constraints.py`). + One shape, not two: it takes a chain step's `frame_snap` field names, and a + chain writes `"frame_snap": "constraint:num_frames"` rather than repeating + the numbers. `snap: "up"` rounds an off-grid value to the next legal one + and warns (at validation *and* through `emit_warning`, so it reaches the + job's `warnings`); without `snap` an off-grid value is refused. The bounds + hold for the value the run will use, matching diffusers' own + `align_num_frames`, which snaps before it range-checks — so 108 is accepted + (it becomes 124) and 346 refused (it would become 362). LTX-2.5's templates + declare the `8 * n + 1` grid with *no* `snap`, because those pipelines floor + an off-grid count rather than raising: rounding up here would be a second + silent change to the length. Checked in `validation_errors` (so + `POST /api/validate`, `validate_workflow` and the pre-queue check all + refuse it at `arguments.` / `variables.`), at run time in + `apply_constraints` before anything loads, and reported beside the default + by `list_workflows` (terse) / `get_workflow(variables_only=true)` — that + last part is what stops the next consumer picking 61 (#96). + `tests/test_variable_constraints.py` sweeps the whole catalog and pins every + declared number to the diffusers symbol it derives from. A constraint reaches + a top-level variable only, not a field inside a `for_each` entry - **Step cache**: a process-wide singleton (`dw/step_cache.py`) consulted by every `Workflow.run`, including server jobs; entries are keyed by `(workflow id, step name)` and validated against the output *root*, never the per-run directory - a run directory is new every execution and would defeat the cache; disabled entirely when the workflow sets no `seed`; a hit reports the earlier run's files with `reused: true` and writes nothing new; `memory clear` drops it. This is why "Run again" on a seeded workflow finishes instantly and generates nothing - the job page says so when every step was reused, and `POST /api/jobs/{id}/rerun` with `{"new_seed": true}` (MCP `rerun_job(new_seed=True)`) draws a fresh seed into the workflow's seed variable, which is the way to get a different image diff --git a/docs/WORKFLOW_GUIDE.md b/docs/WORKFLOW_GUIDE.md index ec184b8b..dd353954 100644 --- a/docs/WORKFLOW_GUIDE.md +++ b/docs/WORKFLOW_GUIDE.md @@ -336,6 +336,55 @@ in braces keeps it a plain string — `"{nf4}"` is the string `nf4`. Getting thi wrong fails at load time, after validation has already passed, so a value that is meant as text under one of those keys must be braced. +### What a variable is allowed to be + +A model's own rule about a value belongs in the workflow, not in engine code +(CLAUDE.md) and not in a consumer's head. `variable_constraints` declares it +per variable, in the same field names a chain step's `frame_snap` uses: + +```json +"variable_constraints": { + "num_frames": { + "modulus": 17, + "remainder": 5, + "min_frames": 124, + "max_frames": 345, + "snap": "up", + "reason": "the video VAE encodes 17 * n + 5 frames, and MiniMax-H3 generates between 5 and 15 seconds at 24 fps" + } +} +``` + +The value has to be `modulus * n + remainder` within `min_frames` to +`max_frames`. With `snap: "up"` an off-grid value is rounded to the next one +the rule accepts and the run *says so* - `130` becomes `141`, reported as a +warning at validation time and again in the job's `warnings`; without `snap` +it is refused. The bounds are checked against the value the run will use, so +they hold for the rounded number: on the rule above `108` is accepted (it +becomes `124`) and `346` is refused (it would become `362`). + +Checked three times, for the reasons the task-argument domains are: in +`validation_errors`, so `POST /api/validate`, `validate_workflow` and the +pre-queue check all refuse a bad value at `arguments.` or +`variables.` for free; at run time before anything loads, which is the +backstop for a value the static pass cannot see (an inline workflow, a value +a parent passed down); and in the catalog, where `list_workflows` and +`get_workflow(variables_only=true)` report the rule beside the default - the +half that stops the next caller picking a number the model refuses. + +State the rule once. Where a template both declares a constraint and snaps a +chain, the chain's `frame_snap` names it rather than repeating the numbers: + +```json +"frame_snap": "constraint:num_frames" +``` + +Two limits, both accepted. A constraint cannot express a bound that depends +on another variable (a maximum that is `fps * seconds` where a template +exposes `fps`), and it reaches a top-level variable only - not a field inside +a list entry, so a `for_each` template whose entries each carry their own +`num_frames` is unconstrained and relies on the run-time check. + ### A workflow takes only the keys the engine reads The workflow object itself, `step`, `task`, `workflow`, @@ -1271,7 +1320,9 @@ joined into a single file: - `frame_snap` — the constraint the pipeline puts on `num_frames`, used to snap the final `match_audio` segment to a valid length. MiniMax H3 accepts `17n+5` frames between 124 and 345: `{ "modulus": 17, "remainder": 5, "min_frames": 124, - "max_frames": 345 }`. + "max_frames": 345 }`. Where the workflow already declares that rule as a + `variable_constraints` entry, write `"frame_snap": "constraint:num_frames"` + instead, so the numbers live in one place (*What a variable is allowed to be*). - `prompts` — optional per-segment prompt list for narrative progression; segment `i` uses `prompts[min(i, len - 1)]`. - `save_segments` — write each completed segment to the output directory as a diff --git a/docs/proposals/preflight-argument-bounds.md b/docs/proposals/preflight-argument-bounds.md index e3b6d839..3fe87f24 100644 --- a/docs/proposals/preflight-argument-bounds.md +++ b/docs/proposals/preflight-argument-bounds.md @@ -129,3 +129,38 @@ Checking a value against a *loaded* pipeline's signature - that is what the existing signature warnings do. This is only about bounds that are constants of the model, which is the class of error that costs two minutes of GPU before it is reported. + +## Accepted limits (recorded at implementation) + +Approved as Option A with C's reporting folded in (dkackman, 2026-09-14) and +implemented on that basis. Three limits were accepted rather than designed +around: + +1. **A constraint cannot depend on another variable.** A bound that is + `fps * seconds` where a template exposes `fps` has no expression here. + H3's fps is fixed and LTX-2.5's templates carry `frame_rate` as a + variable but bound only the VAE grid, which is rate-independent, so + nothing in the catalog needs it today. +2. **A constraint reaches a top-level variable only.** A `for_each` template + whose entries each carry their own `num_frames` - + `templates/minimax/dialogue-short` is the live case - cannot declare the + rule per entry, and relies on the run-time check as before. Reaching + inside a list entry would need a path-shaped constraint key, which is a + second dialect and was ruled out. +3. **The run-time check stays.** It is the backstop for an inline workflow an + agent just wrote, which no declared constraint covers, and for a value a + parent workflow passed down. + +One thing the implementation resolved that the proposal did not anticipate: +the two families round in *opposite* directions. H3's `align_num_frames` +snaps a frame count **up** and then range-checks the aligned value, so the +H3 templates declare `snap: "up"` and the engine's bounds are checked against +the rounded number (108 is accepted, 346 refused). LTX-2.5's pipelines floor +an off-grid count instead - `(num_frames - 1) // 8 * 8 + 1` - so those +templates declare the `8 * n + 1` grid with **no** `snap`: an off-grid count +is refused with the rule stated, rather than rounded up (which would hand +back a longer clip than either the caller or the pipeline chose) or silently +floored (which is the behaviour being fixed). `snap: "down"` was deliberately +not added - one shape, not two. + +*Implementer agent, model `opus` via provider `anthropic`.* diff --git a/dw/schema.py b/dw/schema.py index f2d866e2..66260169 100644 --- a/dw/schema.py +++ b/dw/schema.py @@ -175,9 +175,10 @@ def load_schema(schema_name): "configures", "cost", "variables", + "variable_constraints", "seed", ], - "defs": [], + "defs": ["variable_constraint"], }, "steps": {"properties": ["steps"], "defs": ["step", "workflow_reference"]}, "pipelines": { diff --git a/dw/server/app.py b/dw/server/app.py index dd4cbefe..96127e34 100644 --- a/dw/server/app.py +++ b/dw/server/app.py @@ -64,6 +64,7 @@ resolve_prompt_reference, ) from ..assets import is_asset_reference, resolve_asset_reference +from ..variable_constraints import constraint_errors, constraint_warnings from ..variables import argument_errors from ..workflow import Workflow, workflow_from_definition, workflow_from_file from .enhancers import build_enhance_workflow, preset_descriptions @@ -293,6 +294,9 @@ def workflow_details(sources_by_name): "traits": metadata["traits"], "summary": metadata["summary"], "lists": metadata["lists"], + # What a variable's value is allowed to be, so the rule is + # read rather than guessed at (#96) + "constraints": definition.get("variable_constraints") or {}, "cost": cost if isinstance(cost, list) and cost else None, } except Exception: @@ -307,6 +311,7 @@ def workflow_details(sources_by_name): "traits": [], "summary": "", "lists": {}, + "constraints": {}, "cost": None, } _workflow_detail_cache[path] = (mtime, detail) @@ -933,6 +938,11 @@ def _candidate_for( ) candidate.validate() problems = argument_errors(candidate.workflow_definition, arguments) + # A value outside a rule the workflow declares, refused before the + # job id rather than after the weights are loaded (#96) + problems += constraint_errors( + candidate.workflow_definition, arguments, supplied=set(arguments or {}) + ) if problems: raise ValueError( "; ".join( @@ -1587,6 +1597,10 @@ def validate_workflow( "error": None, "errors": [], "warnings": workflow_argument_warnings(definition) + # A value a declared constraint will round up - the silent half + # of #96: the run changed the caller's frame count and only the + # server's log said so + + constraint_warnings(definition, request.arguments) + entry_field_warnings(definition, request.arguments) # Why `plan.cached_steps` is 0 for a workflow with no seed - the # cache is off, not empty @@ -1901,13 +1915,20 @@ def preview(value, path): values, truncated = {}, [] for variable, value in variables.items(): values[variable] = value if full else preview(value, variable) - return { + answer = { "name": name, "variables": values, "truncated": truncated, "seed": definition.get("seed"), "origin": source.origin, } + # The rule beside the default it constrains: a consumer reading + # `num_frames: 124` with no range picked 61 and paid 138 s of + # loading to be told the rule was 17n + 5 from 124 (#96) + constraints = definition.get("variable_constraints") + if isinstance(constraints, dict) and constraints: + answer["constraints"] = constraints + return answer @app.get("/api/workflows/{name:path}") def get_workflow(name: str, ws: Workspace = Depends(selected_workspace)): diff --git a/dw/server/catalog_shape.py b/dw/server/catalog_shape.py index 2fa557f5..c8994276 100644 --- a/dw/server/catalog_shape.py +++ b/dw/server/catalog_shape.py @@ -341,12 +341,39 @@ def derive_catalog_metadata(definition): "kinds", "variable_names", "lists", + "constraints", "configures", ) # Carried in the compact view only when set: a template has no -# `configures`, and most workflows have no list-driven step -_COMPACT_WHEN_SET = frozenset({"configures", "lists"}) +# `configures`, most workflows have no list-driven step, and most declare +# no bound on a variable +_COMPACT_WHEN_SET = frozenset({"configures", "lists", "constraints"}) + + +def terse_constraint(rule): + """One variable's rule as a phrase, for the compact listing. + + The block itself carries a `reason` in the author's words, which is what + a consumer reading one workflow wants and what the whole catalog cannot + afford - the compact listing has a token budget (#101), and the numbers + are the part that stops the next consumer picking 61 (#96). + """ + if not isinstance(rule, dict): + return rule + parts = [] + if rule.get("modulus"): + parts.append(f"{rule['modulus']}*n+{rule.get('remainder', 0)}") + low, high = rule.get("min_frames"), rule.get("max_frames") + if low is not None and high is not None: + parts.append(f"{low}-{high}") + elif low is not None: + parts.append(f"{low}+") + elif high is not None: + parts.append(f"up to {high}") + if rule.get("snap") == "up": + parts.append("rounds up") + return ", ".join(parts) def project_listing( @@ -399,6 +426,11 @@ def project_listing( for key in COMPACT_FIELDS if key not in _COMPACT_WHEN_SET or detail.get(key) } + if slim.get("constraints"): + slim["constraints"] = { + variable: terse_constraint(rule) + for variable, rule in slim["constraints"].items() + } if detail.get("configures_missing"): slim["configures_missing"] = detail["configures_missing"] projected[name] = slim diff --git a/dw/variable_constraints.py b/dw/variable_constraints.py new file mode 100644 index 00000000..02cb0314 --- /dev/null +++ b/dw/variable_constraints.py @@ -0,0 +1,360 @@ +"""What a variable's value is allowed to be, declared by the author. + +`validate_workflow` passed `num_frames: 61` on an H3 template and answered +`valid: true`, naming `num_frames` in `checked_arguments` - so the answer +claimed to cover the caller's value. The run then spent 138.7 s loading the +weights and the turbo LoRA, entered the text encoder, and failed on a check +against two integers: + + MiniMax-H3 generates between 5.0 and 15.0 seconds at 24 fps, so + `num_frames`, rounded up to the next `17 * n + 5` the video VAE can + encode, must be between 120 and 360, got 61 (rounded up to 73). + +Every term in that message is a property of the model. None of it needed a +loaded pipeline - and nothing reachable over the API stated it, so a consumer +reading `num_frames: 124` with no range had no way to know the rule was +`17n + 5` from 124 (#96). + +The rule is declared in the workflow, where `frame_snap` already puts the +same numbers for a chain step, and the engine holds no model knowledge of its +own (CLAUDE.md). One shape, not two: a `variable_constraints` entry takes +`frame_snap`'s field names, and a chain step's `frame_snap` may be the string +`"constraint:"` so a template states `17n + 5` once rather than +twice in one file. + +Checked in three places for the reasons the task-argument domains are +(dw/task_domains.py): statically in `validation_errors`, so a stored default +or a caller's argument outside the rule is a free refusal at the JSON path it +sits at; at run time before anything loads, where a value that arrived from +somewhere the static pass cannot see is snapped or refused; and reported +beside the variable's default by the catalog, which is the half that stops +the next consumer picking 61. +""" + +import logging +import math +import numbers + +logger = logging.getLogger("dw") + +CONSTRAINTS_KEY = "variable_constraints" +# What a chain step's `frame_snap` writes instead of repeating the numbers +CONSTRAINT_PREFIX = "constraint:" +# The fields a `frame_snap` block carries, which are the fields a constraint +# is checked on - the rest of a constraint says what to do about a violation +SNAP_FIELDS = ("modulus", "remainder", "min_frames", "max_frames") + + +def declared_constraints(definition): + """The workflow's constraint block, or {}.""" + constraints = definition.get(CONSTRAINTS_KEY) + return constraints if isinstance(constraints, dict) else {} + + +def snap_block(constraint): + """The `frame_snap`-shaped part of a constraint - what a chain step needs + when it names the constraint rather than repeating it.""" + return { + key: constraint[key] + for key in SNAP_FIELDS + if isinstance(constraint, dict) and key in constraint + } + + +def _as_integer(value): + """The value as an int, or None when this is not a number to check. + + A string that is still a reference, a null, a list: not a violation, just + not something this pass can answer about. + """ + if isinstance(value, bool): + return None + if isinstance(value, numbers.Integral): + return int(value) + if isinstance(value, numbers.Real) and float(value).is_integer(): + return int(value) + if isinstance(value, str): + try: + return int(value) + except ValueError: + return None + return None + + +def aligned(value, constraint): + """The smallest value on the constraint's grid that is not below `value`, + ignoring the range, or None when there is no grid to align to. + + The range is checked against *this* number rather than against what the + caller wrote, because that is what the pipeline does: diffusers' + `align_num_frames` snaps first and the duration bound then holds for the + aligned count, which is why 346 frames is refused (it becomes 362) and + 108 is accepted (it becomes 124). + """ + number = _as_integer(value) + modulus = constraint.get("modulus") + if number is None or not modulus: + return None + remainder = constraint.get("remainder", 0) % modulus + steps = math.ceil((number - remainder) / modulus) + target = steps * modulus + remainder + while target < number: + target += modulus + return target + + +def effective(value, constraint): + """The number the run will actually use - `value` itself, or what + `snap: "up"` rounds it to. None when this is not a number to check.""" + number = _as_integer(value) + if number is None: + return None + if constraint.get("snap") != "up": + return number + target = aligned(number, constraint) + return number if target is None else target + + +def snapped(value, constraint): + """What `value` becomes, or None when it is left alone. + + Only `snap: "up"` snaps, and only off the grid: rounding is a change to + what the caller asked for, so it happens where the author said it should + and nowhere else. 61 frames on an H3 template is not a request for 124 - + it rounds to 73, which the range refuses; 130 is a request for 141 (#96). + """ + number = _as_integer(value) + target = effective(value, constraint) + if target is None or target == number: + return None + return target + + +def violations(value, constraint): + """Why `value` breaks `constraint`, as phrases, or [] when it does not. + + Judged on the value the run would use, so a constraint that rounds is + only ever refused for a bound the rounded number still breaks. + """ + number = effective(value, constraint) + if number is None: + return [] + problems = [] + modulus = constraint.get("modulus") + remainder = constraint.get("remainder", 0) + if modulus and (number - remainder) % modulus != 0: + problems.append(f"must be {modulus} * n + {remainder}") + minimum = constraint.get("min_frames") + if minimum is not None and number < minimum: + problems.append(f"must be at least {minimum}") + maximum = constraint.get("max_frames") + if maximum is not None and number > maximum: + problems.append(f"must be at most {maximum}") + return problems + + +def _accepted(constraint): + """The rule as a phrase a consumer can act on, which is the half of #96 + that stops the next caller picking 61: the range and the step, together.""" + parts = [] + low, high = constraint.get("min_frames"), constraint.get("max_frames") + if low is not None and high is not None: + parts.append(f"{low} to {high}") + elif low is not None: + parts.append(f"{low} or more") + elif high is not None: + parts.append(f"{high} or less") + modulus = constraint.get("modulus") + if modulus: + parts.append(f"{modulus} * n + {constraint.get('remainder', 0)}") + return ", ".join(parts) + + +def refusal(name, value, constraint): + """The one wording every layer uses for a value outside its constraint.""" + problems = violations(value, constraint) + if not problems: + return None + reason = constraint.get("reason") + because = f" - {reason}" if reason else "" + target = snapped(value, constraint) + rounding = f" (rounds up to {target})" if target is not None else "" + return ( + f"'{name}' is {value}{rounding}, which this workflow does not " + f"accept: {'; '.join(problems)}. Accepted: {_accepted(constraint)}{because}. " + f"The rule is declared on the workflow, so it is checked before " + f"anything loads rather than after the weights are in memory" + ) + + +def snap_notice(name, value, constraint): + """The wording for a value that will be rounded, or None. + + A value the range still refuses after rounding is not a notice - it is a + refusal, and saying both would be two answers to one mistake. + """ + target = snapped(value, constraint) + if target is None or violations(value, constraint): + return None + reason = constraint.get("reason") + because = f" - {reason}" if reason else "" + return ( + f"'{name}' is {value}, which this workflow rounds up to {target}" + f"{because}. The run generates {target}, not {value}" + ) + + +def constraint_errors(definition, arguments=None, supplied=()): + """Every declared constraint a value breaks, as [{path, message}]. + + A value the caller supplied is reported at `arguments.`, where they + wrote it; a stored default at `variables.`. A constraint that + declares `snap: "up"` and can reach a legal value is a warning rather + than an error - the run will round it, and saying so is the point. + """ + constraints = declared_constraints(definition) + if not constraints: + return [] + variables = definition.get("variables") or {} + values = {**variables, **(arguments or {})} + errors = [] + for name in sorted(constraints): + constraint = constraints[name] + if not isinstance(constraint, dict) or name not in values: + continue + message = refusal(name, values[name], constraint) + if message is None: + # Either legal, or legal once rounded - a value the workflow + # rounds is a warning (constraint_warnings), never an error + continue + where = "arguments" if name in (supplied or ()) else "variables" + errors.append({"path": f"{where}.{name}", "message": message}) + return errors + + +def constraint_warnings(definition, arguments=None): + """Every value a declared constraint will round, as messages. + + The silent half of #96: the engine already rounded `num_frames` up to the + VAE's grid with a `logger.warning`, which by the #82 rule does not exist + out there - a caller got a frame count they did not ask for and nothing + said so. + """ + constraints = declared_constraints(definition) + if not constraints: + return [] + values = {**(definition.get("variables") or {}), **(arguments or {})} + notices = [] + for name in sorted(constraints): + constraint = constraints[name] + if not isinstance(constraint, dict) or name not in values: + continue + notice = snap_notice(name, values[name], constraint) + if notice is not None: + notices.append(notice) + return notices + + +def apply_constraints(definition, variables): + """Refuse or round the run's variable values, before anything loads. + + The run-time half, and the backstop for everything the static pass cannot + see - an inline workflow, a value a parent workflow passed down. Rounds in + place and emits a warning for each rounded value, so a frame count the run + changed reaches the job's `warnings` rather than only the log; raises + ValueError for a value no rule can reach. + """ + from .events import emit_warning + + constraints = declared_constraints(definition) + if not constraints or not isinstance(variables, dict): + return + for name in sorted(constraints): + constraint = constraints[name] + if not isinstance(constraint, dict) or name not in variables: + continue + value = variables[name] + notice = snap_notice(name, value, constraint) + if notice is not None: + variables[name] = snapped(value, constraint) + emit_warning( + notice, kind="value_snapped", variable=name, value=variables[name] + ) + continue + message = refusal(name, value, constraint) + if message is not None: + raise ValueError(message) + + +def resolve_constraint_references(definition): + """Replace every `"frame_snap": "constraint:"` with the declared + constraint's numbers, in place. + + So a template states `17n + 5` once - on the variable, where the catalog + reports it and validation checks it - rather than twice, with a chain + step's copy free to drift from it. A name that is not declared is an + error rather than a silently absent constraint: a chain that snapped to + nothing would stitch segments the pipeline refuses. + """ + constraints = declared_constraints(definition) + + def walk(node): + if isinstance(node, list): + for item in node: + walk(item) + return + if not isinstance(node, dict): + return + reference = node.get("frame_snap") + if isinstance(reference, str) and reference.startswith(CONSTRAINT_PREFIX): + name = reference[len(CONSTRAINT_PREFIX) :] + if name not in constraints: + raise ValueError( + f"'frame_snap': '{reference}' names no entry of this " + f"workflow's 'variable_constraints'. Declared: " + + (", ".join(sorted(constraints)) or "") + ) + node["frame_snap"] = snap_block(constraints[name]) + for value in node.values(): + walk(value) + + walk(definition) + return definition + + +def constraint_reference_errors(definition): + """`frame_snap: "constraint:x"` naming nothing declared, as + [{path, message}] - the validation-time form of the error + resolve_constraint_references raises.""" + constraints = declared_constraints(definition) + errors = [] + + def walk(node, path): + if isinstance(node, list): + for index, item in enumerate(node): + walk(item, f"{path}[{index}]") + return + if not isinstance(node, dict): + return + for key, value in node.items(): + where = f"{path}.{key}" if path else key + if ( + key == "frame_snap" + and isinstance(value, str) + and value.startswith(CONSTRAINT_PREFIX) + and value[len(CONSTRAINT_PREFIX) :] not in constraints + ): + errors.append( + { + "path": where, + "message": ( + f"'{value}' names no entry of this workflow's " + f"'variable_constraints'. Declared: " + + (", ".join(sorted(constraints)) or "") + ), + } + ) + walk(value, where) + + walk(definition, "") + return errors diff --git a/dw/workflow.py b/dw/workflow.py index e845bc61..bf9fc08c 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -32,6 +32,12 @@ from .elision import elide_definition, warn_elided from .introspection import task_signature_errors from .task_domains import task_argument_errors +from .variable_constraints import ( + apply_constraints, + constraint_errors, + constraint_reference_errors, + resolve_constraint_references, +) from .subfolders import step_subfolder, subfolder_errors from .step import Step from .step_cache import ( @@ -611,6 +617,13 @@ def validation_errors(self, arguments=None, composing=None): # is the one mistake a free pre-flight most obviously exists for # (dw/introspection.py, #141) + task_signature_errors(expanded, source_indices) + # A value outside a rule the workflow declares - the bound that + # cost 138 s of loading to discover, refused for free at the + # path the value sits at (dw/variable_constraints.py, #96) + + constraint_errors( + self.workflow_definition, arguments, supplied=set(arguments or {}) + ) + + constraint_reference_errors(self.workflow_definition) + self.sub_workflow_errors(expanded, source_indices, composing) ) @@ -694,6 +707,11 @@ def _prepare_definition(self, workflow_def, arguments, base_dir): # first set variable values base don the arguments passed to the workflow # these may come form the command line or form a parent workflow set_variables(arguments, variables) + # A value outside a rule the workflow declares is refused, and + # one the rule rounds is rounded with a warning saying so - + # before anything loads, and before substitution puts the value + # everywhere it is referenced (dw/variable_constraints.py, #96) + apply_constraints(workflow_def, variables) # an entry of a list-valued variable may name another # variable; resolve those before anything inside it is # realized, so a reference type in an entry is a type name @@ -709,6 +727,11 @@ def _prepare_definition(self, workflow_def, arguments, base_dir): # seed, the run id and the realized workflow are computed, so # each covers what actually runs. A ForEachError here fails the # run before anything loads + # A chain step's `frame_snap` may name the declared constraint + # rather than repeating its numbers, so a template states the rule + # once (#96) + resolve_constraint_references(workflow_def) + workflow_def = expand_for_each(workflow_def) # A step nothing after it reads, and which saves no file, does not diff --git a/dw/workflow_schema.json b/dw/workflow_schema.json index e33235ce..2936640f 100644 --- a/dw/workflow_schema.json +++ b/dw/workflow_schema.json @@ -62,6 +62,13 @@ "variables": { "$ref": "#/$defs/arguments" }, + "variable_constraints": { + "description": "What each variable's value is allowed to be, by variable name. A rule the engine could only enforce after the weights were loaded costs minutes to discover; declared here it is a free refusal at validation time and a line in the catalog beside the default.", + "type": "object", + "additionalProperties": { + "$ref": "#/$defs/variable_constraint" + } + }, "seed": { "description": "Default seed for the entire workflow. Accepts a 'variable:' reference.", "type": [ @@ -89,6 +96,47 @@ "steps" ], "$defs": { + "variable_constraint": { + "description": "A rule the author declares for one variable's value, in the same field names a chain step's 'frame_snap' uses, so a model's frame rule is written once. Checked before anything is queued and reported beside the variable's default, so a consumer reads the rule rather than guessing it (#96).", + "type": "object", + "properties": { + "modulus": { + "type": "integer", + "minimum": 1 + }, + "remainder": { + "type": "integer", + "minimum": 0 + }, + "min_frames": { + "type": "integer", + "minimum": 1 + }, + "max_frames": { + "type": "integer", + "minimum": 1 + }, + "snap": { + "description": "What to do with a value the rule refuses. 'up' rounds to the next value the rule accepts, and warns that it did; absent, the value is refused.", + "enum": [ + "up" + ] + }, + "reason": { + "description": "Why the rule exists, in the author's words - quoted in the refusal and reported beside the default.", + "type": "string" + } + }, + "dependentRequired": { + "modulus": [ + "remainder" + ], + "remainder": [ + "modulus" + ] + }, + "additionalProperties": false + }, "image": { "type": "object", "properties": { @@ -294,31 +342,39 @@ "exclusiveMinimum": 0 }, "frame_snap": { - "description": "Constraint the pipeline places on num_frames - counts must equal modulus * n + remainder within the bounds (MiniMax H3: modulus 17, remainder 5, min 124, max 345). Used to snap the final match_audio segment to a valid length.", - "type": "object", - "properties": { - "modulus": { - "type": "integer", - "minimum": 1 - }, - "remainder": { - "type": "integer", - "minimum": 0 - }, - "min_frames": { - "type": "integer", - "minimum": 1 + "description": "Constraint the pipeline places on num_frames - counts must equal modulus * n + remainder within the bounds (MiniMax H3: modulus 17, remainder 5, min 124, max 345). Used to snap the final match_audio segment to a valid length. May instead be the string 'constraint:', naming an entry of the workflow's 'variable_constraints' - so the rule is stated once rather than twice in one file.", + "oneOf": [ + { + "type": "object", + "properties": { + "modulus": { + "type": "integer", + "minimum": 1 + }, + "remainder": { + "type": "integer", + "minimum": 0 + }, + "min_frames": { + "type": "integer", + "minimum": 1 + }, + "max_frames": { + "type": "integer", + "minimum": 1 + } + }, + "required": [ + "modulus", + "remainder" + ], + "additionalProperties": false }, - "max_frames": { - "type": "integer", - "minimum": 1 + { + "type": "string", + "pattern": "^constraint:[a-zA-Z_][a-zA-Z0-9_-]*$" } - }, - "required": [ - "modulus", - "remainder" - ], - "additionalProperties": false + ] }, "prompts": { "description": "Per-segment prompt overrides; segment i uses prompts[min(i, len(prompts) - 1)].", diff --git a/dw_mcp/server.py b/dw_mcp/server.py index aaca5fe4..88a7a9e5 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -179,8 +179,12 @@ def list_workflows( takes, the steps run over it and the default's length; there `cost[].per_entry`, when present, is the measured cost of one entry (`{variable, minutes, entries}`), so a run over a - different-length list can be priced from it. Templates only by - default; `configures=