← Files PixVerseARCHIVED FILE

pvx/pixverse.py

93.6 KB · Oct 4, 2026 · 12:28 UTC

↓ Download file

from __future__ import annotations

import hashlib
import json
import re
import shlex
import subprocess
import sys
import time
from concurrent.futures import Future, ThreadPoolExecutor, TimeoutError as FutureTimeoutError
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any

from .capabilities import load_create_capabilities
from .compatibility import PIXVERSE_CLI_BASELINE_VERSION
from .model_defaults import IMAGE_MODEL, IMAGE_QUALITY, IMAGE_DETAIL, VIDEO_MODEL, VIDEO_QUALITY
from .region import effective_pixverse_region
from .shell import CommandResult, run_json, which
from .state import append_jsonl, read_jsonl, utc_now


# A task in one of these classes was submitted to PixVerse and therefore billed,
# but pvx never saw a terminal status for it. Its outcome is UNKNOWN, not failed.
# `queue reconcile` re-checks exactly these over free, read-only endpoints.
UNRESOLVED_ERROR_CLASSES = {"deadline", "deadline_unresolved"}

CREATE_KINDS = {
    "video",
    "image",
    "transition",
    "voice",
    "music",
    "extend",
    "modify",
    "upscale",
    "reference",
    "motion-control",
    "template",
}

CONCURRENCY_CODES = {429, 500041, 500042, 500044}
INSUFFICIENT_BALANCE_CODES = {500043}
PROMPT_CODES = {400018, 400038}
PARAM_CODES = {400017}
MEMBERSHIP_CODES = {500323, 500342}
SUBMISSION_FEEDBACK_MIN_SECONDS = 1.0
DOWNLOAD_INITIAL_FEEDBACK_SECONDS = 2.0
DOWNLOAD_MAX_FEEDBACK_SECONDS = 30.0
QUEUE_DOWNLOAD_WORKERS = 2


def media_type_for_kind(kind: str) -> str:
    if kind == "image":
        return "image"
    if kind in {"voice", "music"}:
        return "audio"
    return "video"


def poll_type_for_kind(kind: str) -> str:
    return media_type_for_kind(kind)


def stable_key(project: str, task_id: str, command: str) -> str:
    raw = f"{project}:{task_id}:{command}".encode("utf-8")
    return hashlib.sha256(raw).hexdigest()[:32]


def split_command(command: str) -> list[str]:
    argv = shlex.split(command)
    if not argv:
        raise ValueError("empty command")
    if Path(argv[0]).name != "pixverse":
        raise ValueError("PixVerse Agent Plugin queue only runs pixverse commands")
    return argv


def create_kind(argv: list[str]) -> str:
    normalized = _normalize_pixverse_create_globals(argv)
    if len(normalized) >= 3 and normalized[1] == "create":
        kind = normalized[2]
        if kind in CREATE_KINDS:
            return kind
    if argv[1:3] == ["miniapps", "create"]:
        raise ValueError(
            "MiniApps creation is not supported by the pvx paid queue yet; "
            "pixverse miniapps list/info remain available for read-only inspection"
        )
    raise ValueError("queued command must be `pixverse create <kind> ...`")


def apply_generation_defaults(argv: list[str]) -> list[str]:
    """Materialize defaults when authoring a new queue, preserving explicit choices.

    Do not rewrite saved queues at load/run time: their commands may already have
    paid receipts or approval, and changing them would change idempotency keys.
    """
    args = list(argv)
    try:
        kind = create_kind(args)
    except ValueError:
        return args
    model = _argv_option_value(args, "--model", "-m")
    if not model:
        if kind == "image":
            model = IMAGE_MODEL
        elif kind in {"video", "reference"}:
            model = VIDEO_MODEL
        elif kind == "transition" and len(_argv_list_option_values(args, "--images")) == 2:
            model = VIDEO_MODEL
        if model:
            args.extend(["--model", model])
    if model in {IMAGE_MODEL, VIDEO_MODEL} and not _argv_option_value(args, "--quality", "-q"):
        args.extend(["--quality", IMAGE_QUALITY if model == IMAGE_MODEL else VIDEO_QUALITY])
    if model == IMAGE_MODEL and not _argv_option_value(args, "--detail-level"):
        args.extend(["--detail-level", IMAGE_DETAIL])
    return args


def validate_create_argv(argv: list[str]) -> str:
    """Validate release-pinned capability walls that would otherwise spend a failed task."""
    effective_pixverse_region(argv[1:])
    argv = _normalize_pixverse_create_globals(argv)
    kind = create_kind(argv)
    workspace_id = _argv_option_value(argv, "--workspace-id")
    if workspace_id:
        raise ValueError(
            "queued paid generation does not accept --workspace-id; select the intended active "
            "workspace before preflight so account, balance, submission, and reconciliation stay aligned"
        )
    if kind == "upscale":
        quality = _argv_option_value(argv, "--quality", "-q")
        if quality and quality != "2160p":
            raise ValueError(
                f"PixVerse CLI {PIXVERSE_CLI_BASELINE_VERSION}+ video upscale supports only --quality 2160p"
            )
    model = _argv_option_value(argv, "--model", "-m")
    if model == "seedance-2.5":
        if kind not in {"video", "reference", "transition"}:
            raise ValueError(
                "Seedance 2.5 is supported only for video, reference, and two-frame transition generation"
            )
        for flag, message in (
            ("--audio", "Seedance 2.5 does not expose an audio toggle; omit --audio and describe the desired sound in the prompt. This does not mean the model cannot generate audio"),
            ("--no-audio", "Seedance 2.5 does not expose an audio toggle; omit --no-audio and remove the audio track during local export if a silent deliverable is required"),
            ("--multi-shot", "Seedance 2.5 does not support generated multi-shot mode"),
            ("--off-peak", "Seedance 2.5 does not support off-peak generation"),
        ):
            if _argv_has_option(argv, flag):
                raise ValueError(message)
        quality = _argv_option_value(argv, "--quality", "-q")
        if quality and quality not in {"480p", "720p", "1080p"}:
            raise ValueError("Seedance 2.5 quality must be 480p, 720p, or 1080p")
        aspect_ratio = _argv_option_value(argv, "--aspect-ratio")
        if aspect_ratio and kind in {"video", "reference"} and aspect_ratio not in {
            "auto",
            "21:9",
            "16:9",
            "4:3",
            "1:1",
            "3:4",
            "9:16",
        }:
            raise ValueError(f"Seedance 2.5 does not support aspect ratio {aspect_ratio}")
        duration = _argv_option_value(argv, "--duration", "-d")
        if duration and not (kind == "reference" and duration == "auto"):
            try:
                duration_value = float(duration)
            except ValueError as exc:
                raise ValueError("Seedance 2.5 duration must be an integer from 4 to 30") from exc
            if not duration_value.is_integer() or duration_value < 4 or duration_value > 30:
                raise ValueError("Seedance 2.5 duration must be an integer from 4 to 30")
        task_type = _argv_option_value(argv, "--task-type")
        if task_type:
            if kind != "reference":
                raise ValueError("--task-type is available only for Seedance 2.5 reference generation")
            if task_type not in {"auto", "reference", "edit", "extend"}:
                raise ValueError("Seedance 2.5 --task-type must be auto, reference, edit, or extend")
        if kind == "transition" and len(_argv_list_option_values(argv, "--images")) != 2:
            raise ValueError("Seedance 2.5 transition requires exactly two keyframe images")
        if kind == "reference":
            reference_limits = {"--images": 30, "--videos": 10, "--audios": 10}
            reference_counts = {
                option: len(_argv_list_option_values(argv, option))
                for option in reference_limits
            }
            for option, limit in reference_limits.items():
                if reference_counts[option] > limit:
                    raise ValueError(
                        f"Seedance 2.5 reference accepts at most {limit} values for {option}"
                    )
            if sum(reference_counts.values()) > 50:
                raise ValueError("Seedance 2.5 reference accepts at most 50 total media references")
    if kind == "music" and (
        _argv_has_option(argv, "--no-duration-auto") or _argv_option_value(argv, "--duration-seconds")
    ):
        # PixVerse music routes generate at automatic duration; the live service rejects
        # a fixed target (400017 "duration_seconds is reserved"). Trim, loop or fade locally.
        raise ValueError(
            "PixVerse music generation uses automatic duration; drop --duration-seconds / "
            "--no-duration-auto and trim, loop or fade the returned audio locally to the picture"
        )
    _validate_create_capability_argv(argv, kind=kind, model=model)
    return kind


def _normalize_pixverse_create_globals(argv: list[str]) -> list[str]:
    if not argv:
        return []
    command_args = [argv[0]]
    deferred: list[str] = []
    index = 1
    while index < len(argv):
        argument = argv[index]
        if argument == "--":
            command_args.extend(argv[index:])
            break
        if argument in {"--json", "-p"}:
            deferred.append(argument)
            index += 1
            continue
        if argument in {"--workspace-id", "--trace-id", "--region"}:
            deferred.append(argument)
            if index + 1 < len(argv):
                deferred.append(argv[index + 1])
            index += 2
            continue
        if any(
            argument.startswith(f"{name}=")
            for name in ("--workspace-id", "--trace-id", "--region")
        ):
            deferred.append(argument)
            index += 1
            continue
        command_args.append(argument)
        index += 1
    return [*command_args, *deferred]


def _parameter_flags(parameter: dict[str, Any]) -> list[str]:
    declaration = str(parameter.get("flag") or "")
    flags: list[str] = []
    for part in declaration.replace(",", "/").split("/"):
        token = part.strip().split(maxsplit=1)[0] if part.strip() else ""
        if token.startswith("-"):
            flags.append(token)
    return flags


def _parameter_is_present(argv: list[str], parameter: dict[str, Any]) -> bool:
    return any(_argv_has_option(argv, flag) for flag in _parameter_flags(parameter))


def _parameter_value(argv: list[str], parameter: dict[str, Any]) -> str:
    flags = _parameter_flags(parameter)
    return _argv_option_value(argv, *flags) if flags else ""


def _parameter_list_values(argv: list[str], parameter: dict[str, Any]) -> list[str]:
    flags = _parameter_flags(parameter)
    return _argv_list_option_values(argv, flags[0]) if flags else []


def _positive_boolean_requested(argv: list[str], parameter: dict[str, Any]) -> bool:
    flags = _parameter_flags(parameter)
    positive = next((flag for flag in flags if not flag.startswith("--no-")), "")
    return bool(positive and _argv_has_option(argv, positive))


def _validate_integer_parameter(name: str, value: str, contract: dict[str, Any]) -> None:
    try:
        parsed = int(value)
    except (TypeError, ValueError) as exc:
        raise ValueError(f"--{name.replace('_', '-')} must be an integer") from exc
    minimum = contract.get("min")
    maximum = contract.get("max")
    if isinstance(minimum, int) and parsed < minimum:
        raise ValueError(f"--{name.replace('_', '-')} must be at least {minimum}")
    if isinstance(maximum, int) and parsed > maximum:
        raise ValueError(f"--{name.replace('_', '-')} must be at most {maximum}")


def _validate_create_capability_argv(argv: list[str], *, kind: str, model: str) -> None:
    create = load_create_capabilities()
    modes = create.get("modes") if isinstance(create.get("modes"), dict) else {}
    mode = modes.get(kind) if isinstance(modes, dict) else None
    if not isinstance(mode, dict):
        return
    parameters = mode.get("parameters") if isinstance(mode.get("parameters"), dict) else {}
    shared = (
        create.get("shared_parameters")
        if isinstance(create.get("shared_parameters"), dict)
        else {}
    )
    selected_model = model or str(mode.get("default_model") or "")
    model_ids = mode.get("model_ids")
    if model and isinstance(model_ids, list) and model not in model_ids:
        raise ValueError(f"model {model!r} is not supported by pixverse create {kind}")

    # The offline video enum combines T2V and I2V ratios. The CLI's H3
    # builders reject T2V auto and silently force I2V auto, so validate the
    # selected input mode before a paid plan can promise different framing.
    if selected_model in {"minimax-h3", "minimax-h3-max"} and kind == "video":
        aspect = _argv_option_value(argv, "--aspect-ratio")
        image = _argv_option_value(argv, "--image")
        if not image and aspect == "auto":
            raise ValueError(f"{selected_model} text-to-video does not support --aspect-ratio auto")
        if image and aspect and aspect != "auto":
            raise ValueError(
                f"{selected_model} image-to-video forces automatic framing; "
                "omit --aspect-ratio or use auto, or choose reference mode for a fixed ratio"
            )

    for name, parameter in parameters.items():
        if not isinstance(parameter, dict):
            continue
        parameter_type = str(parameter.get("type") or "")
        present = _parameter_is_present(argv, parameter)
        values = _parameter_list_values(argv, parameter) if parameter_type == "media_list" else []
        value = _parameter_value(argv, parameter) if parameter_type not in {"boolean", "media_list"} else ""
        if parameter.get("required") is True and not (present and (values or value or parameter_type == "boolean")):
            flag = _parameter_flags(parameter)[0] if _parameter_flags(parameter) else name
            raise ValueError(f"pixverse create {kind} requires {flag}")
        minimum_count = parameter.get("min_count")
        if isinstance(minimum_count, int) and present and len(values) < minimum_count:
            flag = _parameter_flags(parameter)[0] if _parameter_flags(parameter) else name
            raise ValueError(f"{flag} requires at least {minimum_count} values")
        if present and parameter_type == "integer":
            _validate_integer_parameter(name, value, parameter)

    model_parameters = mode.get("model_parameters")
    for name, by_model in (
        model_parameters.items() if isinstance(model_parameters, dict) else []
    ):
        contract = by_model.get(selected_model) if isinstance(by_model, dict) else None
        parameter = parameters.get(name) if isinstance(parameters.get(name), dict) else shared.get(name)
        if not isinstance(contract, dict) or not isinstance(parameter, dict):
            continue
        present = _parameter_is_present(argv, parameter)
        parameter_type = str(parameter.get("type") or "")
        if contract.get("supported") is False and (
            _positive_boolean_requested(argv, parameter) if parameter_type == "boolean" else present
        ):
            flag = _parameter_flags(parameter)[0] if _parameter_flags(parameter) else name
            raise ValueError(f"{selected_model} does not support {flag} for pixverse create {kind}")
        if not present:
            continue
        values = _parameter_list_values(argv, parameter) if parameter_type == "media_list" else []
        value = _parameter_value(argv, parameter) if parameter_type not in {"boolean", "media_list"} else ""
        if parameter_type == "integer":
            _validate_integer_parameter(name, value, contract)
        allowed = contract.get("enum")
        if isinstance(allowed, list) and value and str(value) not in {str(item) for item in allowed}:
            flag = _parameter_flags(parameter)[0] if _parameter_flags(parameter) else name
            raise ValueError(f"{flag}={value!r} is not supported by {selected_model}")
        maximum_count = contract.get("max_count")
        if isinstance(maximum_count, int) and len(values) > maximum_count:
            flag = _parameter_flags(parameter)[0] if _parameter_flags(parameter) else name
            raise ValueError(f"{selected_model} accepts at most {maximum_count} values for {flag}")

    for name, parameter in shared.items():
        if not isinstance(parameter, dict) or not _parameter_is_present(argv, parameter):
            continue
        if parameter.get("type") == "integer":
            _validate_integer_parameter(name, _parameter_value(argv, parameter), parameter)


    if kind == "reference":
        reference_counts = {
            name: len(_parameter_list_values(argv, parameters[name]))
            for name in ("images", "videos", "audios")
            if isinstance(parameters.get(name), dict)
        }
        if sum(reference_counts.values()) == 0:
            raise ValueError("pixverse create reference requires at least one image, video, or audio reference")
        if reference_counts.get("audios", 0) and not (
            reference_counts.get("images", 0) or reference_counts.get("videos", 0)
        ) and selected_model != "wan-3.0":
            raise ValueError("audio references require an image or video reference for this model")
        metadata = mode.get("model_metadata")
        model_metadata = metadata.get(selected_model) if isinstance(metadata, dict) else None
        total_limit = model_metadata.get("max_reference_items") if isinstance(model_metadata, dict) else None
        if isinstance(total_limit, int) and sum(reference_counts.values()) > total_limit:
            raise ValueError(f"{selected_model} accepts at most {total_limit} total media references")

    if kind == "transition":
        image_parameter = parameters.get("images") if isinstance(parameters.get("images"), dict) else {}
        image_count = len(_parameter_list_values(argv, image_parameter))
        variants = mode.get("variants") if isinstance(mode.get("variants"), dict) else {}
        variant_name = "two_frames" if image_count == 2 else "multi_frame" if image_count >= 3 else ""
        variant = variants.get(variant_name) if isinstance(variants, dict) else None
        allowed_models = variant.get("allowed_models") if isinstance(variant, dict) else None
        if selected_model and isinstance(allowed_models, list) and selected_model not in allowed_models:
            raise ValueError(
                f"{selected_model} is not supported for a {image_count}-frame transition"
            )
        rules = mode.get("rules")
        for rule in rules if isinstance(rules, list) else []:
            if not isinstance(rule, dict) or rule.get("id") != "transition-prompt-required":
                continue
            required_models = rule.get("when", {}).get("model", {}).get("in", [])
            if selected_model in required_models:
                prompt = parameters.get("prompt") if isinstance(parameters.get("prompt"), dict) else {}
                if not _parameter_value(argv, prompt):
                    raise ValueError(f"{selected_model} transition requires --prompt")


def _argv_option_value(argv: list[str], *names: str) -> str:
    for index, value in enumerate(argv):
        if value in names:
            return argv[index + 1] if index + 1 < len(argv) else ""
        for name in names:
            if value.startswith(f"{name}="):
                return value.split("=", 1)[1]
    return ""


def _argv_has_option(argv: list[str], name: str) -> bool:
    return any(value == name or value.startswith(f"{name}=") for value in argv)


def _argv_list_option_values(argv: list[str], name: str) -> list[str]:
    try:
        index = argv.index(name) + 1
    except ValueError:
        return []
    values: list[str] = []
    while index < len(argv) and not argv[index].startswith("-"):
        values.append(argv[index])
        index += 1
    return values


def ensure_async_json_args(argv: list[str], key: str) -> list[str]:
    out = list(argv)
    if "--json" not in out and "-p" not in out:
        out.append("--json")
    if "--no-wait" not in out:
        out.append("--no-wait")
    try:
        kind = create_kind(out)
    except ValueError:
        kind = ""
    if kind in {"voice", "music"}:
        if "--client-request-id" not in out:
            out.extend(["--client-request-id", key])
    elif "--idempotency-key" not in out:
        out.extend(["--idempotency-key", key])
    return out


def extract_task_id(payload: dict[str, Any], kind: str) -> str | None:
    media = media_type_for_kind(kind)
    candidates = [
        f"{media}_id",
        "video_id",
        "image_id",
        "audio_id",
        "voice_id",
        "music_id",
        "id",
    ]
    for key in candidates:
        value = payload.get(key)
        if value:
            return str(value)
    for key in ("video_ids", "image_ids", "audio_ids", "ids"):
        value = payload.get(key)
        if isinstance(value, list) and value:
            return str(value[0])
    return None


def classify_error(result: CommandResult, payload: dict[str, Any] | None = None) -> str:
    if result.returncode == 7:
        return "concurrency"
    payload = payload or {}
    if not payload:
        payload = _parse_error_payload(result)
    text = f"{result.stderr}\n{result.stdout}\n{json.dumps(payload, ensure_ascii=False)}".lower()
    # PixVerse media paths avoid the CLI's URL re-upload ceiling, but a
    # provider-generated image can still exceed a downstream model's dimension
    # wall (for example a 6336 px-wide 21:9 board against a 6000 px maximum).
    # Keep this more specific than the provider's generic 400017 parameter code
    # so the queue can localize the internal asset and let the CLI's existing
    # local-image resize path repair it once, after the server reports the issue.
    if (
        "width and height between 300 and 6000" in text
        or "no larger than 30mb" in text
        or "image dimensions" in text and "6000" in text
    ):
        return "reference_input_invalid"
    if "file too large" in text or "payload too large" in text or "max: 10mb" in text:
        return "input_too_large"
    code = payload.get("code")
    if isinstance(code, str) and code.isdigit():
        code = int(code)
    if code in CONCURRENCY_CODES:
        return "concurrency"
    if code in INSUFFICIENT_BALANCE_CODES:
        return "insufficient_balance"
    if code in PROMPT_CODES:
        return "prompt_invalid"
    if code in PARAM_CODES:
        return "param_invalid"
    if code in MEMBERSHIP_CODES:
        return "membership_required"
    if any(
        phrase in text
        for phrase in (
            "membership required",
            "subscription required",
            "upgrade your plan",
            "upgrade plan",
            "not available for current plan",
            "not available on your plan",
            "insufficient entitlement",
            "user rights insufficient",
            "\u6743\u76ca\u4e0d\u8db3",
            "\u4f1a\u5458\u6743\u76ca",
            "\u9700\u8981\u4f1a\u5458",
        )
    ):
        return "membership_required"
    if "auth" in text or "login" in text:
        return "auth"
    if "quota" in text or "concurrent" in text or "over limit" in text:
        return "concurrency"
    if "credit" in text or "balance" in text:
        return "insufficient_balance"
    if "voice is required" in text or "duration_seconds is reserved" in text:
        return "param_invalid"
    if "timeout" in text:
        return "timeout"
    return "unknown"


def _parse_error_payload(result: CommandResult) -> dict[str, Any]:
    for raw in (result.stdout, result.stderr):
        text = raw.strip()
        if not text.startswith("{"):
            continue
        try:
            payload = json.loads(text)
        except json.JSONDecodeError:
            continue
        if isinstance(payload, dict):
            return payload
    return {}


@dataclass
class TaskSpec:
    id: str
    command: str
    label: str = ""
    depends_on: list[str] = field(default_factory=list)
    # A reused task carries an already generated asset from an earlier run. It is
    # never submitted again; dependents resolve `{{id.path}}` from these values.
    reuse: dict[str, Any] | None = None


REUSE_FIELDS = ("task_id", "path", "url", "cover_url", "local_path")


def reuse_record_from_manifest(manifest: Path, task_id: str) -> dict[str, Any] | None:
    """Find the latest successful manifest row for a queue task id (free, read-only)."""
    latest: dict[str, Any] | None = None
    for row in read_jsonl(manifest):
        if str(row.get("event") or "") != "task.success" or str(row.get("id") or "") != task_id:
            continue
        if not (row.get("path") or row.get("local_path") or row.get("url")):
            continue
        latest = row
    if latest is None:
        return None
    record = {field_name: str(latest.get(field_name) or "") for field_name in REUSE_FIELDS}
    record["command"] = str(latest.get("command") or "")
    record["source_task"] = task_id
    record["completed_at"] = str(latest.get("completed_at") or latest.get("at") or "")
    record["cost_credits"] = latest.get("cost_credits")
    return record


@dataclass
class TaskState:
    spec: TaskSpec
    kind: str
    media_type: str
    command: str
    idempotency_key: str
    task_id: str = ""
    status: str = "pending"
    url: str = ""
    cover_url: str = ""
    path: str = ""
    local_path: str = ""
    local_preview_status: str = ""
    local_preview_error: str = ""
    error_class: str = ""
    error_message: str = ""
    submitted_at: str = ""
    completed_at: str = ""
    cost_credits: int | None = None
    cost_source: str = ""
    status_code: Any = None
    provider_info: dict[str, Any] = field(default_factory=dict)
    raw: dict[str, Any] = field(default_factory=dict)

    def record(self) -> dict[str, Any]:
        record: dict[str, Any] = {
            "id": self.spec.id,
            "label": self.spec.label,
            "kind": self.kind,
            "media_type": self.media_type,
            "command": self.command,
            "task_id": self.task_id,
            "status": self.status,
            "url": self.url,
            "cover_url": self.cover_url,
            "path": self.path,
            "local_path": self.local_path,
            "local_preview_status": self.local_preview_status,
            "local_preview_error": self.local_preview_error,
            "error_class": self.error_class,
            "error_message": self.error_message,
            "submitted_at": self.submitted_at,
            "completed_at": self.completed_at,
        }
        if self.cost_credits is not None:
            record["cost_credits"] = self.cost_credits
            record["cost_source"] = self.cost_source
        if self.status_code is not None:
            record["status_code"] = self.status_code
        return record


def _terminal_event(event: str, state: TaskState, **extra: Any) -> dict[str, Any]:
    """Build a terminal manifest event that keeps the server's own words.

    pvx used to drop `status_code` and the provider payload on `task.failed`,
    so a terminal event recorded *less* than the `task.progress` events that
    preceded it. When v04 produced nine silent `generation_failed` results there
    was nothing on disk to autopsy. Terminal events now carry the raw payload.
    """
    payload: dict[str, Any] = {"event": event, "at": utc_now(), **state.record()}
    if state.provider_info:
        payload["provider_info"] = state.provider_info
    payload.update(extra)
    return payload


def _apply_terminal(
    state: TaskState,
    info: dict[str, Any],
    terminal: str,
    error_class: str,
    error_message: str,
) -> None:
    """Fold a terminal PixVerse status payload into a task state."""
    state.status = terminal
    state.error_class = error_class
    state.error_message = error_message
    state.status_code = info.get("status_code")
    state.provider_info = dict(info)
    state.completed_at = utc_now()
    if terminal != "success":
        return
    state.url, state.cover_url = result_urls(info)
    poll_cost, poll_source = extract_credit_cost(info)
    if poll_cost is not None and state.cost_credits is None:
        state.cost_credits = poll_cost
        state.cost_source = f"task_status.{poll_source}"
    asset_payload = asset_info(state.task_id, state.kind)
    state.path = asset_path_from_info(asset_payload)
    asset_cost, asset_source = extract_credit_cost(asset_payload)
    if asset_cost is not None and state.cost_credits is None:
        state.cost_credits = asset_cost
        state.cost_source = f"asset_info.{asset_source}"
    if asset_payload:
        state.provider_info = {**state.provider_info, "asset_info": asset_payload}


PLACEHOLDER_RE = re.compile(r"\{\{([\w-]+)\.(url|cover_url|path|task_id)\}\}")


def canonicalize_internal_asset_placeholders(command: str) -> str:
    """Keep PixVerse-generated dependencies inside PixVerse's media store.

    A generated image URL can point at a PNG larger than the CLI's 10 MB
    upload ceiling. Feeding that URL back into another create command makes the
    CLI download and re-upload its own asset, which is both unnecessary and can
    fail after the paid upstream task has completed. Internal queue references
    therefore use the provider media path. Existing queue specs that still use
    ``{{task.url}}`` are normalized at load time for backward compatibility.
    User-supplied local paths and external URLs are unaffected because they are
    not placeholders.
    """
    return re.sub(r"\{\{([\w-]+)\.url\}\}", r"{{\1.path}}", command)


def placeholder_refs(command: str) -> list[str]:
    refs: list[str] = []
    for match in PLACEHOLDER_RE.finditer(command):
        ref = match.group(1)
        if ref not in refs:
            refs.append(ref)
    return refs


def substitute_placeholders(command: str, states: dict[str, TaskState]) -> str:
    def replace(match: re.Match[str]) -> str:
        task_name, field = match.group(1), match.group(2)
        state = states.get(task_name)
        if state is None or state.status != "success":
            raise ValueError(f"cannot resolve placeholder {match.group(0)}")
        value = getattr(state, field)
        if not value:
            raise ValueError(f"placeholder {match.group(0)} resolved empty")
        return value

    return PLACEHOLDER_RE.sub(replace, command)


def _localized_internal_image_command(
    *,
    state: TaskState,
    states: dict[str, TaskState],
    project_path: Path,
) -> tuple[str, list[dict[str, str]]] | None:
    """Localize generated image dependencies only after a provider limit error.

    The normal route stays fast and server-side: ``{{image.path}}`` is sent as a
    PixVerse media path. If the downstream endpoint rejects that generated image
    for dimensions or payload size before returning a task id, this helper
    downloads only the affected internal image(s). PixVerse CLI already resizes
    oversized *local* image inputs to its safe 1920x1920 envelope, so the retry
    reuses that established path instead of duplicating image inspection logic in
    pvx. Direct user paths and external URLs are deliberately untouched.
    """
    replacements: list[tuple[str, TaskState]] = []
    for match in PLACEHOLDER_RE.finditer(state.command):
        dependency_id, field = match.group(1), match.group(2)
        upstream = states.get(dependency_id)
        if field != "path" or upstream is None or upstream.kind != "image" or not upstream.task_id:
            continue
        if all(existing_id != dependency_id for existing_id, _ in replacements):
            replacements.append((dependency_id, upstream))
    if not replacements:
        return None

    localized = state.command
    records: list[dict[str, str]] = []
    cache_root = project_path / "assets" / "reference-fallbacks" / state.spec.id
    cache_root.mkdir(parents=True, exist_ok=True)
    for dependency_id, upstream in replacements:
        destination = cache_root / dependency_id
        destination.mkdir(parents=True, exist_ok=True)
        candidates = sorted(
            (
                item
                for item in destination.iterdir()
                if item.is_file() and item.suffix.lower() in {".jpg", ".jpeg", ".png", ".webp"}
            ),
            key=lambda item: item.stat().st_mtime,
            reverse=True,
        )
        local_file = candidates[0].resolve() if candidates else None
        if local_file is None:
            try:
                result, payload = run_json(
                    [
                        "pixverse",
                        "asset",
                        "download",
                        upstream.task_id,
                        "--type",
                        "image",
                        "--dest",
                        str(destination),
                        "--json",
                    ],
                    timeout=180,
                )
            except subprocess.TimeoutExpired:
                return None
            downloaded = Path(str(payload.get("file") or "")).expanduser()
            if result.ok and downloaded.is_file():
                local_file = downloaded.resolve()
            else:
                candidates = sorted(
                    (
                        item
                        for item in destination.iterdir()
                        if item.is_file()
                        and item.suffix.lower() in {".jpg", ".jpeg", ".png", ".webp"}
                    ),
                    key=lambda item: item.stat().st_mtime,
                    reverse=True,
                )
                local_file = candidates[0].resolve() if candidates else None
        if local_file is None:
            return None
        token = f"{{{{{dependency_id}.path}}}}"
        localized = localized.replace(token, shlex.quote(str(local_file)))
        records.append(
            {
                "dependency_id": dependency_id,
                "task_id": upstream.task_id,
                "provider_path": upstream.path,
                "local_file": str(local_file),
            }
        )

    return substitute_placeholders(localized, states), records


def read_slots() -> dict[str, int]:
    fallback = {"image": 1, "video": 1, "audio": 1, "shared_pool": 1, "shared_limit": 1}
    if not which("pixverse"):
        return fallback
    try:
        result, payload = run_json(["pixverse", "account", "slots", "--json"], timeout=20)
    except subprocess.TimeoutExpired:
        return fallback
    if not result.ok or not isinstance(payload, dict):
        return fallback
    out: dict[str, int] = {}
    shared = bool(payload.get("shared_pool"))
    out["shared_pool"] = int(shared)
    for kind in ("image", "video"):
        info = payload.get(kind)
        if isinstance(info, dict):
            if info.get("unlimited"):
                out[kind] = 4
                out[f"{kind}_limit"] = 4
            else:
                out[kind] = max(0, int(info.get("remaining", 0) or 0))
        else:
            out[kind] = 1
    out["audio"] = min(out.get("video", 1), 2) if shared else 1
    out["audio_limit"] = 2 if shared else 1
    if shared and any(f"{kind}_limit" in out for kind in ("image", "video")):
        out["shared_limit"] = 4
    return out


def terminal_from_status(info: dict[str, Any]) -> tuple[str | None, str, str]:
    status = str(info.get("status") or "").lower()
    status_code = info.get("status_code")
    if status in {"completed", "success", "succeeded"}:
        return "success", "", ""
    if status_code in {7, 8} or status in {"failed", "error", "rejected", "not approved"}:
        blob = " ".join(
            str(info.get(key) or "")
            for key in ("error", "error_message", "message", "code", "status")
        )
        if status_code == 7:
            return "failed", "audit_reject", blob.strip()
        code = info.get("code")
        if isinstance(code, str) and code.isdigit():
            code = int(code)
        lowered = blob.lower()
        if code in MEMBERSHIP_CODES or any(
            phrase in lowered
            for phrase in (
                "membership required",
                "subscription required",
                "upgrade your plan",
                "upgrade plan",
                "not available for current plan",
                "not available on your plan",
                "insufficient entitlement",
                "\u6743\u76ca\u4e0d\u8db3",
                "\u4f1a\u5458\u6743\u76ca",
                "\u9700\u8981\u4f1a\u5458",
            )
        ):
            return "failed", "membership_required", blob.strip()
        if "not approved" in lowered:
            return "failed", "audit_reject", blob.strip()
        if code in INSUFFICIENT_BALANCE_CODES or "insufficient balance" in lowered or "insufficient credit" in lowered:
            return "failed", "insufficient_balance", blob.strip()
        if "auth" in lowered or "login" in lowered or "\u767b\u5f55" in lowered:
            return "failed", "auth", blob.strip()
        return "failed", "generation_failed", blob.strip()
    return None, "", ""


def result_urls(info: dict[str, Any]) -> tuple[str, str]:
    url = str(
        info.get("video_url")
        or info.get("image_url")
        or info.get("audio_url")
        or info.get("url")
        or ""
    )
    cover = str(info.get("cover_url") or info.get("thumbnail_url") or "")
    return url, cover


def asset_info(task_id: str, kind: str) -> dict[str, Any]:
    try:
        result, payload = run_json(
            ["pixverse", "asset", "info", task_id, "--type", poll_type_for_kind(kind), "--json"],
            timeout=30,
        )
    except subprocess.TimeoutExpired:
        return {}
    if not result.ok or not isinstance(payload, dict):
        return {}
    return payload


def asset_path_from_info(info: dict[str, Any]) -> str:
    return str(info.get("video_path") or info.get("image_path") or info.get("audio_path") or "")


def ensure_local_asset(
    *,
    task_id: str,
    media_type: str,
    project_path: Path,
    existing_path: str = "",
    progress: bool = False,
    label: str = "",
    expected_bytes: int | None = None,
) -> dict[str, str]:
    """Download a generated asset to a stable project-local preview path.

    PixVerse status and asset-info payloads expose public URLs and provider media
    paths, neither of which Codex can reliably render inline. Keep those provider
    values for audit and queue dependencies, but always give user-facing preview
    and project QA a real local file. A task-specific directory also makes this
    operation idempotent when a queue is resumed or reconciled.
    """
    if existing_path:
        existing = Path(existing_path).expanduser()
        if _usable_local_file(existing):
            return {
                "local_path": str(existing.resolve()),
                "status": "ready",
                "error": "",
            }

    normalized_type = media_type if media_type in {"image", "video", "audio"} else ""
    directory_name = {"image": "images", "video": "videos", "audio": "audio"}.get(
        normalized_type,
        "media",
    )
    safe_task_id = re.sub(r"[^A-Za-z0-9._-]+", "-", task_id).strip("-.") or "unknown-task"
    destination = (project_path / "assets" / directory_name / safe_task_id).resolve()

    cached = _latest_local_asset(destination, normalized_type)
    if cached is not None:
        return {"local_path": str(cached), "status": "ready", "error": ""}
    if not task_id:
        return {
            "local_path": "",
            "status": "task_id_missing",
            "error": "Cannot download a generated asset without its PixVerse task id.",
        }

    destination.mkdir(parents=True, exist_ok=True)
    argv = ["pixverse", "asset", "download", task_id]
    if normalized_type:
        argv.extend(["--type", normalized_type])
    argv.extend(["--dest", str(destination), "--json"])
    try:
        result, payload = _run_json_with_download_feedback(
            argv,
            timeout=5 * 60,
            progress=progress,
            destination=destination,
            label=label or task_id,
            expected_bytes=expected_bytes,
        )
    except subprocess.TimeoutExpired:
        return {
            "local_path": "",
            "status": "download_timeout",
            "error": "Local preview download timed out; the generated provider asset is unchanged.",
        }

    raw_file = str(payload.get("file") or "") if isinstance(payload, dict) else ""
    candidates: list[Path] = []
    if raw_file:
        reported = Path(raw_file).expanduser()
        candidates.append(reported)
        if not reported.is_absolute():
            candidates.append(destination / reported)
    candidates.extend(
        path
        for path in [_latest_local_asset(destination, normalized_type)]
        if path is not None
    )
    downloaded = next((path.resolve() for path in candidates if _usable_local_file(path)), None)
    if result.ok and downloaded is not None:
        return {"local_path": str(downloaded), "status": "ready", "error": ""}
    return {
        "local_path": "",
        "status": "download_failed",
        "error": result.stderr or result.stdout or "PixVerse asset download returned no local file.",
    }


def _latest_local_asset(destination: Path, media_type: str = "") -> Path | None:
    if not destination.is_dir():
        return None
    suffixes = {
        "image": {".gif", ".jpeg", ".jpg", ".png", ".webp"},
        "video": {".m4v", ".mov", ".mp4", ".webm"},
        "audio": {".aac", ".flac", ".m4a", ".mp3", ".ogg", ".wav"},
    }.get(media_type, set())
    candidates = sorted(
        (
            path
            for path in destination.rglob("*")
            if not path.name.startswith(".")
            and (not suffixes or path.suffix.lower() in suffixes)
            and _usable_local_file(path)
        ),
        key=lambda path: path.stat().st_mtime,
        reverse=True,
    )
    return candidates[0].resolve() if candidates else None


def _usable_local_file(path: Path) -> bool:
    try:
        return path.is_file() and path.stat().st_size > 0
    except OSError:
        return False


def asset_size_from_info(info: dict[str, Any]) -> int | None:
    for key in ("file_size_bytes", "size_bytes", "content_length", "file_size", "size"):
        value = info.get(key)
        if isinstance(value, bool):
            continue
        if isinstance(value, int) and value > 0:
            return value
        if isinstance(value, str) and value.isdigit() and int(value) > 0:
            return int(value)
    return None


def _localize_successful_state(
    state: TaskState,
    *,
    project_path: Path,
    progress: bool,
) -> None:
    if state.status != "success":
        return
    _emit_progress(
        progress,
        f"{state.spec.id} generated successfully; downloading the {state.media_type} for a local Codex preview.",
    )
    localized = ensure_local_asset(
        task_id=state.task_id,
        media_type=state.media_type,
        project_path=project_path,
        existing_path=state.local_path,
        progress=progress,
        label=state.spec.label or state.spec.id,
        expected_bytes=asset_size_from_info(
            state.provider_info.get("asset_info")
            if isinstance(state.provider_info.get("asset_info"), dict)
            else {}
        ),
    )
    state.local_path = localized["local_path"]
    state.local_preview_status = localized["status"]
    state.local_preview_error = localized["error"]
    if state.local_path:
        _emit_progress(
            progress,
            f"Local preview ready for {state.spec.label or state.spec.id}: {state.local_path}",
        )
    else:
        _emit_progress(
            progress,
            f"{state.spec.id} succeeded, but its local preview is unavailable "
            f"({state.local_preview_status}). The provider asset was not regenerated.",
        )


def extract_credit_cost(payload: dict[str, Any]) -> tuple[int | None, str]:
    candidates = (
        ("cost_credits", "cost_credits"),
        ("credits", "credits"),
        ("credit", "credit"),
        ("consume_credits", "consume_credits"),
        ("used_credits", "used_credits"),
    )
    for key, source in candidates:
        value = payload.get(key)
        if isinstance(value, bool):
            continue
        if isinstance(value, int) and value >= 0:
            return value, source
        if isinstance(value, float) and value >= 0 and value.is_integer():
            return int(value), source
        if isinstance(value, str) and value.isdigit():
            return int(value), source
    return None, ""


def load_task_specs(path: Path) -> tuple[str, list[TaskSpec]]:
    payload = json.loads(path.read_text(encoding="utf-8"))
    project = str(payload.get("project") or path.stem)
    tasks_raw = payload.get("tasks")
    if not isinstance(tasks_raw, list) or not tasks_raw:
        raise ValueError("queue spec must contain a non-empty `tasks` list")
    tasks: list[TaskSpec] = []
    seen: set[str] = set()
    for idx, item in enumerate(tasks_raw, start=1):
        if not isinstance(item, dict):
            raise ValueError(f"tasks[{idx}] must be an object")
        task_id = str(item.get("id") or f"task-{idx}")
        if task_id in seen:
            raise ValueError(f"duplicate task id: {task_id}")
        seen.add(task_id)
        reuse_raw = item.get("reuse")
        reuse: dict[str, Any] | None = None
        if reuse_raw is not None:
            if not isinstance(reuse_raw, dict):
                raise ValueError(f"task {task_id} reuse must be an object")
            reuse = {key: str(reuse_raw.get(key) or "") for key in REUSE_FIELDS}
            reuse["source_task"] = str(reuse_raw.get("source_task") or task_id)
            reuse["project"] = str(reuse_raw.get("project") or "")
            if not reuse["task_id"] or not (reuse["path"] or reuse["local_path"] or reuse["url"]):
                raise ValueError(
                    f"task {task_id} reuse needs task_id and one of path/local_path/url "
                    "(use `queue append --reuse <project>:<task-id>` to fill them from a manifest)"
                )
        command = canonicalize_internal_asset_placeholders(
            str(item.get("cmd") or item.get("command") or (reuse_raw or {}).get("command") or "").strip()
        )
        if not command:
            raise ValueError(f"task {task_id} missing command")
        try:
            validate_create_argv(split_command(command))
        except ValueError as exc:
            raise ValueError(f"task {task_id} invalid command: {exc}") from exc
        deps = item.get("depends_on") or []
        if isinstance(deps, str):
            deps = [deps]
        deps = [str(dep) for dep in deps]
        if reuse is None:
            for ref in placeholder_refs(command):
                if ref not in deps:
                    deps.append(ref)
        tasks.append(
            TaskSpec(
                id=task_id,
                command=command,
                label=str(item.get("label") or task_id),
                depends_on=deps,
                reuse=reuse,
            )
        )
    for task in tasks:
        for dep in task.depends_on:
            if dep not in seen:
                raise ValueError(f"task {task.id} depends on unknown task {dep}")
    return project, tasks


def _restore_queue_states_from_manifest(
    states: dict[str, TaskState], manifest: Path
) -> dict[str, int]:
    """Resume an unchanged queue without recreating completed audio or media.

    Video and image commands have provider idempotency keys, but current voice
    and music commands only expose a client request id. A host crash followed by
    an unchanged queue rerun could therefore create duplicate paid audio. Reuse
    exact successful commands and resume exact in-flight commands from the
    append-only manifest before considering any new submission.
    """
    counts = {"success": 0, "running": 0, "unresolved": 0, "failed": 0}
    restored_ids: set[str] = set()
    for row in read_jsonl(manifest):
        event = str(row.get("event") or "")
        spec_id = str(row.get("id") or "")
        state = states.get(spec_id)
        if state is None or not event.startswith("task."):
            continue
        row_command = str(row.get("command") or "").strip()
        canonical_row_command = canonicalize_internal_asset_placeholders(row_command)
        command_matches = bool(row_command) and canonical_row_command == state.command
        exact_command_matches = bool(row_command) and row_command == state.command

        if (
            event == "task.localized"
            and state.status == "success"
            and str(row.get("task_id") or "") == state.task_id
        ):
            state.local_path = str(row.get("local_path") or state.local_path)
            state.local_preview_status = str(
                row.get("local_preview_status") or state.local_preview_status
            )
            state.local_preview_error = str(
                row.get("local_preview_error") or state.local_preview_error
            )
            continue
        if not command_matches:
            continue

        if event == "task.success":
            state.status = "success"
            state.task_id = str(row.get("task_id") or "")
            state.url = str(row.get("url") or "")
            state.cover_url = str(row.get("cover_url") or "")
            state.path = str(row.get("path") or "")
            state.local_path = str(row.get("local_path") or "")
            state.local_preview_status = str(row.get("local_preview_status") or "")
            state.local_preview_error = str(row.get("local_preview_error") or "")
            state.submitted_at = str(row.get("submitted_at") or "")
            state.completed_at = str(row.get("completed_at") or "")
            state.status_code = row.get("status_code")
            if isinstance(row.get("cost_credits"), int):
                state.cost_credits = int(row["cost_credits"])
                state.cost_source = str(row.get("cost_source") or "manifest")
            restored_ids.add(spec_id)
            continue

        if state.status == "success":
            continue
        task_id = str(row.get("task_id") or "")
        if event == "task.submitted" and task_id:
            state.status = "running"
            state.task_id = task_id
            state.submitted_at = str(row.get("submitted_at") or row.get("at") or "")
            restored_ids.add(spec_id)
        elif event == "task.unresolved":
            state.task_id = task_id
            state.error_class = str(row.get("error_class") or "")
            state.error_message = str(row.get("error_message") or "")
            state.status = "running" if task_id else "unresolved"
            restored_ids.add(spec_id)
        elif event == "task.failed":
            if exact_command_matches:
                # A genuine failure for the exact command is terminal.
                state.status = "failed"
                state.task_id = task_id
                state.error_class = str(row.get("error_class") or "")
                state.error_message = str(row.get("error_message") or "")
                state.completed_at = str(row.get("completed_at") or row.get("at") or "")
                restored_ids.add(spec_id)
            else:
                # Historic URL transport was canonicalized to media path. That
                # is a materially repaired command, so do not keep an earlier
                # URL submission or failure alive in the new run.
                state.status = "pending"
                state.task_id = ""
                state.error_class = ""
                state.error_message = ""
                restored_ids.discard(spec_id)

    for spec_id in restored_ids:
        status = states[spec_id].status
        if status in counts:
            counts[status] += 1
    return counts


def _apply_reused_states(states: dict[str, TaskState], manifest: Path) -> list[str]:
    """Mark reused tasks successful from their recorded asset instead of submitting them."""
    applied: list[str] = []
    for state in states.values():
        reuse = state.spec.reuse
        if not reuse or state.status == "success":
            continue
        state.status = "success"
        state.task_id = str(reuse.get("task_id") or "")
        state.url = str(reuse.get("url") or "")
        state.cover_url = str(reuse.get("cover_url") or "")
        state.path = str(reuse.get("path") or "")
        local_path = str(reuse.get("local_path") or "")
        if local_path and Path(local_path).is_file():
            state.local_path = local_path
            state.local_preview_status = "ready"
        state.cost_credits = 0
        state.cost_source = "reused"
        state.completed_at = utc_now()
        state.provider_info = {"reused_from": {k: v for k, v in reuse.items() if v}}
        append_jsonl(manifest, _terminal_event("task.success", state, reused=True))
        applied.append(state.spec.id)
    return applied


def _emit_progress(enabled: bool, message: str) -> None:
    if enabled:
        print(f"[pixverse-agent] {message}", file=sys.stderr, flush=True)


def _run_json_with_wait_feedback(
    argv: list[str],
    *,
    timeout: float,
    progress: bool,
    feedback_interval: float,
    message: str,
) -> tuple[CommandResult, dict[str, Any]]:
    """Keep a slow provider submission conversational without retrying it."""
    interval = max(SUBMISSION_FEEDBACK_MIN_SECONDS, min(feedback_interval, 30.0))
    with ThreadPoolExecutor(max_workers=1, thread_name_prefix="pvx-submit") as executor:
        future = executor.submit(run_json, argv, timeout)
        while True:
            try:
                return future.result(timeout=interval)
            except FutureTimeoutError:
                _emit_progress(progress, message)


def _run_json_with_download_feedback(
    argv: list[str],
    *,
    timeout: float,
    progress: bool,
    destination: Path,
    label: str,
    expected_bytes: int | None = None,
) -> tuple[CommandResult, dict[str, Any]]:
    """Keep local delivery responsive without spamming long downloads.

    Fast downloads return immediately. For slower ones, sample byte growth in
    the task-specific destination, estimate recent throughput, and adapt the
    next heartbeat from a quick initial response to a 30-second large-file cap.
    """
    started = time.monotonic()
    last_sample_at = started
    last_bytes = _downloaded_bytes(destination)
    smoothed_rate = 0.0
    feedback_interval = DOWNLOAD_INITIAL_FEEDBACK_SECONDS
    with ThreadPoolExecutor(max_workers=1, thread_name_prefix="pvx-download") as executor:
        future = executor.submit(run_json, argv, timeout)
        while True:
            try:
                return future.result(timeout=feedback_interval)
            except FutureTimeoutError:
                now = time.monotonic()
                downloaded = _downloaded_bytes(destination)
                sample_seconds = max(0.001, now - last_sample_at)
                instant_rate = max(0, downloaded - last_bytes) / sample_seconds
                if instant_rate > 0:
                    smoothed_rate = (
                        instant_rate
                        if smoothed_rate <= 0
                        else smoothed_rate * 0.4 + instant_rate * 0.6
                    )
                elif smoothed_rate > 0:
                    smoothed_rate *= 0.5
                elapsed = now - started
                _emit_progress(
                    progress,
                    _download_feedback_message(
                        label=label,
                        downloaded_bytes=downloaded,
                        bytes_per_second=smoothed_rate,
                        expected_bytes=expected_bytes,
                        elapsed_seconds=elapsed,
                    ),
                )
                feedback_interval = _next_download_feedback_interval(
                    elapsed_seconds=elapsed,
                    downloaded_bytes=downloaded,
                    bytes_per_second=smoothed_rate,
                )
                last_sample_at = now
                last_bytes = downloaded


def _downloaded_bytes(destination: Path) -> int:
    if not destination.is_dir():
        return 0
    total = 0
    for path in destination.rglob("*"):
        try:
            if path.is_file():
                total += path.stat().st_size
        except OSError:
            continue
    return total


def _next_download_feedback_interval(
    *,
    elapsed_seconds: float,
    downloaded_bytes: int,
    bytes_per_second: float,
) -> float:
    mib = 1024 * 1024
    if downloaded_bytes >= 256 * mib:
        return DOWNLOAD_MAX_FEEDBACK_SECONDS
    if downloaded_bytes >= 64 * mib:
        return 20.0
    if bytes_per_second >= 8 * mib:
        return 15.0
    if bytes_per_second >= 1 * mib:
        return 10.0
    if bytes_per_second > 0:
        return 5.0
    if elapsed_seconds < 5:
        return 3.0
    if elapsed_seconds < 15:
        return 5.0
    if elapsed_seconds < 35:
        return 10.0
    return DOWNLOAD_MAX_FEEDBACK_SECONDS


def _download_feedback_message(
    *,
    label: str,
    downloaded_bytes: int,
    bytes_per_second: float,
    expected_bytes: int | None,
    elapsed_seconds: float,
) -> str:
    details: list[str] = []
    if downloaded_bytes > 0:
        details.append(f"downloaded {_format_bytes(downloaded_bytes)}")
    if expected_bytes and expected_bytes > 0:
        percent = min(100.0, downloaded_bytes / expected_bytes * 100)
        details.append(f"{percent:.0f}%")
    if bytes_per_second > 0:
        details.append(f"recent speed {_format_bytes(bytes_per_second)}/s")
    if (
        expected_bytes
        and expected_bytes > downloaded_bytes
        and bytes_per_second > 0
    ):
        eta = (expected_bytes - downloaded_bytes) / bytes_per_second
        details.append(f"ETA {_format_elapsed(eta)}")
    progress_text = ", ".join(details) if details else "waiting for local file data"
    return (
        f"Still downloading the local preview for {label}: {progress_text}; "
        f"elapsed {_format_elapsed(elapsed_seconds)}. Generation is complete and no regeneration was submitted."
    )


def _format_bytes(value: float | int) -> str:
    amount = float(max(0, value))
    for unit in ("B", "KiB", "MiB", "GiB"):
        if amount < 1024 or unit == "GiB":
            return f"{amount:.1f} {unit}" if unit != "B" else f"{amount:.0f} {unit}"
        amount /= 1024
    return f"{amount:.1f} GiB"


def _resolve_states(
    states: list[TaskState],
    *,
    manifest: Path,
    project_path: Path,
    progress: bool = False,
) -> list[TaskState]:
    """Re-check already-submitted tasks once, over free read-only endpoints.

    A submitted task has already been billed. Whatever pvx thinks, its outcome is
    unknown -- never "failed" -- until PixVerse itself reports a terminal status.
    Every recovered result is APPENDED to the manifest; no existing line is ever
    rewritten. Tasks still running are listed in `pending-reconcile.jsonl`.

    Returns the states that are still not terminal.
    """
    still_open: list[TaskState] = []
    by_type: dict[str, list[TaskState]] = {}
    for state in states:
        if state.task_id:
            by_type.setdefault(poll_type_for_kind(state.kind), []).append(state)

    for poll_type, group in by_type.items():
        ids = ",".join(state.task_id for state in group)
        try:
            result, payload = run_json(
                ["pixverse", "task", "status", "--ids", ids, "--type", poll_type, "--json"],
                timeout=40,
            )
        except subprocess.TimeoutExpired:
            _emit_progress(
                progress,
                f"Status re-check timed out for {ids}; keeping the billed task(s) unresolved for later reconcile.",
            )
            still_open.extend(group)
            continue
        if not result.ok or not isinstance(payload, dict):
            still_open.extend(group)
            continue
        by_id = {state.task_id: state for state in group}
        for task_id, info in payload.items():
            state = by_id.get(str(task_id))
            if state is None or not isinstance(info, dict):
                continue
            terminal, error_class, error_message = terminal_from_status(info)
            if terminal is None:
                continue
            previous = state.error_class or "orphan"
            _apply_terminal(state, info, terminal, error_class, error_message)
            if terminal == "success":
                _localize_successful_state(
                    state,
                    project_path=project_path,
                    progress=progress,
                )
            append_jsonl(
                manifest,
                _terminal_event(
                    f"task.{terminal}",
                    state,
                    reconciled=True,
                    reconciled_from=previous,
                ),
            )
            _emit_progress(
                progress,
                f"{state.spec.id} re-checked after the deadline: PixVerse reports {terminal} "
                f"(task {state.task_id}).",
            )
        still_open.extend(state for state in group if state.status not in {"success", "failed"})

    for state in still_open:
        append_jsonl(
            project_path / "pending-reconcile.jsonl",
            {
                "at": utc_now(),
                "id": state.spec.id,
                "task_id": state.task_id,
                "kind": state.kind,
                "media_type": state.media_type,
                "status": "unresolved",
                "error_class": state.error_class or "deadline_unresolved",
                "note": (
                    "Submitted and billed; PixVerse had not reached a terminal status at re-check "
                    "time. This is not a failure. Re-check later with `pvx queue reconcile "
                    "<project-slug>`, which is free and read-only."
                ),
            },
        )
    return still_open


def _state_from_manifest_row(row: dict[str, Any]) -> TaskState:
    """Rebuild a minimal TaskState from an append-only manifest event."""
    kind = str(row.get("kind") or "")
    if kind not in CREATE_KINDS:
        kind = "video"
    command = str(row.get("command") or "")
    return TaskState(
        spec=TaskSpec(
            id=str(row.get("id") or ""),
            command=command,
            label=str(row.get("label") or ""),
        ),
        kind=kind,
        media_type=str(row.get("media_type") or media_type_for_kind(kind)),
        command=command,
        idempotency_key="",
        task_id=str(row.get("task_id") or ""),
        status="unresolved",
        error_class=str(row.get("error_class") or ""),
        submitted_at=str(row.get("submitted_at") or ""),
    )


def unresolved_manifest_entries(manifest: Path) -> list[dict[str, Any]]:
    """Find submitted-but-unresolved tasks in an append-only manifest.

    Two ways a paid task ends up stranded:

    1. No terminal event at all -- the process died mid-poll. This is the v05
       ghost-directory case: 4 x `task.submitted`, 112 x `task.progress`, zero
       terminal events, four completed videos nobody ever downloaded.
    2. A terminal event whose `error_class` is in the deadline family -- the
       queue gave up on a task PixVerse was still working on.

    A task that PixVerse itself reported as failed (`generation_failed`,
    `audit_reject`, ...) is genuinely finished and is NOT re-checked.
    """
    latest: dict[str, dict[str, Any]] = {}
    resolved: set[str] = set()
    genuinely_failed: set[str] = set()
    for row in read_jsonl(manifest):
        event = str(row.get("event") or "")
        spec_id = str(row.get("id") or "")
        if not spec_id or not event.startswith("task."):
            continue
        if str(row.get("task_id") or ""):
            latest[spec_id] = row
        if event == "task.success":
            resolved.add(spec_id)
        elif event == "task.failed" and str(row.get("error_class") or "") not in UNRESOLVED_ERROR_CLASSES:
            genuinely_failed.add(spec_id)
    return [
        row
        for spec_id, row in latest.items()
        if spec_id not in resolved and spec_id not in genuinely_failed
    ]


def reconcile_project(*, project_path: Path, dry_run: bool = False) -> dict[str, Any]:
    """Recover billed-but-unrecorded tasks for a project. Free and read-only.

    Fixes both the deadline mis-accounting and the orphaned-task path (a queue
    process that died before its tasks finished).
    """
    manifest = project_path / "manifest.jsonl"
    candidates = [_state_from_manifest_row(row) for row in unresolved_manifest_entries(manifest)]
    checked = [
        {"id": state.spec.id, "task_id": state.task_id, "error_class": state.error_class or "orphan"}
        for state in candidates
    ]
    if dry_run or not candidates:
        return {
            "dry_run": dry_run,
            "checked": checked,
            "recovered": [],
            "confirmed_failed": [],
            "still_pending": checked if dry_run else [],
        }

    _resolve_states(candidates, manifest=manifest, project_path=project_path)
    return {
        "dry_run": False,
        "checked": checked,
        "recovered": [
            {
                "id": state.spec.id,
                "task_id": state.task_id,
                "url": state.url,
                "path": state.path,
                "local_path": state.local_path,
                "local_preview_status": state.local_preview_status,
            }
            for state in candidates
            if state.status == "success"
        ],
        "confirmed_failed": [
            {"id": state.spec.id, "task_id": state.task_id, "error_class": state.error_class}
            for state in candidates
            if state.status == "failed"
        ],
        "still_pending": [
            {"id": state.spec.id, "task_id": state.task_id}
            for state in candidates
            if state.status not in {"success", "failed"}
        ],
    }


def run_queue(
    *,
    spec_path: Path,
    project_path: Path,
    dry_run: bool = False,
    poll_interval: float = 15.0,
    status_interval: float = 30.0,
    deadline_seconds: float = 25 * 60,
    progress: bool = True,
) -> list[dict[str, Any]]:
    # Downloading completed media must not prevent filling free generation
    # slots or checking other paid tasks. Keep local delivery bounded and let
    # the queue's main thread own manifest writes and the final result.
    with ThreadPoolExecutor(
        max_workers=QUEUE_DOWNLOAD_WORKERS, thread_name_prefix="pvx-local-preview"
    ) as download_executor:
        return _run_queue(
            spec_path=spec_path,
            project_path=project_path,
            download_executor=download_executor,
            dry_run=dry_run,
            poll_interval=poll_interval,
            status_interval=status_interval,
            deadline_seconds=deadline_seconds,
            progress=progress,
        )


def _run_queue(
    *,
    spec_path: Path,
    project_path: Path,
    download_executor: ThreadPoolExecutor,
    dry_run: bool = False,
    poll_interval: float = 15.0,
    status_interval: float = 30.0,
    deadline_seconds: float = 25 * 60,
    progress: bool = True,
) -> list[dict[str, Any]]:
    project, specs = load_task_specs(spec_path)
    manifest = project_path / "manifest.jsonl"
    states: dict[str, TaskState] = {}
    for spec in specs:
        command = spec.command
        for ref in placeholder_refs(command):
            if ref not in spec.depends_on:
                spec.depends_on.append(ref)
        argv = split_command(command)
        kind = create_kind(argv)
        key = stable_key(project, spec.id, command)
        states[spec.id] = TaskState(
            spec=spec,
            kind=kind,
            media_type=media_type_for_kind(kind),
            command=command,
            idempotency_key=key,
        )
    if dry_run:
        return [state.record() | {"status": "dry_run"} for state in states.values()]

    restored = _restore_queue_states_from_manifest(states, manifest)
    reused_ids = _apply_reused_states(states, manifest)
    if reused_ids:
        _emit_progress(
            progress,
            f"Reusing {len(reused_ids)} accepted asset(s) without new generation: {', '.join(reused_ids)}.",
        )

    for state in states.values():
        if state.status != "success":
            continue
        previous_localization = (
            state.local_path,
            state.local_preview_status,
            state.local_preview_error,
        )
        _localize_successful_state(
            state,
            project_path=project_path,
            progress=progress,
        )
        if (
            state.local_path,
            state.local_preview_status,
            state.local_preview_error,
        ) != previous_localization:
            append_jsonl(
                manifest,
                {"event": "task.localized", "at": utc_now(), **state.record()},
            )

    _emit_progress(
        progress,
        f"Queue ready: {len(states)} task(s), polling every {poll_interval:g}s, deadline {deadline_seconds:g}s.",
    )
    restored_total = sum(restored.values())
    if restored_total:
        _emit_progress(
            progress,
            "Recovered unchanged queue state from the project manifest: "
            f"{restored['success']} ready, {restored['running']} still running, "
            f"{restored['unresolved']} unresolved, {restored['failed']} failed; "
            "no duplicate submissions were made for those tasks.",
        )
    progress_cache: dict[str, str] = {}
    periodic_cache: dict[str, tuple[str, float]] = {}
    task_progress_signatures: dict[str, str] = {}
    submitted_monotonic: dict[str, float] = {}

    def progress_once(key: str, message: str) -> None:
        if progress_cache.get(key) == message:
            return
        progress_cache[key] = message
        _emit_progress(progress, message)

    def progress_periodic(key: str, signature: str, message: str) -> None:
        now = time.monotonic()
        previous = periodic_cache.get(key)
        interval = max(status_interval, poll_interval, 1.0)
        if previous is None or previous[0] != signature or now - previous[1] >= interval:
            periodic_cache[key] = (signature, now)
            _emit_progress(progress, message)

    download_jobs: dict[str, Future[None]] = {}

    def collect_downloads(*, wait: bool = False) -> None:
        for spec_id, future in list(download_jobs.items()):
            if not wait and not future.done():
                continue
            state = states[spec_id]
            try:
                future.result()
            except Exception as exc:
                # Delivery failures never change the provider's successful
                # generation or authorize a new paid attempt.
                state.local_preview_status = "download_failed"
                state.local_preview_error = str(exc)
                _emit_progress(progress, f"{spec_id} generated successfully, but local delivery failed: {exc}")
            append_jsonl(manifest, _terminal_event("task.success", state))
            del download_jobs[spec_id]
            finished = sum(
                item.status in {"success", "failed", "unresolved"} for item in states.values()
            )
            preview = "is ready" if state.local_path else "has no usable local preview yet"
            progress_once(
                f"success:{spec_id}",
                f"{spec_id} completed successfully — {state.spec.label or spec_id} {preview} "
                f"({finished}/{len(states)} generated or terminal).",
            )

    def stop_pending_for_account_failure() -> None:
        blocker = next((item for item in states.values()
                        if item.error_class in {"membership_required", "auth", "insufficient_balance"}), None)
        if blocker is None:
            return
        for pending in states.values():
            if pending.status != "pending":
                continue
            pending.status = "failed"
            pending.error_class = "account_blocked"
            pending.error_message = f"Not submitted: {blocker.spec.id} reported {blocker.error_class}. Wait for the user's account/route choice and a fresh preflight."
            pending.completed_at = utc_now()
            append_jsonl(manifest, _terminal_event("task.failed", pending))
        progress_once("account-blocked", "Account access stopped new submissions. Keep existing task IDs; show the subscription link and wait for upgrade or explicit fallback choice.")

    deadline = time.monotonic() + deadline_seconds
    while True:
        collect_downloads()
        stop_pending_for_account_failure()
        progressed = False
        # Capacity only affects new submissions. Polling an existing paid task
        # (or waiting for its dependent input) needs no account-slots request.
        ready_to_submit = any(
            state.status == "pending"
            and all(states[dep].status == "success" for dep in state.spec.depends_on)
            for state in states.values()
        )
        slots = read_slots() if ready_to_submit else {}
        inflight_by_type: dict[str, int] = {"image": 0, "video": 0, "audio": 0}
        for state in states.values():
            if state.status == "running":
                inflight_by_type[state.media_type] = inflight_by_type.get(state.media_type, 0) + 1

        # `remaining` already excludes work occupying the account's slots.
        # Treat it as a budget for NEW submissions, not a total inflight limit.
        # Only the explicit local safeguards (unknown/unlimited capacity and
        # audio) need to subtract this queue's known inflight tasks.
        submission_budget = {
            kind: max(0, slots.get(kind, 0)) for kind in ("image", "video", "audio")
        }
        for kind in submission_budget:
            local_limit = slots.get(f"{kind}_limit")
            if local_limit is not None:
                submission_budget[kind] = min(
                    submission_budget[kind], max(0, local_limit - inflight_by_type[kind])
                )
        shared_pool = bool(slots.get("shared_pool"))
        shared_budget = min(slots.get("image", 0), slots.get("video", 0))
        if "shared_limit" in slots:
            shared_budget = min(
                shared_budget, max(0, slots["shared_limit"] - sum(inflight_by_type.values()))
            )

        for state in states.values():
            if state.status != "pending":
                continue
            if any(states[dep].status != "success" for dep in state.spec.depends_on):
                blocking_statuses = {states[dep].status for dep in state.spec.depends_on}
                if blocking_statuses & {"failed", "unresolved"}:
                    state.status = "failed"
                    state.error_class = (
                        "upstream_unresolved" if "unresolved" in blocking_statuses else "upstream_failed"
                    )
                    state.completed_at = utc_now()
                    append_jsonl(manifest, _terminal_event("task.failed", state))
                    progress_once(
                        f"failed:{state.spec.id}",
                        f"{state.spec.id} was not submitted because an upstream task did not finish safely.",
                    )
                    progressed = True
                continue
            available = submission_budget[state.media_type]
            if shared_pool:
                available = min(available, shared_budget)
            if available <= 0:
                progress_once(
                    f"slot:{state.spec.id}",
                    f"{state.spec.id} is waiting for a {state.media_type} slot "
                    f"({inflight_by_type.get(state.media_type, 0)} awaiting status, {available} available"
                    f"{' in the shared pool' if shared_pool else ''}).",
                )
                continue
            try:
                resolved = substitute_placeholders(state.command, states)
            except ValueError as exc:
                state.status = "failed"
                state.error_class = "upstream_asset_unavailable"
                state.error_message = str(exc)
                state.completed_at = utc_now()
                append_jsonl(manifest, _terminal_event("task.failed", state))
                progress_once(
                    f"failed:{state.spec.id}",
                    f"{state.spec.id} was not submitted because its upstream PixVerse media path is unavailable.",
                )
                progressed = True
                continue
            argv = ensure_async_json_args(split_command(resolved), state.idempotency_key)
            progress_once(
                f"submit:{state.spec.id}",
                f"Submitting {state.spec.id} ({state.media_type}): {state.spec.label or state.spec.id}.",
            )
            try:
                result, payload = _run_json_with_wait_feedback(
                    argv,
                    timeout=90,
                    progress=progress,
                    feedback_interval=status_interval,
                    message=(
                        f"Still submitting {state.spec.id}; PixVerse has not returned a task id yet. "
                        "The request is still in flight and no duplicate retry has been sent."
                    ),
                )
            except subprocess.TimeoutExpired:
                state.status = "unresolved"
                state.error_class = "submit_timeout_unknown"
                state.error_message = (
                    "PixVerse CLI did not return a task id before the submission timeout. "
                    "The request may have reached the provider, so pvx did not retry or claim that no credits were spent."
                )
                state.completed_at = utc_now()
                append_jsonl(manifest, _terminal_event("task.unresolved", state))
                progress_once(
                    f"unresolved:{state.spec.id}",
                    f"{state.spec.id} submission outcome is unknown; no task id returned and no duplicate retry was sent.",
                )
                progressed = True
                continue
            if not result.ok:
                state.error_class = classify_error(result, payload)
                # Do not inspect every reference up front. Only a definite
                # pre-submission image-limit rejection activates this fallback,
                # and only for generated image placeholders with a known task id.
                # The retry is safe because the rejected response contained no
                # provider task id; a distinct key prevents a cached rejection
                # from shadowing the repaired local-input request.
                if (
                    state.error_class in {"reference_input_invalid", "input_too_large"}
                    and not extract_task_id(payload, state.kind)
                ):
                    localized = _localized_internal_image_command(
                        state=state,
                        states=states,
                        project_path=project_path,
                    )
                    if localized is not None:
                        localized_command, localized_assets = localized
                        append_jsonl(
                            manifest,
                            {
                                "event": "task.reference_fallback",
                                "at": utc_now(),
                                **state.record(),
                                "trigger_error_class": state.error_class,
                                "localized_assets": localized_assets,
                                "note": (
                                    "The provider rejected an internal media path before task submission. "
                                    "pvx localized the generated image and retried once through PixVerse "
                                    "CLI's built-in local-image resize path."
                                ),
                            },
                        )
                        progress_once(
                            f"reference-fallback:{state.spec.id}",
                            f"{state.spec.id} hit a reference-image limit before submission; "
                            "localized the internal image and retrying once through the CLI's safe resize path.",
                        )
                        retry_key = stable_key(
                            project,
                            state.spec.id,
                            f"{state.command}:localized-reference-fallback:v1",
                        )
                        retry_argv = ensure_async_json_args(
                            split_command(localized_command),
                            retry_key,
                        )
                        try:
                            result, payload = _run_json_with_wait_feedback(
                                retry_argv,
                                timeout=90,
                                progress=progress,
                                feedback_interval=status_interval,
                                message=(
                                    f"Still submitting the repaired {state.spec.id} request; "
                                    "no further retry will be sent without a provider result."
                                ),
                            )
                        except subprocess.TimeoutExpired:
                            state.status = "unresolved"
                            state.error_class = "submit_timeout_unknown"
                            state.error_message = (
                                "The one-time localized-reference retry did not return a task id before "
                                "the submission timeout. It may have reached PixVerse, so pvx did not retry again."
                            )
                            state.completed_at = utc_now()
                            append_jsonl(manifest, _terminal_event("task.unresolved", state))
                            progress_once(
                                f"unresolved:{state.spec.id}",
                                f"{state.spec.id} localized-reference retry outcome is unknown; no duplicate retry was sent.",
                            )
                            progressed = True
                            continue
                        state.error_class = "" if result.ok else classify_error(result, payload)
                if not result.ok:
                    if state.error_class == "concurrency":
                        # A definite full-pool response invalidates this pass's
                        # capacity snapshot. Do not probe every pending task;
                        # poll existing work, then obtain fresh capacity.
                        submission_budget[state.media_type] = 0
                        progress_once(
                            f"concurrency:{state.spec.id}",
                            f"{state.spec.id} hit a PixVerse concurrency response; waiting instead of resubmitting.",
                        )
                        if shared_pool:
                            shared_budget = 0
                            break
                        continue
                    state.status = "failed"
                    state.error_message = result.stderr or result.stdout
                    state.completed_at = utc_now()
                    state.provider_info = payload if isinstance(payload, dict) else {}
                    state.status_code = state.provider_info.get("code")
                    append_jsonl(manifest, _terminal_event("task.failed", state))
                    progress_once(f"failed:{state.spec.id}", f"{state.spec.id} failed: {state.error_class}.")
                    stop_pending_for_account_failure()
                    progressed = True
                    continue
            task_id = extract_task_id(payload, state.kind)
            if not task_id:
                state.status = "failed"
                state.error_class = "missing_task_id"
                state.error_message = json.dumps(payload, ensure_ascii=False)
                state.completed_at = utc_now()
                state.provider_info = payload if isinstance(payload, dict) else {}
                append_jsonl(manifest, _terminal_event("task.failed", state))
                progress_once(f"failed:{state.spec.id}", f"{state.spec.id} failed: missing PixVerse task id.")
                progressed = True
                continue
            state.task_id = task_id
            state.status = "running"
            state.submitted_at = utc_now()
            submitted_monotonic[state.spec.id] = time.monotonic()
            state.raw = payload
            state.cost_credits, source = extract_credit_cost(payload)
            state.cost_source = f"submit.{source}" if source else ""
            inflight_by_type[state.media_type] = inflight_by_type.get(state.media_type, 0) + 1
            submission_budget[state.media_type] -= 1
            if shared_pool:
                shared_budget -= 1
            append_jsonl(manifest, {"event": "task.submitted", "at": utc_now(), **state.record()})
            progress_once(f"submitted:{state.spec.id}", f"Submitted {state.spec.id}: task {state.task_id}.")
            progressed = True

        running = [state for state in states.values() if state.status == "running"]
        if running:
            by_type: dict[str, list[TaskState]] = {}
            for state in running:
                by_type.setdefault(poll_type_for_kind(state.kind), []).append(state)
            for poll_type, group in by_type.items():
                ids = ",".join(state.task_id for state in group)
                progress_once(f"poll:{poll_type}:{ids}", f"Polling {len(group)} {poll_type} task(s): {ids}.")
                try:
                    result, payload = run_json(
                        ["pixverse", "task", "status", "--ids", ids, "--type", poll_type, "--json"],
                        timeout=40,
                    )
                except subprocess.TimeoutExpired:
                    progress_once(
                        f"poll-timeout:{poll_type}:{ids}",
                        f"Status check timed out for {ids}; the paid task stays running and will be checked again.",
                    )
                    continue
                if not result.ok or not isinstance(payload, dict):
                    continue
                by_id = {state.task_id: state for state in group}
                for task_id, info in payload.items():
                    if not isinstance(info, dict):
                        continue
                    state = by_id.get(str(task_id))
                    if state is None:
                        continue
                    terminal, error_class, error_message = terminal_from_status(info)
                    if terminal is None:
                        status_text = str(info.get("status") or info.get("status_code") or "running")
                        percent = info.get("progress_percent")
                        percent_suffix = f", {percent}%" if percent not in (None, "") else ""
                        code = info.get("status_code")
                        code_suffix = f", code {code}" if code not in (None, "") else ""
                        elapsed = _format_elapsed(time.monotonic() - submitted_monotonic.get(state.spec.id, time.monotonic()))
                        signature = f"{status_text}:{code}:{percent}"
                        # Provider state is often unchanged across dozens of polls.
                        # Emit and persist only real state/percent changes; the aggregate
                        # studio heartbeat owns elapsed-time reassurance.
                        if task_progress_signatures.get(state.spec.id) != signature:
                            task_progress_signatures[state.spec.id] = signature
                            _emit_progress(
                                progress,
                                f"{state.spec.id} still {status_text}{percent_suffix}{code_suffix}; "
                                f"elapsed {elapsed}; task {state.task_id}.",
                            )
                            append_jsonl(
                                manifest,
                                {
                                    "event": "task.progress",
                                    "at": utc_now(),
                                    "id": state.spec.id,
                                    "task_id": state.task_id,
                                    "status": info.get("status"),
                                    "status_code": info.get("status_code"),
                                    "progress_percent": info.get("progress_percent"),
                                },
                            )
                        continue
                    _apply_terminal(state, info, terminal, error_class, error_message)
                    if terminal == "success":
                        state.local_preview_status = "downloading"
                        download_jobs[state.spec.id] = download_executor.submit(
                            _localize_successful_state,
                            state,
                            project_path=project_path,
                            progress=progress,
                        )
                    else:
                        progress_once(f"failed:{state.spec.id}", f"{state.spec.id} failed: {error_class}.")
                        append_jsonl(manifest, _terminal_event(f"task.{terminal}", state))
                    progressed = True

        collect_downloads()
        status_counts = {
            status: sum(state.status == status for state in states.values())
            for status in ("success", "failed", "unresolved", "running", "pending")
        }
        active_labels = [
            state.spec.label or state.spec.id
            for state in states.values()
            if state.status == "running"
        ][:3]
        local_ready = sum(state.local_preview_status == "ready" for state in states.values())
        heartbeat_signature = ":".join(
            f"{task_id}={state.status}:{state.local_preview_status}" for task_id, state in states.items()
        )
        active_text = ", ".join(active_labels) if active_labels else "none"
        progress_periodic(
            "queue-heartbeat",
            heartbeat_signature,
            "Studio heartbeat: "
            f"{status_counts['success']}/{len(states)} generated, {local_ready} local previews ready, "
            f"{len(download_jobs)} downloading or queued for download, "
            f"{status_counts['running']} awaiting provider status, {status_counts['pending']} waiting, "
            f"{status_counts['failed']} failed, {status_counts['unresolved']} unresolved; active: {active_text}.",
        )

        if all(state.status in {"success", "failed", "unresolved"} for state in states.values()):
            collect_downloads(wait=True)
            _emit_progress(
                progress,
                f"Queue finished — {status_counts['success']}/{len(states)} asset task(s) ready, "
                f"{status_counts['failed']} failed, {status_counts['unresolved']} unresolved.",
            )
            return [state.record() for state in states.values()]
        if time.monotonic() >= deadline:
            # The deadline is one wall-clock budget shared by the whole queue, not a
            # per-task timeout. Reaching it says nothing about any individual task.
            # It used to mark every non-terminal task `failed / deadline` with no
            # final poll -- including tasks PixVerse had already completed and
            # billed. Split by whether a task was ever submitted, then re-check.
            unresolved: list[TaskState] = []
            for state in states.values():
                if state.status in {"success", "failed", "unresolved"}:
                    continue
                state.completed_at = utc_now()
                if not state.task_id:
                    # Never submitted: no PixVerse task exists and nothing was billed.
                    state.status = "failed"
                    state.error_class = "deadline_not_submitted"
                    state.error_message = (
                        "Queue deadline reached before this task was submitted. "
                        "No PixVerse task was created and no credits were spent."
                    )
                    append_jsonl(manifest, _terminal_event("task.failed", state))
                    progress_once(
                        f"deadline:{state.spec.id}",
                        f"{state.spec.id} was never submitted before the queue deadline; no credits were spent.",
                    )
                    continue
                # Submitted means billed. Outcome unknown, not failed.
                state.status = "unresolved"
                state.error_class = "deadline_unresolved"
                state.error_message = (
                    "Queue deadline reached while PixVerse was still working on this task. "
                    "It was submitted and billed; its outcome is unknown, not failed."
                )
                append_jsonl(manifest, _terminal_event("task.unresolved", state))
                unresolved.append(state)
            if unresolved:
                _emit_progress(
                    progress,
                    f"Deadline reached with {len(unresolved)} submitted task(s) still running. "
                    "They are already billed, so re-checking them once before reporting.",
                )
                _resolve_states(
                    unresolved,
                    manifest=manifest,
                    project_path=project_path,
                    progress=progress,
                )
            collect_downloads(wait=True)
            _emit_progress(progress, "Queue finished at deadline.")
            return [state.record() for state in states.values()]
        # A submission is progress, but does not justify immediately polling
        # the same still-running job again. Advance ready dependencies now;
        # otherwise let the normal polling interval elapse.
        can_advance_pending = any(
            state.status == "pending"
            and all(states[dep].status not in {"pending", "running"} for dep in state.spec.depends_on)
            for state in states.values()
        )
        if not progressed or not can_advance_pending:
            time.sleep(min(max(poll_interval, 1.0), max(1.0, deadline - time.monotonic())))


def _format_elapsed(seconds: float) -> str:
    seconds = max(0, int(seconds))
    minutes, remainder = divmod(seconds, 60)
    hours, minutes = divmod(minutes, 60)
    if hours:
        return f"{hours}h{minutes:02d}m{remainder:02d}s"
    if minutes:
        return f"{minutes}m{remainder:02d}s"
    return f"{remainder}s"

SHA-256: 633a5b94d51376aa797d4648f49aae9d8a019afc45d2040316f70e3597c3a4b1