From d6443625f5a2f91ce5495371194d4499dfbc5e06 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Tue, 22 Sep 2026 12:56:28 -0500 Subject: [PATCH 1/5] Give every run a version, so one of several can be named A file's name is per step, not per run, so four runs of one workflow write four files called AcornWarsCutAndScore-film.7-0.0.mp4 and the gallery drew four identical captions. The run id told them apart but is not something anyone says out loud, so there was no way for an agent to point a person at one of them and no way for the person to find it. Every run now takes an ordinal. assign_run_version runs when Workflow.run opens the run directory and records it as `version` in manifest.json; run_versions reads it back. Assigned once and never recomputed, which is the point - deleting a middle run leaves a gap rather than sliding every later number down, so a number quoted today still means the same run tomorrow. Assignment is max(recorded) + 1 over every sibling manifest rather than one past the newest: run ids are chronological only to the second, and within one second the spec digest decides the sort, which is exactly what three quick reruns hit. A run with no recorded number - made before this field, or killed before its manifest landed - is backfilled by rank, the unrecorded runs older than every recorded one taking the numbers beneath the lowest. GET /api/gallery and the metadata route carry `version` and `run_id` (run_versions read once per identity per listing, not per file). list_gallery teaches the vocabulary over MCP, which costs 44 tokens the surface budget now takes deliberately. The web UI reads the field only: a v4 chip ahead of the caption, the run id in the detail pane. Nothing on disk is renamed, so `output:` references, the step cache and keep_output are untouched. Two limits taken deliberately: deleting the newest run frees its number for reuse, and the flat layout has no runs to number. Co-Authored-By: Claude Opus 5 (1M context) --- CLAUDE.md | 25 +++++++ docs/MCP.md | 2 +- dw/runs.py | 104 +++++++++++++++++++++++++++ dw/server/app.py | 40 +++++++++-- dw/workflow.py | 22 +++++- dw_mcp/catalog.py | 8 ++- dw_mcp/server.py | 5 ++ tests/test_events.py | 25 +++++++ tests/test_mcp_server.py | 8 ++- tests/test_runs.py | 104 +++++++++++++++++++++++++++ tests/test_server.py | 73 +++++++++++++++++++ ui/CLAUDE.md | 7 ++ ui/src/lib/pages/GalleryPage.svelte | 33 ++++++++- ui/src/lib/pages/GalleryPage.test.ts | 45 ++++++++++++ ui/src/lib/proofs.test.ts | 2 + ui/src/lib/types.ts | 7 ++ 16 files changed, 501 insertions(+), 9 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index a8b7efef..337c4d7f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -368,6 +368,31 @@ same reason - default setup cannot load a pack. `JobManager.realized` finds the file. `exports` is a reserved workspace name: `POST /api/jobs/{id}/export` gathers one finished job into `/exports//` and `GET /exports/.zip` streams it. +- **A run has a number, and it is not derived from the listing** - a file's name + is per *step*, so four runs of one workflow write four files called + `AcornWarsCutAndScore-film.7-0.0.mp4` and the gallery drew four identical + captions: the run id told them apart but is not something anyone says out + loud, so an agent had no way to name one of them to a person. Every run now + takes an ordinal, `assign_run_version` (`dw/runs.py`) at the moment + `Workflow.run` opens the run directory, recorded as `version` in + `manifest.json` and read back by `run_versions`. Assigned once and never + recomputed, which is the point: deleting a middle run leaves a gap rather + than sliding every later number down, so "version 5" still means the same + run tomorrow. Assignment is `max(recorded) + 1` over *every* sibling + manifest, not one past the newest - run ids are chronological only to the + second, and within one second the spec digest decides the sort, which is + exactly what three quick reruns hit. A run with no recorded number (made + before the field, or killed before its manifest landed) is backfilled by + rank: the unrecorded runs older than every recorded one take the numbers + beneath the lowest, later ones continue from the run before. `GET + /api/gallery` and the metadata route carry `version` and `run_id` + (`run_versions` read once per identity per listing, not per file), MCP + `list_gallery` teaches the vocabulary, and the web UI reads the field only - + a `v4` chip on the card, the run id in the detail pane. Nothing on disk is + renamed, so `output:` references, the step cache and `keep_output` are + untouched. Two limits taken deliberately: deleting the *newest* run frees + its number for reuse (the high-water mark lived in the manifest that went + with it), and the flat layout has no runs, so `version` is null there. - **Result subfolders**: a step's `result.subfolder` (`dw/subfolders.py`) puts its files in a subfolder of the run directory - `/final/x.mp4` - by convention `final` or `intermediate`; the engine treats no name specially and there is no default. diff --git a/docs/MCP.md b/docs/MCP.md index a9ee2f1f..58f7762d 100644 --- a/docs/MCP.md +++ b/docs/MCP.md @@ -226,7 +226,7 @@ when no single workflow covers it. | `get_health()` | — | Check that the server is alive, and which machine answered: `version`, `device`, whether a model process is currently resident (`worker_alive`), the job running now and the queue depth. `worker_alive: false` is the normal idle state on a server that has not run a job since startup or the last memory clear - not a fault - the on-demand worker starts with the next job (#206) | | `get_server_info()` | — | What this installation can do and where it keeps things: `device` (the accelerator a run will use), `version`, the `workspace` this session is working in and the workflow/asset/output/prompt `directories` of *that* workspace, the bind address and port, whether a token is required, and whether MCP is mounted. Check the device before authoring - a CUDA-only choice (bitsandbytes, `torch.compile`, flash attention) is not available on an `mps` or `cpu` server. `runtime` (#222) reports the Python version, torch version and the CUDA version torch was built against, the NVIDIA driver version (when `nvidia-smi` is reachable), and the installed versions of diffusers, transformers, accelerate, bitsandbytes, peft, safetensors and sentencepiece (`null` for one not installed) - for diagnosing an environment mismatch between boxes without shelling in | | `list_jobs(limit=20, status=None, workspace=None)` | optional `limit` (newest N), `status` (one state or a comma-separated set of `queued`, `running`, `succeeded`, `failed`, `cancelled`), `workspace` | List queued, running and recent jobs, **newest first**. Bounded by default: the unbounded listing was over a client's tool-result limit on a server with a few months of history, which made it a tool that could not be called at all. `total` says how many matched and `truncated`/`next` say so when the answer was cut - raise `limit` or narrow with `status`. Without `workspace`, a named workspace lists its own jobs and the default one lists every job the server holds | -| `list_gallery(limit=50, subfolder=None, only_orphans=False, workspace=None)` | `limit`, `subfolder`, `only_orphans`, `workspace` | List generated output files, newest first. A name is `//`, where `` may sit in the subfolder the step chose (`final/episode.mp4`); each entry carries `folder` (the workflow) and `subfolder` (by convention `final` or `intermediate`, `''` when the step chose none, any path the workflow wrote otherwise), and `subfolder="final"` lists only deliverables. Each entry also carries a ready-made `url`, already scoped to the workspace that made it - a hand-built `/outputs/` URL 404s for anything but the default workspace. `only_orphans=True` inverts the call: instead of files, it returns run directories with no media anywhere under them (`runs`, each `{name, mtime}`) - a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything; `subfolder` does not apply in this mode, and `name` is exactly what `delete_output` accepts (#170). `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | +| `list_gallery(limit=50, subfolder=None, only_orphans=False, workspace=None)` | `limit`, `subfolder`, `only_orphans`, `workspace` | List generated output files, newest first. A name is `//`, where `` may sit in the subfolder the step chose (`final/episode.mp4`); each entry carries `folder` (the workflow) and `subfolder` (by convention `final` or `intermediate`, `''` when the step chose none, any path the workflow wrote otherwise), and `subfolder="final"` lists only deliverables. Each entry carries `run_id` and `version` - that run's ordinal among the workflow's runs, which is how one of several runs that wrote the same basename is named to a person: the web UI labels the same file `v5`. The number is assigned when the run opens and never renumbered, so deleting a run leaves a gap rather than sliding the rest down, and it is `null` under the flat output layout, which has no runs. Tools take `name`, never `version`. Each entry also carries a ready-made `url`, already scoped to the workspace that made it - a hand-built `/outputs/` URL 404s for anything but the default workspace. `only_orphans=True` inverts the call: instead of files, it returns run directories with no media anywhere under them (`runs`, each `{name, mtime}`) - a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything; `subfolder` does not apply in this mode, and `name` is exactly what `delete_output` accepts (#170). `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | | `get_gallery_metadata(name, envelope=False, workspace=None)` | `name`, `workspace` | Get the metadata embedded in a generated file — or, when `name` is an `asset:` reference, what an *input* asset holds (`source` says which; `job` is null for an asset). Reading an input's duration, frame count, fps and sample rate before a run is how a caller learns the `total_frames`, `fps` and `sample_rate` a workflow expects it to supply: the exact workflow and arguments that produced it, and, for audio/video, a `media` block (duration, rate, channels, fps, size, peak/mean dBFS). `envelope=true` adds `media.envelope` — `rms_dbfs` and `peak_dbfs` one entry per second — which is what locates something in a track rather than measuring the whole of it. `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | ### Media diff --git a/dw/runs.py b/dw/runs.py index ac42818b..7fa11c46 100644 --- a/dw/runs.py +++ b/dw/runs.py @@ -345,6 +345,110 @@ def strip_run_id(relative_path): return split_run_path(relative_path)[0] +# The key a run's ordinal is recorded under in its manifest. It is assigned +# once, when the run directory is opened, and never recomputed - which is the +# whole point: a number quoted in conversation has to still mean the same run +# after a sibling is deleted. Deleting a middle run leaves a gap +RUN_VERSION_KEY = "version" + + +def _recorded_version(run_dir): + """The ordinal a run recorded for itself, or None. + + None covers every way the number can be missing: a run made before this + field existed, one killed before its manifest landed, and one whose + manifest cannot be parsed. All three are ranked rather than trusted. + """ + try: + with open(os.path.join(run_dir, MANIFEST_FILE_NAME)) as file: + manifest = json.load(file) + except (OSError, ValueError): + return None + version = manifest.get(RUN_VERSION_KEY) if isinstance(manifest, dict) else None + return version if isinstance(version, int) and version > 0 else None + + +def _run_ids(identity_dir): + """Every run directory under one workflow identity, oldest first. + + Run ids sort by their UTC timestamp, so lexical order is chronological + to the second - the same property `latest` relies on. Within one second + the spec digest decides, which is arbitrary but stable; nothing here + needs finer ordering than that. + """ + try: + entries = os.listdir(identity_dir) + except OSError: + return [] + return sorted( + name + for name in entries + if is_run_id(name) and os.path.isdir(os.path.join(identity_dir, name)) + ) + + +def run_versions(identity_dir): + """Every run of one workflow mapped to its ordinal: {run id: version}. + + A run that recorded a version keeps it verbatim - that is what makes the + number survive a sibling being deleted. A run that recorded none (made + before the field existed, or killed before its manifest landed) is + ranked into the sequence around it: the unrecorded runs *older* than + every recorded one take the numbers just beneath the lowest recorded + one, so history that predates the field lands where it belongs, and an + unrecorded run anywhere later simply continues from the run before it. + Ordering is by run id, which is chronological. + """ + run_ids = _run_ids(identity_dir) + recorded = { + run_id: _recorded_version(os.path.join(identity_dir, run_id)) + for run_id in run_ids + } + # Room beneath the lowest recorded number for the unrecorded runs that + # come before it. Where there is not enough room the sequence starts at + # 1 and the recorded numbers stand: a duplicate is better than + # renumbering a run someone has already been told the number of + leading = 0 + for run_id in run_ids: + if recorded[run_id] is not None: + break + leading += 1 + next_number = 1 + if leading < len(run_ids): + next_number = max(1, recorded[run_ids[leading]] - leading) + versions = {} + for run_id in run_ids: + if recorded[run_id] is not None: + versions[run_id] = recorded[run_id] + next_number = recorded[run_id] + 1 + else: + versions[run_id] = next_number + next_number += 1 + return versions + + +def assign_run_version(output_dir, identity): + """The ordinal the run about to open under `identity` takes. + + One past the highest ordinal any sibling holds - not one past the newest + run's, because run ids are chronological only across seconds: two runs + started in the same second are ordered by their spec digest, so the last + id is not reliably the highest number. Three quick reruns are exactly + that case. + + It reads every sibling manifest, which is a small JSON file per run of + one workflow, once, against a run measured in minutes. Sharing + `run_versions` rather than deriving the maximum separately is what keeps + the number assigned here and the number the gallery reports from + drifting apart. + + Best effort, like everything else that writes a run's bookkeeping: a + directory that cannot be read yields 1 rather than failing the run. + """ + versions = run_versions(os.path.join(output_dir, identity)) + return max(versions.values(), default=0) + 1 + + def run_directory(output_dir, file_spec, workflow_id, run_id): """Where one execution writes: //. diff --git a/dw/server/app.py b/dw/server/app.py index a384c6a5..19511e1f 100644 --- a/dw/server/app.py +++ b/dw/server/app.py @@ -99,6 +99,7 @@ is_output_reference, is_run_id, resolve_output_reference, + run_versions, split_run_path, ) from ..workspace import ( @@ -2701,13 +2702,14 @@ def _iter_gallery_files(root, group_runs=True): continue relative_name = name if not directory else f"{directory}/{name}" if group_runs: - folder, _run_id, subfolder = split_run_path(relative_name) + folder, run_id, subfolder = split_run_path(relative_name) else: - folder, subfolder = directory, "" + folder, subfolder, run_id = directory, "", "" yield ( relative_name, folder, subfolder, + run_id, kind, os.path.join(current, name), ) @@ -2718,7 +2720,19 @@ def _gallery_entries(root, ws): files = list(_iter_gallery_files(root)) except OSError: files = [] - for relative_name, folder, subfolder, kind, path in files: + # One read of each workflow's run ordinals per listing, not per file: + # a run of fifty files would otherwise re-read the same manifests + # fifty times + versions_by_folder = {} + + def _version(folder, run_id): + if not run_id: + return None + if folder not in versions_by_folder: + versions_by_folder[folder] = run_versions(os.path.join(root, folder)) + return versions_by_folder[folder].get(run_id) + + for relative_name, folder, subfolder, run_id, kind, path in files: try: stat = os.stat(path) except OSError: @@ -2736,6 +2750,14 @@ def _gallery_entries(root, ws): "name": relative_name, "folder": folder, "subfolder": subfolder, + # Which run wrote it, and that run's ordinal among this + # workflow's runs - the 'v4' a person sees in the grid + # and an agent says out loud. Two runs write the same + # basename, so `label` cannot tell them apart and + # `name` is too long to quote. None under the flat + # layout, which has no runs to number + "run_id": run_id, + "version": _version(folder, run_id), # Quoted (slashes kept literal): a name carrying '#', '?' # or '%' would otherwise break the src the gallery # renders it into. The mtime still rides along for cache @@ -2884,6 +2906,7 @@ def gallery_metadata( the only way to read a wav's length was to run a job that copied it into the output directory. `job` is null for an asset (nothing here produced it) and `source` says which of the two roots answered.""" + run_id, version = "", None if is_asset_reference(name): path = _asset_file(name, ws) source, job = "asset", None @@ -2897,6 +2920,13 @@ def gallery_metadata( job = manager.history.job_for_file(name, workspace=ws.name) except Exception: job = None + # Which run wrote it, and that run's ordinal - the same 'v4' the + # listing reports. After "look at version 3" this is the next + # call, so it confirms the right file was reached rather than + # sending the caller back to the listing + folder, run_id, _subfolder = split_run_path(name) + if run_id: + version = run_versions(os.path.join(ws.outputs, folder)).get(run_id) metadata = read_embedded_metadata(path) extension = os.path.splitext(path)[1].lower() media = ( @@ -2909,6 +2939,8 @@ def gallery_metadata( "source": source, "metadata": metadata, "job": job, + "run_id": run_id, + "version": version, "media": media, } @@ -3597,7 +3629,7 @@ def list_assets(ws: Workspace = Depends(selected_workspace)): except OSError: files = [] origin = _asset_origin(ws, root) - for relative, folder, _subfolder, kind, path in files: + for relative, folder, _subfolder, _run_id, kind, path in files: try: stat = os.stat(path) except OSError: diff --git a/dw/workflow.py b/dw/workflow.py index 7077e7f1..4a99c34e 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -58,6 +58,7 @@ FLAT_LAYOUT, REALIZED_FILE_NAME, activate_output_root, + assign_run_version, deactivate_output_root, workflow_identity, manifest_relative_files, @@ -315,6 +316,11 @@ class Workflow: # leaves no manifest of its own, since its steps are already rolled up # into the parent's _run_dir_inherited = False + # That directory's ordinal among this workflow's runs - what the gallery + # shows as 'v4'. None in the flat layout, for a sub-workflow (which is + # part of the parent's run, not a run of its own), and before a run + # starts + _run_version = None # Where the parent step that delegated to this workflow sits in the # run the caller queued: {"step", "index", "total_steps"}. A child # counts its own steps from zero, so without this a composed run @@ -1113,7 +1119,18 @@ def run( self._run_dir = run_directory( self.output_dir, self.file_spec, workflow_id, run_id ) - logger.debug(f"Run directory: {self._run_dir}") + # The run's ordinal among this workflow's runs, taken + # once here and carried into the manifest. Assigning it + # at run time rather than deriving it when the gallery + # asks is what lets a sibling be deleted without + # renumbering the runs that outlive it + self._run_version = assign_run_version( + self.output_dir, + workflow_identity(self.file_spec, workflow_id), + ) + logger.debug( + f"Run directory: {self._run_dir} (v{self._run_version})" + ) # The record of what actually ran, written before the first step # so a crash or a cancel still leaves it. A sub-workflow inherits @@ -1483,6 +1500,9 @@ def _write_run_manifest( self._run_dir, { "run_id": run_id, + # This run's ordinal among the workflow's runs - 'v4' in the + # gallery. Recorded, never recomputed + "version": self._run_version, "status": status, "started_at": started_at, "finished_at": datetime.now(timezone.utc).isoformat(), diff --git a/dw_mcp/catalog.py b/dw_mcp/catalog.py index 8da6ede1..98555712 100644 --- a/dw_mcp/catalog.py +++ b/dw_mcp/catalog.py @@ -230,7 +230,13 @@ def list_gallery(client, limit=50, subfolder=None, only_orphans=False, workspace Each file entry also carries `label`, a bare display basename for a UI grid - it is not a valid reference on its own (two runs can write the same basename) and is not accepted by `get_gallery_metadata` or - `delete_output`. Pass `name` to those, not `label`.""" + `delete_output`. Pass `name` to those, not `label`. + + `version` is that run's ordinal among the workflow's runs, and `run_id` + the run it came from. The version is what to quote to a person - the web + UI labels the same file `v5` - and is stable: it is assigned when the + run opens and a deleted sibling leaves a gap rather than renumbering + what is left. Null under the flat output layout, which has no runs.""" params = {"limit": limit} if subfolder is not None: params["subfolder"] = subfolder diff --git a/dw_mcp/server.py b/dw_mcp/server.py index e00806c7..44e02089 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -410,6 +410,11 @@ def list_gallery( file over HTTP, already scoped to the right workspace; use it as given rather than composing one from the name. + Entries also carry `run_id` and `version`, that run's ordinal among + the workflow's runs - stable, never renumbered. Quote the version + to a person: the web UI labels the same file `v5`. Tools still take + `name`. + `only_orphans=True` inverts the call: instead of files, it returns run directories holding nothing but their own bookkeeping (manifest.json, workflow.json, job.json) as `runs`, each diff --git a/tests/test_events.py b/tests/test_events.py index 4931f41f..8985d425 100644 --- a/tests/test_events.py +++ b/tests/test_events.py @@ -461,3 +461,28 @@ def test_watchdog_event_carries_the_required_fields(): assert stall["phase"] == "saving" assert isinstance(stall["seconds_since_phase_start"], (int, float)) assert "message" in stall + + +def test_each_run_records_its_own_version(tmp_path): + """Consecutive runs of one workflow number themselves 1, 2, 3 - the + ordinal the gallery shows as 'v2' and an agent quotes.""" + + def mock_load(self, shared_components): + self.pipeline = FakePipeline() + + versions = [] + for _ in range(3): + workflow_def = _workflow_def() + workflow_def["steps"][0]["result"] = {"content_type": "image/png"} + workflow = Workflow(workflow_def, str(tmp_path), "test.json") + with patch.object(Pipeline, "load", mock_load): + with patch("dw.workflow.empty_device_cache"): + workflow.run({}, previous_pipelines={}) + # The run's own directory, not the one its files came from: a + # cached step reports the earlier run's files while still being a + # run of its own with its own number + run_dir = pathlib.Path(workflow._run_dir) + manifest = json.loads((run_dir / "manifest.json").read_text()) + versions.append(manifest["version"]) + + assert versions == [1, 2, 3] diff --git a/tests/test_mcp_server.py b/tests/test_mcp_server.py index 86f034bd..b4a45d5e 100644 --- a/tests/test_mcp_server.py +++ b/tests/test_mcp_server.py @@ -1435,7 +1435,13 @@ def test_the_stated_tool_count_is_the_registered_one(): # downscale" without the old space-filling clause, which pushed descriptions # over budget first (13_820). Measured 2026-09-21 at 13_790.8 (9_025.0 / # 3_751.8 / 1_014.0). 9.2 tokens of headroom left. -SURFACE_BUDGET = 13_800 +# Run versions added four sentences to list_gallery teaching `version` and +# `run_id` - the handle for naming one of several runs that wrote the same +# basename, which is the one thing the surface could not say before. Written +# as tightly as it can be said and still 44 tokens over, so the budget takes +# them deliberately rather than the sentence being cut to nothing. Measured +# 2026-09-22 at 13_844.0 (9_078.0 / 3_752.0 / 1_014.0). 6 tokens of headroom. +SURFACE_BUDGET = 13_850 @pytest.mark.asyncio diff --git a/tests/test_runs.py b/tests/test_runs.py index 18713d9a..069b2538 100644 --- a/tests/test_runs.py +++ b/tests/test_runs.py @@ -10,6 +10,7 @@ from dw.runs import ( FLAT_LAYOUT, + assign_run_version, OUTPUT_LAYOUT_ENV_VAR, RUN_LAYOUT, is_run_id, @@ -17,6 +18,7 @@ new_run_id, output_layout, resolve_output_reference, + run_versions, split_run_path, strip_run_id, workflow_identity, @@ -733,3 +735,105 @@ def test_a_chain_spill_lands_in_the_steps_subfolder(self, tmp_path, fake_pipelin ) segments = sorted((tmp_path / "Gyre" / "run" / "final").glob("*.segment-*.mp4")) assert len(segments) == 2 + + +class TestRunVersions: + """A run's ordinal among the runs of its workflow - the number a person + sees as 'v4' in the gallery and an agent says out loud. Assigned once, + recorded in the manifest, and never renumbered when a sibling is + deleted.""" + + @staticmethod + def _run(identity_dir, run_id, version=None): + """A run directory holding a manifest, with or without a version.""" + run_dir = os.path.join(identity_dir, run_id) + os.makedirs(run_dir, exist_ok=True) + manifest = {"run_id": run_id, "status": "completed"} + if version is not None: + manifest["version"] = version + with open(os.path.join(run_dir, "manifest.json"), "w") as file: + json.dump(manifest, file) + return run_dir + + def test_the_first_run_of_a_workflow_is_version_one(self, tmp_path): + assert assign_run_version(str(tmp_path), "ltx2/Gyre") == 1 + + def test_the_next_run_takes_the_number_after_the_newest(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + self._run(str(identity), "20260902-120000-bbbbbbbb", version=2) + assert assign_run_version(str(tmp_path), "ltx2/Gyre") == 3 + + def test_a_deleted_middle_run_leaves_a_gap_rather_than_renumbering(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + self._run(str(identity), "20260903-120000-cccccccc", version=3) + versions = run_versions(str(identity)) + # v2 is gone; v3 is still v3, and the next run is v4 + assert versions == { + "20260901-120000-aaaaaaaa": 1, + "20260903-120000-cccccccc": 3, + } + assert assign_run_version(str(tmp_path), "ltx2/Gyre") == 4 + + def test_runs_started_in_the_same_second_still_number_upward(self, tmp_path): + # Run ids are chronological only across seconds - within one second + # the spec digest decides the sort, so the highest number is not + # necessarily the last id. Three quick reruns are exactly this case + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-dddddddd", version=1) + self._run(str(identity), "20260901-120000-aaaaaaaa", version=2) + assert assign_run_version(str(tmp_path), "ltx2/Gyre") == 3 + + def test_runs_predating_the_field_are_ranked_by_run_id(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260903-120000-cccccccc") + self._run(str(identity), "20260901-120000-aaaaaaaa") + self._run(str(identity), "20260902-120000-bbbbbbbb") + assert run_versions(str(identity)) == { + "20260901-120000-aaaaaaaa": 1, + "20260902-120000-bbbbbbbb": 2, + "20260903-120000-cccccccc": 3, + } + + def test_backfilled_runs_sit_below_the_lowest_recorded_number(self, tmp_path): + # Two runs made before the field existed, then one that records it. + # The recorded number is authoritative; the older two are ranked + # beneath it so nothing collides + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa") + self._run(str(identity), "20260902-120000-bbbbbbbb") + self._run(str(identity), "20260903-120000-cccccccc", version=3) + assert run_versions(str(identity)) == { + "20260901-120000-aaaaaaaa": 1, + "20260902-120000-bbbbbbbb": 2, + "20260903-120000-cccccccc": 3, + } + + def test_a_run_whose_manifest_never_landed_is_still_numbered(self, tmp_path): + # Killed hard enough to write nothing: the directory is there, the + # manifest is not, and it is ranked like any pre-field run + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + os.makedirs(str(identity / "20260902-120000-bbbbbbbb")) + assert run_versions(str(identity)) == { + "20260901-120000-aaaaaaaa": 1, + "20260902-120000-bbbbbbbb": 2, + } + + def test_a_directory_that_is_not_a_run_is_ignored(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + os.makedirs(str(identity / "not-a-run-id")) + assert run_versions(str(identity)) == {"20260901-120000-aaaaaaaa": 1} + + def test_an_unreadable_manifest_does_not_lose_the_run(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + run_dir = identity / "20260901-120000-aaaaaaaa" + os.makedirs(str(run_dir)) + with open(os.path.join(str(run_dir), "manifest.json"), "w") as file: + file.write("{ not json") + assert run_versions(str(identity)) == {"20260901-120000-aaaaaaaa": 1} + + def test_a_workflow_with_no_runs_yet_has_none(self, tmp_path): + assert run_versions(str(tmp_path / "never" / "ran")) == {} diff --git a/tests/test_server.py b/tests/test_server.py index 56f1b165..5dc08882 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -1993,6 +1993,79 @@ def test_gallery_reports_and_filters_by_subfolder(server, tmp_path): assert finals["subfolders"] == full["subfolders"] +def test_gallery_reports_each_run_version(server, tmp_path): + """Four runs of one workflow write the same basename, so the gallery + label alone cannot tell them apart. Every entry carries the run it came + from and that run's ordinal - 'v4' - which is the handle an agent quotes + and a person finds in the grid.""" + import json as _json + + from PIL import Image + + def _run(identity, run_id, version=None): + run = tmp_path / "outputs" / identity / run_id + (run / "final").mkdir(parents=True) + Image.new("RGB", (2, 2)).save(run / "final" / "film.7-0.0.png") + manifest = {"run_id": run_id} + if version is not None: + manifest["version"] = version + (run / "manifest.json").write_text(_json.dumps(manifest)) + return run + + with server(success_script) as client: + _run("acorn/cut", "20260901-120000-aaaaaaaa", version=1) + # v2 was deleted; v3 keeps its number rather than sliding down + _run("acorn/cut", "20260903-120000-cccccccc", version=3) + # a run from before the field existed is ranked, not dropped + _run("acorn/score", "20260902-120000-bbbbbbbb") + # the flat layout has no runs at all + (tmp_path / "outputs" / "ltx").mkdir() + Image.new("RGB", (2, 2)).save(tmp_path / "outputs" / "ltx" / "flat.png") + + by_name = {f["name"]: f for f in client.get("/api/gallery").json()["files"]} + + first = by_name["acorn/cut/20260901-120000-aaaaaaaa/final/film.7-0.0.png"] + assert first["version"] == 1 + assert first["run_id"] == "20260901-120000-aaaaaaaa" + third = by_name["acorn/cut/20260903-120000-cccccccc/final/film.7-0.0.png"] + assert third["version"] == 3 + # the two runs are indistinguishable by label alone - which is the + # whole reason the version is here + assert first["label"] == third["label"] == "film.7-0.0.png" + # numbering is per workflow identity, so another workflow's first + # run is its own v1 + assert ( + by_name["acorn/score/20260902-120000-bbbbbbbb/final/film.7-0.0.png"][ + "version" + ] + == 1 + ) + assert by_name["ltx/flat.png"]["version"] is None + assert by_name["ltx/flat.png"]["run_id"] == "" + + +def test_gallery_metadata_names_the_run_and_its_version(server, tmp_path): + """After "look at version 3", the next call is usually this one - so it + answers with the run and the ordinal rather than making the caller go + back to the listing to confirm it read the right file.""" + import json as _json + + from PIL import Image + + with server(success_script) as client: + run = tmp_path / "outputs" / "acorn/cut" / "20260903-120000-cccccccc" + (run / "final").mkdir(parents=True) + Image.new("RGB", (2, 2)).save(run / "final" / "film.7-0.0.png") + (run / "manifest.json").write_text(_json.dumps({"version": 3})) + + body = client.get( + "/api/gallery/acorn/cut/20260903-120000-cccccccc" + "/final/film.7-0.0.png/metadata" + ).json() + assert body["run_id"] == "20260903-120000-cccccccc" + assert body["version"] == 3 + + def test_gallery_only_orphans_lists_media_less_run_directories(server, tmp_path): """#170: a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything, has no diff --git a/ui/CLAUDE.md b/ui/CLAUDE.md index 58aceef6..979067f3 100644 --- a/ui/CLAUDE.md +++ b/ui/CLAUDE.md @@ -61,6 +61,13 @@ workflow that has never run gets no frame at all rather than a grey placeholder (a fresh workspace would otherwise be a wall of empty plates). Within a folder, workflows that have produced something sort first. +Stripping the run id is also what makes four runs of one workflow four +identical captions, since a file's name is per step rather than per run. +The gallery entry carries `version` - the run's ordinal, assigned by the +engine and never renumbered - and the grid draws it as a `v4` chip ahead +of the label, with the run id in the detail pane beside it. The UI reads +the field only; nothing here computes or orders a version. + Every picture in the app sits in the global `.frame` (app.css): the media fills it edge to edge, with no inner padding and no rounding of its own, which is what makes it read as a proof on a sheet rather than as another diff --git a/ui/src/lib/pages/GalleryPage.svelte b/ui/src/lib/pages/GalleryPage.svelte index e9f0b1ac..3c6be3e3 100644 --- a/ui/src/lib/pages/GalleryPage.svelte +++ b/ui/src/lib/pages/GalleryPage.svelte @@ -319,7 +319,15 @@ {:else} ♪ {file.label} {/if} - {file.label} + + {#if file.version} + v{file.version} + {/if}{file.label} {/snippet} @@ -352,6 +360,15 @@ class="muted" title="open the file itself in a new tab">open file + {#if selected.version} + + version {selected.version} · {selected.run_id} + {/if} {formatBytes(selected.size)} · {formatMtime(selected.mtime)} @@ -524,6 +541,20 @@ white-space: normal; word-break: break-all; } + /* The one part of the caption that must not be broken or clamped away: + with four runs writing the same name it is the only thing on the card + that differs. Inline-block so word-break: break-all cannot split 'v10' + across lines */ + .version { + display: inline-block; + margin-right: 0.35rem; + padding: 0 0.3rem; + border-radius: 0.2rem; + background: var(--chip, rgb(255 255 255 / 0.08)); + color: var(--fg); + font-weight: 600; + word-break: keep-all; + } .detail { position: sticky; bottom: 1rem; diff --git a/ui/src/lib/pages/GalleryPage.test.ts b/ui/src/lib/pages/GalleryPage.test.ts index 16c2f495..68a859d6 100644 --- a/ui/src/lib/pages/GalleryPage.test.ts +++ b/ui/src/lib/pages/GalleryPage.test.ts @@ -6,6 +6,7 @@ import { within, } from '@testing-library/svelte' import { afterEach, beforeEach, expect, it, vi } from 'vitest' +import { fireEvent } from '@testing-library/dom' // Hoisted above the imports so the static import of the component below - // itself hoisted - sees an initialized mock. Importing the component inside // the test instead would charge its (multi-second) compile to the test timeout @@ -18,6 +19,8 @@ const file = (name: string, subfolder = ''): GalleryFile => ({ name, folder: name.includes('/') ? name.split('/')[0] : '', subfolder, + run_id: '', + version: null, url: `/outputs/${name}`, kind: 'image', size: 1024, @@ -383,3 +386,45 @@ it('leaves the detail open when Escape answers a confirm dialog', async () => { screen.getByLabelText('delete this file from the output directory'), ).toBeTruthy() }) + +it('marks each file with the version of the run that wrote it', async () => { + // Two runs of one workflow write the same basename - the case where the + // label alone tells a person nothing about which is which + listing.files = [ + { ...file('acorn/r1/film.mp4'), label: 'film.mp4', version: 1 }, + { ...file('acorn/r2/film.mp4'), label: 'film.mp4', version: 4 }, + ] + render(GalleryPage) + await waitFor(() => expect(screen.getAllByText('film.mp4')).toHaveLength(2)) + expect(screen.getByText('v1')).toBeTruthy() + expect(screen.getByText('v4')).toBeTruthy() +}) + +it('shows no version for a flat-layout file, which belongs to no run', async () => { + // The flat layout has no runs to number, and a card must not read + // 'vnull' or 'vundefined' because of it + listing.files = [{ ...file('ltx/flat.png'), label: 'flat.png' }] + render(GalleryPage) + await waitFor(() => expect(screen.getByText('flat.png')).toBeTruthy()) + expect(screen.queryByText(/^v\S+$/)).toBeNull() +}) + +it('names the run and its version in the details of the selected file', async () => { + // "Look at version 4" ends here: the pane says which run it reached, so + // the number in the grid can be checked against the one quoted + listing.files = [ + { + ...file('acorn/r2/film.mp4'), + label: 'film.mp4', + run_id: 'r2', + version: 4, + }, + ] + render(GalleryPage) + await waitFor(() => expect(screen.getByText('film.mp4')).toBeTruthy()) + await fireEvent.click(screen.getByText('film.mp4')) + const detail = document.querySelector('.detail') as HTMLElement + // One span holding both forms, so the query is over its whole text + expect(detail.textContent).toContain('version 4') + expect(detail.querySelector('code')?.textContent).toBe('r2') +}) diff --git a/ui/src/lib/proofs.test.ts b/ui/src/lib/proofs.test.ts index 895d00c2..59029b4b 100644 --- a/ui/src/lib/proofs.test.ts +++ b/ui/src/lib/proofs.test.ts @@ -10,6 +10,8 @@ const file = ( name, folder, subfolder: '', + run_id: '', + version: null, url: '/' + name, kind, size: 1, diff --git a/ui/src/lib/types.ts b/ui/src/lib/types.ts index d1dabb7c..357fcbaf 100644 --- a/ui/src/lib/types.ts +++ b/ui/src/lib/types.ts @@ -262,6 +262,13 @@ export interface GalleryFile { /** What followed the run id in the file's path - the `final` / * `intermediate` a step's `result.subfolder` chose, `''` for none. */ subfolder: string + /** The run that wrote the file, `''` under the flat layout. */ + run_id: string + /** That run's ordinal among the workflow's runs - what the grid shows as + * `v4`. Two runs write the same `label`, so this is what tells them + * apart at a glance. Assigned when the run opens and never renumbered, + * so a deleted sibling leaves a gap. Null when there is no run. */ + version: number | null url: string kind: 'image' | 'video' | 'audio' size: number From 7240ede218af4fa30321dd6be2193d347b34e28d Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Tue, 22 Sep 2026 16:06:04 -0500 Subject: [PATCH 2/5] feat: implement run versioning and manifest updates for workflow executions --- CLAUDE.md | 17 +++- docs/MCP.md | 2 +- dw/runs.py | 146 ++++++++++++++++++++++----- dw/server/app.py | 7 ++ dw/workflow.py | 21 +++- dw_mcp/catalog.py | 4 +- tests/test_events.py | 30 ++++++ tests/test_runs.py | 79 +++++++++++++++ tests/test_server.py | 38 +++++++ ui/src/lib/pages/GalleryPage.svelte | 9 +- ui/src/lib/pages/GalleryPage.test.ts | 6 +- 11 files changed, 320 insertions(+), 39 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 337c4d7f..f6849a83 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -381,10 +381,19 @@ same reason - default setup cannot load a pack. run tomorrow. Assignment is `max(recorded) + 1` over *every* sibling manifest, not one past the newest - run ids are chronological only to the second, and within one second the spec digest decides the sort, which is - exactly what three quick reruns hit. A run with no recorded number (made - before the field, or killed before its manifest landed) is backfilled by - rank: the unrecorded runs older than every recorded one take the numbers - beneath the lowest, later ones continue from the run before. `GET + exactly what three quick reruns hit. The number is on disk from the moment + the run opens - a `status: "running"` manifest is written before the first + step and rewritten in full at the end - so a hard kill does not lose it and + a second process opening a run of the same workflow sees it. A run with no + recorded number (made before the field, or killed before even that first + manifest) is ranked: the unrecorded runs older than every recorded one take + the numbers beneath the lowest, later ones continue from the highest before + them. A ranked number would move when an older sibling is deleted, so + `record_run_versions` writes it into the manifest on the two write paths - + a run opening and a run directory being deleted; the listing never writes, + and a run with no manifest at all is left ranked. A gap in the numbers is + not only a deletion: a failed run or a fully cached rerun takes a number + and may have no media for the gallery to show under it. `GET /api/gallery` and the metadata route carry `version` and `run_id` (`run_versions` read once per identity per listing, not per file), MCP `list_gallery` teaches the vocabulary, and the web UI reads the field only - diff --git a/docs/MCP.md b/docs/MCP.md index 58f7762d..f06956cb 100644 --- a/docs/MCP.md +++ b/docs/MCP.md @@ -226,7 +226,7 @@ when no single workflow covers it. | `get_health()` | — | Check that the server is alive, and which machine answered: `version`, `device`, whether a model process is currently resident (`worker_alive`), the job running now and the queue depth. `worker_alive: false` is the normal idle state on a server that has not run a job since startup or the last memory clear - not a fault - the on-demand worker starts with the next job (#206) | | `get_server_info()` | — | What this installation can do and where it keeps things: `device` (the accelerator a run will use), `version`, the `workspace` this session is working in and the workflow/asset/output/prompt `directories` of *that* workspace, the bind address and port, whether a token is required, and whether MCP is mounted. Check the device before authoring - a CUDA-only choice (bitsandbytes, `torch.compile`, flash attention) is not available on an `mps` or `cpu` server. `runtime` (#222) reports the Python version, torch version and the CUDA version torch was built against, the NVIDIA driver version (when `nvidia-smi` is reachable), and the installed versions of diffusers, transformers, accelerate, bitsandbytes, peft, safetensors and sentencepiece (`null` for one not installed) - for diagnosing an environment mismatch between boxes without shelling in | | `list_jobs(limit=20, status=None, workspace=None)` | optional `limit` (newest N), `status` (one state or a comma-separated set of `queued`, `running`, `succeeded`, `failed`, `cancelled`), `workspace` | List queued, running and recent jobs, **newest first**. Bounded by default: the unbounded listing was over a client's tool-result limit on a server with a few months of history, which made it a tool that could not be called at all. `total` says how many matched and `truncated`/`next` say so when the answer was cut - raise `limit` or narrow with `status`. Without `workspace`, a named workspace lists its own jobs and the default one lists every job the server holds | -| `list_gallery(limit=50, subfolder=None, only_orphans=False, workspace=None)` | `limit`, `subfolder`, `only_orphans`, `workspace` | List generated output files, newest first. A name is `//`, where `` may sit in the subfolder the step chose (`final/episode.mp4`); each entry carries `folder` (the workflow) and `subfolder` (by convention `final` or `intermediate`, `''` when the step chose none, any path the workflow wrote otherwise), and `subfolder="final"` lists only deliverables. Each entry carries `run_id` and `version` - that run's ordinal among the workflow's runs, which is how one of several runs that wrote the same basename is named to a person: the web UI labels the same file `v5`. The number is assigned when the run opens and never renumbered, so deleting a run leaves a gap rather than sliding the rest down, and it is `null` under the flat output layout, which has no runs. Tools take `name`, never `version`. Each entry also carries a ready-made `url`, already scoped to the workspace that made it - a hand-built `/outputs/` URL 404s for anything but the default workspace. `only_orphans=True` inverts the call: instead of files, it returns run directories with no media anywhere under them (`runs`, each `{name, mtime}`) - a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything; `subfolder` does not apply in this mode, and `name` is exactly what `delete_output` accepts (#170). `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | +| `list_gallery(limit=50, subfolder=None, only_orphans=False, workspace=None)` | `limit`, `subfolder`, `only_orphans`, `workspace` | List generated output files, newest first. A name is `//`, where `` may sit in the subfolder the step chose (`final/episode.mp4`); each entry carries `folder` (the workflow) and `subfolder` (by convention `final` or `intermediate`, `''` when the step chose none, any path the workflow wrote otherwise), and `subfolder="final"` lists only deliverables. Each entry carries `run_id` and `version` - that run's ordinal among the workflow's runs, which is how one of several runs that wrote the same basename is named to a person: the web UI labels the same file `v5`. The number is assigned when the run opens and never renumbered, so deleting a run leaves a gap rather than sliding the rest down (a failed run, or a rerun that reused every step, leaves one too - it took a number and may have nothing to list), and it is `null` under the flat output layout, which has no runs. Tools take `name`, never `version`. Each entry also carries a ready-made `url`, already scoped to the workspace that made it - a hand-built `/outputs/` URL 404s for anything but the default workspace. `only_orphans=True` inverts the call: instead of files, it returns run directories with no media anywhere under them (`runs`, each `{name, mtime}`) - a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything; `subfolder` does not apply in this mode, and `name` is exactly what `delete_output` accepts (#170). `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | | `get_gallery_metadata(name, envelope=False, workspace=None)` | `name`, `workspace` | Get the metadata embedded in a generated file — or, when `name` is an `asset:` reference, what an *input* asset holds (`source` says which; `job` is null for an asset). Reading an input's duration, frame count, fps and sample rate before a run is how a caller learns the `total_frames`, `fps` and `sample_rate` a workflow expects it to supply: the exact workflow and arguments that produced it, and, for audio/video, a `media` block (duration, rate, channels, fps, size, peak/mean dBFS). `envelope=true` adds `media.envelope` — `rms_dbfs` and `peak_dbfs` one entry per second — which is what locates something in a track rather than measuring the whole of it. `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | ### Media diff --git a/dw/runs.py b/dw/runs.py index 7fa11c46..797369ca 100644 --- a/dw/runs.py +++ b/dw/runs.py @@ -118,6 +118,7 @@ def _runs_newest_first(directory): for name in os.listdir(directory) if is_run_id(name) and os.path.isdir(os.path.join(directory, name)) ), + key=run_id_sort_key, reverse=True, ) except OSError: @@ -348,9 +349,44 @@ def strip_run_id(relative_path): # The key a run's ordinal is recorded under in its manifest. It is assigned # once, when the run directory is opened, and never recomputed - which is the # whole point: a number quoted in conversation has to still mean the same run -# after a sibling is deleted. Deleting a middle run leaves a gap +# after a sibling is deleted. Deleting a middle run leaves a gap, and so does +# a run that wrote no media (it failed, or every step was reused from the +# cache): it took a number and has nothing in the gallery to show under it RUN_VERSION_KEY = "version" +# The length of a run id before any '-N' counter a same-second rerun takes +_RUN_ID_BASE_LENGTH = len("20260101-000000-00000000") + +# Recorded ordinals by manifest path, keyed on the manifest's stat so an +# edited or replaced manifest is read again. A recorded number never changes, +# so this is what keeps a gallery listing from parsing every manifest under +# the output root on every call +_recorded_versions = {} + + +def run_id_sort_key(run_id): + """Order run ids oldest first, with a rerun's '-N' counter compared as a + number - lexically '-10' would sort before '-2'.""" + base, counter = run_id[:_RUN_ID_BASE_LENGTH], run_id[_RUN_ID_BASE_LENGTH + 1 :] + return (base, int(counter) if counter.isdigit() else 1) + + +def _read_manifest(run_dir): + """A run's manifest as a dict, or None when it is missing or unreadable.""" + try: + with open(os.path.join(run_dir, MANIFEST_FILE_NAME)) as file: + manifest = json.load(file) + except (OSError, ValueError): + return None + return manifest if isinstance(manifest, dict) else None + + +def _valid_version(version): + # bool is an int subclass, and True is not version 1 + if isinstance(version, bool) or not isinstance(version, int): + return None + return version if version > 0 else None + def _recorded_version(run_dir): """The ordinal a run recorded for itself, or None. @@ -359,19 +395,26 @@ def _recorded_version(run_dir): field existed, one killed before its manifest landed, and one whose manifest cannot be parsed. All three are ranked rather than trusted. """ + path = os.path.join(run_dir, MANIFEST_FILE_NAME) try: - with open(os.path.join(run_dir, MANIFEST_FILE_NAME)) as file: - manifest = json.load(file) - except (OSError, ValueError): + stat = os.stat(path) + except OSError: + _recorded_versions.pop(path, None) return None - version = manifest.get(RUN_VERSION_KEY) if isinstance(manifest, dict) else None - return version if isinstance(version, int) and version > 0 else None + signature = (stat.st_mtime_ns, stat.st_size, stat.st_ino) + cached = _recorded_versions.get(path) + if cached is not None and cached[0] == signature: + return cached[1] + manifest = _read_manifest(run_dir) + version = _valid_version(manifest.get(RUN_VERSION_KEY)) if manifest else None + _recorded_versions[path] = (signature, version) + return version def _run_ids(identity_dir): """Every run directory under one workflow identity, oldest first. - Run ids sort by their UTC timestamp, so lexical order is chronological + Run ids sort by their UTC timestamp, so this order is chronological to the second - the same property `latest` relies on. Within one second the spec digest decides, which is arbitrary but stable; nothing here needs finer ordering than that. @@ -381,24 +424,17 @@ def _run_ids(identity_dir): except OSError: return [] return sorted( - name - for name in entries - if is_run_id(name) and os.path.isdir(os.path.join(identity_dir, name)) + ( + name + for name in entries + if is_run_id(name) and os.path.isdir(os.path.join(identity_dir, name)) + ), + key=run_id_sort_key, ) -def run_versions(identity_dir): - """Every run of one workflow mapped to its ordinal: {run id: version}. - - A run that recorded a version keeps it verbatim - that is what makes the - number survive a sibling being deleted. A run that recorded none (made - before the field existed, or killed before its manifest landed) is - ranked into the sequence around it: the unrecorded runs *older* than - every recorded one take the numbers just beneath the lowest recorded - one, so history that predates the field lands where it belongs, and an - unrecorded run anywhere later simply continues from the run before it. - Ordering is by run id, which is chronological. - """ +def _ranked_versions(identity_dir): + """({run id: version}, {run id: recorded version or None}).""" run_ids = _run_ids(identity_dir) recorded = { run_id: _recorded_version(os.path.join(identity_dir, run_id)) @@ -420,10 +456,64 @@ def run_versions(identity_dir): for run_id in run_ids: if recorded[run_id] is not None: versions[run_id] = recorded[run_id] - next_number = recorded[run_id] + 1 + # Never backwards: two runs of one second can sort in the + # opposite order to their numbers, and an unrecorded run after + # them must not take a number the higher one already holds + next_number = max(next_number, recorded[run_id] + 1) else: versions[run_id] = next_number next_number += 1 + return versions, recorded + + +def run_versions(identity_dir): + """Every run of one workflow mapped to its ordinal: {run id: version}. + + A run that recorded a version keeps it verbatim - that is what makes the + number survive a sibling being deleted. A run that recorded none (made + before the field existed, or killed before its manifest landed) is + ranked into the sequence around it: the unrecorded runs *older* than + every recorded one take the numbers just beneath the lowest recorded + one, so history that predates the field lands where it belongs, and an + unrecorded run anywhere later continues from the highest number before + it. Ordering is by run id, which is chronological. + + Read only. A ranked number is only as stable as its neighbours until + `record_run_versions` writes it down. + """ + return _ranked_versions(identity_dir)[0] + + +def record_run_versions(identity_dir): + """Write each ranked number into the manifest of a run that has one but + records no version, and return every run's ordinal. + + A ranked number moves when an older unrecorded sibling is deleted, so + runs made before the field existed are pinned the first time anything + writes under their workflow: a new run opening, or a run directory + being deleted. The listing never writes. A run with no manifest at all + is left alone - writing one would invent a record of a run nobody + recorded - and stays ranked. + + Best effort: a manifest that cannot be rewritten keeps its ranked number. + """ + versions, recorded = _ranked_versions(identity_dir) + for run_id, version in versions.items(): + if recorded[run_id] is not None: + continue + run_dir = os.path.join(identity_dir, run_id) + manifest = _read_manifest(run_dir) + if manifest is None: + continue + manifest[RUN_VERSION_KEY] = version + path = os.path.join(run_dir, MANIFEST_FILE_NAME) + partial = f"{path}.partial" + try: + with open(partial, "w") as file: + json.dump(manifest, file, indent=2, default=str) + os.replace(partial, path) + except OSError as e: + logger.warning(f"Could not record version {version} in {path}: {e}") return versions @@ -436,16 +526,16 @@ def assign_run_version(output_dir, identity): id is not reliably the highest number. Three quick reruns are exactly that case. - It reads every sibling manifest, which is a small JSON file per run of - one workflow, once, against a run measured in minutes. Sharing - `run_versions` rather than deriving the maximum separately is what keeps - the number assigned here and the number the gallery reports from + Pins the ranked numbers of older runs on the way (`record_run_versions`), + so history that predates the field stops moving once a new run joins it. + Sharing that ranking rather than deriving the maximum separately is what + keeps the number assigned here and the number the gallery reports from drifting apart. Best effort, like everything else that writes a run's bookkeeping: a directory that cannot be read yields 1 rather than failing the run. """ - versions = run_versions(os.path.join(output_dir, identity)) + versions = record_run_versions(os.path.join(output_dir, identity)) return max(versions.values(), default=0) + 1 diff --git a/dw/server/app.py b/dw/server/app.py index 19511e1f..28ee5929 100644 --- a/dw/server/app.py +++ b/dw/server/app.py @@ -98,6 +98,7 @@ REALIZED_FILE_NAME, is_output_reference, is_run_id, + record_run_versions, resolve_output_reference, run_versions, split_run_path, @@ -3387,6 +3388,9 @@ def _prune_empty_run_directory(name, root): continue return None + # Pin the siblings' numbers first: a run that predates versions is + # ranked, and removing one ahead of it would renumber it + record_run_versions(os.path.dirname(run_dir)) shutil.rmtree(run_dir, ignore_errors=True) # And the identity folders above it, while they are empty - a swept # workspace should not keep one directory per workflow it once ran @@ -3429,6 +3433,9 @@ def delete_output(name: str, ws: Workspace = Depends(selected_workspace)): """ run_dir = _run_directory(name, ws.outputs) if run_dir is not None: + # As in _prune_empty_run_directory: pin the siblings' numbers + # before one of them goes + record_run_versions(os.path.dirname(run_dir)) shutil.rmtree(run_dir, ignore_errors=True) parent = os.path.dirname(run_dir) while os.path.normpath(parent) != os.path.normpath(ws.outputs): diff --git a/dw/workflow.py b/dw/workflow.py index 4a99c34e..0845237e 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -1152,6 +1152,20 @@ def run( # Never fatal: the record is worth less than the run logger.warning(f"Could not realize workflow {workflow_id}: {e}") + # A manifest now, rewritten in full when the run ends: the + # version held only in memory until then was lost to a hard + # kill, and a second process opening a run of this workflow + # meanwhile could not see it and took the same number + self._write_run_manifest( + run_id, + "running", + started_at, + arguments, + resolved_seed, + realized_name, + annotations, + ) + # Which run this is, so a server job can find the directory # it wrote. Emitted even when the realized file did not land: # the manifest is still there, and so are the files @@ -1505,7 +1519,12 @@ def _write_run_manifest( "version": self._run_version, "status": status, "started_at": started_at, - "finished_at": datetime.now(timezone.utc).isoformat(), + # None on the manifest written as the run opens + "finished_at": ( + None + if status == "running" + else datetime.now(timezone.utc).isoformat() + ), "dw_version": __version__, "device": str(get_device()), "workflow": { diff --git a/dw_mcp/catalog.py b/dw_mcp/catalog.py index 98555712..4e1c70bb 100644 --- a/dw_mcp/catalog.py +++ b/dw_mcp/catalog.py @@ -236,7 +236,9 @@ def list_gallery(client, limit=50, subfolder=None, only_orphans=False, workspace the run it came from. The version is what to quote to a person - the web UI labels the same file `v5` - and is stable: it is assigned when the run opens and a deleted sibling leaves a gap rather than renumbering - what is left. Null under the flat output layout, which has no runs.""" + what is left - as does a run that failed, or reused every step from + the cache, and so wrote nothing to list. Null under the flat output + layout, which has no runs.""" params = {"limit": limit} if subfolder is not None: params["subfolder"] = subfolder diff --git a/tests/test_events.py b/tests/test_events.py index 8985d425..1a60baa8 100644 --- a/tests/test_events.py +++ b/tests/test_events.py @@ -486,3 +486,33 @@ def mock_load(self, shared_components): versions.append(manifest["version"]) assert versions == [1, 2, 3] + + +def test_the_version_is_on_disk_before_the_first_step_runs(tmp_path): + """A run killed mid-step - which is how a stuck server gets restarted - + never reaches the closing manifest, so the number has to land when the + run opens. Also what lets a second process opening a run of the same + workflow see this one's number rather than taking it too.""" + seen = {} + + def mock_load(self, shared_components): + manifest_path = pathlib.Path(workflow._run_dir) / "manifest.json" + seen.update(json.loads(manifest_path.read_text())) + self.pipeline = FakePipeline() + + workflow_def = _workflow_def() + workflow_def["steps"][0]["result"] = {"content_type": "image/png"} + workflow = Workflow(workflow_def, str(tmp_path), "test.json") + with patch.object(Pipeline, "load", mock_load): + with patch("dw.workflow.empty_device_cache"): + workflow.run({}, previous_pipelines={}) + + assert seen["version"] == 1 + assert seen["status"] == "running" + assert seen["finished_at"] is None + closing = json.loads( + (pathlib.Path(workflow._run_dir) / "manifest.json").read_text() + ) + assert closing["status"] == "completed" + assert closing["version"] == 1 + assert closing["finished_at"] is not None diff --git a/tests/test_runs.py b/tests/test_runs.py index 069b2538..8fe7f208 100644 --- a/tests/test_runs.py +++ b/tests/test_runs.py @@ -11,6 +11,7 @@ from dw.runs import ( FLAT_LAYOUT, assign_run_version, + record_run_versions, OUTPUT_LAYOUT_ENV_VAR, RUN_LAYOUT, is_run_id, @@ -837,3 +838,81 @@ def test_an_unreadable_manifest_does_not_lose_the_run(self, tmp_path): def test_a_workflow_with_no_runs_yet_has_none(self, tmp_path): assert run_versions(str(tmp_path / "never" / "ran")) == {} + + def test_a_run_after_an_out_of_order_second_takes_no_held_number(self, tmp_path): + # Two runs of one second sort opposite to their numbers (v6 then + # v5); a killed run after them must not be ranked back down onto 6 + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-00000000", version=6) + self._run(str(identity), "20260901-120000-ffffffff", version=5) + os.makedirs(str(identity / "20260901-120005-aaaaaaaa")) + assert run_versions(str(identity)) == { + "20260901-120000-00000000": 6, + "20260901-120000-ffffffff": 5, + "20260901-120005-aaaaaaaa": 7, + } + assert assign_run_version(str(tmp_path), "ltx2/Gyre") == 8 + + def test_a_rerun_counter_sorts_as_a_number(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + for run_id in ( + "20260901-120000-aaaaaaaa-10", + "20260901-120000-aaaaaaaa", + "20260901-120000-aaaaaaaa-2", + ): + self._run(str(identity), run_id) + assert list(run_versions(str(identity)).items()) == [ + ("20260901-120000-aaaaaaaa", 1), + ("20260901-120000-aaaaaaaa-2", 2), + ("20260901-120000-aaaaaaaa-10", 3), + ] + + def test_a_boolean_is_not_a_recorded_version(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=True) + self._run(str(identity), "20260902-120000-bbbbbbbb", version=5) + assert run_versions(str(identity))["20260901-120000-aaaaaaaa"] == 4 + + def test_an_edited_manifest_is_read_again(self, tmp_path): + # Recorded numbers are cached against the manifest's stat, so a + # rewrite - the run's closing manifest, or a backfill - is seen + identity = tmp_path / "ltx2" / "Gyre" + run_dir = self._run(str(identity), "20260901-120000-aaaaaaaa") + assert run_versions(str(identity)) == {"20260901-120000-aaaaaaaa": 1} + with open(os.path.join(run_dir, "manifest.json"), "w") as file: + json.dump({"version": 9, "padding": "changes the size"}, file) + assert run_versions(str(identity)) == {"20260901-120000-aaaaaaaa": 9} + + def test_opening_a_run_pins_the_numbers_of_older_runs(self, tmp_path): + # History from before the field is ranked, and a ranked number moves + # when an older sibling goes - until a new run writes it down + identity = tmp_path / "ltx2" / "Gyre" + for day in (1, 2, 3): + self._run(str(identity), f"2026090{day}-120000-aaaaaaaa") + assert assign_run_version(str(tmp_path), "ltx2/Gyre") == 4 + for day, version in ((1, 1), (2, 2), (3, 3)): + manifest_path = identity / f"2026090{day}-120000-aaaaaaaa" / "manifest.json" + manifest = json.loads(manifest_path.read_text()) + assert manifest["version"] == version + # the rest of the record is untouched + assert manifest["status"] == "completed" + # now deleting the oldest renumbers nothing + import shutil + + shutil.rmtree(str(identity / "20260901-120000-aaaaaaaa")) + assert run_versions(str(identity)) == { + "20260902-120000-aaaaaaaa": 2, + "20260903-120000-aaaaaaaa": 3, + } + + def test_pinning_leaves_a_run_with_no_manifest_alone(self, tmp_path): + # Writing one would invent a record of a run nobody recorded + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + killed = identity / "20260902-120000-bbbbbbbb" + os.makedirs(str(killed)) + assert record_run_versions(str(identity)) == { + "20260901-120000-aaaaaaaa": 1, + "20260902-120000-bbbbbbbb": 2, + } + assert not (killed / "manifest.json").exists() diff --git a/tests/test_server.py b/tests/test_server.py index 5dc08882..cdccac7c 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -2066,6 +2066,44 @@ def test_gallery_metadata_names_the_run_and_its_version(server, tmp_path): assert body["version"] == 3 +def test_deleting_an_older_run_renumbers_none_of_its_siblings(server, tmp_path): + """Runs from before versions existed are ranked, so removing the oldest + would slide every later one down a number. The delete pins the + siblings' numbers into their manifests first - both for a whole run + directory and for the last file of a run, which sweeps the directory.""" + import json as _json + + from PIL import Image + + identity = tmp_path / "outputs" / "acorn" / "cut" + run_ids = [f"2026090{day}-120000-aaaaaaaa" for day in (1, 2, 3, 4)] + with server(success_script) as client: + for run_id in run_ids: + (identity / run_id).mkdir(parents=True) + Image.new("RGB", (2, 2)).save(identity / run_id / "film.png") + (identity / run_id / "manifest.json").write_text( + _json.dumps({"run_id": run_id}) + ) + + def versions(): + return { + f["run_id"]: f["version"] + for f in client.get("/api/gallery").json()["files"] + } + + assert versions() == dict(zip(run_ids, (1, 2, 3, 4))) + # the whole run directory + assert client.delete(f"/api/gallery/acorn/cut/{run_ids[0]}").status_code == 200 + assert versions() == dict(zip(run_ids[1:], (2, 3, 4))) + # the last file of a run, which takes its directory with it + assert ( + client.delete(f"/api/gallery/acorn/cut/{run_ids[1]}/film.png").status_code + == 200 + ) + assert not (identity / run_ids[1]).exists() + assert versions() == dict(zip(run_ids[2:], (3, 4))) + + def test_gallery_only_orphans_lists_media_less_run_directories(server, tmp_path): """#170: a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything, has no diff --git a/ui/src/lib/pages/GalleryPage.svelte b/ui/src/lib/pages/GalleryPage.svelte index 3c6be3e3..32507caf 100644 --- a/ui/src/lib/pages/GalleryPage.svelte +++ b/ui/src/lib/pages/GalleryPage.svelte @@ -326,6 +326,10 @@ title="version {file.version} of this workflow" >v{file.version} + + {' '} {/if}{file.label} @@ -547,11 +551,10 @@ across lines */ .version { display: inline-block; - margin-right: 0.35rem; padding: 0 0.3rem; border-radius: 0.2rem; - background: var(--chip, rgb(255 255 255 / 0.08)); - color: var(--fg); + background: var(--line); + color: var(--ink); font-weight: 600; word-break: keep-all; } diff --git a/ui/src/lib/pages/GalleryPage.test.ts b/ui/src/lib/pages/GalleryPage.test.ts index 68a859d6..3d9d6339 100644 --- a/ui/src/lib/pages/GalleryPage.test.ts +++ b/ui/src/lib/pages/GalleryPage.test.ts @@ -1,12 +1,12 @@ import { cleanup, + fireEvent, render, screen, waitFor, within, } from '@testing-library/svelte' import { afterEach, beforeEach, expect, it, vi } from 'vitest' -import { fireEvent } from '@testing-library/dom' // Hoisted above the imports so the static import of the component below - // itself hoisted - sees an initialized mock. Importing the component inside // the test instead would charge its (multi-second) compile to the test timeout @@ -398,6 +398,10 @@ it('marks each file with the version of the run that wrote it', async () => { await waitFor(() => expect(screen.getAllByText('film.mp4')).toHaveLength(2)) expect(screen.getByText('v1')).toBeTruthy() expect(screen.getByText('v4')).toBeTruthy() + // One space between chip and label, so a screen reader does not run + // 'v4' into the file name + const caption = screen.getByText('v4').closest('.caption') as HTMLElement + expect(caption.textContent?.replace(/\s+/g, ' ').trim()).toBe('v4 film.mp4') }) it('shows no version for a flat-layout file, which belongs to no run', async () => { From 667642e4e720d9a5da429b8ce7d2e57fa7456b3d Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Tue, 22 Sep 2026 16:12:40 -0500 Subject: [PATCH 3/5] feat: add versioning to output references and gallery listings - Enhanced the `list_gallery` function to support filtering by `folder` and `version`, allowing users to specify a particular run's output. - Introduced a `version` parameter in various API endpoints and functions to track the ordinal of runs, displayed as `v` in the gallery. - Updated the `Job` model to include `run_version`, reflecting the version of the run associated with each job. - Modified the realization process to correctly handle output references that specify a version, ensuring that the correct run is referenced. - Adjusted the UI components to display the run version alongside job details, improving user clarity on which version of a workflow is being referenced. - Added tests to verify the correct behavior of versioned output references and ensure that the system behaves as expected when handling versions. --- CLAUDE.md | 15 +++++++--- docs/MCP.md | 2 +- docs/SERVER.md | 3 +- docs/WORKFLOW_GUIDE.md | 13 ++++++--- dw/realize.py | 18 ++++++++---- dw/runs.py | 48 ++++++++++++++++++++++++++++---- dw/server/app.py | 29 ++++++++++++++++++- dw/server/exports.py | 12 ++++++++ dw/server/jobs.py | 24 ++++++++++++---- dw/workflow.py | 1 + dw_mcp/catalog.py | 18 ++++++++++-- dw_mcp/server.py | 9 ++++-- tests/test_events.py | 3 ++ tests/test_mcp_catalog.py | 13 +++++++++ tests/test_realize.py | 15 ++++++++++ tests/test_runs.py | 41 +++++++++++++++++++++++++++ tests/test_server.py | 12 ++++++++ tests/test_server_exports.py | 9 ++++++ tests/test_server_jobs.py | 8 ++++++ ui/CLAUDE.md | 6 ++-- ui/src/lib/pages/JobPage.svelte | 16 +++++++++++ ui/src/lib/pages/JobPage.test.ts | 18 ++++++++++++ ui/src/lib/pages/JobsPage.svelte | 7 +++++ ui/src/lib/types.ts | 6 ++++ 24 files changed, 313 insertions(+), 33 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index f6849a83..06e89aa8 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -153,7 +153,9 @@ references) are documented above in *Workflow sources* and *Type System*. `"output:ltx2/Gyre/latest/still.png"`. The name is `//` under the output root, and `latest` in the run-id position picks the newest run that holds the file (run ids sort by their UTC timestamp; a failed or fully-cached run holds - only a manifest and is skipped). Resolved in `realize_args` beside `asset:` (`dw/runs.py`), + only a manifest and is skipped), and `v` there picks the run whose version is N + (below) - exactly that run, with no fallback to an older one. Either is a selector only + where run directories are, and the realized workflow pins both to the run id. Resolved in `realize_args` beside `asset:` (`dw/runs.py`), against the output root `Workflow.run` activates, and confined to it - A generated file becomes a stable input with `POST /api/assets/keep` (gallery "Keep as asset", MCP `keep_output`): it is hard-linked, else copied, from the workspace's outputs @@ -395,9 +397,14 @@ same reason - default setup cannot load a pack. not only a deletion: a failed run or a fully cached rerun takes a number and may have no media for the gallery to show under it. `GET /api/gallery` and the metadata route carry `version` and `run_id` - (`run_versions` read once per identity per listing, not per file), MCP - `list_gallery` teaches the vocabulary, and the web UI reads the field only - - a `v4` chip on the card, the run id in the detail pane. Nothing on disk is + (`run_versions` read once per identity per listing, not per file), and + `?folder=&version=` lists one run's files. The number is also a name: + `output:/v4/`. The `run_start` event carries it, the job + records it (`run_version`, a `jobs.sqlite` column) and the export README and + zip download name (`-v4-.zip`) carry it too. MCP + `list_gallery` teaches the vocabulary and takes `folder`/`version`, and the + web UI reads the field only - a `v4` chip on the gallery card, the jobs list + and the job page, the run id in the gallery's detail pane. Nothing on disk is renamed, so `output:` references, the step cache and `keep_output` are untouched. Two limits taken deliberately: deleting the *newest* run frees its number for reuse (the high-water mark lived in the manifest that went diff --git a/docs/MCP.md b/docs/MCP.md index f06956cb..b3ed1c69 100644 --- a/docs/MCP.md +++ b/docs/MCP.md @@ -226,7 +226,7 @@ when no single workflow covers it. | `get_health()` | — | Check that the server is alive, and which machine answered: `version`, `device`, whether a model process is currently resident (`worker_alive`), the job running now and the queue depth. `worker_alive: false` is the normal idle state on a server that has not run a job since startup or the last memory clear - not a fault - the on-demand worker starts with the next job (#206) | | `get_server_info()` | — | What this installation can do and where it keeps things: `device` (the accelerator a run will use), `version`, the `workspace` this session is working in and the workflow/asset/output/prompt `directories` of *that* workspace, the bind address and port, whether a token is required, and whether MCP is mounted. Check the device before authoring - a CUDA-only choice (bitsandbytes, `torch.compile`, flash attention) is not available on an `mps` or `cpu` server. `runtime` (#222) reports the Python version, torch version and the CUDA version torch was built against, the NVIDIA driver version (when `nvidia-smi` is reachable), and the installed versions of diffusers, transformers, accelerate, bitsandbytes, peft, safetensors and sentencepiece (`null` for one not installed) - for diagnosing an environment mismatch between boxes without shelling in | | `list_jobs(limit=20, status=None, workspace=None)` | optional `limit` (newest N), `status` (one state or a comma-separated set of `queued`, `running`, `succeeded`, `failed`, `cancelled`), `workspace` | List queued, running and recent jobs, **newest first**. Bounded by default: the unbounded listing was over a client's tool-result limit on a server with a few months of history, which made it a tool that could not be called at all. `total` says how many matched and `truncated`/`next` say so when the answer was cut - raise `limit` or narrow with `status`. Without `workspace`, a named workspace lists its own jobs and the default one lists every job the server holds | -| `list_gallery(limit=50, subfolder=None, only_orphans=False, workspace=None)` | `limit`, `subfolder`, `only_orphans`, `workspace` | List generated output files, newest first. A name is `//`, where `` may sit in the subfolder the step chose (`final/episode.mp4`); each entry carries `folder` (the workflow) and `subfolder` (by convention `final` or `intermediate`, `''` when the step chose none, any path the workflow wrote otherwise), and `subfolder="final"` lists only deliverables. Each entry carries `run_id` and `version` - that run's ordinal among the workflow's runs, which is how one of several runs that wrote the same basename is named to a person: the web UI labels the same file `v5`. The number is assigned when the run opens and never renumbered, so deleting a run leaves a gap rather than sliding the rest down (a failed run, or a rerun that reused every step, leaves one too - it took a number and may have nothing to list), and it is `null` under the flat output layout, which has no runs. Tools take `name`, never `version`. Each entry also carries a ready-made `url`, already scoped to the workspace that made it - a hand-built `/outputs/` URL 404s for anything but the default workspace. `only_orphans=True` inverts the call: instead of files, it returns run directories with no media anywhere under them (`runs`, each `{name, mtime}`) - a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything; `subfolder` does not apply in this mode, and `name` is exactly what `delete_output` accepts (#170). `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | +| `list_gallery(limit=50, subfolder=None, only_orphans=False, workspace=None, folder=None, version=None)` | `limit`, `subfolder`, `only_orphans`, `workspace`, `folder`, `version` | List generated output files, newest first. A name is `//`, where `` may sit in the subfolder the step chose (`final/episode.mp4`); each entry carries `folder` (the workflow) and `subfolder` (by convention `final` or `intermediate`, `''` when the step chose none, any path the workflow wrote otherwise), and `subfolder="final"` lists only deliverables. Each entry carries `run_id` and `version` - that run's ordinal among the workflow's runs, which is how one of several runs that wrote the same basename is named to a person: the web UI labels the same file `v5`. The number is assigned when the run opens and never renumbered, so deleting a run leaves a gap rather than sliding the rest down (a failed run, or a rerun that reused every step, leaves one too - it took a number and may have nothing to list), and it is `null` under the flat output layout, which has no runs. `folder` with `version` lists that one run's files, and `output:/v5/` names one in a workflow; every other tool takes `name`. Each entry also carries a ready-made `url`, already scoped to the workspace that made it - a hand-built `/outputs/` URL 404s for anything but the default workspace. `only_orphans=True` inverts the call: instead of files, it returns run directories with no media anywhere under them (`runs`, each `{name, mtime}`) - a run whose output was deleted before `delete_output` could remove it by name, or one that failed before writing anything; `subfolder` does not apply in this mode, and `name` is exactly what `delete_output` accepts (#170). `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | | `get_gallery_metadata(name, envelope=False, workspace=None)` | `name`, `workspace` | Get the metadata embedded in a generated file — or, when `name` is an `asset:` reference, what an *input* asset holds (`source` says which; `job` is null for an asset). Reading an input's duration, frame count, fps and sample rate before a run is how a caller learns the `total_frames`, `fps` and `sample_rate` a workflow expects it to supply: the exact workflow and arguments that produced it, and, for audio/video, a `media` block (duration, rate, channels, fps, size, peak/mean dBFS). `envelope=true` adds `media.envelope` — `rms_dbfs` and `peak_dbfs` one entry per second — which is what locates something in a track rather than measuring the whole of it. `workspace` names the workspace for this one call without switching the session to it - the same pin `run_workflow` takes, so a job run into another workspace stays reachable from the session that queued it | ### Media diff --git a/docs/SERVER.md b/docs/SERVER.md index f357d87c..30afe061 100644 --- a/docs/SERVER.md +++ b/docs/SERVER.md @@ -441,7 +441,8 @@ The editor's forms come from these; they are just as usable from scripts: gallery entry carries `folder` (the workflow identity, the run id dropped) and `subfolder` (what followed the run id - the `final`/`intermediate` a step's `result.subfolder` chose, `''` when it chose none); `?folder=` and - `?subfolder=` filter independently, and the reply's `folders` and + `?subfolder=` filter independently (`?version=` too - with `?folder=`, + the one run the gallery labels `v4`), and the reply's `folders` and `subfolders` list every distinct value over the whole tree, `''` always a member of each so root-level files stay selectable - `GET /api/gallery/{name:path}/download` — download an output file diff --git a/docs/WORKFLOW_GUIDE.md b/docs/WORKFLOW_GUIDE.md index 53fbeac6..895b34cd 100644 --- a/docs/WORKFLOW_GUIDE.md +++ b/docs/WORKFLOW_GUIDE.md @@ -289,7 +289,8 @@ for existence. to the same file. - `output:` — `output://` is a file an earlier run wrote, under the output root and confined to it. `latest` in the run-id position - picks the newest run that holds that file. A run id is not stable against + picks the newest run that holds that file; `v` picks the run the gallery labels + `v` (`list_gallery`'s `version`), and only that run. A run id is not stable against pruning: to depend on a generated file, promote it with `keep_output` and reference the `asset:` name instead. - `prompt:` — `prompt:name` or `prompt:folder/name` is a stored prompt's @@ -1459,7 +1460,7 @@ Beside that manifest the run also writes `workflow.json` — the *realized* workflow, meaning the one that actually ran. Every mutable input is pinned into it: the caller's `arguments` folded into the `variables` defaults, the seed the run used, each `prompt:` reference replaced by the stored text, and each -`output:/latest/` rewritten to the run id it resolved to. +`output:/latest/` (or `/v/`) rewritten to the run id it resolved to. `asset:`, `constant:`, `previous_result:` and `builtin:` are kept as written — each already names something pinned by the asset library or by the manifest's `dw_version` — and a sub-workflow named by local path is kept with its file's @@ -1657,8 +1658,12 @@ second-stage workflow name the first stage's product without being edited after run - and keeps working when the newest run failed part way, or reused every step from the cache and so wrote nothing of its own but a manifest. Runs sort by their id, which starts with a UTC timestamp, so "newest" needs no file timestamps and survives a -directory being copied. `latest` only selects a run where run directories are; a -workflow or file that happens to be called `latest` is still named as itself. +directory being copied. `v` in the same position names the run whose version is N - +the `v4` the gallery labels its files with - so the number a person was told is a name +a workflow can take. Unlike `latest` it picks exactly one run: `v4` not holding the file +is an error, not a reason to try `v3`. `latest` and `v` only select a run where run +directories are; a workflow or file that happens to be called either is still named as +itself. Like `asset:`, a reference resolves to a path and then whatever loads paths loads it, so it works under `image`, `video`, a `from_file`, or a list of them. The audio tasks take diff --git a/dw/realize.py b/dw/realize.py index c970aad2..f00f402d 100644 --- a/dw/realize.py +++ b/dw/realize.py @@ -29,6 +29,7 @@ is_output_reference, output_root as default_output_root, resolve_output_reference, + version_selector, ) from .security import SecurityError, validate_workflow_path from .workflow_sources import resolve_sub_workflow, SubWorkflowNotFound @@ -170,14 +171,19 @@ def _inline_prompt(reference, annotations, prompt_dir, base_dir): def _pin_output(reference, output_root): - """'output:/latest/' rewritten to the run it resolved to. - - An explicit run id is already pinned, so it is returned untouched without - touching the disk - realizing must not fail on a reference the run has - not reached yet. + """'output:/latest/' - or '/v4/' - rewritten to the run it + resolved to. + + A version is stable, but deleting the newest run frees its number for + reuse, so the realized copy names the run id either way. An explicit run + id is already pinned, so it is returned untouched without touching the + disk - realizing must not fail on a reference the run has not reached + yet. """ name = reference.removeprefix(OUTPUT_PREFIX).strip() - if LATEST not in name.split("/"): + if not any( + part == LATEST or version_selector(part) is not None for part in name.split("/") + ): return reference root = output_root or default_output_root() try: diff --git a/dw/runs.py b/dw/runs.py index 797369ca..919a3437 100644 --- a/dw/runs.py +++ b/dw/runs.py @@ -57,6 +57,17 @@ # newest run directory LATEST = "latest" +# 'v4' in the run-id position of an 'output:' reference: the run whose +# ordinal is 4 - the number the gallery shows and an agent quotes +_VERSION_SELECTOR = re.compile(r"^v([1-9][0-9]*)$") + + +def version_selector(segment): + """The ordinal a 'v' segment names, or None for any other segment.""" + match = _VERSION_SELECTOR.match(segment) + return int(match.group(1)) if match else None + + # What a run id looks like: a UTC timestamp and a short digest of the spec. # The pattern is not only documentation - the gallery reads it to group a # workflow's runs under one folder rather than listing every run separately @@ -126,8 +137,8 @@ def _runs_newest_first(directory): def _resolve_segments(directory, parts, reference, root): - """Build the path a name stands for, expanding 'latest' where it names - a run. + """Build the path a name stands for, expanding 'latest' or 'v' where + it names a run. 'latest' means the newest run *that has the file*, not the newest run directory: a run that failed part way, or one whose every step was a @@ -136,9 +147,14 @@ def _resolve_segments(directory, parts, reference, root): stage before it plainly produced something. So the runs are tried newest first and the first one holding the rest of the name wins. + 'v' means the run whose recorded ordinal is N - the 'v4' the gallery + shows - so the number quoted to a person is also a name a workflow can + take. Unlike 'latest' it picks exactly one run: a v4 that did not write + the file is an error, not a reason to try v3. + Only a segment standing where run directories are is a run selector. A - 'latest' segment in a directory that holds no runs is a name like any - other, so a workflow or a file called 'latest' stays reachable. + 'latest' or 'v4' segment in a directory that holds no runs is a name + like any other, so a workflow or a file called either stays reachable. Returns the path, or None when runs were found and none of them holds the file. @@ -162,6 +178,27 @@ def _resolve_segments(directory, parts, reference, root): f"'{reference}' names the newest run of a workflow that " f"has not produced one" ) + wanted = version_selector(part) + if wanted is not None and _runs_newest_first(directory): + versions = run_versions(directory) + matching = [run for run, version in versions.items() if version == wanted] + if not matching: + held = ", ".join(f"v{v}" for v in sorted(set(versions.values()))) + raise ValueError( + f"No run v{wanted} under {os.path.relpath(directory, root)} - " + f"'{reference}' names a run by its version, and the runs " + f"there are {held}" + ) + # Normally one; two only where history predating versions could + # not be ranked beneath the first recorded number. Newest first, + # as 'latest' would try them + for run in sorted(matching, key=run_id_sort_key, reverse=True): + candidate = _resolve_segments( + os.path.join(directory, run), rest, reference, root + ) + if candidate and os.path.isfile(candidate): + return candidate + return None return _resolve_segments(os.path.join(directory, part), rest, reference, root) @@ -170,7 +207,8 @@ def resolve_output_reference(reference, root=None): The name is a path under the output directory - '//' - and the run id may be written as 'latest', which resolves - to the newest run of that workflow that holds the file. That is what + to the newest run of that workflow that holds the file, or as 'v', + the run whose version is N. 'latest' is what lets a second-stage workflow name the first stage's product without being edited after every run, and without breaking when the newest run failed or reused cached files and so wrote none of its own. diff --git a/dw/server/app.py b/dw/server/app.py index 28ee5929..cb98172f 100644 --- a/dw/server/app.py +++ b/dw/server/app.py @@ -14,6 +14,7 @@ import tempfile import copy import json +import re import uuid import asyncio import logging @@ -2824,6 +2825,7 @@ def gallery( folder: Optional[str] = None, subfolder: Optional[str] = None, only_orphans: bool = False, + version: Optional[int] = None, ws: Workspace = Depends(selected_workspace), ): """A page of media files in the output directory, newest first. @@ -2838,6 +2840,8 @@ def gallery( way: the in-run subfolders steps wrote into ('final', 'intermediate'), '' for files at a run's root. `folder` and `subfolder` filter independently and intersect when both are given. + `version` narrows to the runs holding that ordinal - with `folder`, + the one run "v4" names; without it, that run of every workflow. `only_orphans=true` inverts the whole call: instead of media files, it returns run directories holding nothing but their own @@ -2868,6 +2872,8 @@ def gallery( entries = [e for e in entries if e["folder"] == folder] if subfolder is not None: entries = [e for e in entries if e["subfolder"] == subfolder] + if version is not None: + entries = [e for e in entries if e["version"] == version] offset = max(0, offset) limit = max(0, limit) page = entries[offset : offset + limit] @@ -4165,6 +4171,27 @@ async def input_file( files = _static_files_for(roots[0]) return await files.get_response(name, request.scope) + def _export_download_name(directory, job_id): + """'-v4-.zip' when the exported manifest says which + run it was, else '.zip'. Only the saved file's name: the + URL and the entries inside keep the job id, so nothing that already + names an export changes.""" + try: + with open(os.path.join(directory, MANIFEST_FILE_NAME)) as file: + manifest = json.load(file) + except (OSError, ValueError): + return f"{job_id}.zip" + if not isinstance(manifest, dict): + return f"{job_id}.zip" + version = manifest.get("version") + identity = (manifest.get("workflow") or {}).get("identity") + if not isinstance(version, int) or isinstance(version, bool): + return f"{job_id}.zip" + if not isinstance(identity, str) or not identity: + return f"v{version}-{job_id}.zip" + slug = re.sub(r"[^A-Za-z0-9_.-]+", "-", identity).strip("-.") + return f"{slug}-v{version}-{job_id}.zip" if slug else f"v{version}-{job_id}.zip" + # Ungated for the same reason the two above are: a download link cannot # attach an Authorization header either @app.get("/exports/{job_id}.zip") @@ -4186,7 +4213,7 @@ def export_zip(job_id: str, ws: Workspace = Depends(selected_workspace)): path = os.path.join(current, name) entry = os.path.relpath(path, directory).replace(os.sep, "/") entries.append((f"{job_id}/{entry}", path)) - return _zip_download(entries, f"{job_id}.zip") + return _zip_download(entries, _export_download_name(directory, job_id)) # ---------------------------------------------------------------- the UI diff --git a/dw/server/exports.py b/dw/server/exports.py index 7953f7a9..bf573a45 100644 --- a/dw/server/exports.py +++ b/dw/server/exports.py @@ -77,6 +77,7 @@ "error", "run_id", "run_dir", + "run_version", ) README_TEMPLATE = """# {workflow_name} - job {job_id} @@ -90,6 +91,7 @@ | Job | `{job_id}` | | Workflow | `{workflow_name}` | | Catalog entry | {catalog_name} | +| Run | {run} | | Status | {status} | | Started | {started_at} | | Finished | {finished_at} | @@ -394,6 +396,16 @@ def _readme(job_id, detail, manifest, workflow, summary, realized): "predates run tracking, so its arguments and prompts are not " "pinned into it." ), + run=( + f"`{detail['run_id']}`" + + ( + f" - version {manifest['version']}" + if isinstance(manifest.get("version"), int) + else "" + ) + if detail.get("run_id") + else "not recorded" + ), status=detail.get("status"), started_at=detail.get("started_at"), finished_at=detail.get("finished_at"), diff --git a/dw/server/jobs.py b/dw/server/jobs.py index c01b6950..69fee683 100644 --- a/dw/server/jobs.py +++ b/dw/server/jobs.py @@ -134,6 +134,11 @@ def __init__(self, db_path): connection.execute("ALTER TABLE jobs ADD COLUMN run_id TEXT") if "run_dir" not in columns: connection.execute("ALTER TABLE jobs ADD COLUMN run_dir TEXT") + # That run's ordinal among the workflow's runs - the 'v4' the + # gallery shows. NULL before the column, and for a job that + # never opened a run + if "run_version" not in columns: + connection.execute("ALTER TABLE jobs ADD COLUMN run_version INTEGER") # Which form of cost acknowledgement queued the job. Rows before # the column are 'none' - nothing recorded is nothing recorded if "acknowledged" not in columns: @@ -178,8 +183,8 @@ def record(self, job): " started_at, finished_at, arguments, spec, manifest, warnings," " error, events, workspace, workflow_name, run_id, run_dir," " acknowledged, host_memory_peak_rss_mb," - " host_memory_job_peak_rss_mb) VALUES" - " (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", + " host_memory_job_peak_rss_mb, run_version) VALUES" + " (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", ( job.id, job.workflow_name, @@ -204,6 +209,7 @@ def record(self, job): # itself allows getattr(job, "host_memory_peak_rss_mb", None), getattr(job, "host_memory_job_peak_rss_mb", None), + getattr(job, "run_version", None), ), ) @@ -219,7 +225,8 @@ def recent_summaries(self, limit=200, workspace=None, statuses=None): """ query = ( "SELECT id, workflow, status, created_at, started_at, finished_at," - " workspace, workflow_name, run_id, acknowledged FROM jobs" + " workspace, workflow_name, run_id, acknowledged, run_version" + " FROM jobs" ) params = [] clauses = [] @@ -249,6 +256,7 @@ def recent_summaries(self, limit=200, workspace=None, statuses=None): "workflow_name": row[7], "run_id": row[8], "acknowledged": row[9] or ACK_NONE, + "run_version": row[10], "historical": True, } for row in rows @@ -259,8 +267,8 @@ def get(self, job_id): row = connection.execute( "SELECT id, workflow, status, created_at, started_at, finished_at," " arguments, spec, manifest, warnings, error, workspace," - " workflow_name, run_id, run_dir, acknowledged, events FROM jobs" - " WHERE id = ?", + " workflow_name, run_id, run_dir, acknowledged, events," + " run_version FROM jobs WHERE id = ?", (job_id,), ).fetchone() return self._to_detail(row) if row else None @@ -479,6 +487,7 @@ def parse(text, fallback): "workflow_name": row[12], "run_id": row[13], "run_dir": row[14], + "run_version": row[17], "acknowledged": row[15] or ACK_NONE, "acknowledged_cost": (spec or {}).get("acknowledged_cost"), "traceback": None, @@ -510,6 +519,7 @@ def __init__(self, spec): # never got that far self.run_id = None self.run_dir = None + self.run_version = None # Which form of cost acknowledgement queued this job (#85) self.acknowledged = spec.get("acknowledged") or ACK_NONE # The worker's own high-water mark for this run, from its final @@ -685,6 +695,9 @@ def summary(self): # so it defaults the same way history's column does "workspace": self.spec.get("workspace") or DEFAULT_WORKSPACE_NAME, "run_id": self.run_id, + # The run's ordinal - 'v4' - so the job that just ran can be + # named the way the gallery will name it + "run_version": self.run_version, "acknowledged": self.acknowledged, } @@ -1303,6 +1316,7 @@ def _consume_results(self, job): if event.get("event") == "run_start": job.run_id = event.get("run_id") job.run_dir = event.get("run_dir") + job.run_version = event.get("version") job.add_event(event) elif message_type in ("output", "workflow_loaded"): text = message.get("message") or message.get("workflow_name", "") diff --git a/dw/workflow.py b/dw/workflow.py index 0845237e..e65e173b 100644 --- a/dw/workflow.py +++ b/dw/workflow.py @@ -1172,6 +1172,7 @@ def run( run_context.emit( "run_start", run_id=run_id, + version=self._run_version, identity=workflow_identity(self.file_spec, workflow_id), run_dir=os.path.relpath(self._run_dir, self.output_dir).replace( os.sep, "/" diff --git a/dw_mcp/catalog.py b/dw_mcp/catalog.py index 4e1c70bb..669fd502 100644 --- a/dw_mcp/catalog.py +++ b/dw_mcp/catalog.py @@ -213,10 +213,20 @@ def list_jobs(client, limit=20, status=None, workspace=None): return answer -def list_gallery(client, limit=50, subfolder=None, only_orphans=False, workspace=None): +def list_gallery( + client, + limit=50, + subfolder=None, + only_orphans=False, + workspace=None, + folder=None, + version=None, +): """Generated media in the output directory, newest first. `subfolder` narrows to one in-run subfolder ('final', 'intermediate', '' for files - at a run's root); None means every file. + at a run's root); None means every file. `folder` narrows to one + workflow and `version` to one run's ordinal, so the two together list + the run a person calls "v4". `only_orphans=True` inverts the call: instead of files, it returns run directories holding nothing but their own bookkeeping (manifest.json, @@ -242,6 +252,10 @@ def list_gallery(client, limit=50, subfolder=None, only_orphans=False, workspace params = {"limit": limit} if subfolder is not None: params["subfolder"] = subfolder + if folder is not None: + params["folder"] = folder + if version is not None: + params["version"] = version if only_orphans: params["only_orphans"] = "true" return client.get_json("/api/gallery", params=params, workspace=workspace) diff --git a/dw_mcp/server.py b/dw_mcp/server.py index 44e02089..4272c9ba 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -395,6 +395,8 @@ def list_gallery( subfolder: str | None = None, only_orphans: bool = False, workspace: str | None = None, + folder: str | None = None, + version: int | None = None, ) -> dict: """List generated output files, newest first. A name is //, where may itself sit in a @@ -412,8 +414,9 @@ def list_gallery( Entries also carry `run_id` and `version`, that run's ordinal among the workflow's runs - stable, never renumbered. Quote the version - to a person: the web UI labels the same file `v5`. Tools still take - `name`. + to a person: the web UI labels the same file `v5`. `folder=` with + `version=` lists that one run; "output:/v5/" names it + in a workflow. Other tools still take `name`. `only_orphans=True` inverts the call: instead of files, it returns run directories holding nothing but their own bookkeeping @@ -438,6 +441,8 @@ def list_gallery( subfolder=subfolder, only_orphans=only_orphans, workspace=workspace, + folder=folder, + version=version, ) def get_gallery_metadata( diff --git a/tests/test_events.py b/tests/test_events.py index 1a60baa8..b97f7e3e 100644 --- a/tests/test_events.py +++ b/tests/test_events.py @@ -89,6 +89,9 @@ def test_progress_event_sequence(): names = [event["event"] for event in events] assert names[0] == "run_start" + # the run's ordinal, so a job can name its run the way the gallery will + # (the output root here is shared, so only its shape is fixed) + assert isinstance(events[0]["version"], int) and events[0]["version"] >= 1 assert names[1] == "workflow_start" assert names[-1] == "workflow_end" assert "step_start" in names and "step_end" in names diff --git a/tests/test_mcp_catalog.py b/tests/test_mcp_catalog.py index 2d28cb48..b7b698cb 100644 --- a/tests/test_mcp_catalog.py +++ b/tests/test_mcp_catalog.py @@ -117,6 +117,19 @@ def test_list_gallery_sends_a_subfolder_only_when_given(): assert seen["params"]["subfolder"] == "" +def test_list_gallery_sends_folder_and_version_only_when_given(): + client, seen = recording_client() + catalog.list_gallery(client, limit=7) + assert "folder" not in seen["params"] + assert "version" not in seen["params"] + + # together they name one run - what a person calls "v4" + client, seen = recording_client() + catalog.list_gallery(client, folder="acorn/cut", version=4) + assert seen["params"]["folder"] == "acorn/cut" + assert seen["params"]["version"] == "4" + + def test_list_gallery_sends_only_orphans_only_when_true(): client, seen = recording_client() catalog.list_gallery(client, limit=7) diff --git a/tests/test_realize.py b/tests/test_realize.py index cec60148..656399b6 100644 --- a/tests/test_realize.py +++ b/tests/test_realize.py @@ -136,6 +136,21 @@ def test_latest_is_pinned_to_the_run_it_resolved_to(self, output_root): f"output:ltx2/Gyre/{run_id}/still.png" ) + def test_a_version_is_pinned_to_the_run_it_named(self, output_root): + # A version is stable, but deleting the newest run frees its number, + # so the realized copy names the run id as it does for 'latest' + root, run_id = output_root + source = definition() + source["steps"][0]["pipeline"]["arguments"]["image"] = ( + "output:ltx2/Gyre/v1/still.png" + ) + + realized, _ = realize_workflow(source, {}, 7, output_root=root) + + assert realized["steps"][0]["pipeline"]["arguments"]["image"] == ( + f"output:ltx2/Gyre/{run_id}/still.png" + ) + def test_an_explicit_run_id_is_kept_as_written(self, output_root): root, run_id = output_root written = f"output:ltx2/Gyre/{run_id}/still.png" diff --git a/tests/test_runs.py b/tests/test_runs.py index 8fe7f208..360355a7 100644 --- a/tests/test_runs.py +++ b/tests/test_runs.py @@ -905,6 +905,47 @@ def test_opening_a_run_pins_the_numbers_of_older_runs(self, tmp_path): "20260903-120000-aaaaaaaa": 3, } + def test_an_output_reference_can_name_a_run_by_its_version(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + first = self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + self._run(str(identity), "20260903-120000-cccccccc", version=3) + with open(os.path.join(first, "still.png"), "wb") as file: + file.write(b"png") + + resolved = resolve_output_reference( + "output:ltx2/Gyre/v1/still.png", str(tmp_path) + ) + assert resolved == os.path.realpath(os.path.join(first, "still.png")) + + def test_a_version_that_did_not_write_the_file_does_not_fall_back(self, tmp_path): + # Unlike 'latest', 'v3' picks one run: v3 lacking the file is an + # error, not a reason to hand back v1's + identity = tmp_path / "ltx2" / "Gyre" + first = self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + self._run(str(identity), "20260903-120000-cccccccc", version=3) + with open(os.path.join(first, "still.png"), "wb") as file: + file.write(b"png") + + with pytest.raises(ValueError, match="not found"): + resolve_output_reference("output:ltx2/Gyre/v3/still.png", str(tmp_path)) + + def test_a_missing_version_names_the_ones_there_are(self, tmp_path): + identity = tmp_path / "ltx2" / "Gyre" + self._run(str(identity), "20260901-120000-aaaaaaaa", version=1) + self._run(str(identity), "20260903-120000-cccccccc", version=3) + + with pytest.raises(ValueError, match="No run v2 .* v1, v3"): + resolve_output_reference("output:ltx2/Gyre/v2/still.png", str(tmp_path)) + + def test_a_v_segment_where_no_runs_are_is_an_ordinary_name(self, tmp_path): + # A workflow folder called 'v2' stays reachable, as 'latest' does + target = tmp_path / "flat" / "v2" + target.mkdir(parents=True) + (target / "still.png").write_bytes(b"png") + + resolved = resolve_output_reference("output:flat/v2/still.png", str(tmp_path)) + assert resolved == os.path.realpath(str(target / "still.png")) + def test_pinning_leaves_a_run_with_no_manifest_alone(self, tmp_path): # Writing one would invent a record of a run nobody recorded identity = tmp_path / "ltx2" / "Gyre" diff --git a/tests/test_server.py b/tests/test_server.py index cdccac7c..b808d253 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -2043,6 +2043,18 @@ def _run(identity, run_id, version=None): assert by_name["ltx/flat.png"]["version"] is None assert by_name["ltx/flat.png"]["run_id"] == "" + # "show me v3": folder and version together list exactly that run + names = [ + f["name"] + for f in client.get( + "/api/gallery", params={"folder": "acorn/cut", "version": 3} + ).json()["files"] + ] + assert names == ["acorn/cut/20260903-120000-cccccccc/final/film.7-0.0.png"] + # version alone spans workflows: each one's v1 + v1 = client.get("/api/gallery", params={"version": 1}).json()["files"] + assert {f["folder"] for f in v1} == {"acorn/cut", "acorn/score"} + def test_gallery_metadata_names_the_run_and_its_version(server, tmp_path): """After "look at version 3", the next call is usually this one - so it diff --git a/tests/test_server_exports.py b/tests/test_server_exports.py index 95163a52..35f11140 100644 --- a/tests/test_server_exports.py +++ b/tests/test_server_exports.py @@ -56,8 +56,10 @@ def exporting_script(command): json.dump( { "run_id": RUN_ID, + "version": 4, "status": "completed", "seed": 7, + "workflow": {"identity": "server_test"}, "steps": [{"step": "gen", "files": ["still.png"]}], }, file, @@ -66,6 +68,7 @@ def exporting_script(command): "type": "progress", "event": "run_start", "run_id": RUN_ID, + "version": 4, "identity": "server_test", "run_dir": RUN_DIR, } @@ -282,6 +285,8 @@ def test_the_readme_names_the_job_and_says_how_to_run_it(self, server): assert job_id in readme assert "python -m dw.run workflow.json" in readme assert "Git LFS" in readme + # which run, in the form the gallery labels it + assert f"`{RUN_ID}` - version 4" in readme def test_the_job_s_own_asset_dir_is_used_not_the_export_s_workspace( self, server, workspace_root @@ -376,6 +381,10 @@ def test_it_lists_the_same_entries_as_the_directory(self, server): response = client.get(f"/exports/{job_id}.zip") assert response.status_code == 200 + # the saved file says which run it is; the URL and the entries + # inside keep the job id + disposition = response.headers["content-disposition"] + assert f"server_test-v4-{job_id}.zip" in disposition archive = zipfile.ZipFile(io.BytesIO(response.content)) assert sorted(archive.namelist()) == sorted( f"{job_id}/{entry['path']}" for entry in body["files"] diff --git a/tests/test_server_jobs.py b/tests/test_server_jobs.py index d5f4d2ef..0ab4347d 100644 --- a/tests/test_server_jobs.py +++ b/tests/test_server_jobs.py @@ -28,6 +28,7 @@ def tracked_script(command): "run_id": RUN_ID, "identity": "server_test", "run_dir": RUN_DIR, + "version": 4, } yield {"type": "success", "message": "ok", "run_count": 1, "manifest": []} @@ -59,6 +60,9 @@ def test_run_start_populates_the_job(manager): assert job.run_dir == RUN_DIR assert job.summary()["run_id"] == RUN_ID assert job.detail()["run_dir"] == RUN_DIR + # the ordinal the gallery shows for this run's files + assert job.run_version == 4 + assert job.summary()["run_version"] == 4 def test_both_persist_and_read_back(manager): @@ -66,6 +70,10 @@ def test_both_persist_and_read_back(manager): historical = manager.history.get(job.id) assert historical["run_id"] == RUN_ID assert historical["run_dir"] == RUN_DIR + assert historical["run_version"] == 4 + # and in the polled list, not only the detail + (summary,) = manager.history.recent_summaries() + assert summary["run_version"] == 4 def test_realized_reads_the_file_the_run_wrote(manager, tmp_path): diff --git a/ui/CLAUDE.md b/ui/CLAUDE.md index 979067f3..47005503 100644 --- a/ui/CLAUDE.md +++ b/ui/CLAUDE.md @@ -65,8 +65,10 @@ Stripping the run id is also what makes four runs of one workflow four identical captions, since a file's name is per step rather than per run. The gallery entry carries `version` - the run's ordinal, assigned by the engine and never renumbered - and the grid draws it as a `v4` chip ahead -of the label, with the run id in the detail pane beside it. The UI reads -the field only; nothing here computes or orders a version. +of the label, with the run id in the detail pane beside it. A job carries +the same number as `run_version`, drawn as `v4` in the jobs list and on the +job page (from the `run_start` event while the job runs). The UI reads the +field only; nothing here computes or orders a version. Every picture in the app sits in the global `.frame` (app.css): the media fills it edge to edge, with no inner padding and no rounding of its own, diff --git a/ui/src/lib/pages/JobPage.svelte b/ui/src/lib/pages/JobPage.svelte index f8a7c07d..2cf92e1f 100644 --- a/ui/src/lib/pages/JobPage.svelte +++ b/ui/src/lib/pages/JobPage.svelte @@ -215,6 +215,13 @@ events.find((e) => e.event === 'workflow_start')?.seed as number | undefined, ) + // The record's number once the job has one; while it runs, the run_start + // event says it first - so the page names the run the moment it opens + const runVersion = $derived( + job?.run_version ?? + (events.find((e) => e.event === 'run_start')?.version as + number | undefined), + ) const etaSeconds = $derived.by(() => { if (!denoise?.total_steps || stepTimes.length < 3) return null const window = stepTimes.slice(-6) @@ -345,6 +352,15 @@ >seed {seed} {/if} + {#if runVersion} + + v{runVersion} + {/if} {#if job.acknowledged === 'bound'} diff --git a/ui/src/lib/pages/JobPage.test.ts b/ui/src/lib/pages/JobPage.test.ts index 6e0f90c7..94fb8e4d 100644 --- a/ui/src/lib/pages/JobPage.test.ts +++ b/ui/src/lib/pages/JobPage.test.ts @@ -407,3 +407,21 @@ it("corrects the URL to the job's own workspace", async () => { render(JobPage, { jobId: 'j1' }) await waitFor(() => expect(location.hash).toBe('#/ws/studio/jobs/j1')) }) + +it('names the run by the version the gallery labels its files with', async () => { + detail.job = { + ...job([]), + run_id: '20260922-120000-aaaaaaaa', + run_version: 4, + } + render(JobPage, { jobId: 'j1' }) + const chip = await waitFor(() => screen.getByText('v4')) + expect(chip.getAttribute('title')).toContain('20260922-120000-aaaaaaaa') +}) + +it('shows no version for a job that never opened a run', async () => { + detail.job = { ...job([]), run_version: null } + render(JobPage, { jobId: 'j1' }) + await waitFor(() => expect(screen.getByText('j1')).toBeTruthy()) + expect(screen.queryByText(/^v\d+$/)).toBeNull() +}) diff --git a/ui/src/lib/pages/JobsPage.svelte b/ui/src/lib/pages/JobsPage.svelte index 3f494da2..9c94b050 100644 --- a/ui/src/lib/pages/JobsPage.svelte +++ b/ui/src/lib/pages/JobsPage.svelte @@ -131,6 +131,13 @@ {job.status} {job.workflow} + {#if job.run_version} + v{job.run_version} + {/if} {#if scope === 'all' && (workspace.names?.length ?? 0) > 1} {job.workspace} {/if} diff --git a/ui/src/lib/types.ts b/ui/src/lib/types.ts index 357fcbaf..a1e586cf 100644 --- a/ui/src/lib/types.ts +++ b/ui/src/lib/types.ts @@ -13,6 +13,12 @@ export interface JobSummary { * and any caller that sent nothing), a bare boolean, or one bound to the * plan a validate answered with. Absent on rows from older servers. */ acknowledged?: 'none' | 'boolean' | 'bound' + /** The run this job opened - null until it opens one, and for a job + * recorded before runs were tracked. */ + run_id?: string | null + /** That run's ordinal among the workflow's runs - the `v4` the gallery + * shows for its files. Null until the run opens, and for older rows. */ + run_version?: number | null } export interface ManifestEntry { From 7132590c1bb578561f057dc03705f8258927b4cf Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Tue, 22 Sep 2026 16:17:03 -0500 Subject: [PATCH 4/5] Pay for list_gallery's new parameters; find npm on a non-interactive PATH The folder/version parameters pushed the MCP surface 69 tokens over budget. The docstring now costs less than it did before the change (descriptions 9_078 -> 9_068); the two schema entries (+49) are taken by the budget deliberately, 13_850 -> 13_890, since no docstring can pay for a schema. deploy.sh ran `npm run build` on whatever PATH the caller's shell gave it. `ssh lem deploy.sh` is non-interactive, and lem's ~/.bashrc puts ~/.local/node/bin on PATH only after its interactive-only guard, so npm that works at a prompt was missing. find_npm now looks in DW_NODE_BIN and the usual per-user install locations (and nvm) before the UI build, and fails naming DW_NODE_BIN rather than with a bare command-not-found. Co-Authored-By: Claude Opus 5.5 (1M context) --- dw_mcp/server.py | 16 +++++++--------- scripts/deploy.sh | 30 +++++++++++++++++++++++++++++- tests/test_mcp_server.py | 8 +++++++- 3 files changed, 43 insertions(+), 11 deletions(-) diff --git a/dw_mcp/server.py b/dw_mcp/server.py index 4272c9ba..093e7bb0 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -408,15 +408,13 @@ def list_gallery( by convention `final` is the deliverable and `intermediate` the scratch work, '' when the step chose none); `subfolder=` filters on the latter, so `subfolder="final"` is "what did these runs - deliver". Each entry also carries a ready-made `url` for viewing the - file over HTTP, already scoped to the right workspace; use it as - given rather than composing one from the name. - - Entries also carry `run_id` and `version`, that run's ordinal among - the workflow's runs - stable, never renumbered. Quote the version - to a person: the web UI labels the same file `v5`. `folder=` with - `version=` lists that one run; "output:/v5/" names it - in a workflow. Other tools still take `name`. + deliver". Each entry's `url` is already scoped to its workspace; + use it as given rather than composing one from the name. + + Entries also carry `run_id` and `version`, the run's stable ordinal + (the web UI shows `v5`) - quote the version to a person. `folder=` + plus `version=` lists that run; "output:/v5/" names + it. Other tools take `name`. `only_orphans=True` inverts the call: instead of files, it returns run directories holding nothing but their own bookkeeping diff --git a/scripts/deploy.sh b/scripts/deploy.sh index e90ac842..289cf005 100755 --- a/scripts/deploy.sh +++ b/scripts/deploy.sh @@ -30,7 +30,9 @@ # # Environment overrides, all optional: # DW_DIR (checkout, default ~/diffusers-workflow), DW_TOKEN (default xyz), -# DW_PORT (8765), DW_WORKSPACE (~/diffusers-workspace), DW_HOST (0.0.0.0). +# DW_PORT (8765), DW_WORKSPACE (~/diffusers-workspace), DW_HOST (0.0.0.0), +# DW_NODE_BIN (the directory holding npm, when it is not on a +# non-interactive PATH and not in one of the places find_npm looks). set -euo pipefail DW_DIR="${DW_DIR:-$HOME/diffusers-workflow}" @@ -57,6 +59,31 @@ say() { echo "[deploy $(ts)] $*"; } health() { curl -s -m 5 -H "Authorization: Bearer $DW_TOKEN" "$HEALTH" 2>/dev/null; } server_pids() { pgrep -f 'python -m dw\.serve' || true; } +# `ssh lem deploy.sh` runs a non-interactive shell, and a per-user node +# install is usually put on PATH by ~/.bashrc *after* its "not interactive, +# stop here" guard - so npm that works at a prompt is missing here. Look in +# the usual per-user places rather than depend on the caller's shell +find_npm() { + command -v npm >/dev/null 2>&1 && return 0 + local dir + for dir in "${DW_NODE_BIN:-}" "$HOME/.local/node/bin" "$HOME/.volta/bin" \ + "$HOME/.local/share/fnm/aliases/default/bin" "$HOME/.local/bin" /usr/local/bin; do + if [ -n "$dir" ] && [ -x "$dir/npm" ]; then + PATH="$dir:$PATH"; export PATH + say "npm not on PATH; using $dir" + return 0 + fi + done + # nvm is a shell function, not a directory on PATH + if [ -s "${NVM_DIR:-$HOME/.nvm}/nvm.sh" ]; then + # shellcheck disable=SC1091 + . "${NVM_DIR:-$HOME/.nvm}/nvm.sh" >/dev/null 2>&1 && command -v npm >/dev/null 2>&1 \ + && { say "npm from nvm: $(command -v npm)"; return 0; } + fi + say "npm not found (PATH=$PATH); set DW_NODE_BIN to the directory holding npm" + exit 1 +} + cd "$DW_DIR" [ -z "$branch" ] && branch="$(git branch --show-current)" @@ -80,6 +107,7 @@ if echo "$changed" | grep -qx 'pyproject.toml'; then venv/bin/pip install -q -e . fi if echo "$changed" | grep -q '^ui/' || [ ! -d ui/dist ]; then + find_npm if echo "$changed" | grep -qx 'ui/package-lock.json' || [ ! -d ui/node_modules ]; then say "ui/package-lock.json changed; npm ci" (cd ui && npm ci --silent --no-audit --no-fund) diff --git a/tests/test_mcp_server.py b/tests/test_mcp_server.py index b4a45d5e..4248fedb 100644 --- a/tests/test_mcp_server.py +++ b/tests/test_mcp_server.py @@ -1441,7 +1441,13 @@ def test_the_stated_tool_count_is_the_registered_one(): # as tightly as it can be said and still 44 tokens over, so the budget takes # them deliberately rather than the sentence being cut to nothing. Measured # 2026-09-22 at 13_844.0 (9_078.0 / 3_752.0 / 1_014.0). 6 tokens of headroom. -SURFACE_BUDGET = 13_850 +# Then list_gallery took `folder` and `version`, so "show me v5" is one call +# rather than a scan of the listing. The docstring paid for its own new +# sentence and then some (descriptions 9_078 -> 9_068, the `url` sentence +# said in fewer words); the two schema entries (+49) are what the budget +# takes, since no docstring can pay for a parameter's schema. Measured +# 2026-09-22 at 13_883.0 (9_068.0 / 3_801.0 / 1_014.0). 7 tokens of headroom. +SURFACE_BUDGET = 13_890 @pytest.mark.asyncio From b2fbf6e7825f050abb91bb6e3556b6e090a499fc Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Tue, 22 Sep 2026 16:28:33 -0500 Subject: [PATCH 5/5] wait_for_job's slim job carries run_version The slim projection named run_id but not run_version, so an agent that had just waited on a run could not say "that was v5" - the label the gallery puts on its files - without a second call to get_job. One key in _SLIM_KEYS; no docstring change, so the MCP surface budget is untouched. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/MCP.md | 2 +- dw_mcp/diagnose.py | 3 +++ tests/test_mcp_diagnose.py | 18 ++++++++++++++++++ 3 files changed, 22 insertions(+), 1 deletion(-) diff --git a/docs/MCP.md b/docs/MCP.md index b3ed1c69..1ddac9ee 100644 --- a/docs/MCP.md +++ b/docs/MCP.md @@ -306,7 +306,7 @@ references written in the same session. | `get_job_workflow(job_id)` | `job_id` | The REST equivalent is `GET /api/jobs/{id}/workflow` (see [SERVER.md](SERVER.md#jobs-api)). The workflow the job actually ran. `realized: true` means every mutable input is pinned (arguments, seed, prompts, `output:latest`); `false` means the job predates run tracking and this is the definition as submitted. Pass it to `save_workflow` to keep it under a name | | `export_job(job_id, overwrite=False)` | `job_id`, `overwrite` | Gather one finished job into `/exports//` on the server: the realized workflow, the run's manifest, the job row, a README, and copies of the assets, earlier-run inputs and outputs. Returns the directory, a zip URL, the file list with sizes and the total. The three JSON files are in the zip, not repeated here - get_job_workflow and get_job serve them individually. **The directory is on the machine running the server**, like `download_output`'s destination - fetch the zip URL and unpack it into `exports/` under the session's working directory (a deliverable, not a temp file); the archive already unpacks into one folder named after the job id | | `get_job_events(job_id, after=-1, limit=200)` | `job_id`, `after`, `limit` | Get a page of a job's progress events | -| `wait_for_job(job_id, timeout_seconds=20)` | `job_id`, `timeout_seconds` | Block until a job reaches a terminal status, or `timeout_seconds` elapses. **One call blocks for at most 55 seconds** — a larger `timeout_seconds` is clamped, not honoured, because no MCP client holds a tool call open for a generation's real runtime, so budget one call per ~55s of the job. Every reply carries `waited_seconds`, `timeout_requested_seconds`, `timeout_applied_seconds` and `timeout_capped`, so a capped return is distinguishable from an elapsed one. Use instead of hand-polling `get_job`/`get_job_events` in a loop; if it returns `still_running: true`, call it again. Returns a slim job - status, warnings, error, and the manifest once finished - without the arguments; `get_job` has those. A running job also carries `progress` (below) | +| `wait_for_job(job_id, timeout_seconds=20)` | `job_id`, `timeout_seconds` | Block until a job reaches a terminal status, or `timeout_seconds` elapses. **One call blocks for at most 55 seconds** — a larger `timeout_seconds` is clamped, not honoured, because no MCP client holds a tool call open for a generation's real runtime, so budget one call per ~55s of the job. Every reply carries `waited_seconds`, `timeout_requested_seconds`, `timeout_applied_seconds` and `timeout_capped`, so a capped return is distinguishable from an elapsed one. Use instead of hand-polling `get_job`/`get_job_events` in a loop; if it returns `still_running: true`, call it again. Returns a slim job - status, warnings, error, `run_id` and `run_version` (the run's `v5`, as the gallery labels it), and the manifest once finished - without the arguments; `get_job` has those. A running job also carries `progress` (below) | | `cancel_job(job_id)` | `job_id` | Ask a queued or running job to stop | | `clear_memory()` | — | Drop every loaded pipeline and the step cache, freeing VRAM/RAM immediately instead of waiting for the next job to evict one model for another. Also drops the step cache, so a seeded workflow that would otherwise reuse cached results regenerates on its next run. Refused with a 409 while a job is running or queued - the queue is FIFO, so wait for it to finish and retry rather than expecting this call to block until it does (#221) | | `rerun_job(job_id, acknowledged_cost=False, new_seed=False)` | `job_id`, `acknowledged_cost`, `new_seed` | Queue a fresh job from a previous job's stored specification. Costs GPU time, so it passes the same gate as `run_workflow`. `new_seed=true` draws a fresh seed into the workflow's seed variable — without it a seeded workflow's rerun repeats its arguments exactly and the step cache serves the whole run from the earlier one's files (`reused: true`), generating nothing. `get_job_workflow`'s `seed_variable` says whether there is one - `acknowledged_cost` is `true` or the bound `{fingerprint, minutes, downloads}` from the validate plan; a 409 means the plan changed and the message carries the new estimate | diff --git a/dw_mcp/diagnose.py b/dw_mcp/diagnose.py index 66d7bd79..ade6ca36 100644 --- a/dw_mcp/diagnose.py +++ b/dw_mcp/diagnose.py @@ -198,6 +198,9 @@ def get_job_events(client, job_id, after=-1, limit=200): "finished_at", "workspace", "run_id", + # The run's ordinal - the 'v5' the gallery labels its files with - so + # the caller can name the run it just waited on without another call + "run_version", "queue_position", "warnings", "error", diff --git a/tests/test_mcp_diagnose.py b/tests/test_mcp_diagnose.py index 01192f1c..98e2474a 100644 --- a/tests/test_mcp_diagnose.py +++ b/tests/test_mcp_diagnose.py @@ -417,6 +417,24 @@ def test_wait_for_job_keeps_the_manifest_and_error_once_terminal(monkeypatch): assert "get_job" in result["next"] +def test_wait_for_job_names_the_run_the_way_the_gallery_will(monkeypatch): + """The job that just finished is the one a person asks about next, and + the gallery labels its files 'v5' - so the slim job carries the number + beside the run id rather than sending the caller to get_job for it.""" + done = { + **FAT_JOB, + "status": "succeeded", + "run_id": "20260922-212005-cd189c68", + "run_version": 5, + } + client, _ = scripted({("GET", "/api/jobs/job-1"): (200, done)}) + + result = diagnose.wait_for_job(client, "job-1") + + assert result["job"]["run_id"] == "20260922-212005-cd189c68" + assert result["job"]["run_version"] == 5 + + def test_wait_for_job_reports_queue_position_for_a_still_queued_job(monkeypatch): monkeypatch.setattr(diagnose, "MAX_WAIT_SECONDS", 0) queued = {**FAT_JOB, "status": "queued", "queue_position": 2}