From 7f40f23af925e367b99cd788e1ea5580ee6d1af1 Mon Sep 17 00:00:00 2001 From: Don Kackman Date: Mon, 21 Sep 2026 17:45:11 -0500 Subject: [PATCH] =?UTF-8?q?feat(mcp):=20run=5Fworkflow=20wait=5Fseconds=20?= =?UTF-8?q?and=20delete=5Foutput=20by=20job=5Fid=20=E2=80=94=20fold=20the?= =?UTF-8?q?=20wait=20turn=20into=20the=20run,=20delete=20a=20run=20by=20th?= =?UTF-8?q?e=20job=20that=20wrote=20it?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Measured over ~1,400 regression/tester cases: 604 run_workflow calls were followed by 776 wait_for_job calls (682 at the 55 s cap), and delete_output was called 647 times, 274 of them per file. For an unattended agent each of those is a tool turn, and cost is turns x context. Two additive changes: - run_workflow(..., wait_seconds=0): above 0, after queuing, block exactly as wait_for_job(job_id, timeout_seconds=wait_seconds) would - the same loop, not a second one, so the MAX_WAIT_SECONDS clamp and the waited_seconds / timeout_* / still_running fields are identical - and return the queued-job fields plus that wait's slim job. The cost gate is untouched and a refused queue (409) waits on nothing. Default 0 is the old behaviour byte for byte. - delete_output(name=None, workspace=None, job_id=None): exactly one of name / job_id. By job_id the job record's run_dir - the / the run-directory delete already accepts - is read and the whole run deleted; the reply adds job_id and the resolved run_dir. A job with no run_dir, or an unknown job, errors before anything is deleted. With no workspace pin the delete goes to the workspace the job ran in, since that is where the run directory is. The two schema parameters and their sentences first measured the resident tool surface at 14_097 against the 13_800 ceiling; paid for inside the three docstrings they touch (run_workflow, wait_for_job, delete_output) with every pinned phrase and fact kept, landing at 13_796.75. docs/MCP.md rows and dw_mcp/CLAUDE.md updated; tests for both paths in test_mcp_diagnose / test_mcp_media, schema exposure in test_mcp_server. Co-Authored-By: Claude Opus 5 --- docs/MCP.md | 4 +- dw_mcp/CLAUDE.md | 8 ++- dw_mcp/diagnose.py | 22 ++++++- dw_mcp/media.py | 39 +++++++++++- dw_mcp/server.py | 124 ++++++++++++++++++++---------------- tests/test_mcp_diagnose.py | 126 +++++++++++++++++++++++++++++++++++++ tests/test_mcp_media.py | 108 +++++++++++++++++++++++++++++++ tests/test_mcp_server.py | 30 +++++++++ 8 files changed, 399 insertions(+), 62 deletions(-) diff --git a/docs/MCP.md b/docs/MCP.md index f04779ea..fb253e37 100644 --- a/docs/MCP.md +++ b/docs/MCP.md @@ -238,7 +238,7 @@ when no single workflow covers it. | `get_output_frames(name, at=None, seams=None, count=None, boundaries=None, names=None, max_dimension=512, hear=None, workspace=None)` | `name`, `at`, `seams`, `count`, `boundaries`, `names`, `max_dimension`, `hear`, `workspace` | See a generated video as frames, since there is no video content type over MCP (#193). One selector per call: `count` for an evenly spaced contact sheet, `at` for moments (seconds or `"frame:N"`), `seams` (true, or seam numbers from 1) for the last frame before and first frame after each join side by side (each pair carries `difference`, the mean pixel change across the join, 0-255 - rank seams by it and look at the worst) - `boundaries` is each later shot's first frame - the running sum of the shots' `frame_count` from `get_gallery_metadata` on their own `intermediate/` files - and `names` names them. Tiles are fitted to `max_dimension` and, when the set would exceed the 4MB budget, shrunk together rather than dropped; the text part lists each tile and says so. `hear=N` also returns N seconds of soundtrack centred on each `at` moment, after its image - the way to check a hit point or lip-sync without reconciling two clocks; a mute clip keeps its frames and says `no soundtrack` | | `get_output_text(name, max_characters=20000, workspace=None)` | `name`, `max_characters`, `workspace` | Read a text output — a prompt enhancement, or any step whose result is `text/plain` or JSON. Reports the file's real length and whether it was truncated. `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 | | `download_output(name, destination=None, overwrite=False, workspace=None)` | `name`, `destination`, `overwrite`, `workspace` | Save one output file to local disk, of any content type. `destination` may be a full path, a directory, or omitted to save under the output's own name in the current working directory; `~` expands and missing parent directories are created. `overwrite=True` is required to replace a file already at the resolved path. Over a `dw.serve --mcp` endpoint the file lands on the server, so the destination is confined to that workspace and a relative one is joined onto it. Returns nothing to the conversation but where the file landed — unlike the other media tools, the point is a file on disk, not a payload in context. Writes on the machine running the MCP server - over `dw.serve --mcp` that is the GPU box. A write that fails there (a path that exists only on the client, for instance) comes back as an error naming the server-side write and the client-side alternatives, not as an anonymous tool failure. `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 | -| `delete_output(name, workspace=None)` | `name`, `workspace` | Permanently remove one generated file from the output directory. `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 | +| `delete_output(name=None, workspace=None, job_id=None)` | exactly one of `name` / `job_id`, `workspace` | Permanently remove one generated file from the output directory, or - with `name` a `/` run directory, or with `job_id` - a whole run. By `job_id` the run directory is read from the job record (`run_dir`) and the reply adds `job_id` and the resolved `run_dir` to the usual `name` / `deleted` / `run_swept`; a job that never wrote a run directory, or an unknown one, is an error. `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; a `job_id` delete with no `workspace` goes to the workspace the job ran in | ### Authoring, assets and workspaces @@ -284,7 +284,7 @@ references written in the same session. | Tool | Arguments | Purpose | | --- | --- | --- | -| `run_workflow(workflow_path=None, inline_workflow=None, arguments=None, acknowledged_cost=False, workspace=None)` | exactly one of `workflow_path` (a catalog name from `list_workflows`, with or without `.json`, or a path to a workflow file on the server) or `inline_workflow`, optional `arguments`, `acknowledged_cost`, `workspace` | Queue a workflow for generation. Returns as soon as the job is queued. `workspace` names the workspace for this one call without switching the session to it - use it to pin a job whose `output:` or `asset:` references live in a workspace other than the session's - `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 | +| `run_workflow(workflow_path=None, inline_workflow=None, arguments=None, acknowledged_cost=False, workspace=None, wait_seconds=0)` | exactly one of `workflow_path` (a catalog name from `list_workflows`, with or without `.json`, or a path to a workflow file on the server) or `inline_workflow`, optional `arguments`, `acknowledged_cost`, `workspace`, `wait_seconds` | Queue a workflow for generation. Returns as soon as the job is queued - unless `wait_seconds` is above 0, in which case the call then waits on the queued job exactly as `wait_for_job(job_id, timeout_seconds=wait_seconds)` would (same 55 s cap per call, clamped not honoured) and the result carries the queued-job fields plus that wait's (`status`, `still_running`, `waited_seconds`, `timeout_requested_seconds`, `timeout_applied_seconds`, `timeout_capped`, the slim `job`); when the cap covers the job's runtime one call is the run and the wait, and a `still_running: true` result is followed with `wait_for_job` as before. `workspace` names the workspace for this one call without switching the session to it - use it to pin a job whose `output:` or `asset:` references live in a workspace other than the session's - `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, and nothing is waited on | | `get_job(job_id)` | `job_id` | Get a job's status, warnings, output manifest, error and traceback; each manifest entry's `subfolder` is the in-run subfolder the step declared - by convention `final` for the deliverable, `intermediate` for scratch, `''` for none. A running job also carries `progress` (below) | | `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 | diff --git a/dw_mcp/CLAUDE.md b/dw_mcp/CLAUDE.md index d8ccc7d7..50a53f2d 100644 --- a/dw_mcp/CLAUDE.md +++ b/dw_mcp/CLAUDE.md @@ -71,6 +71,12 @@ forwarded verbatim, which the server refuses with 409 when the run's shape changed since (`_acknowledgement_body` in `diagnose.py`; the 409 is rendered with the new estimate by `DwClient._format_detail`). The three job-queuing tools return as soon as the job is queued, since a generation outlasts any client's tool-call -timeout. Authoring has two halves: `get_schema` describes a workflow and +timeout; `run_workflow(wait_seconds=N)` then folds the first `wait_for_job` +into the same call (same `MAX_WAIT_SECONDS` clamp, same budget fields), because +measured over ~1,400 agent-driven cases almost every run was followed by a +wait turn of its own. `delete_output(job_id=...)` is the same economy for +cleanup: the job record's `run_dir` is the `/` the +run-directory delete already accepts, so a whole run goes in one call without +a gallery listing to find its name. Authoring has two halves: `get_schema` describes a workflow and `get_prompt_schema` a stored prompt, which a workflow reaches by `"prompt:name"`. See docs/MCP.md. diff --git a/dw_mcp/diagnose.py b/dw_mcp/diagnose.py index 8a1940f9..66d7bd79 100644 --- a/dw_mcp/diagnose.py +++ b/dw_mcp/diagnose.py @@ -75,6 +75,7 @@ def run_workflow( arguments=None, acknowledged_cost=False, workspace=None, + wait_seconds=0, ): """Queue a workflow. `workflow_path` (or `name` - the same thing `validate_workflow` calls it) is either a catalog name from @@ -85,6 +86,17 @@ def run_workflow( spelling. Returns as soon as it is queued - it does not wait for the job to finish. Poll `get_job_events` for progress. + `wait_seconds` folds the first `wait_for_job` into this call: when it + is above 0 the queued job is waited on exactly as + `wait_for_job(job_id, timeout_seconds=wait_seconds)` would - same clamp + to MAX_WAIT_SECONDS, same `waited_seconds` / `timeout_*` / + `still_running` fields - and the answer carries the queued-job fields + plus that wait's slim job. Almost every run is followed by a wait, and + an unattended agent pays a whole tool turn for it; where the cap covers + the job's runtime this one call is the run and the wait. The gate is + untouched: queuing is refused before anything is waited on, and a + refused queue (409 on a stale plan) returns nothing extra. + `acknowledged_cost` is true or, better, the plan it was quoted from: {fingerprint, minutes, downloads} from `validate_workflow` - see COST_REFUSAL. A bound one the server checks; a 409 means the run's @@ -120,13 +132,21 @@ def run_workflow( # output: references in the wrong root fails after it was queued params = {"workspace": workspace} if workspace else None job = client.post_json("/api/jobs", payload, params=params) - return { + queued = { "job_id": job.get("id"), "status": job.get("status"), "queue_position": job.get("queue_position"), "next": "Poll get_job_events(job_id) for progress, then get_job(job_id) " "for the manifest or the error.", } + if not wait_seconds or float(wait_seconds) <= 0 or queued["job_id"] is None: + return queued + # The same loop wait_for_job runs, not a second one: its clamp, its + # budget fields and its `next` are what a caller already paces against. + # The wait's status and next overwrite the queued ones, since the job + # has moved on from "queued" by the time either is read + waited = wait_for_job(client, queued["job_id"], timeout_seconds=wait_seconds) + return {**queued, **waited} def get_job(client, job_id): diff --git a/dw_mcp/media.py b/dw_mcp/media.py index e243d373..26f1021c 100644 --- a/dw_mcp/media.py +++ b/dw_mcp/media.py @@ -387,14 +387,47 @@ def is_text(content_type): } -def delete_output(client, name, workspace=None): +def delete_output(client, name=None, workspace=None, job_id=None): """Remove one file from the output directory. The gallery is the output directory read back, so this is where a delete belongs. The run directory goes too once its last media file is gone, sidecars included, and a `/` name removes a whole run - what a - failed run, which has a manifest and nothing else, needs (#134).""" - return client.delete_json(api_path("api", "gallery", name), workspace=workspace) + failed run, which has a manifest and nothing else, needs (#134). + + `job_id` is the other handle on a whole run: the job record carries + the `/` its run wrote (`run_dir`, relative to the + output root), so the run is deleted without the caller listing the + gallery to find the name. Exactly one of `name` / `job_id`. A job that + never wrote a run directory - refused before it started, or from + before run tracking - has nothing to delete and says so. `workspace` + pins the call as it always has; without one, a job's delete goes to + the workspace the job itself ran in, since that is where its run + directory is.""" + if (name is None) == (job_id is None): + raise DwApiError( + "Provide exactly one of `name` (a gallery name or a " + "`/` run directory) or `job_id` (the run that " + "job wrote, deleted whole)." + ) + if job_id is None: + return client.delete_json(api_path("api", "gallery", name), workspace=workspace) + + job = client.get_json(api_path("api", "jobs", job_id)) + run_dir = job.get("run_dir") + if not run_dir: + raise DwApiError( + f"Job {job_id} ({job.get('status') or 'unknown status'}) has no " + "run directory to delete - it never started a run, or predates " + "run tracking. If it left files, list_gallery(only_orphans=True) " + "finds the run directory by name." + ) + # The record's path is slash-separated relative to the output root - + # exactly the run-directory form the gallery route accepts + run_dir = "/".join(part for part in str(run_dir).split("/") if part) + target = workspace or job.get("workspace") or None + deleted = client.delete_json(api_path("api", "gallery", run_dir), workspace=target) + return {**deleted, "job_id": job_id, "run_dir": run_dir} def _remote_root(client): diff --git a/dw_mcp/server.py b/dw_mcp/server.py index 721ed029..cd9ebb95 100644 --- a/dw_mcp/server.py +++ b/dw_mcp/server.py @@ -680,23 +680,27 @@ def get_output_text( client, name, max_characters=max_characters, workspace=workspace ) - def delete_output(name: str, workspace: str | None = None) -> dict: + def delete_output( + name: str | None = None, + workspace: str | None = None, + job_id: str | None = None, + ) -> dict: """Permanently remove one generated file from the output directory. - Not recoverable: rerunning the job that made it is the only way - back, and any "output:" reference pointing at it stops resolving. - Prefer `keep_output` first if it is worth keeping. When it was the - last media file of its run, the run directory goes with it - - `manifest.json` and `workflow.json` included - so deleting what you - made leaves the workspace as you found it. `name` may also be a run + Not recoverable (rerun the job to get it back), and any "output:" + reference to it stops resolving; prefer `keep_output` if it is + worth keeping. When it was the last media file of its run, the run + directory goes with it, sidecars included. `name` may also be a run directory ("/", the first two parts of a gallery - name), which removes the whole run: the only way to clear a run that - failed before it wrote any media. + name), which removes the whole run - the only handle on a run that + failed before writing any media - or give `job_id` instead: the run + that job wrote is removed whole, and the reply adds `job_id` and + the resolved `run_dir`. Exactly one of the two; a job with no run + directory, or unknown, is an error. - `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 is reachable from - here without leaving this one (#99).""" - return media.delete_output(client, name, workspace=workspace) + `workspace` pins this call to another workspace without switching + the session (#99); a `job_id` delete with no `workspace` goes to + the workspace the job ran in.""" + return media.delete_output(client, name, workspace=workspace, job_id=job_id) def download_output( name: str, @@ -1092,30 +1096,33 @@ def run_workflow( arguments: dict | None = None, acknowledged_cost: bool | dict = False, workspace: str | None = None, + wait_seconds: int = 0, ) -> dict: """Queue a workflow for generation. THIS COSTS GPU TIME: a run occupies the machine for minutes and the engine runs one job at a time. Tell the user what will run and get their go-ahead, then pass - acknowledged_cost=true. Returns as soon as the job is queued - a - generation outlasts any tool-call timeout - so follow it with - `wait_for_job`, then `get_job` for the manifest. Give exactly one of - `workflow_path` - a catalog name from `list_workflows`, with or - without .json, or a path on the server - or `inline_workflow`, a - full definition for a request nothing stored covers. `validate_workflow` - calls these same two concepts `name` and `workflow`; both tools - accept both spellings, so a document just validated can be run - without renaming a key. `arguments` overrides the workflow's - variables by name, which is how one stored workflow serves many - requests without being edited or copied. `workspace` names the - workspace for this one call without switching the session to it - - use it to pin a job whose `output:` or `asset:` references live in a - workspace other than the session's. + acknowledged_cost=true. Returns as soon as the job is queued; + follow it with `wait_for_job`, then `get_job` for the manifest - or + fold that first wait in with `wait_seconds` above 0, which waits on + the job exactly as `wait_for_job(job_id, + timeout_seconds=wait_seconds)` would ({cap}s cap per call) and adds + its fields to the result (`still_running`, `waited_seconds`, + `timeout_*`, the slim `job`). If the cap covers the job's + runtime one call is enough; on `still_running: true` call + `wait_for_job` as before. Give exactly one of `workflow_path` - a + catalog name from `list_workflows`, with or without .json, or a + path on the server - or `inline_workflow`, a full definition + nothing stored covers; `validate_workflow` calls these `name` and + `workflow`, and both tools accept both spellings. `arguments` + overrides the workflow's variables by name. `workspace` pins this + call to another workspace without switching the session (where its + `output:`/`asset:` references live). Bind the acknowledgement to what you quoted: pass {"fingerprint": plan.fingerprint, "minutes": plan.estimate.minutes, - "downloads": [...the non-null repos in plan.downloads_required]} from the - validate answer, and the server refuses with 409 - naming the new - plan - if the run's shape changed since; bare true is for a plan + "downloads": [...the non-null repos in plan.downloads_required]} from + the validate plan; the server refuses with 409, naming the new + plan, if the run's shape changed since. Bare true is for a plan that was null.""" return diagnose.run_workflow( client, @@ -1126,6 +1133,15 @@ def run_workflow( arguments=arguments, acknowledged_cost=acknowledged_cost, workspace=workspace, + wait_seconds=wait_seconds, + ) + + # The cap is a number a caller paces against, so the description + # states it (as wait_for_job's does, below). replace rather than + # format: the docstring spells out a literal {fingerprint, ...} dict. + if run_workflow.__doc__: # absent under python -OO + run_workflow.__doc__ = run_workflow.__doc__.replace( + "{cap}", str(diagnose.MAX_WAIT_SECONDS) ) def get_job(job_id: str) -> dict: @@ -1166,31 +1182,29 @@ def get_job_events(job_id: str, after: int = -1, limit: int = 200) -> dict: def wait_for_job(job_id: str, timeout_seconds: int = 20) -> dict: """Block until a job finishes, instead of polling get_job or - get_job_events by hand. Returns as soon as the job's status is - succeeded, failed or cancelled, or - if timeout_seconds elapses - first - returns its current status with still_running: true so you - can call again. Does not queue anything, so no acknowledged_cost. - - One call blocks for at most {cap} seconds, no matter what - timeout_seconds asks for - this deployment's cap, set for the tool - call budget the MCP client actually holds open. A larger value is - not honoured, it is clamped, so budget roughly one call per {cap}s - of the job - if {cap} covers the job's whole runtime, one call is - enough. Every reply says which happened: waited_seconds, - timeout_requested_seconds, timeout_applied_seconds and - timeout_capped. - - 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`: the step, phase, and + get_job_events by hand: returns as soon as its status is succeeded, + failed or cancelled, or with still_running: true when + timeout_seconds elapses first, so you can call again. Queues + nothing, so no acknowledged_cost. + + One call blocks for at most {cap} seconds, whatever timeout_seconds + asks for - this deployment's cap, set for the tool-call budget the + client holds open; a larger value is clamped, not honoured, so + budget one call per {cap}s of the job, and one call is enough when + {cap} covers its runtime. Every reply says which happened: + waited_seconds, timeout_requested_seconds, timeout_applied_seconds + and timeout_capped. + + Returns a slim job - status, warnings, error, the manifest once + finished - without the arguments (get_job has those). A running job + also carries `progress`: step, phase, and `denoise_step`/`denoise_total_steps`, null until the denoise loop - starts. Judge a slow run against a stuck one by whether - `denoise_step` has moved since a poll minutes ago, not by silence - past a fixed threshold - a video reference's lead-in can run many - minutes emitting nothing, and denoise gaps are uneven under a - transformer block cache; both are normal. The full diagnosis, and - why `denoise_total_steps` sometimes reads one less than what was - asked for, are in WORKFLOW_GUIDE's "The loop", step 5.""" + starts. Tell a slow run from a stuck one by whether + `denoise_step` has moved since a poll minutes ago, not by silence: + a video reference's lead-in can run many minutes emitting nothing, + and denoise gaps are uneven under a transformer block cache - both + normal. Full diagnosis, and why `denoise_total_steps` can read one + less than asked, in WORKFLOW_GUIDE's "The loop", step 5.""" return diagnose.wait_for_job(client, job_id, timeout_seconds=timeout_seconds) # The cap is a number a caller paces against, so the description states diff --git a/tests/test_mcp_diagnose.py b/tests/test_mcp_diagnose.py index 9cc759a2..01192f1c 100644 --- a/tests/test_mcp_diagnose.py +++ b/tests/test_mcp_diagnose.py @@ -427,6 +427,132 @@ def test_wait_for_job_reports_queue_position_for_a_still_queued_job(monkeypatch) assert result["job"]["queue_position"] == 2 +def running_then_done(submit_body=SUBMITTED, bodies=None): + """POST /api/jobs queues; GET /api/jobs/job-1 walks `bodies`, repeating + the last one - a run that is then polled.""" + seen = [] + bodies = bodies or [{"id": "job-1", "status": "succeeded", "manifest": ["x"]}] + + def handler(request): + key = (request.method, request.url.path) + seen.append({"key": key, "params": dict(request.url.params)}) + if key == ("POST", "/api/jobs"): + return httpx.Response(201, json=submit_body) + if key == ("GET", "/api/jobs/job-1"): + polls = sum(1 for entry in seen if entry["key"] == key) + return httpx.Response(200, json=bodies[min(polls - 1, len(bodies) - 1)]) + return httpx.Response(404, json={"detail": f"unrouted {key}"}) + + return DwClient(transport=httpx.MockTransport(handler)), seen + + +def test_run_with_wait_seconds_zero_is_the_plain_queue(): + """The default: one POST, the queued-job shape, nothing polled - what + every caller before `wait_seconds` existed gets unchanged.""" + client, seen = running_then_done() + + result = diagnose.run_workflow( + client, workflow_path="w.json", acknowledged_cost=True, wait_seconds=0 + ) + + assert [entry["key"] for entry in seen] == [("POST", "/api/jobs")] + assert result["job_id"] == "job-1" + assert result["status"] == "queued" + assert "still_running" not in result + assert "waited_seconds" not in result + + +def test_run_with_wait_seconds_folds_the_first_wait_in(monkeypatch): + """Queue, then wait on the same call: the answer is the queued-job + fields plus everything wait_for_job would have returned - status, + budget fields, the slim job with its manifest.""" + monkeypatch.setattr(diagnose, "WAIT_POLL_SECONDS", 0.01) + client, seen = running_then_done( + bodies=[ + {"id": "job-1", "status": "running", "progress": {"step": "s"}}, + {"id": "job-1", "status": "succeeded", "manifest": ["out.png"]}, + ] + ) + + result = diagnose.run_workflow( + client, workflow_path="w.json", acknowledged_cost=True, wait_seconds=5 + ) + + assert [entry["key"] for entry in seen] == [ + ("POST", "/api/jobs"), + ("GET", "/api/jobs/job-1"), + ("GET", "/api/jobs/job-1"), + ] + assert result["job_id"] == "job-1" + assert result["queue_position"] == 2, "the queued-job fields survive" + assert result["status"] == "succeeded", "the wait's status wins" + assert result["still_running"] is False + assert result["job"]["manifest"] == ["out.png"] + assert result["timeout_requested_seconds"] == 5.0 + assert result["timeout_applied_seconds"] == 5.0 + assert result["timeout_capped"] is False + assert "waited_seconds" in result + assert "get_job" in result["next"] + + +def test_run_with_wait_seconds_reports_still_running_at_the_budget(monkeypatch): + monkeypatch.setattr(diagnose, "WAIT_POLL_SECONDS", 0.01) + client, _seen = running_then_done( + bodies=[{**FAT_JOB, "status": "running"}], + ) + + result = diagnose.run_workflow( + client, workflow_path="w.json", acknowledged_cost=True, wait_seconds=0.03 + ) + + assert result["status"] == "running" + assert result["still_running"] is True + assert "arguments" not in result["job"], "the slim job, as wait_for_job's" + assert "wait_for_job" in result["next"] + + +def test_run_with_wait_seconds_is_clamped_like_wait_for_job(monkeypatch): + """The cap is the deployment's, not the caller's: asking for 600 here + gets the same clamp and the same honest budget fields as wait_for_job.""" + monkeypatch.setattr(diagnose, "WAIT_POLL_SECONDS", 0.01) + monkeypatch.setattr(diagnose, "MAX_WAIT_SECONDS", 0.03) + client, _seen = running_then_done(bodies=[{"id": "job-1", "status": "running"}]) + + result = diagnose.run_workflow( + client, workflow_path="w.json", acknowledged_cost=True, wait_seconds=600 + ) + + assert result["still_running"] is True + assert result["timeout_capped"] is True + assert result["timeout_requested_seconds"] == 600 + assert result["timeout_applied_seconds"] == 0.0 + assert "600" in result["next"] + + +def test_run_with_wait_seconds_still_refuses_without_an_acknowledged_cost(): + """Folding the wait in changes nothing about the gate.""" + client, seen = running_then_done() + + with pytest.raises(DwApiError, match="acknowledged_cost"): + diagnose.run_workflow(client, workflow_path="w.json", wait_seconds=30) + + assert seen == [] + + +def test_run_with_wait_seconds_waits_on_nothing_when_the_queue_is_refused(): + """A 409 (stale plan) surfaces as before and no poll follows it.""" + client, seen = scripted( + {("POST", "/api/jobs"): (409, {"detail": {"message": "plan changed"}})} + ) + + with pytest.raises(DwApiError, match="plan changed"): + diagnose.run_workflow( + client, workflow_path="w.json", acknowledged_cost=True, wait_seconds=30 + ) + + assert [entry["key"] for entry in seen] == [("POST", "/api/jobs")] + + class TestGetJobWorkflow: def test_a_realized_workflow_comes_back_with_the_flag_set(self): client, seen = scripted( diff --git a/tests/test_mcp_media.py b/tests/test_mcp_media.py index 874dc5d3..8966041b 100644 --- a/tests/test_mcp_media.py +++ b/tests/test_mcp_media.py @@ -562,6 +562,114 @@ def handler(request): media.delete_output(client, "ghost.png") +def deleting_by_job(job): + """GET /api/jobs/job-1 answers `job`; DELETE on the gallery answers as + the server's run-directory form does.""" + seen = [] + + def handler(request): + seen.append((request.method, request.url.path, dict(request.url.params))) + if request.method == "GET" and request.url.path == "/api/jobs/job-1": + return httpx.Response(200, json=job) + if request.method == "GET": + return httpx.Response(404, json={"detail": "Unknown job"}) + name = request.url.path.removeprefix("/api/gallery/") + return httpx.Response( + 200, json={"name": name, "deleted": True, "run_swept": name.split("/")[-1]} + ) + + return DwClient(transport=httpx.MockTransport(handler)), seen + + +def test_delete_output_by_job_id_deletes_the_run_directory_the_job_wrote(): + client, seen = deleting_by_job( + { + "id": "job-1", + "status": "succeeded", + "run_dir": "ltx2/Gyre", + "workspace": "default", + } + ) + + result = media.delete_output(client, job_id="job-1") + + assert [entry[:2] for entry in seen] == [ + ("GET", "/api/jobs/job-1"), + ("DELETE", "/api/gallery/ltx2/Gyre"), + ] + assert result == { + "name": "ltx2/Gyre", + "deleted": True, + "run_swept": "Gyre", + "job_id": "job-1", + "run_dir": "ltx2/Gyre", + } + + +def test_delete_output_by_job_id_goes_to_the_workspace_the_job_ran_in(): + """The run directory is wherever the job wrote it, so with no pin the + delete follows the job's own workspace rather than the session's.""" + client, seen = deleting_by_job( + {"id": "job-1", "status": "failed", "run_dir": "w/Run", "workspace": "shots"} + ) + client.workspace = "elsewhere" + + media.delete_output(client, job_id="job-1") + + assert seen[-1][1] == "/api/gallery/w/Run" + assert seen[-1][2] == {"workspace": "shots"} + + +def test_delete_output_by_job_id_honours_an_explicit_workspace(): + client, seen = deleting_by_job( + {"id": "job-1", "status": "failed", "run_dir": "w/Run", "workspace": "shots"} + ) + + media.delete_output(client, job_id="job-1", workspace="pinned") + + assert seen[-1][2] == {"workspace": "pinned"} + + +def test_delete_output_refuses_neither_name_nor_job_id(): + client, seen = deleting_by_job({}) + + with pytest.raises(DwApiError, match="exactly one"): + media.delete_output(client) + + assert seen == [] + + +def test_delete_output_refuses_both_name_and_job_id(): + client, seen = deleting_by_job({}) + + with pytest.raises(DwApiError, match="exactly one"): + media.delete_output(client, "out.png", job_id="job-1") + + assert seen == [] + + +def test_delete_output_by_job_id_surfaces_an_unknown_job(): + client, seen = deleting_by_job({}) + + with pytest.raises(DwApiError, match="Unknown job"): + media.delete_output(client, job_id="ghost") + + assert [entry[0] for entry in seen] == ["GET"], "nothing is deleted" + + +def test_delete_output_by_job_id_refuses_a_job_with_no_run_directory(): + """A job refused before it started (or one from before run tracking) + has no run_dir; that is an error, not a delete of nothing.""" + client, seen = deleting_by_job( + {"id": "job-1", "status": "failed", "run_dir": None, "workspace": "default"} + ) + + with pytest.raises(DwApiError, match="no run directory"): + media.delete_output(client, job_id="job-1") + + assert [entry[0] for entry in seen] == ["GET"], "nothing is deleted" + + # --------------------------------------------------------- output download diff --git a/tests/test_mcp_server.py b/tests/test_mcp_server.py index 2e4e8256..ae73aa7b 100644 --- a/tests/test_mcp_server.py +++ b/tests/test_mcp_server.py @@ -239,6 +239,26 @@ async def test_run_workflow_takes_an_acknowledged_cost_flag(): assert "acknowledged_cost" in tools["run_workflow"].input_schema["properties"] +@pytest.mark.asyncio +async def test_run_workflow_and_delete_output_take_the_turn_saving_parameters(): + """Almost every run is followed by a wait, and most deletes are of the + run a job just wrote; each is a whole tool turn for an unattended agent. + `wait_seconds` folds the first wait into the run, `job_id` deletes the + run without a gallery listing to find its name, and neither is required + - a caller from before either existed sees the same defaults.""" + tools = await tools_of(server_over(ok({}))) + + run = tools["run_workflow"].input_schema + assert run["properties"]["wait_seconds"]["default"] == 0 + assert "wait_seconds" not in run.get("required", []) + assert "wait_for_job" in tools["run_workflow"].description + + delete = tools["delete_output"].input_schema + assert "job_id" in delete["properties"] + assert "name" not in delete.get("required", []) + assert "job_id" in tools["delete_output"].description + + @pytest.mark.asyncio async def test_run_workflow_advertises_its_cost(): """The description is what an agent reads before spending GPU minutes.""" @@ -1389,6 +1409,16 @@ def test_the_stated_tool_count_is_the_registered_one(): # sentence to the terse "pins this call" form, and the envelope and asset # paragraphs said in fewer words with every fact kept. Measured 2026-09-21 # at 13_760.5 (9_075.0 / 3_671.5 / 1_014.0). 39.5 tokens of headroom left. +# run_workflow's `wait_seconds` and delete_output's `job_id` (two schema +# parameters, ~50 tokens, plus the sentences that explain them) first +# measured at 14_097. Paid for inside the three docstrings they touch: +# run_workflow's wait paragraph said once and tersely, its spelling / +# arguments / workspace sentences shortened with every fact kept, +# wait_for_job's cap paragraph and stall sentence said in fewer words +# (the pinned phrases stay), delete_output's "leaves the workspace as you +# found it" clause dropped. Measured 2026-09-21 at 13_796.75 (9_061.5 / +# 3_721.25 / 1_014.0). 3.25 tokens of headroom left; measure again before +# the next docstring change. SURFACE_BUDGET = 13_800