diff --git a/CHANGELOG.md b/CHANGELOG.md index 04c05879a..ac43c47b5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,6 +24,17 @@ history. `--json` they go to stderr so stdout stays the single envelope. The rate is measured over the last ten seconds on a timer, so a stalled upload reports a falling rate instead of going quiet. Schema: `build_push_event.json`. +- `comfy deploy up --watch` and `comfy deploy status` show where a deployment that + is coming up has got to: the step, and while models are copied onto its storage + the model, bytes done of the total, rate and time left ("Staging models: model 1 + of 2 sd_xl_base_1.0.safetensors, 3.5 GB of 7.3 GB, 44.2 MB/s, 1m 25s left"). The + numbers are the deploy service's own `progress` object, which `status --json` + and `up --json` now carry while the status is `provisioning` or `starting` and + omit otherwise. Under `--watch`, `--json-stream` emits a `deploy_progress` event + per new sample on stdout and `--json` puts the same lines on stderr. Ctrl-C + during `--watch` stops the watching and nothing else, and prints the command + that re-attaches. A service that sends no `progress` prints what it printed + before. Schema: `deploy_progress_event.json`. - `comfy build push` prints every warning a save returns, and `--release` cuts no release while one says a deployment could not download a model link (`build_release_held`); `--release-despite-warnings` cuts anyway. @@ -50,6 +61,18 @@ history. carries that error code and no pick is marked. The flag is off by default because it fetches uncached template workflows and calls the local server. +### Changed + +- `comfy deploy up` now follows the deployment until it settles, instead of + returning as soon as the deploy service accepts it. A script that relied on + `up` returning at once passes `--no-watch`. +- A watch (`up`, or `status --watch`) now stops at `unhealthy` instead of + waiting for `ready`: the status only ever follows `ready`, so the wait could + last as long as the endpoint stayed degraded, with nothing printed. `up` reports + an unhealthy deployment as not ok (`deploy_status_terminal`, exit 1), with + or without the watch, since it is billing without serving; `status` still + reports it as recoverable. + ### Fixed - `insert_workflow` (`comfy workflow insert-workflow`) now rebases the inserted diff --git a/comfy_cli/command/deploy.py b/comfy_cli/command/deploy.py index e0ecb9324..b674b691b 100644 --- a/comfy_cli/command/deploy.py +++ b/comfy_cli/command/deploy.py @@ -18,6 +18,7 @@ from comfy_cli.command.build_spec import BuildSpecInvalidError from comfy_cli.command.deploy_compute import prompt_gpu as _prompt_gpu from comfy_cli.command.deploy_compute import prompt_region as _prompt_region +from comfy_cli.command.deploy_progress import DeployWatchReporter from comfy_cli.command.deploy_resolve import DeployResolveError from comfy_cli.command.deploy_runtime import command_clients as _command_clients from comfy_cli.command.deploy_runtime import poll_deployment as _poll_deployment @@ -202,6 +203,8 @@ def status_cmd( typer.Argument(help="ComfyUI install directory or build spec path. Default: the current directory."), ] = None, deployment_id: DeploymentOption = None, + # `status` answers a question and exits; watching is the caller asking to + # stay, so here it stays opt-in. `up` starts the wait, so there it is on. watch: Annotated[bool, typer.Option("--watch", help="Poll until the deployment reaches a terminal state.")] = False, ) -> None: _run_status(path, deployment_id=deployment_id, watch=watch) @@ -234,7 +237,16 @@ def up_cmd( ] = None, release: Annotated[str | None, typer.Option("--release", help="Deploy this release id.")] = None, deployment_id: DeploymentOption = None, - watch: Annotated[bool, typer.Option("--watch", help="Poll until the deployment reaches a terminal state.")] = False, + # Watching is what someone who just asked for a deployment wants: the command + # that starts a several-minute wait should say how the wait is going. Ctrl-C + # and --no-watch both leave the deploy running and print how to re-attach. + watch: Annotated[ + bool, + typer.Option( + "--watch/--no-watch", + help="Follow the deployment until it settles. Use --no-watch to return as soon as it is accepted.", + ), + ] = True, ) -> None: renderer = get_renderer() _require_paired_bounds(renderer, minimum, maximum) @@ -270,7 +282,19 @@ def up_cmd( ) result = reconcile_up(builder, client, replace(request, gpu=selected_gpu, region=selected_region)) if watch: - watched = _poll_deployment(client, _required_string(result.deployment, "id"), _sleep) + watched_id = _required_string(result.deployment, "id") + reporter = DeployWatchReporter(renderer, watched_id) + try: + watched = _poll_deployment(client, watched_id, _sleep, reporter.snapshot) + except KeyboardInterrupt: + # The deploy runs on the service's side and never needed this + # process: say so, report where it had got to, and leave it be. + reporter.interrupted() + if reporter.last is not None: + _render_result(renderer, replace(result, deployment=reporter.last), watch=False) + raise typer.Exit(code=130) from None + finally: + reporter.close() result = replace(result, deployment=watched) _render_result(renderer, result, watch=watch) except (BuildSpecNotFoundError, BuildSpecInvalidError) as error: diff --git a/comfy_cli/command/deploy_progress.py b/comfy_cli/command/deploy_progress.py new file mode 100644 index 000000000..852484c56 --- /dev/null +++ b/comfy_cli/command/deploy_progress.py @@ -0,0 +1,370 @@ +"""Progress for a deployment that is coming up. + +While a deployment's status is ``provisioning`` or ``starting`` the deploy +service carries a ``progress`` object on the deployment read: the step it is in +and, while models are copied onto its storage, models and bytes done against the +total, the model in flight, the rate and the time left. The service computes the +rate and the time left itself, so this module only words what it is given: the +CLI, the portal and an agent reading the events cannot disagree about a number. + +Three audiences, the same split ``comfy build push`` makes for an upload: + +- a person at a terminal gets a single redrawn Rich progress line; +- a person whose output is piped gets a plain line each time the service writes + a new sample, with no carriage-return redraws to litter the file; +- an agent gets ``deploy_progress`` events through ``Renderer.progress_event``: + on stdout under ``--json-stream``, on stderr under ``--json``, where stdout + stays the single envelope. The ``progress`` object rides through unchanged. + +A server that sends no ``progress`` (an older one, or any status but the two +above) produces nothing here, so those commands print what they always printed. +""" + +from __future__ import annotations + +import json +from collections.abc import Callable +from datetime import datetime, timezone +from typing import Any, Final + +from comfy_cli.command.build_spec import JsonObject +from comfy_cli.output.progress import Surface, human_bytes, human_seconds, surface_for +from comfy_cli.output.sanitize import sanitize_markup +from comfy_cli.utils import parse_rfc3339 + +EVENT_PROGRESS: Final = "deploy_progress" + +# The service rewrites the object every few seconds in every step while the +# deploy is alive, the two that count nothing included, and its writes are +# best-effort. A minute of silence is a sample worth doubting; one or two +# dropped writes are ordinary. How long a step has run is `startedAt`'s to say: +# `updatedAt` only says how fresh the sample is. +STALE_SECONDS: Final = 60.0 + +STAGING_STEP: Final = "staging_models" + +_STEP_LABELS: Final = { + STAGING_STEP: "Staging models", + "creating_endpoint": "Creating the endpoint", + "waiting_for_worker": "Waiting for the first worker", +} + +Now = Callable[[], datetime] + + +def _utcnow() -> datetime: + return datetime.now(timezone.utc) + + +COMING_UP = frozenset({"provisioning", "starting"}) + + +def progress_of(deployment: JsonObject) -> JsonObject | None: + """The deployment's progress object, or None where it sent none. + + Lenient where the rest of the deploy parsing is strict: progress narrates a + deploy and must never be what fails a command, so anything that is not an + object naming a step reads as no progress at all. Only a deployment coming + up has any: a settled one still carrying an object is off the contract. + """ + if deployment.get("status") not in COMING_UP: + return None + progress = deployment.get("progress") + if not isinstance(progress, dict) or not isinstance(progress.get("step"), str): + return None + return progress + + +def _sample_key(progress: JsonObject, *, by_stamp: bool) -> str: + """What makes one sample the same as the last. + + By its stamp for an agent, which is told about every write. By what it says + for a person reading piped lines: a step with nothing to count is re-stamped + every few seconds, and a line per re-stamp would bury the ones that matter. + """ + updated_at = progress.get("updatedAt") + if by_stamp and isinstance(updated_at, str): + return updated_at + content = {key: value for key, value in progress.items() if key != "updatedAt"} + return json.dumps(content, sort_keys=True, separators=(",", ":"), default=str) + + +def _number(progress: JsonObject, key: str) -> int | None: + value = progress.get(key) + if isinstance(value, bool) or not isinstance(value, (int, float)): + return None + return int(value) + + +def step_label(progress: JsonObject) -> str: + step = str(progress.get("step")) + # A step this version has never heard of still names itself. + return _STEP_LABELS.get(step, step.replace("_", " ")) + + +def _seconds_since(progress: JsonObject, key: str, now: datetime) -> float | None: + stamp = progress.get(key) + if not isinstance(stamp, str): + return None + try: + then = parse_rfc3339(stamp) + except ValueError: + return None + return max((now - then).total_seconds(), 0.0) + + +def seconds_since_update(progress: JsonObject, now: datetime) -> float | None: + return _seconds_since(progress, "updatedAt", now) + + +def is_stale(progress: JsonObject, now: datetime) -> bool: + age = seconds_since_update(progress, now) + return age is not None and age > STALE_SECONDS + + +def _model_part(progress: JsonObject) -> str | None: + model = progress.get("currentModel") + if not isinstance(model, str) or not model: + return None + name = model.rsplit("/", 1)[-1] + done, total = _number(progress, "modelsDone"), _number(progress, "modelsTotal") + if done is None or not total: + return name + return f"model {min(done + 1, total)} of {total} {name}" + + +def _bytes_part(progress: JsonObject) -> str | None: + done, total = _number(progress, "bytesDone"), _number(progress, "bytesTotal") + if done is None: + return None + if total is None: + return f"{human_bytes(done)} copied" + of = "of at least" if progress.get("bytesTotalIsFloor") is True else "of" + return f"{human_bytes(done)} {of} {human_bytes(total)}" + + +def _parts(progress: JsonObject, now: datetime) -> tuple[str, str | None, list[str]]: + """The step's name, the model it is on, and the numbers, kept apart. + + The sentence and the live bar want the same facts in different shapes: one + joins them with commas, the other spreads them across columns and has to + know which piece may be truncated when the terminal is narrow. + """ + label = step_label(progress) + model: str | None = None + parts: list[str] = [] + if progress.get("step") != STAGING_STEP: + # Nothing here is measured, so the only honest number is how long the + # step has been running. Without it the line never changes and a wait + # that is working reads exactly like one that has died. + waited = _seconds_since(progress, "startedAt", now) + if waited is not None and waited >= 1: + parts.append(f"{human_seconds(waited)} so far") + else: + if _number(progress, "bytesTotal") == 0: + parts.append("every model is already in place") + else: + model = _model_part(progress) + rate, left = _number(progress, "bytesPerSecond"), _number(progress, "etaSeconds") + parts.extend( + part + for part in ( + _bytes_part(progress), + None if not rate else f"{human_bytes(rate)}/s", + None if left is None else f"{human_seconds(left)} left", + ) + if part is not None + ) + return label, model, parts + + +def _notes(progress: JsonObject, now: datetime, *, short: bool = False) -> list[str]: + """What a reader must know about the sample itself: a restart, an old read.""" + notes = [] + attempt = _number(progress, "attempt") + if attempt is not None and attempt > 1: + notes.append(f"attempt {attempt}" if short else f"attempt {attempt}, this step was restarted") + # The same test the events' `stale` flag uses, so a line and an event + # read at the same moment never disagree about the sample. + age = seconds_since_update(progress, now) + if age is not None and is_stale(progress, now): + waited = human_seconds(age) + notes.append(f"no update for {waited}" if short else f"last update {waited} ago, so these numbers may be stale") + return notes + + +def describe(progress: JsonObject, *, now: datetime) -> str: + """One line saying where the deployment is, from one progress object.""" + label, model, parts = _parts(progress, now) + if model is not None: + parts.insert(0, model) + line = label if not parts else f"{label}: {', '.join(parts)}" + return line + "".join(f" ({note})" for note in _notes(progress, now)) + + +def live_line(progress: JsonObject, *, now: datetime) -> str: + """The same facts for the redrawn line, ordered by what may be cut. + + The terminal crops from the right. The notes ride on the label in their short + form, because a restart or a silent service changes how every number after + them reads. The model's name goes last: it is the one piece long enough to + need cutting and the only one a reader can lose without losing a number. + """ + label, model, parts = _parts(progress, now) + notes = _notes(progress, now, short=True) + if notes: + label = f"{label} ({', '.join(notes)})" + if model is not None: + parts.append(model) + return label if not parts else f"{label}: {', '.join(parts)}" + + +def reattach_hint(deployment_id: str) -> str: + return f"comfy deploy status --deployment {deployment_id} --watch" + + +class DeployWatchReporter: + """Reports each new sample of a watched deployment's progress. + + Fed every snapshot the poll reads. The poll runs every two seconds and the + service writes about every three, so a sample is reported once, when it + changes, and once more if it then goes stale: a reader is never shown + movement that did not happen. + """ + + def __init__(self, renderer: Any, deployment_id: str, *, now: Now = _utcnow) -> None: + self._renderer = renderer + self._deployment_id = deployment_id + self._now = now + self._reported: tuple[str, str, bool] | None = None + self._live: Any = None + self._live_task: Any = None + # The last snapshot read, for the envelope an interrupted watch still owes. + self.last: JsonObject | None = None + # Progress is never worth failing a deploy command over: the first write + # the stream refuses turns reporting off for the rest of the watch. + self._muted = False + self._surface: Surface = surface_for(renderer) + + def snapshot(self, deployment: JsonObject) -> None: + self.last = deployment + progress = progress_of(deployment) + if progress is None: + self.close() + return + now = self._now() + stale = is_stale(progress, now) + if self._surface == "live": + self._update_live(progress, now) + return + # Keyed on the sample, not on the wording: the stale suffix and the + # time so far count seconds, and a line per poll is what this exists to + # avoid. A sample with no stamp is keyed on its content, so a number that + # moved is still reported and a repeat still is not. + by_stamp = self._surface == "events" + key = (str(deployment.get("status")), _sample_key(progress, by_stamp=by_stamp), stale) + if key == self._reported: + return + self._reported = key + if self._surface == "events": + self._event(deployment, progress, stale) + else: + self._say(describe(progress, now=now)) + + def close(self) -> None: + live, self._live, self._live_task = self._live, None, None + if live is None: + return + try: + live.stop() + except OSError: + self._muted = True + + def interrupted(self) -> None: + """Say what Ctrl-C did not do: the deploy runs on the service's side.""" + self.close() + self._say( + f"Stopped watching. Deployment {self._deployment_id} keeps coming up on our side.", + hint=f"run `{reattach_hint(self._deployment_id)}` to watch it again", + ) + + # ----- internals ----- + + def _event(self, deployment: JsonObject, progress: JsonObject, stale: bool) -> None: + if self._muted: + return + try: + self._renderer.progress_event( + EVENT_PROGRESS, + deployment_id=self._deployment_id, + status=deployment.get("status"), + stale=stale, + progress=progress, + ) + except OSError: + self._muted = True + + def _say(self, message: str, *, hint: str | None = None) -> None: + if self._muted: + return + try: + self._renderer.info(message, hint=hint) + except OSError: + self._muted = True + + def _open_live(self) -> None: + from rich.progress import BarColumn, Progress, SpinnerColumn, TextColumn + from rich.table import Column + + # The spinner turns on Rich's own clock, not on the service's writes, so + # a step that counts nothing still shows the command is alive. + # One text column holding the whole sentence, after a short fixed bar. + # Spread over several columns, Rich's table gave the narrow terminal's + # shortfall to whichever column it chose: at 80 columns it cut the label + # and the time left and wrapped the bar onto a line of its own. A single + # column that never wraps is cropped from its right end, and + # `live_line` puts the model's name there. + live = Progress( + SpinnerColumn(), + BarColumn(bar_width=10), + TextColumn("{task.fields[line]}", table_column=Column(no_wrap=True, overflow="ellipsis", ratio=1)), + console=self._renderer.console(), + transient=True, + expand=True, + ) + # Adding the task redraws the display, so it can refuse the stream just + # as starting it can; both stay inside one boundary and nothing is + # published until both have landed. + try: + live.start() + task = live.add_task("", total=None, line="") + except OSError: + self._muted = True + try: + live.stop() + except OSError: + pass + return + self._live = live + self._live_task = task + + def _update_live(self, progress: JsonObject, now: datetime) -> None: + if self._muted: + return + if self._live is None: + self._open_live() + if self._live is None: + return + total = _number(progress, "bytesTotal") if progress.get("step") == STAGING_STEP else None + done = _number(progress, "bytesDone") or 0 + try: + # No total means no fraction to draw, and Rich pulses the bar instead: + # the step is moving, and how far along it is was never measured. + self._live.update( + self._live_task, + total=total if total else None, + completed=min(done, total) if total else 0, + line=sanitize_markup(live_line(progress, now=now)), + ) + except OSError: + self._muted = True diff --git a/comfy_cli/command/deploy_runtime.py b/comfy_cli/command/deploy_runtime.py index 96be6472e..83e2da1b0 100644 --- a/comfy_cli/command/deploy_runtime.py +++ b/comfy_cli/command/deploy_runtime.py @@ -27,7 +27,11 @@ from comfy_cli.output.renderer import Renderer DEPLOY_POLL_SECONDS: Final = 2.0 -_WATCH_TERMINAL: Final = frozenset({"ready", "failed", "stopped", "stop_failed"}) +# Where a watch stops. `unhealthy` is here although the service can still move +# it back to `ready`: it only ever follows `ready`, so a deployment in it has +# already come up, and a watch that waited on it would wait silently for as +# long as the endpoint stays degraded. +_WATCH_TERMINAL: Final = frozenset({"ready", "unhealthy", "failed", "stopped", "stop_failed"}) @dataclass(frozen=True, slots=True) @@ -80,9 +84,21 @@ def terminal_status_error(deployment_id: str, status: str) -> JsonObject: } -def poll_deployment(client: DeployUpClient, deployment_id: str, sleep_fn: Callable[[float], None]) -> JsonObject: +def poll_deployment( + client: DeployUpClient, + deployment_id: str, + sleep_fn: Callable[[float], None], + on_snapshot: Callable[[JsonObject], None] | None = None, +) -> JsonObject: + """Read the deployment until it settles, handing each read to ``on_snapshot``. + + The progress a watcher shows rides the same read the loop already makes, so + watching costs the service nothing it was not already answering. + """ while True: snapshot = client.get_deployment(deployment_id) + if on_snapshot is not None: + on_snapshot(snapshot) status = required_string(snapshot, "status") if status in _WATCH_TERMINAL: return snapshot diff --git a/comfy_cli/command/deploy_status.py b/comfy_cli/command/deploy_status.py index ced1a4ac2..017cbb50e 100644 --- a/comfy_cli/command/deploy_status.py +++ b/comfy_cli/command/deploy_status.py @@ -4,6 +4,7 @@ import urllib.error from dataclasses import dataclass, replace +from datetime import datetime, timezone from typing import Final import typer @@ -11,6 +12,8 @@ from comfy_cli.builder_api import BuilderAuthError from comfy_cli.command.build_paths import BuildSpecNotFoundError, resolve_build_paths from comfy_cli.command.build_spec import BuildSpecInvalidError, JsonObject, read_build_spec +from comfy_cli.command.deploy_progress import DeployWatchReporter, progress_of +from comfy_cli.command.deploy_progress import describe as describe_progress from comfy_cli.command.deploy_resolve import ( _STATUS_RANK, BuilderReleaseClient, @@ -60,14 +63,21 @@ class StatusResult: deployment: JsonObject | None release: JsonObject | None serving: JsonObject | None + # Where a deployment that is coming up has got to, as the service sent it. + progress: JsonObject | None = None def payload(self) -> JsonObject: - return { + payload: JsonObject = { "build": {"id": self.build_id, "name": self.build_name}, "deployment": self.deployment, "release": self.release, "serving": self.serving, } + # Absent rather than null once the deployment has settled, and from a + # service that never sends it. + if self.progress is not None: + payload["progress"] = self.progress + return payload def _nullable_string(value: JsonObject, key: str) -> str | None: @@ -161,9 +171,34 @@ def status_result(builder: BuilderReleaseClient, target: StatusTarget) -> Status _normalized_deployment(deployment), release, _normalized_serving(deployment), + progress_of(deployment), ) +def _interrupted_result(builder: BuilderReleaseClient, target: StatusTarget) -> StatusResult: + """The last read, for a watch the person stopped. + + Its release summary needs one more read of the Build's releases, and a + person often presses Ctrl-C because the network went away. That read is not + worth turning an interrupt (130) into a server error (1): without it the + envelope still carries the deployment and its progress, with no release. + """ + try: + return status_result(builder, target) + except (DeployAPIError, BuilderAuthError, ResponseTooLarge, TimeoutError, urllib.error.URLError, KeyError): + deployment = target.deployment + if deployment is None: + return StatusResult(target.build_id, target.build_name, None, None, None) + return StatusResult( + target.build_id, + target.build_name, + _normalized_deployment(deployment), + None, + _normalized_serving(deployment), + progress_of(deployment), + ) + + def _render_serving(renderer: Renderer, serving: JsonObject | None) -> None: if serving is None: renderer.info("Serving: not sampled yet.") @@ -233,6 +268,8 @@ def render_status(renderer: Renderer, result: StatusResult) -> None: status = _render_deployment(renderer, deployment) if renderer.is_pretty(): + if result.progress is not None: + renderer.info(describe_progress(result.progress, now=datetime.now(timezone.utc))) _render_stop_reason(renderer, deployment) _render_serving(renderer, result.serving) release = result.release @@ -263,7 +300,19 @@ def run_status(path: str | None, *, deployment_id: str | None = None, watch: boo builder, client = _command_clients() target = resolve_status(builder, client, path, deployment_id) if watch and target.deployment is not None: - watched = poll_deployment(client, required_string(target.deployment, "id"), _sleep) + watched_id = required_string(target.deployment, "id") + reporter = DeployWatchReporter(renderer, watched_id) + try: + watched = poll_deployment(client, watched_id, _sleep, reporter.snapshot) + except KeyboardInterrupt: + # Watching is all this command does, so Ctrl-C stops the watching + # and nothing else: the deployment is the service's to bring up. + reporter.interrupted() + if reporter.last is not None: + render_status(renderer, _interrupted_result(builder, replace(target, deployment=reporter.last))) + raise typer.Exit(code=130) from None + finally: + reporter.close() target = replace(target, deployment=watched) render_status(renderer, status_result(builder, target)) except (BuildSpecNotFoundError, BuildSpecInvalidError) as error: diff --git a/comfy_cli/command/deploy_types.py b/comfy_cli/command/deploy_types.py index ded698081..d338b494e 100644 --- a/comfy_cli/command/deploy_types.py +++ b/comfy_cli/command/deploy_types.py @@ -6,6 +6,7 @@ from typing import Protocol from comfy_cli.command.build_spec import JsonObject, JsonValue +from comfy_cli.command.deploy_progress import progress_of from comfy_cli.deploy_api_errors import DeployAPIError @@ -61,12 +62,19 @@ def payload(self) -> JsonObject: "status": required_string(self.deployment, "status"), "created": self.created, } - return { + payload: JsonObject = { "deployment": deployment, "release": self.release, "computeConfig": self.compute_config, "supersedes": supersedes, } + # Present only while the deployment is coming up, and absent rather than + # null otherwise: an older service never sends it, and a settled + # deployment has nothing left to narrate. + progress = progress_of(self.deployment) + if progress is not None: + payload["progress"] = progress + return payload class ComputeRequiredError(Exception): diff --git a/comfy_cli/command/deploy_up.py b/comfy_cli/command/deploy_up.py index bfdd26f57..124beeeb3 100644 --- a/comfy_cli/command/deploy_up.py +++ b/comfy_cli/command/deploy_up.py @@ -194,6 +194,15 @@ def _render_result(renderer, result: UpResult, *, watch: bool) -> None: f"Deployment {deployment_id} could not stop and may still be billing.", hint=f"run `comfy deploy stop --deployment {deployment_id}` again", ) + elif status == "unhealthy": + # `up` leaves an unhealthy deployment as it is, so saying nothing would + # read as success for one that is billing and not serving. + renderer.warn( + f"Deployment {deployment_id} is unhealthy: it came up, then its endpoint degraded. " + "It is still billing, and the service moves it back to ready if it recovers.", + hint=f"run `comfy deploy logs --deployment {deployment_id}` to see why, " + f"or `comfy deploy stop --deployment {deployment_id}` to stop billing", + ) elif watch and status in {"failed", "stopped"}: renderer.warn(f"Deployment {deployment_id} reached terminal status {status}.") if result.dropped_bounds: @@ -208,7 +217,7 @@ def _render_result(renderer, result: UpResult, *, watch: bool) -> None: if status == "stop_failed" else f"run `comfy deploy scale --deployment {deployment_id} --min --max ` to change them", ) - terminal = status in {"failed", "stopped", "stop_failed"} + terminal = status in {"failed", "stopped", "stop_failed", "unhealthy"} renderer.emit( result.payload(), command="deploy up", diff --git a/comfy_cli/discovery.py b/comfy_cli/discovery.py index c92f74bd6..0f404d66e 100644 --- a/comfy_cli/discovery.py +++ b/comfy_cli/discovery.py @@ -216,6 +216,9 @@ "comfy jobs watch": "run_event", # upload progress; under plain --json the same lines go to stderr "comfy build push": "build_push_event", + # the deployment coming up; same rule about stderr + "comfy deploy up": "deploy_progress_event", + "comfy deploy status": "deploy_progress_event", } diff --git a/comfy_cli/error_codes.py b/comfy_cli/error_codes.py index 351c2e5bc..db7881f6c 100644 --- a/comfy_cli/error_codes.py +++ b/comfy_cli/error_codes.py @@ -1440,9 +1440,11 @@ class ErrorCode: "`details.status` names the state. The two commands differ, deliberately: `comfy deploy status` " "reports only `failed` and `stop_failed`, since a `stopped` deployment is a normal thing to be " "asked about; `comfy deploy up` adds `stopped` (with or without `--watch`), because a deployment it was " - "asked to bring up and that is stopped did not come up.", + "asked to bring up and that is stopped did not come up, and `unhealthy`, because one that came up and " + "then degraded is billing without serving and `up` does not change it.", "for `failed`, inspect `comfy deploy logs` and redeploy with `comfy deploy up`; for `stop_failed`, " - "re-run `comfy deploy stop` -- it may still be billing; for `stopped`, `comfy deploy start`", + "re-run `comfy deploy stop` -- it may still be billing; for `stopped`, `comfy deploy start`; for " + "`unhealthy`, inspect `comfy deploy logs`, or `comfy deploy stop` to stop billing", ), ErrorCode( "deploy_delete_needs_confirm", diff --git a/comfy_cli/schemas/deploy_progress_event.json b/comfy_cli/schemas/deploy_progress_event.json new file mode 100644 index 000000000..768a8a3f7 --- /dev/null +++ b/comfy_cli/schemas/deploy_progress_event.json @@ -0,0 +1,94 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://comfy.org/schemas/deploy_progress_event.json", + "title": "comfy deploy up / status --watch progress event", + "description": "One JSON object per line while a watched deployment comes up. On stdout under --json-stream, ahead of the envelope; on stderr under --json, where stdout stays the single envelope. One event per new sample from the deploy service (it rewrites the object about every three seconds in every step), and one more with stale: true if a sample then goes unrefreshed for a minute. A service that sends no progress object produces no events.", + "type": "object", + "required": [ + "schema", + "type", + "deployment_id", + "status", + "stale", + "progress" + ], + "additionalProperties": true, + "properties": { + "schema": { + "const": "event/1" + }, + "type": { + "const": "deploy_progress" + }, + "deployment_id": { + "type": "string", + "minLength": 1 + }, + "status": { + "type": "string", + "description": "The deployment status the sample was read beside: provisioning or starting." + }, + "stale": { + "type": "boolean", + "description": "True when progress.updatedAt is more than 60 seconds old. The numbers are the last the service wrote, not evidence that the deploy stopped." + }, + "progress": { + "type": "object", + "description": "Where a deployment that is coming up has got to, exactly as the deploy service sent it. Present only while the status is provisioning or starting; absent otherwise and from a service that does not send it. The byte fields, bytesPerSecond and etaSeconds are present only while step is staging_models. bytesTotal is absent for a release whose models were never measured and 0 when every model is already in place. New fields may appear.", + "required": [ + "step" + ], + "additionalProperties": true, + "properties": { + "step": { + "type": "string", + "description": "staging_models, creating_endpoint or waiting_for_worker today; read an unknown value as a step name." + }, + "modelsDone": { + "type": "integer", + "minimum": 0 + }, + "modelsTotal": { + "type": "integer", + "minimum": 0 + }, + "bytesDone": { + "type": "integer", + "minimum": 0 + }, + "bytesTotal": { + "type": "integer", + "minimum": 0 + }, + "bytesTotalIsFloor": { + "type": "boolean" + }, + "currentModel": { + "type": "string" + }, + "bytesPerSecond": { + "type": "integer", + "minimum": 0 + }, + "etaSeconds": { + "type": "integer", + "minimum": 0 + }, + "attempt": { + "type": "integer", + "minimum": 1 + }, + "startedAt": { + "type": "string", + "format": "date-time", + "description": "When this attempt of the step began. How long the step has run is measured from here." + }, + "updatedAt": { + "type": "string", + "format": "date-time", + "description": "When the service last wrote this object. It advances every few seconds in every step, so it says how fresh the sample is, not how long the step has run." + } + } + } + } +} diff --git a/comfy_cli/schemas/deploy_status.json b/comfy_cli/schemas/deploy_status.json index e0c7873a5..5e10e8d67 100644 --- a/comfy_cli/schemas/deploy_status.json +++ b/comfy_cli/schemas/deploy_status.json @@ -80,6 +80,26 @@ "jobsInQueue": {"type": "integer", "minimum": 0}, "sampledAt": {"type": "string", "format": "date-time"} } + }, + "progress": { + "type": "object", + "description": "Where a deployment that is coming up has got to, exactly as the deploy service sent it. Present only while the status is provisioning or starting; absent otherwise, and from a service that does not send it. The byte fields, bytesPerSecond and etaSeconds are present only while step is staging_models. bytesTotal is absent for a release whose models were never measured and 0 when every model is already in place. New fields may appear.", + "required": ["step"], + "additionalProperties": true, + "properties": { + "step": {"type": "string", "description": "staging_models, creating_endpoint or waiting_for_worker today; read an unknown value as a step name."}, + "modelsDone": {"type": "integer", "minimum": 0}, + "modelsTotal": {"type": "integer", "minimum": 0}, + "bytesDone": {"type": "integer", "minimum": 0}, + "bytesTotal": {"type": "integer", "minimum": 0}, + "bytesTotalIsFloor": {"type": "boolean"}, + "currentModel": {"type": "string"}, + "bytesPerSecond": {"type": "integer", "minimum": 0}, + "etaSeconds": {"type": "integer", "minimum": 0}, + "attempt": {"type": "integer", "minimum": 1}, + "startedAt": {"type": "string", "format": "date-time", "description": "When this attempt of the step began. How long the step has run is measured from here."}, + "updatedAt": {"type": "string", "format": "date-time", "description": "When the service last wrote this object. It advances every few seconds in every step, so it says how fresh the sample is, not how long the step has run."} + } } } } diff --git a/comfy_cli/schemas/deploy_up.json b/comfy_cli/schemas/deploy_up.json index 82b838fac..c497d1747 100644 --- a/comfy_cli/schemas/deploy_up.json +++ b/comfy_cli/schemas/deploy_up.json @@ -55,6 +55,26 @@ } } } + }, + "progress": { + "type": "object", + "description": "Where a deployment that is coming up has got to, exactly as the deploy service sent it. Present only while the status is provisioning or starting; absent otherwise, and from a service that does not send it. The byte fields, bytesPerSecond and etaSeconds are present only while step is staging_models. bytesTotal is absent for a release whose models were never measured and 0 when every model is already in place. New fields may appear.", + "required": ["step"], + "additionalProperties": true, + "properties": { + "step": {"type": "string", "description": "staging_models, creating_endpoint or waiting_for_worker today; read an unknown value as a step name."}, + "modelsDone": {"type": "integer", "minimum": 0}, + "modelsTotal": {"type": "integer", "minimum": 0}, + "bytesDone": {"type": "integer", "minimum": 0}, + "bytesTotal": {"type": "integer", "minimum": 0}, + "bytesTotalIsFloor": {"type": "boolean"}, + "currentModel": {"type": "string"}, + "bytesPerSecond": {"type": "integer", "minimum": 0}, + "etaSeconds": {"type": "integer", "minimum": 0}, + "attempt": {"type": "integer", "minimum": 1}, + "startedAt": {"type": "string", "format": "date-time", "description": "When this attempt of the step began. How long the step has run is measured from here."}, + "updatedAt": {"type": "string", "format": "date-time", "description": "When the service last wrote this object. It advances every few seconds in every step, so it says how fresh the sample is, not how long the step has run."} + } } } } diff --git a/comfy_cli/skills/comfy-deploy/SKILL.md b/comfy_cli/skills/comfy-deploy/SKILL.md index 7690522ef..ee9d08e6c 100644 --- a/comfy_cli/skills/comfy-deploy/SKILL.md +++ b/comfy_cli/skills/comfy-deploy/SKILL.md @@ -54,7 +54,7 @@ refs compute — deployable regions and GPU classes with availability. ```shell comfy build release show # confirm deployable: true comfy deploy refs compute # read the real GPU/region pairs -comfy deploy up --gpu --region --min 0 --max 1 --watch +comfy deploy up --gpu --region --min 0 --max 1 comfy deploy status comfy deploy run --workflow .json comfy deploy stop # when the user is done @@ -77,7 +77,7 @@ release version. **Read it and act on it.** An empty array means nothing else is running; a non-empty one is a bill the user has not agreed to. ```shell -comfy deploy up --watch # note `supersedes` in the output +comfy deploy up # note `supersedes` in the output comfy deploy stop --deployment ``` @@ -109,8 +109,39 @@ as a pair, `--min` accepts 0–20 and `--max` 1–20. | `stop_failed` | **maybe** | Stop did not take; retry it | | `failed` | no | Permanent failure | -`--watch` on `up` and `status` polls until `ready`, `failed`, `stopped` or -`stop_failed`. The other five are transitional and it keeps waiting. +`up` follows the deployment by default; pass `--no-watch` to return as soon as it +is accepted. `status` waits only when asked, with `--watch`. Either way the wait +ends at `ready`, `unhealthy`, `failed`, `stopped` or `stop_failed`. `unhealthy` +only ever follows `ready`, so the deployment already came up: `up` reports it as +not ok (`deploy_status_terminal`, exit 1) because it is billing without serving, +and `status` reports it as recoverable. Through `queued`, `provisioning`, +`starting` and `stopping` it keeps waiting. + +While the status is `provisioning` or `starting` the deployment carries a +`progress` object, and `status --json` returns it as `data.progress`: `step` +(`staging_models`, `creating_endpoint`, `waiting_for_worker`) and, while models +are copied onto the deployment's storage, `modelsDone` / `modelsTotal`, +`bytesDone` / `bytesTotal`, `currentModel`, `bytesPerSecond` and `etaSeconds`. +Under `--watch` the same object arrives as `deploy_progress` events, one per new +sample (the service rewrites it about every three seconds in every step): on **stderr** under `--json`, +on stdout under `--json-stream`. Relay those numbers instead of "still +provisioning". Things to read correctly: + +- `bytesTotal` absent means nobody measured the release, so there is no time + left to quote; `0` means every model was already in place. +- `bytesTotalIsFloor: true` means "at least this much", with no `etaSeconds`. +- `etaSeconds` covers staging only. Creating the endpoint and the first worker's + cold start come after it. +- `attempt` above 1 means the step was restarted, and `bytesDone` started again. +- `stale: true` on an event means the service has not rewritten the sample for a + minute. Its writes are best-effort, so that is not evidence the deploy stopped; + the status is still the verdict. +- How long a step has run is now minus `startedAt`. `updatedAt` only says how + fresh the sample is, and it moves every few seconds in every step. +- No `progress` at all is an older service, or a status other than the two above. + +Interrupting the wait leaves the deployment coming up on the service's side; +`comfy deploy status --deployment --watch` attaches again. `status` also reports **why** a deployment stopped, as `stopReason`: `user`, `credits`, or `policy`. `credits` is a billing problem and not something a retry @@ -120,7 +151,7 @@ fixes — say so rather than restarting into the same wall. ```shell comfy deploy up [PATH] --gpu --region [--min N --max N] - [--release ] [--deployment ] [--watch] + [--release ] [--deployment ] [--no-watch] ``` - **It selects the newest deployable release of the Build** unless `--release` diff --git a/docs/json-output.md b/docs/json-output.md index 1d3904a86..90c9ae5b9 100644 --- a/docs/json-output.md +++ b/docs/json-output.md @@ -432,6 +432,46 @@ connection, which runs slightly ahead of bytes the server has acknowledged. At a terminal the same numbers render as one redrawn progress line; with output piped under `--no-json` they are plain lines every five seconds, with no carriage returns. +### `deploy_progress` + +Emitted by `comfy deploy up --watch` and `comfy deploy status --watch` while the +deployment's status is `provisioning` or `starting`. It is part of those +commands, not the `run` stream, and validates against +`deploy_progress_event.json`. + +Where it goes depends on the mode, the same way upload progress does. Under +`--json-stream` it is on stdout, ahead of the envelope. Under plain `--json` +(what a caller with a piped stdout resolves to, so every agent) it is on +**stderr**, because stdout in that mode is exactly one envelope. + +```json +{"schema": "event/1", "type": "deploy_progress", "deployment_id": "dep-7edc1262", "status": "provisioning", "stale": false, "progress": {"step": "staging_models", "modelsDone": 0, "modelsTotal": 2, "bytesDone": 3536540667, "bytesTotal": 7272719498, "currentModel": "models/checkpoints/sd_xl_base_1.0.safetensors", "bytesPerSecond": 44205525, "etaSeconds": 85, "attempt": 1, "startedAt": "2026-09-20T01:59:55Z", "updatedAt": "2026-09-20T02:01:15Z"}} +{"schema": "event/1", "type": "deploy_progress", "deployment_id": "dep-7edc1262", "status": "starting", "stale": false, "progress": {"step": "waiting_for_worker", "attempt": 1, "startedAt": "2026-09-20T02:02:46Z", "updatedAt": "2026-09-20T02:02:46Z"}} +``` + +The event's own fields are snake_case like every other event. `progress` is the +deploy service's object, passed through unchanged, so its fields are the API's +camelCase, the same as `computeConfig` and `endpointUrl` in the envelopes. The +service computes `bytesPerSecond` and `etaSeconds`; the CLI computes nothing, so +the portal and the CLI show the same numbers. + +One event per new sample: the CLI polls every two seconds and the service writes +about every three, and a sample is keyed on its `updatedAt` (on its whole content +where a sample carries no `updatedAt`). If a sample then goes +a minute without being rewritten, one more event carries it with `stale: true`. +The service's writes are best-effort, so a stale sample is not evidence that the +deploy stopped; `status` remains the verdict. A service that sends no `progress` +object produces no events, and the command behaves as it did before. + +`comfy deploy status --json` (with or without `--watch`) and `comfy deploy up +--json` carry the same object as `data.progress` while the deployment is coming +up, and omit the key otherwise. + +At a terminal the numbers render as one redrawn progress line; with output piped +under `--no-json` they are plain lines, one per new sample, with no carriage +returns. Ctrl-C during `--watch` exits 130 after printing that the deployment +keeps coming up and the command that re-attaches, and in the JSON modes still +writes the envelope for the last state it read. ## Success envelope diff --git a/tests/comfy_cli/command/test_deploy_progress.py b/tests/comfy_cli/command/test_deploy_progress.py new file mode 100644 index 000000000..d50fcd6b0 --- /dev/null +++ b/tests/comfy_cli/command/test_deploy_progress.py @@ -0,0 +1,727 @@ +"""`comfy deploy up --watch` and `comfy deploy status` show a deployment's progress.""" + +from __future__ import annotations + +import copy +import errno +import importlib +import inspect +import io +import json +import re +import urllib.error +from datetime import datetime, timedelta, timezone +from pathlib import Path + +import jsonschema +import pytest +from deploy_up_support import FakeBuilder, FakeDeploy, deployment, write_spec +from typer.testing import CliRunner + +from comfy_cli.cmdline import app +from comfy_cli.command.build_spec import JsonObject +from comfy_cli.command.deploy_progress import ( + EVENT_PROGRESS, + STALE_SECONDS, + DeployWatchReporter, + describe, + is_stale, + progress_of, +) +from comfy_cli.command.deploy_runtime import poll_deployment +from comfy_cli.output.renderer import OutputMode, Renderer + +GB = 1 << 30 +NOW = datetime(2026, 9, 20, 2, 1, 20, tzinfo=timezone.utc) + + +def _staging(**over) -> JsonObject: + progress: JsonObject = { + "step": "staging_models", + "modelsDone": 0, + "modelsTotal": 2, + "bytesDone": 3 * GB, + "bytesTotal": 7 * GB, + "currentModel": "models/checkpoints/sd_xl_base_1.0.safetensors", + "bytesPerSecond": 44 * (1 << 20), + "etaSeconds": 85, + "attempt": 1, + "startedAt": "2026-09-20T01:59:55Z", + "updatedAt": "2026-09-20T02:01:15Z", + } + progress.update(over) + return {key: value for key, value in progress.items() if value is not None} + + +def _step(step: str, **over) -> JsonObject: + # The service re-stamps a running step every few seconds, the two that count + # nothing included, so a live sample's updatedAt sits just behind NOW + # whatever its startedAt says. How long the step has run is startedAt's. + return { + "step": step, + "attempt": 1, + "startedAt": "2026-09-20T02:01:20Z", + "updatedAt": "2026-09-20T02:01:18Z", + **over, + } + + +# ----- the wording ----- + + +@pytest.mark.parametrize( + ("progress", "now", "expected"), + [ + pytest.param( + _staging(), + NOW, + "Staging models: model 1 of 2 sd_xl_base_1.0.safetensors, 3.0 GB of 7.0 GB, 44.0 MB/s, 1m 25s left", + id="a model in flight", + ), + pytest.param( + _staging(bytesTotal=None, etaSeconds=None), + NOW, + "Staging models: model 1 of 2 sd_xl_base_1.0.safetensors, 3.0 GB copied, 44.0 MB/s", + id="a release nobody measured has no total and no time left", + ), + pytest.param( + _staging(bytesTotalIsFloor=True, etaSeconds=None), + NOW, + "Staging models: model 1 of 2 sd_xl_base_1.0.safetensors, 3.0 GB of at least 7.0 GB, 44.0 MB/s", + id="a floor says at least", + ), + pytest.param( + _staging(bytesDone=0, bytesTotal=0, currentModel=None, bytesPerSecond=None, etaSeconds=None, modelsDone=2), + NOW, + "Staging models: every model is already in place", + id="nothing to copy", + ), + pytest.param( + _staging(bytesPerSecond=None, etaSeconds=None, bytesDone=0), + NOW, + "Staging models: model 1 of 2 sd_xl_base_1.0.safetensors, 0 B of 7.0 GB", + id="the first seconds have no rate", + ), + pytest.param( + _step("creating_endpoint"), + NOW, + "Creating the endpoint", + id="a step with nothing to count, the moment it begins", + ), + pytest.param( + _step("waiting_for_worker", startedAt="2026-09-20T01:59:05Z"), + NOW, + "Waiting for the first worker: 2m 15s so far", + id="a step with nothing to count says how long it has run, from its start, not its last write", + ), + pytest.param( + _step("waiting_for_worker", startedAt="2026-09-20T01:58:00Z"), + NOW, + "Waiting for the first worker: 3m 20s so far", + id="a long wait the service keeps re-stamping is long, not stale", + ), + pytest.param( + _step("waiting_for_worker", startedAt="2026-09-20T01:58:00Z", updatedAt="2026-09-20T01:59:00Z"), + NOW, + "Waiting for the first worker: 3m 20s so far (last update 2m 20s ago, so these numbers may be stale)", + id="a step the service stopped re-stamping goes stale like any other", + ), + pytest.param( + _step("waiting_for_worker", startedAt="2026-09-20T01:59:05.43745Z", updatedAt="2026-09-20T02:01:18.5Z"), + NOW, + "Waiting for the first worker: 2m 15s so far", + id="a stamp with five fractional digits, which Python 3.10 cannot read alone", + ), + pytest.param( + _step("waiting_for_worker", attempt=2), + NOW, + "Waiting for the first worker (attempt 2, this step was restarted)", + id="a retried step says so", + ), + pytest.param( + _step("warming_cache"), + NOW, + "warming cache", + id="a step this version never heard of", + ), + pytest.param( + _staging(), + NOW + timedelta(seconds=STALE_SECONDS + 35), + "Staging models: model 1 of 2 sd_xl_base_1.0.safetensors, 3.0 GB of 7.0 GB, 44.0 MB/s, 1m 25s left" + " (last update 1m 40s ago, so these numbers may be stale)", + id="an old sample is called old, and its numbers are not moved", + ), + ], +) +def test_describe(progress, now, expected) -> None: + assert describe(progress, now=now) == expected + + +@pytest.mark.parametrize("value", [None, "staging", 3, [], {}, {"step": 7}]) +def test_anything_that_is_not_a_progress_object_reads_as_none(value) -> None: + """Progress narrates a deploy. A shape this build cannot read must cost the + narration, never the command.""" + row = deployment("dep-1") + row["progress"] = value + assert progress_of(row) is None + + +def test_a_deployment_that_sends_no_progress_has_none() -> None: + assert progress_of(deployment("dep-1")) is None + + +@pytest.mark.parametrize("status", ["ready", "failed", "stopped", "unhealthy", None]) +def test_a_settled_deployment_has_no_progress_whatever_it_carries(status) -> None: + """The service sends progress only while a deployment is provisioning or + starting. One that settled and still carries an object must not narrate.""" + row = deployment("dep-1", status="provisioning") + row["progress"] = _staging() + assert progress_of(row) is not None + row["status"] = status + assert progress_of(row) is None + + +# ----- the reporter ----- + + +class _Clock: + def __init__(self) -> None: + self.now = NOW + + def __call__(self) -> datetime: + return self.now + + +def _renderer(mode: OutputMode, tmp_path: Path): + out, err = (tmp_path / "out").open("w+"), (tmp_path / "err").open("w+") + renderer = Renderer(mode=mode) + renderer.pretty_stream = out if mode is OutputMode.PRETTY else err + renderer.machine_stream = out + return renderer, out, err + + +def _read(handle) -> str: + handle.flush() + handle.seek(0) + return handle.read() + + +def _snapshot(status: str, progress: JsonObject | None) -> JsonObject: + row = deployment("dep-1", status=status) + if progress is not None: + row["progress"] = progress + return row + + +def test_a_sample_with_no_stamp_is_reported_when_its_numbers_move_and_not_when_they_repeat(tmp_path) -> None: + """`updatedAt` is optional to the lenient parser; a sample without one must + not be mistaken for the last one just because both have no stamp.""" + # Given samples the service never stamped + renderer, out, _ = _renderer(OutputMode.NDJSON, tmp_path) + reporter = DeployWatchReporter(renderer, "dep-1", now=_Clock()) + first = _staging(updatedAt=None) + second = _staging(bytesDone=4 * GB, updatedAt=None) + + # When + for progress in (first, first, second, second, first): + reporter.snapshot(_snapshot("provisioning", progress)) + + # Then: one event per change, none per repeat + events = [json.loads(line) for line in _read(out).splitlines()] + assert [event["progress"]["bytesDone"] for event in events] == [3 * GB, 4 * GB, 3 * GB] + + +class _LiveRenderer: + """A pretty renderer on a terminal, so the reporter picks the live surface.""" + + class _Console: + is_terminal = True + + def __init__(self) -> None: + self.said: list[str] = [] + + def is_pretty(self) -> bool: + return True + + def console(self) -> _LiveRenderer._Console: + return self._Console() + + def info(self, message: str, *, hint: str | None = None) -> None: + self.said.append(message) + + +class _DisplayThatRefusesTheTask: + """Rich's Progress with a stream that fails on the first redraw after start.""" + + stopped = 0 + + def __init__(self, *columns: object, **options: object) -> None: + pass + + def start(self) -> None: + pass + + def add_task(self, description: str, **fields: object) -> int: + raise OSError(errno.EIO, "broken terminal") + + def stop(self) -> None: + type(self).stopped += 1 + + +def test_a_display_refused_at_the_first_redraw_mutes_the_watch_instead_of_raising(monkeypatch) -> None: + # Given a terminal whose stream refuses the redraw that adding the task makes + import rich.progress + + monkeypatch.setattr(rich.progress, "Progress", _DisplayThatRefusesTheTask) + _DisplayThatRefusesTheTask.stopped = 0 + renderer = _LiveRenderer() + reporter = DeployWatchReporter(renderer, "dep-1", now=_Clock()) + + # When: two samples arrive and nothing raises + reporter.snapshot(_snapshot("provisioning", _staging())) + reporter.snapshot(_snapshot("provisioning", _staging(bytesDone=4 * GB))) + reporter.close() + + # Then: the half-opened display was closed once, never reopened, and the + # reporter stayed quiet on the stream it could not write to + assert reporter._muted is True + assert reporter._live is None and reporter._live_task is None + assert _DisplayThatRefusesTheTask.stopped == 1 + assert renderer.said == [] + + +def test_an_ndjson_watch_emits_one_event_per_new_sample_and_carries_the_object_unchanged(tmp_path) -> None: + # Given a poll five times as fast as the service writes + renderer, out, _ = _renderer(OutputMode.NDJSON, tmp_path) + reporter = DeployWatchReporter(renderer, "dep-1", now=_Clock()) + first, second = _staging(), _staging(bytesDone=4 * GB, updatedAt="2026-09-20T02:01:25Z") + + # When + for progress in (first, first, first, second, second): + reporter.snapshot(_snapshot("provisioning", progress)) + + # Then + events = [json.loads(line) for line in _read(out).splitlines()] + assert [event["progress"] for event in events] == [first, second] + assert events[0] == { + "schema": "event/1", + "type": EVENT_PROGRESS, + "deployment_id": "dep-1", + "status": "provisioning", + "stale": False, + "progress": first, + } + schema_path = Path(__file__).parents[3] / "comfy_cli" / "schemas" / "deploy_progress_event.json" + validator = jsonschema.Draft202012Validator(json.loads(schema_path.read_text(encoding="utf-8"))) + for event in events: + validator.validate(event) + + +def test_a_sample_that_goes_stale_is_reported_once_more_and_never_again(tmp_path) -> None: + # Given + renderer, out, _ = _renderer(OutputMode.NDJSON, tmp_path) + clock = _Clock() + reporter = DeployWatchReporter(renderer, "dep-1", now=clock) + sample = _snapshot("provisioning", _staging()) + + # When the service stops writing for two minutes + reporter.snapshot(sample) + clock.now = NOW + timedelta(seconds=120) + reporter.snapshot(sample) + clock.now = NOW + timedelta(seconds=122) + reporter.snapshot(sample) + + # Then + events = [json.loads(line) for line in _read(out).splitlines()] + assert [event["stale"] for event in events] == [False, True] + assert events[0]["progress"] == events[1]["progress"], "a stale sample's numbers are not moved" + + +def test_a_json_watch_puts_events_on_stderr_so_stdout_stays_one_envelope(tmp_path, capsys) -> None: + # Given + renderer = Renderer(mode=OutputMode.JSON) + reporter = DeployWatchReporter(renderer, "dep-1", now=_Clock()) + + # When + reporter.snapshot(_snapshot("provisioning", _staging())) + + # Then + captured = capsys.readouterr() + assert captured.out == "" + assert json.loads(captured.err)["type"] == EVENT_PROGRESS + + +def test_piped_pretty_output_is_plain_lines_with_no_carriage_returns(tmp_path) -> None: + # Given + renderer, out, _ = _renderer(OutputMode.PRETTY, tmp_path) + reporter = DeployWatchReporter(renderer, "dep-1", now=_Clock()) + + # When the service re-stamps the last step twice with nothing else moving + reporter.snapshot(_snapshot("provisioning", _staging())) + reporter.snapshot(_snapshot("provisioning", _staging())) + reporter.snapshot(_snapshot("starting", _step("waiting_for_worker"))) + reporter.snapshot(_snapshot("starting", _step("waiting_for_worker", updatedAt="2026-09-20T02:01:19Z"))) + reporter.snapshot(_snapshot("starting", _step("waiting_for_worker", updatedAt="2026-09-20T02:01:20Z"))) + reporter.snapshot(_snapshot("ready", None)) + reporter.close() + + # Then + text = _read(out) + assert "\r" not in text + assert text.count("Staging models:") == 1 + assert text.count("Waiting for the first worker") == 1 + + +def test_an_older_service_that_sends_no_progress_reports_nothing(tmp_path) -> None: + # Given + renderer, out, err = _renderer(OutputMode.NDJSON, tmp_path) + reporter = DeployWatchReporter(renderer, "dep-1", now=_Clock()) + + # When + for status in ("queued", "provisioning", "starting", "ready"): + reporter.snapshot(_snapshot(status, None)) + + # Then + assert _read(out) == "" and _read(err) == "" + + +def test_a_stream_that_refuses_a_write_mutes_the_reporter_instead_of_failing_the_watch(tmp_path) -> None: + # Given + class Refusing(Renderer): + def progress_event(self, type: str, **fields) -> None: + raise BrokenPipeError + + reporter = DeployWatchReporter(Refusing(mode=OutputMode.NDJSON), "dep-1", now=_Clock()) + + # When / Then + reporter.snapshot(_snapshot("provisioning", _staging())) + reporter.snapshot(_snapshot("provisioning", _staging(updatedAt="2026-09-20T02:01:25Z"))) + + +def test_poll_deployment_hands_every_read_to_the_watcher_the_terminal_one_included() -> None: + # Given + client = FakeDeploy([deployment("dep-1", status="queued")], get_statuses=["provisioning", "starting", "ready"]) + seen: list[str] = [] + + # When + final = poll_deployment(client, "dep-1", lambda _: None, lambda row: seen.append(str(row["status"]))) + + # Then + assert seen == ["provisioning", "starting", "ready"] + assert final["status"] == "ready" + + +# ----- the commands ----- + + +class ProgressDeploy(FakeDeploy): + """Serves one (status, progress) pair per read, then repeats the last.""" + + def __init__(self, rows: list[JsonObject], reads: list[tuple[str, JsonObject | None]]) -> None: + super().__init__(rows) + self.reads = list(reads) + + def get_deployment(self, deployment_id: str) -> JsonObject: + with self._lock: + row = self.rows[deployment_id] + if self.reads: + status, progress = self.reads.pop(0) + row["status"] = status + row.pop("progress", None) + if progress is not None: + row["progress"] = copy.deepcopy(progress) + return copy.deepcopy(row) + + +def _status_row(status: str) -> JsonObject: + row = deployment("dep-status", release_id="release-5", status=status, maximum=2) + row.update({"releaseId": "release-5", "endpointUrl": None, "error": None, "serving": None, "stopReason": None}) + return row + + +def _release(version: int) -> JsonObject: + return {"id": f"release-{version}", "buildId": "build-1", "version": version, "deployable": True} + + +def _install_status(monkeypatch, client: FakeDeploy, sleep) -> None: + module = importlib.import_module("comfy_cli.command.deploy_status") + monkeypatch.setattr(module, "_command_clients", lambda: (FakeBuilder([_release(5)]), client)) + monkeypatch.setattr(module, "_sleep", sleep) + + +def _envelope(result) -> dict: + return json.loads([line for line in result.stdout.splitlines() if line.strip()][-1]) + + +def _schema(name: str) -> JsonObject: + path = Path(__file__).parents[3] / "comfy_cli" / "schemas" / name + return json.loads(path.read_text(encoding="utf-8")) + + +def test_status_carries_the_progress_object_while_the_deployment_comes_up(tmp_path, monkeypatch) -> None: + # Given + row = _status_row("provisioning") + row["progress"] = _staging() + _install_status(monkeypatch, FakeDeploy([row]), lambda _: None) + + # When + result = CliRunner().invoke(app, ["--json", "deploy", "status", str(write_spec(tmp_path))]) + + # Then + assert result.exit_code == 0, result.stderr + data = _envelope(result)["data"] + assert data["progress"] == _staging() + jsonschema.Draft202012Validator(_schema("deploy_status.json")).validate(data) + + +def test_status_omits_progress_once_the_deployment_is_ready(tmp_path, monkeypatch) -> None: + # Given + _install_status(monkeypatch, FakeDeploy([_status_row("ready")]), lambda _: None) + + # When + result = CliRunner().invoke(app, ["--json", "deploy", "status", str(write_spec(tmp_path))]) + + # Then + assert result.exit_code == 0, result.stderr + assert "progress" not in _envelope(result)["data"] + + +def test_status_prints_the_progress_line_for_a_person(tmp_path, monkeypatch) -> None: + # Given + row = _status_row("provisioning") + row["progress"] = _staging(updatedAt=datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")) + _install_status(monkeypatch, FakeDeploy([row]), lambda _: None) + + # When + result = CliRunner().invoke( + app, ["--no-json", "deploy", "status", str(write_spec(tmp_path))], env={"COLUMNS": "400"} + ) + + # Then + assert result.exit_code == 0 + assert "Deployment dep-status: provisioning" in result.stdout + assert "Staging models: model 1 of 2 sd_xl_base_1.0.safetensors, 3.0 GB of 7.0 GB" in result.stdout + assert "may be stale" not in result.stdout + + +def test_status_watch_streams_progress_then_settles_on_an_envelope_without_it(tmp_path, monkeypatch) -> None: + # Given + reads = [ + ("provisioning", _staging()), + ("provisioning", _staging(bytesDone=6 * GB, updatedAt="2026-09-20T02:01:25Z")), + ("starting", _step("waiting_for_worker")), + ("ready", None), + ] + _install_status(monkeypatch, ProgressDeploy([_status_row("queued")], reads), lambda _: None) + + # When + result = CliRunner().invoke(app, ["--json-stream", "deploy", "status", str(write_spec(tmp_path)), "--watch"]) + + # Then + assert result.exit_code == 0, result.stderr + lines = [json.loads(line) for line in result.stdout.splitlines() if line.strip()] + events = [line for line in lines if line.get("type") == EVENT_PROGRESS] + assert [event["progress"]["step"] for event in events] == [ + "staging_models", + "staging_models", + "waiting_for_worker", + ] + assert lines[-1]["type"] == "envelope" + assert lines[-1]["data"]["deployment"]["status"] == "ready" + assert "progress" not in lines[-1]["data"] + + +def test_interrupting_a_status_watch_leaves_the_deployment_alone_and_says_how_to_reattach( + tmp_path, monkeypatch +) -> None: + # Given a watch whose first wait is interrupted + client = ProgressDeploy([_status_row("queued")], [("provisioning", _staging())]) + + def interrupt(_: float) -> None: + raise KeyboardInterrupt + + _install_status(monkeypatch, client, interrupt) + + # When + result = CliRunner().invoke( + app, ["--json", "deploy", "status", str(write_spec(tmp_path)), "--watch"], env={"COLUMNS": "400"} + ) + + # Then + assert result.exit_code == 130 + assert "keeps coming up on our side" in result.stderr + assert "comfy deploy status --deployment dep-status --watch" in result.stderr + data = _envelope(result)["data"] + assert data["deployment"]["status"] == "provisioning" + assert data["progress"]["step"] == "staging_models" + assert client.update_calls == [] and client.start_calls == [] + + +def test_interrupting_a_status_watch_with_the_network_gone_still_exits_130(tmp_path, monkeypatch) -> None: + """A person often presses Ctrl-C because the network went away. The release + summary the envelope would carry needs one more read, and losing it must not + turn the interrupt into a server error.""" + + # Given a builder that stops answering once the watch has begun + class BuilderGoneOffline(FakeBuilder): + def list_releases(self, build_id: str) -> list[JsonObject]: + if any(call[0] == "list_releases" for call in self.calls): + raise urllib.error.URLError("network is unreachable") + return super().list_releases(build_id) + + client = ProgressDeploy([_status_row("queued")], [("provisioning", _staging())]) + + def interrupt(_: float) -> None: + raise KeyboardInterrupt + + module = importlib.import_module("comfy_cli.command.deploy_status") + monkeypatch.setattr(module, "_command_clients", lambda: (BuilderGoneOffline([_release(5)]), client)) + monkeypatch.setattr(module, "_sleep", interrupt) + + # When + result = CliRunner().invoke( + app, ["--json", "deploy", "status", str(write_spec(tmp_path)), "--watch"], env={"COLUMNS": "400"} + ) + + # Then + assert result.exit_code == 130, result.stderr + data = _envelope(result)["data"] + assert data["deployment"]["status"] == "provisioning" + assert data["release"] is None + assert data["progress"]["step"] == "staging_models" + jsonschema.Draft202012Validator(_schema("deploy_status.json")).validate(data) + + +def test_up_watch_streams_progress_and_interrupting_it_stops_nothing(tmp_path, monkeypatch) -> None: + # Given + module = importlib.import_module("comfy_cli.command.deploy") + + class CreatedThenStaging(ProgressDeploy): + def get_deployment(self, deployment_id: str) -> JsonObject: + # The create path reads once before any watch; leave that read alone. + if not self.confirmed: + self.confirmed = True + return FakeDeploy.get_deployment(self, deployment_id) + return super().get_deployment(deployment_id) + + client = CreatedThenStaging([], [("provisioning", _staging())]) + client.confirmed = False + monkeypatch.setattr(module, "_command_clients", lambda: (FakeBuilder(), client)) + + def interrupt(_: float) -> None: + raise KeyboardInterrupt + + monkeypatch.setattr(module, "_sleep", interrupt) + + # When + result = CliRunner().invoke( + app, + ["--json", "deploy", "up", str(write_spec(tmp_path)), "--gpu", "l4", "--region", "US-MO-2", "--watch"], + env={"COLUMNS": "400"}, + ) + + # Then + assert result.exit_code == 130 + events = [json.loads(line) for line in result.stderr.splitlines() if line.startswith("{")] + assert [event["type"] for event in events] == [EVENT_PROGRESS] + assert "comfy deploy status --deployment dep-1 --watch" in result.stderr + data = _envelope(result)["data"] + assert data["deployment"]["status"] == "provisioning" + assert data["progress"]["bytesDone"] == 3 * GB + jsonschema.Draft202012Validator(_schema("deploy_up.json")).validate(data) + + +def test_every_step_goes_stale_when_the_service_stops_rewriting_it() -> None: + """The service rewrites the object every few seconds in every step, so a + minute of silence in any of them is a read worth doubting. A step that + never went stale would show a dead wait as a live one for ever.""" + # Given a sample from each step, none rewritten for two minutes + old = NOW + timedelta(seconds=120) + + # Then + assert is_stale(_staging(), old) is True + assert is_stale(_step("creating_endpoint"), old) is True + assert is_stale(_step("waiting_for_worker"), old) is True + assert "may be stale" in describe(_step("waiting_for_worker"), now=old) + assert is_stale(_step("waiting_for_worker"), NOW) is False + + +class _TerminalRenderer: + """A pretty renderer drawing on a real Rich console of a given width.""" + + def __init__(self, width: int) -> None: + from rich.console import Console + + self._console = Console(file=io.StringIO(), width=width, force_terminal=True, color_system=None) + + def is_pretty(self) -> bool: + return True + + def console(self): + return self._console + + def info(self, message: str, *, hint: str | None = None) -> None: + raise AssertionError("the live surface draws, it does not print lines") + + def screen(self) -> list[str]: + text = re.sub(r"\x1b\[[0-9;?]*[A-Za-z]", "", self._console.file.getvalue()).replace("\r", "\n") + return [line for line in text.splitlines() if line.strip()] + + +def _draw(progress: JsonObject, *, width: int, now: datetime) -> str: + renderer = _TerminalRenderer(width) + reporter = DeployWatchReporter(renderer, "dep-1", now=lambda: now) + reporter.snapshot(_snapshot("provisioning", progress)) + reporter._live.refresh() + reporter.close() + return renderer.screen()[-1] + + +def test_at_80_columns_the_live_line_keeps_the_time_left_and_cuts_the_model_name() -> None: + # Given a model name longer than the space left over + progress = _staging(currentModel="models/checkpoints/sd_xl_base_1.0_with_a_long_finetune_name.safetensors") + + # When + line = _draw(progress, width=80, now=NOW) + + # Then: one line, the step and every number on it, the name cut instead + assert len(line) <= 80 + assert "Staging models: 3.0 GB of 7.0 GB, 44.0 MB/s, 1m 25s left" in line + assert "safetensors" not in line + + +def test_the_live_line_says_a_step_restarted_and_a_sample_went_quiet() -> None: + # Given a restarted step the service has not rewritten for four minutes + progress = _staging(attempt=2) + + # When + line = _draw(progress, width=80, now=NOW + timedelta(minutes=4)) + + # Then: the notes ride on the label, where a narrow terminal keeps them + assert "Staging models (attempt 2, no update for 4m 05s): 3.0 GB of 7.0 GB" in line + + +def test_deploy_up_watches_without_being_asked(tmp_path) -> None: + """`up` starts a wait of several minutes, so it follows it by default. + + Asserted through the help a person actually reads, not through Typer's + parameter objects: what matters is that the pair is offered and that the + default shown is to watch. + """ + # Given + from comfy_cli.command import deploy as deploy_module + + # When + help_text = CliRunner().invoke(app, ["deploy", "up", "--help"]).stdout + + # Then + # Rich draws the help in a box, puts the two spellings of the flag in + # separate cells, and wraps at whatever width it detects (80 on CI), so a + # flag can straddle a line break. Compare with every escape code, box + # character and space removed. + flat = re.sub(r"\s+|[\u2500-\u257f]", "", re.sub(r"\x1b\[[0-9;]*m", "", help_text)) + assert "--watch" in flat and "--no-watch" in flat + assert "[default:watch]" in flat + assert inspect.signature(deploy_module.up_cmd).parameters["watch"].default is True + # `status` answers a question and exits, so there watching stays opt-in. + assert inspect.signature(deploy_module.status_cmd).parameters["watch"].default is False diff --git a/tests/comfy_cli/command/test_deploy_status.py b/tests/comfy_cli/command/test_deploy_status.py index c67157015..793be4291 100644 --- a/tests/comfy_cli/command/test_deploy_status.py +++ b/tests/comfy_cli/command/test_deploy_status.py @@ -10,7 +10,6 @@ from comfy_cli.cmdline import app from comfy_cli.command.build_spec import JsonObject -from comfy_cli.command.deploy_runtime import DEPLOY_POLL_SECONDS class RecordingDeploy(FakeDeploy): @@ -273,7 +272,9 @@ def test_credit_stop_is_not_attributed_to_the_user(tmp_path, monkeypatch) -> Non assert "user-initiated" not in rendered -def test_watch_continues_after_first_unhealthy_sample(tmp_path, monkeypatch) -> None: +def test_watch_stops_at_unhealthy_and_calls_it_recoverable(tmp_path, monkeypatch) -> None: + """`unhealthy` only follows `ready`: the deployment came up. Waiting on it + for `ready` would wait silently for as long as the endpoint is degraded.""" # Given client = RecordingDeploy([_status_deployment(status="queued")], get_statuses=["unhealthy", "ready"]) sleeps: list[float] = [] @@ -284,9 +285,9 @@ def test_watch_continues_after_first_unhealthy_sample(tmp_path, monkeypatch) -> # Then assert result.exit_code == 0 - assert _json_envelope(result)["data"]["deployment"]["status"] == "ready" - assert client.get_calls == ["dep-status", "dep-status"] - assert sleeps == [DEPLOY_POLL_SECONDS] + assert _json_envelope(result)["data"]["deployment"]["status"] == "unhealthy" + assert client.get_calls == ["dep-status"] + assert sleeps == [] def test_watch_exits_promptly_on_stop_failed_with_retry_stop_hint(tmp_path, monkeypatch) -> None: diff --git a/tests/comfy_cli/command/test_deploy_up.py b/tests/comfy_cli/command/test_deploy_up.py index 9b5437717..401809afd 100644 --- a/tests/comfy_cli/command/test_deploy_up.py +++ b/tests/comfy_cli/command/test_deploy_up.py @@ -14,7 +14,6 @@ from comfy_cli.caller import Caller from comfy_cli.cmdline import app -from comfy_cli.command.deploy_runtime import DEPLOY_POLL_SECONDS from comfy_cli.deploy_api import _validate_compute_config from comfy_cli.deploy_api_errors import DeployAPIError @@ -443,7 +442,12 @@ def test_the_dropped_bound_warning_reaches_a_json_caller_on_stderr(tmp_path, mon monkeypatch.setattr(module, "_command_clients", lambda: (FakeBuilder(), client)) # When - result = CliRunner().invoke(app, ["--json", "deploy", "up", str(write_spec(tmp_path)), "--min", "3", "--max", "8"]) + # --no-watch because this is about the warning, not the wait. Restarting a + # stopped deployment leaves it `queued`, and watch (now the default) polls a + # fake that never leaves that status, so the run would never end. + result = CliRunner().invoke( + app, ["--json", "deploy", "up", str(write_spec(tmp_path)), "--min", "3", "--max", "8", "--no-watch"] + ) # Then assert "--min had no effect" in result.stderr @@ -579,8 +583,8 @@ def test_watch_exits_immediately_on_stop_failed_with_stop_remedy(tmp_path, monke assert sleeps == [] -def test_watch_continues_through_unhealthy_until_ready(tmp_path, monkeypatch) -> None: - # Given +def test_watch_stops_at_unhealthy_and_reports_it_as_not_ok(tmp_path, monkeypatch) -> None: + # Given a deployment the watch reads as unhealthy module = _deploy() client = FakeDeploy(get_statuses=["queued", "unhealthy", "ready"]) monkeypatch.setattr(module, "_command_clients", lambda: (FakeBuilder(), client)) @@ -591,9 +595,43 @@ def test_watch_continues_through_unhealthy_until_ready(tmp_path, monkeypatch) -> result = CliRunner().invoke( app, ["--json", "deploy", "up", str(write_spec(tmp_path)), "--gpu", "l4", "--region", "US-MO-2", "--watch"], + env={"COLUMNS": "400"}, ) + # Then: `unhealthy` only follows `ready`, so the watch ends there, loudly + assert result.exit_code == 1 + envelope = _json_envelope(result) + assert envelope["data"]["deployment"]["status"] == "unhealthy" + assert envelope["error"]["code"] == "deploy_status_terminal" + assert envelope["error"]["details"]["status"] == "unhealthy" + assert "still billing" in result.stderr + assert "comfy deploy stop --deployment dep-1" in result.stderr + assert sleeps == [] + + +def test_up_on_a_deployment_already_unhealthy_does_not_wait_for_ever(tmp_path, monkeypatch) -> None: + """`up` leaves an unhealthy deployment as it is, so a watch that waited for + `ready` would poll for as long as the endpoint stays degraded, saying + nothing.""" + # Given the release's deployment is unhealthy and stays so + module = _deploy() + client = FakeDeploy([deployment("dep-1", status="unhealthy")]) + monkeypatch.setattr(module, "_command_clients", lambda: (FakeBuilder(), client)) + sleeps: list[float] = [] + + def sleep(seconds: float) -> None: + sleeps.append(seconds) + if len(sleeps) > 5: + raise AssertionError("the watch kept polling an unhealthy deployment") + + monkeypatch.setattr(module, "_sleep", sleep) + + # When + result = CliRunner().invoke(app, ["--no-json", "deploy", "up", str(write_spec(tmp_path))], env={"COLUMNS": "400"}) + # Then - assert result.exit_code == 0 - assert _json_envelope(result)["data"]["deployment"]["status"] == "ready" - assert sleeps == [DEPLOY_POLL_SECONDS] + assert result.exception is None or isinstance(result.exception, SystemExit), result.exception + assert result.exit_code == 1 + assert sleeps == [] + assert "Deployment dep-1 is unhealthy" in result.stdout + result.stderr + assert client.start_calls == [] and client.update_calls == [] diff --git a/tests/comfy_cli/output/test_envelope_schemas.py b/tests/comfy_cli/output/test_envelope_schemas.py index 409ce6e59..bcd387916 100644 --- a/tests/comfy_cli/output/test_envelope_schemas.py +++ b/tests/comfy_cli/output/test_envelope_schemas.py @@ -63,6 +63,9 @@ def _validator_for(name: str) -> jsonschema.Validator: # one of those unions can't ship. "cloud_status.json", "knowledge.json", + # The progress object is the deploy service's own and may grow fields, + # so the event stays open where the envelopes around it are closed. + "deploy_progress_event.json", ], ) def test_schemas_are_well_formed(schema_name):