Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions docs/MCP.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<workflow>/<run id>` 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

Expand Down Expand Up @@ -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 `<workspace>/exports/<job id>/` 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 |
Expand Down
8 changes: 7 additions & 1 deletion dw_mcp/CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<workflow>/<run id>` 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.
22 changes: 21 additions & 1 deletion dw_mcp/diagnose.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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):
Expand Down
39 changes: 36 additions & 3 deletions dw_mcp/media.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<workflow>/<run id>` 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 `<workflow>/<run id>` 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 "
"`<workflow>/<run id>` 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):
Expand Down
Loading
Loading