Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
5cb86ad
feat(studio): admit recoverable image commands in canonical tasks
IAnMove Sep 9, 2026
921149f
Merge remote-tracking branch 'origin/development' into feat/shared-ge…
IAnMove Sep 9, 2026
d9cdf52
fix(commands): harden recovery and simplify receipt validation
IAnMove Sep 9, 2026
5db543c
feat(studio): share durable full image commands across Studio Wizard …
IAnMove Sep 9, 2026
141d630
Merge remote-tracking branch 'origin/development' into feat/shared-ge…
IAnMove Sep 9, 2026
d08853a
fix(commands): preserve source references and freeze native image def…
IAnMove Sep 9, 2026
f0e2ef4
fix(studio): retain image receipts and handle lazy submission failures
IAnMove Sep 9, 2026
93b8bc4
test(wizard): retain admitted image receipt through capability reporting
IAnMove Sep 9, 2026
fe46f41
test(architecture): classify Wizard receipt behavior test
IAnMove Sep 9, 2026
254d6aa
fix(wizard): preserve literal image requests through Suspense present…
IAnMove Sep 9, 2026
cc0d760
fix(studio): wait for an active presentation receiver across Suspense
IAnMove Sep 9, 2026
d256f89
fix(studio): omit leftover form fields from image commands
cursoragent Sep 9, 2026
3a7988e
fix(studio): project known native form leftovers safely
IAnMove Sep 9, 2026
7aebf60
fix(commands): isolate stale image leftovers from queue recovery
cursoragent Sep 9, 2026
57c980c
fix(commands): preserve recovery on storage failure and reject malfor…
IAnMove Sep 9, 2026
b0fc2dd
fix(commands): report recovery discard write failures as unavailable
IAnMove Sep 9, 2026
65ed285
fix(commands): skip invalid listed workspaces during queue recovery
cursoragent Sep 9, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 37 additions & 9 deletions app/_launch_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -608,6 +608,7 @@ def _new_generation_job(
recovered: bool = False,
reserve_generation: bool = True,
provenance: dict | None = None,
publish_task: bool = True,
) -> dict:
frozen_params = copy.deepcopy(params)
execution_mode.validate_generation(workspace)
Expand Down Expand Up @@ -679,7 +680,7 @@ def _new_generation_job(
if reserve_generation:
register_generation_job(_gen_lock, job)
publisher = globals().get("_publish_generation_task")
if callable(publisher):
if publish_task and callable(publisher):
try:
task = publisher(job)
if isinstance(task, dict):
Expand Down Expand Up @@ -10868,7 +10869,8 @@ async def generate(request: Request):
try:
prepare_generation_inputs(body, _generation_model_def, requested_workspace,
uploads_dir=os.path.join(os.getcwd(), "uploads"),
workspace_dir=_workspace_dir(requested_workspace))
workspace_dir=_workspace_dir(requested_workspace),
prepared_images=getattr(request, "prepared_studio_images", False) is True)
except ValueError as error:
raise HTTPException(status_code=400, detail=str(error)) from error
try:
Expand Down Expand Up @@ -11309,6 +11311,9 @@ async def generate(request: Request):

# Capture workspace at submission time — NOT at execution time
workspace = body.pop("workspace", None) or _get_active_workspace()
admission = getattr(request, "admit_generation_command", None)
if callable(admission):
return admission(body, workspace, provenance)
job_out_dir = _workspace_dir(workspace)

h3_preplan_pending = isinstance(
Expand Down Expand Up @@ -26297,7 +26302,10 @@ def _recovery_job_summary(record: dict) -> dict:
@api.get("/api/v1/jobs/recovery")
def get_generation_queue_recovery():
"""Return crash leftovers that are not active in this server process."""
candidates = _durable_generation_queue.list(exclude_ids=_jobs.keys())
with _queue_recovery_lock:
_image_generation_commands.restore_recovery(item["name"] for item in _list_workspaces())
candidates = _image_generation_commands.filter_recovery(
_durable_generation_queue.list(exclude_ids=_jobs.keys()))
return {"jobs": [_recovery_job_summary(record) for record in candidates]}


Expand All @@ -26311,7 +26319,9 @@ def resume_generation_queue():
resumed: list[dict] = []
threads: list[threading.Thread] = []
with _queue_recovery_lock:
candidates = _durable_generation_queue.list(exclude_ids=_jobs.keys())
_image_generation_commands.restore_recovery(item["name"] for item in _list_workspaces())
candidates = _image_generation_commands.filter_recovery(
_durable_generation_queue.list(exclude_ids=_jobs.keys()))
for record in candidates:
job_id = str(record.get("id") or "").strip()
params = record.get("params")
Expand Down Expand Up @@ -26357,6 +26367,9 @@ def resume_generation_queue():
def discard_generation_queue():
"""Clear only inactive recovery candidates; never cancel live work."""
with _queue_recovery_lock:
_image_generation_commands.restore_recovery(item["name"] for item in _list_workspaces())
_image_generation_commands.discard_recovery(
_durable_generation_queue.list(exclude_ids=_jobs.keys()))
removed = _durable_generation_queue.discard(exclude_ids=_jobs.keys())
return {"discarded": removed}

Expand Down Expand Up @@ -36150,7 +36163,7 @@ def _upsert_canonical_task(
return existing


def _publish_generation_task(job: dict) -> dict:
def _generation_task_fields(job: dict) -> dict:
legacy_id = str(job.get("id") or "")
workspace = str(job.get("workspace") or "default")
details = _public_generation_details(job.get("params"))
Expand Down Expand Up @@ -36214,9 +36227,9 @@ def _publish_generation_task(job: dict) -> dict:
task_title = "Tools · Upscale"
elif str(provenance.get("capability") or "") == "revoice":
task_title = "Tools · Revoice"
return _upsert_canonical_task(
workspace,
task_id,
return dict(
workspace=workspace,
id=task_id,
root_id=root_task_id,
parent_id=parent_task_id,
kind=mode,
Expand Down Expand Up @@ -36249,6 +36262,12 @@ def _publish_generation_task(job: dict) -> dict:
)


def _publish_generation_task(job: dict) -> dict:
fields = _generation_task_fields(job)
workspace, task_id = fields.pop("workspace"), fields.pop("id")
return _upsert_canonical_task(workspace, task_id, **fields)


def _observe_generation_job_state(record: dict) -> None:
"""Forward atomic lifecycle changes to the canonical task event stream."""
job = dict(record)
Expand Down Expand Up @@ -36831,11 +36850,20 @@ def _classic_redirect():
# Optional external agents use exactly the same admission endpoints and task IDs.
from routers.wangp_mcp import create_wangp_mcp_router
from services.wangp_agent_adapters import application_handlers as wangp_agent_handlers
from services.image_generation_runtime import create_image_generation_commands
from routers.image_generation_commands import (
create_image_generation_commands_router, image_command_catalog, image_command_handlers,
)
from services.workspace_commands import catalog as workspace_command_catalog

_image_generation_commands = create_image_generation_commands(globals())
api.include_router(create_image_generation_commands_router(_image_generation_commands))
api.include_router(create_wangp_mcp_router(
handlers={"models": lambda args: get_model_options(args['model_type']) if args.get('model_type') else list_models(), "processors": wangp_capabilities, "status": get_status,
"generate": generate, "recast": recast_endpoint, "upscale": tools_upscale,
**wangp_agent_handlers(api)},
**wangp_agent_handlers(api), **image_command_handlers(_image_generation_commands)},
journal_path=os.path.join(os.path.dirname(__file__), "settings", "wangp-mcp-requests.sqlite3"),
command_operations=[*workspace_command_catalog()["operations"], *image_command_catalog()],
))

# ============================================================================
Expand Down
125 changes: 125 additions & 0 deletions app/routers/image_generation_commands.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
"""HTTP and MCP projections of the executable image command contract."""
from __future__ import annotations

from fastapi import APIRouter, Request
from pydantic import BaseModel, ConfigDict, Field, StrictStr, ValidationError, field_validator
from services.image_generation_spec import image_generation_schema
from services.studio_image_spec import studio_image_schema
from services.image_generation_commands import command_error


class ReferenceResolutionInput(BaseModel):
model_config = ConfigDict(extra="forbid", strict=True)
references: list[StrictStr] = Field(min_length=1, max_length=64)


class UISubmissionContext(BaseModel):
"""Optional attribution, never a source of permissions or target IDs."""
model_config = ConfigDict(extra="forbid", strict=True)
workflowId: StrictStr | None = Field(default=None, min_length=1, max_length=200)
runId: StrictStr | None = Field(default=None, min_length=1, max_length=200)

@field_validator("workflowId", "runId")
@classmethod
def nonblank(cls, value):
if value is not None and (not value.strip() or value != value.strip()):
raise ValueError("Use an exact non-blank context ID")
return value


def _ui_context(request):
raw = request.headers.get("X-Hocus-UI-Context", "{}")
if len(raw) > 2048:
raise command_error(422, "invalid_ui_context", "Submission context is too long")
try:
return UISubmissionContext.model_validate_json(raw).model_dump(exclude_none=True)
except ValidationError as error:
raise command_error(422, "invalid_ui_context", "Use only exact workflowId and runId attribution") from error


def image_command_catalog():
spec = image_generation_schema()
studio = studio_image_schema()
studio_input = dict(studio["input"])
definitions = studio_input.pop("$defs", {})
definitions["StudioCommandInput"] = studio_input
envelope = {"type": "object", "additionalProperties": False,
"properties": {"version": {"type": "integer", "enum": [1, 2]},
"operation": {"const": "generation.image"},
"intent_id": spec["intent_id"], "input": {"type": "object"}},
"required": ["version", "operation", "intent_id", "input"],
"$defs": definitions,
"oneOf": [
{"properties": {"version": {"const": 1}, "input": spec["input"]}},
{"properties": {"version": {"const": 2}, "input": {"$ref": "#/$defs/StudioCommandInput"}}},
]}
receipt_input = {"type": "object", "additionalProperties": False,
"properties": {"workspace": spec["input"]["properties"]["workspace"],
"intent_id": spec["intent_id"]}, "required": ["workspace", "intent_id"]}
return [{"name": "generation.image", "version": 2, "supportedVersions": [1, 2], "domain": "studio", "mutation": True,
"description": "Admit an image job with an installed model and explicit output workspace. Version 1 is a single text-to-image request; version 2 accepts the complete typed Studio image parameters, canonical references, LoRAs and image processors. Preserve literal prompts and reuse intent_id only for retries. The receipt proves admission; inspect its task for completion.",
"inputSchema": envelope},
{"name": "generation.receipt", "version": 1, "domain": "studio", "mutation": False,
"description": "Read an immutable image admission and its current canonical task in the exact original output workspace.",
"inputSchema": {"type": "object", "additionalProperties": False,
"properties": {"version": {"type": "integer", "const": 1},
"operation": {"const": "generation.receipt"}, "input": receipt_input},
"required": ["version", "operation", "input"]}}]


def image_command_handlers(service):
async def submit(arguments):
if not isinstance(arguments, dict) or set(arguments) != {"version", "intent_id", "input"}:
raise command_error(422, "invalid_command", "Use version, intent_id and input for the image tool")
return await service.submit({**arguments, "operation": "generation.image"}, trusted_tool="external_agent")

def receipt(arguments):
if (not isinstance(arguments, dict) or set(arguments) != {"version", "input"}
or type(arguments.get("version")) is not int or arguments["version"] != 1
or not isinstance(arguments["input"], dict)
or set(arguments["input"]) != {"workspace", "intent_id"}):
raise command_error(422, "invalid_command", "Use version 1 with workspace and intent_id")
return service.receipt(**arguments["input"])

return {"generation.image": submit, "generation.receipt": receipt}


def create_image_generation_commands_router(service):
router = APIRouter()

@router.get("/api/v1/generation/commands")
def catalog():
return {"version": 2, "operations": image_command_catalog()}

@router.post("/api/v1/generation/commands")
async def submit(request: Request):
try:
command = await request.json()
except ValueError as error:
raise command_error(422, "invalid_command", "Command must be valid JSON") from error
surface = request.headers.get("X-Hocus-UI-Surface", "studio")
if surface not in {"studio", "wizard"}:
raise command_error(422, "invalid_ui_surface", "Choose a known initiating UI surface")
# Declared UI attribution, as on the native Studio endpoint. This is
# never used for authorization. MCP supplies its own external context.
return await service.submit(command, trusted_tool="wizard" if surface == "wizard" else None,
submission_context=_ui_context(request))

@router.get("/api/v1/generation/commands/receipt")
def receipt(workspace: str, intent_id: str):
return service.receipt(workspace, intent_id)

@router.post("/api/v1/generation/commands/references")
def references(body: ReferenceResolutionInput):
"""Read-only migration of exact legacy UI paths into canonical URLs."""
resolve = getattr(service, "canonicalize_reference", None)
if not callable(resolve):
raise command_error(503, "reference_resolution_unavailable", "Reference resolution is unavailable")
if any(not 1 <= len(value) <= 8192 for value in body.references):
raise command_error(422, "invalid_reference", "An exact bounded media reference is required")
try:
return {"references": [resolve(value) for value in body.references]}
except (ValueError, OSError) as error:
raise command_error(422, "invalid_reference", str(error)) from error

return router
Loading
Loading