diff --git a/CLAUDE.md b/CLAUDE.md index a8b7efef..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 @@ -368,6 +370,45 @@ 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. 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), 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 + 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..1ddac9ee 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, 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 @@ -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/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 ac42818b..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 @@ -118,6 +129,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: @@ -125,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 @@ -135,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. @@ -161,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) @@ -169,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. @@ -345,6 +384,199 @@ 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, 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. + + 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. + """ + path = os.path.join(run_dir, MANIFEST_FILE_NAME) + try: + stat = os.stat(path) + except OSError: + _recorded_versions.pop(path, None) + return 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 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. + """ + 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)) + ), + key=run_id_sort_key, + ) + + +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)) + 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] + # 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 + + +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. + + 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 = record_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..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 @@ -98,7 +99,9 @@ REALIZED_FILE_NAME, is_output_reference, is_run_id, + record_run_versions, resolve_output_reference, + run_versions, split_run_path, ) from ..workspace import ( @@ -2701,13 +2704,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 +2722,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 +2752,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 @@ -2801,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. @@ -2815,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 @@ -2845,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] @@ -2884,6 +2913,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 +2927,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 +2946,8 @@ def gallery_metadata( "source": source, "metadata": metadata, "job": job, + "run_id": run_id, + "version": version, "media": media, } @@ -3355,6 +3394,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 @@ -3397,6 +3439,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): @@ -3597,7 +3642,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: @@ -4126,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") @@ -4147,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 7077e7f1..e65e173b 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 @@ -1135,12 +1152,27 @@ 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 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, "/" @@ -1483,9 +1515,17 @@ 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(), + # 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 8da6ede1..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, @@ -230,10 +240,22 @@ 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 - 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 + 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/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/dw_mcp/server.py b/dw_mcp/server.py index e00806c7..093e7bb0 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 @@ -406,9 +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. + 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 @@ -433,6 +439,8 @@ def list_gallery( subfolder=subfolder, only_orphans=only_orphans, workspace=workspace, + folder=folder, + version=version, ) def get_gallery_metadata( 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_events.py b/tests/test_events.py index 4931f41f..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 @@ -461,3 +464,58 @@ 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] + + +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_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_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} diff --git a/tests/test_mcp_server.py b/tests/test_mcp_server.py index 86f034bd..4248fedb 100644 --- a/tests/test_mcp_server.py +++ b/tests/test_mcp_server.py @@ -1435,7 +1435,19 @@ 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. +# 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 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 18713d9a..360355a7 100644 --- a/tests/test_runs.py +++ b/tests/test_runs.py @@ -10,6 +10,8 @@ from dw.runs import ( FLAT_LAYOUT, + assign_run_version, + record_run_versions, OUTPUT_LAYOUT_ENV_VAR, RUN_LAYOUT, is_run_id, @@ -17,6 +19,7 @@ new_run_id, output_layout, resolve_output_reference, + run_versions, split_run_path, strip_run_id, workflow_identity, @@ -733,3 +736,224 @@ 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")) == {} + + 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_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" + 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 56f1b165..b808d253 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -1993,6 +1993,129 @@ 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"] == "" + + # "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 + 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_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/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 58aceef6..47005503 100644 --- a/ui/CLAUDE.md +++ b/ui/CLAUDE.md @@ -61,6 +61,15 @@ 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. 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, 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..32507caf 100644 --- a/ui/src/lib/pages/GalleryPage.svelte +++ b/ui/src/lib/pages/GalleryPage.svelte @@ -319,7 +319,19 @@ {:else} ♪ {file.label} {/if} - {file.label} + + {#if file.version} + v{file.version} + + {' '} + {/if}{file.label} {/snippet} @@ -352,6 +364,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 +545,19 @@ 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; + padding: 0 0.3rem; + border-radius: 0.2rem; + background: var(--line); + color: var(--ink); + 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..3d9d6339 100644 --- a/ui/src/lib/pages/GalleryPage.test.ts +++ b/ui/src/lib/pages/GalleryPage.test.ts @@ -1,5 +1,6 @@ import { cleanup, + fireEvent, render, screen, waitFor, @@ -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,49 @@ 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() + // 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 () => { + // 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/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/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..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 { @@ -262,6 +268,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