diff --git a/app/_launch_runtime.py b/app/_launch_runtime.py index a077eae0..9a88461e 100644 --- a/app/_launch_runtime.py +++ b/app/_launch_runtime.py @@ -449,7 +449,7 @@ async def trace_user_mutations(request: Request, call_next): unregister_abort_state, update_job, ) -from services.asset_manifest import publish_generation_sidecar +from services.asset_manifest import publish_generation_sidecar, publish_generation_sidecar_best_effort from services import resource_scheduler _jobs: dict = {} @@ -552,6 +552,7 @@ def _persist_generation_job(job: dict) -> None: "created_at": job.get("created_at", time.time()), "params": copy.deepcopy(job.get("params") or {}), "workspace": job.get("workspace") or "default", + "provenance": copy.deepcopy(job.get("provenance") or {}), }) except Exception as exc: # Persistence should protect a generation, never prevent it from @@ -576,6 +577,7 @@ def _new_generation_job( created_at: float | None = None, recovered: bool = False, reserve_generation: bool = True, + provenance: dict | None = None, ) -> dict: frozen_params = copy.deepcopy(params) execution_mode.validate_generation(workspace) @@ -610,8 +612,19 @@ def _new_generation_job( f"{len(timeline['intervals'])} timed interval(s) across " f"{timeline['duration_seconds']:.3f}s." ) + resolved_job_id = job_id or uuid.uuid4().hex[:8] + canonical_task_id = f"task-generation-{resolved_job_id}" + owner_id = str(frozen_params.get("_director_pipeline_id") or "") + if owner_id.startswith("series:"): + canonical_root_task_id = f"task-series-render-{owner_id.split(':', 1)[1]}" + elif owner_id: + canonical_root_task_id = f"task-director-{owner_id}" + else: + canonical_root_task_id = canonical_task_id job = { - "id": job_id or uuid.uuid4().hex[:8], + "id": resolved_job_id, + "task_id": canonical_task_id, + "root_task_id": canonical_root_task_id, "status": "queued", "progress": 0, "step": 0, @@ -627,6 +640,7 @@ def _new_generation_job( "error": None, "workspace": workspace, "out_dir": _workspace_dir(workspace), + "provenance": copy.deepcopy(provenance or {}), "recovered": recovered, } # Reserve FIFO order synchronously. Starting one thread per request is @@ -639,13 +653,54 @@ def _new_generation_job( try: task = publisher(job) if isinstance(task, dict): - job["task_id"] = task.get("id") - job["root_task_id"] = task.get("root_id") + job["task_id"] = task.get("id") or job["task_id"] + job["root_task_id"] = ( + task.get("root_id") or task.get("id") or job["root_task_id"] + ) except Exception as exc: print(f"[Task registry] Could not publish generation {job['id']}: {exc}") return job +def _publish_generation_sidecar_for_studio_job( + job: dict, + output_path: str, + sidecar: dict, + *, + tool: str = "studio", +) -> None: + """Publish one Studio result with its initiating command and real location.""" + provenance = job.get("provenance") if isinstance(job.get("provenance"), dict) else {} + command = provenance.get("command") if isinstance(provenance.get("command"), dict) else {} + payload = dict(sidecar) + params = dict(payload.get("params") or {}) + model_type = str(params.get("model_type") or "") + params.setdefault("provider", "minimax" if model_type.startswith("minimax:") else "local") + payload["params"] = params + if not payload.get("job_id"): + payload["job_id"] = job.get("id") + if not payload.get("task_id"): + payload["task_id"] = job.get("task_id") + if not payload.get("root_task_id"): + payload["root_task_id"] = job.get("root_task_id") or job.get("task_id") + payload["created_at"] = job.get("created_at") or payload.get("created_at") or time.time() + payload["queued_at"] = job.get("created_at") or payload.get("queued_at") + payload["started_at"] = job.get("started_at") or payload.get("started_at") + payload["completed_at"] = job.get("finished_at") or time.time() + for key in ("command_id", "workflow_id", "run_id"): + if command.get(key): + payload.setdefault(key, command[key]) + publish_generation_sidecar_best_effort( + output_path, + payload, + workspace_id=provenance.get("workspace_id"), + output_folder=job.get("workspace"), + tool=tool, + actor=provenance.get("actor"), + capability=provenance.get("capability"), + ) + + def _register_manual_generation_job(job: dict) -> dict: """Publish older edit/tool jobs before their background thread starts.""" job_id = str(job.get("id") or "") @@ -7439,7 +7494,13 @@ def model3d_capabilities(): @api.post("/api/v1/model3d/generate") async def generate_model3d(request: Request): from services import model3d_service + from services.generation_provenance import normalize_submission_provenance + body = await request.json() + body["provenance"] = normalize_submission_provenance(body.pop("provenance", None)) + collection_id = body["provenance"].get("workspace_id") + if collection_id and not _workspace_collection_registry.get(collection_id): + raise HTTPException(status_code=400, detail="Unknown Workspace collection") workspace = body.get("workspace") if "workspace" in body else _get_active_workspace() _workspace_dir(workspace) try: @@ -11150,7 +11211,13 @@ def _plan_windows(): @api.post("/api/v1/generate") async def generate(request: Request): """Submit a generation job. Returns immediately with a job_id.""" + from services.generation_provenance import normalize_submission_provenance + body = await request.json() + provenance = normalize_submission_provenance(body.pop("provenance", None)) + collection_id = provenance.get("workspace_id") + if collection_id and not _workspace_collection_registry.get(collection_id): + raise HTTPException(status_code=400, detail="Unknown Workspace collection") # Execution mode is a boot-time trust boundary, never a request option. # Discard spoofed private fields before validating the captured workspace. body.pop("_execution_mode", None) @@ -11619,6 +11686,7 @@ async def generate(request: Request): body, workspace, reserve_generation=not h3_preplan_pending, + provenance=provenance, ) job_id = job["id"] job["out_dir"] = job_out_dir @@ -22842,15 +22910,9 @@ def _run_sfx_generation(job: dict, raw_params: dict, start_time: float): "generation_time": round(elapsed), "created_at": time.time(), } - try: - publish_generation_sidecar( - os.path.join(out_dir, fname), - sidecar, - workspace_id=job.get("workspace"), - tool="studio-sfx", - ) - except Exception: - pass + _publish_generation_sidecar_for_studio_job( + job, os.path.join(out_dir, fname), sidecar, + ) completed = finish_job( job, @@ -23647,12 +23709,7 @@ def publish_progress(message: str, value: int, step: int, total: int) -> None: "simulated": True, "execution_mode": "simulate", } - publish_generation_sidecar( - generated_path, - sidecar, - workspace_id=job.get("workspace"), - tool="studio", - ) + _publish_generation_sidecar_for_studio_job(job, generated_path, sidecar) if not finalize: return update_job( job, @@ -23852,12 +23909,7 @@ def _legacy_h3_progress( for path in generated: file_sidecar = dict(sidecar) file_sidecar["output_filename"] = os.path.basename(path) - publish_generation_sidecar( - path, - file_sidecar, - workspace_id=job.get("workspace"), - tool="studio-h3-legacy", - ) + _publish_generation_sidecar_for_studio_job(job, path, file_sidecar) if not finalize: return update_job( @@ -24484,15 +24536,9 @@ def _write_output_sidecars(file_names): else: file_sidecar.pop("director_clip_index", None) file_sidecar["output_filename"] = fname - try: - publish_generation_sidecar( - os.path.join(out_dir, fname), - file_sidecar, - workspace_id=job.get("workspace"), - tool="studio", - ) - except Exception: - pass + _publish_generation_sidecar_for_studio_job( + job, os.path.join(out_dir, fname), file_sidecar, + ) is_multiclip = total_tasks > 1 and any(t.get('params', {}).get('multi_clip_info') for t in queue) @@ -26403,6 +26449,7 @@ def resume_generation_queue(): reserve_generation=not isinstance( params.get("_h3_window_plan_pending"), dict, ), + provenance=record.get("provenance") if isinstance(record.get("provenance"), dict) else None, ) _jobs[job_id] = job _persist_generation_job(job) @@ -36154,6 +36201,8 @@ def _publish_generation_task(job: dict) -> dict: task_id = f"task-generation-{legacy_id}" params = job.get("params") if isinstance(job.get("params"), dict) else {} owner_id = str(params.get("_director_pipeline_id") or "") + provenance = job.get("provenance") if isinstance(job.get("provenance"), dict) else {} + command = provenance.get("command") if isinstance(provenance.get("command"), dict) else {} if owner_id.startswith("series:"): series_job_id = owner_id.split(":", 1)[1] parent_task_id = f"task-series-render-{series_job_id}" @@ -36215,6 +36264,12 @@ def _publish_generation_task(job: dict) -> dict: metadata={ "adapter": "generation", "generation_details": details, "owner_pipeline_id": owner_id, + "actor": provenance.get("actor") or "unknown", + "tool": provenance.get("tool") or "studio", + "capability": provenance.get("capability"), + "command_id": command.get("command_id"), + "workflow_id": command.get("workflow_id"), + "run_id": command.get("run_id"), }, ) @@ -36364,6 +36419,8 @@ def _publish_generic_legacy_task(record: dict, adapter: str) -> dict | None: or "" ) or None request_body = record.get("request") if isinstance(record.get("request"), dict) else {} + provenance = record.get("provenance") if isinstance(record.get("provenance"), dict) else {} + command = provenance.get("command") if isinstance(provenance.get("command"), dict) else {} provider = str( record.get("provider") or request_body.get("writingProvider") @@ -36425,11 +36482,20 @@ def _publish_generic_legacy_task(record: dict, adapter: str) -> dict | None: and adapter != "comic-plan"), resumable=resumable, recoverable=resumable, error=({"message": str(record.get("error")), "retryable": resumable} if record.get("error") else None), - result_refs=list(record.get("output_files") or ([record["output"]] if record.get("output") else [])), + result_refs=list(record.get("output_files") or ( + [record["output"]] if record.get("output") else + [record["filename"]] if record.get("filename") else [] + )), metadata={ "adapter": adapter, "cancel_mode": record.get("cancel_mode"), "safe_boundary": record.get("safe_boundary"), + "actor": provenance.get("actor") or "unknown", + "tool": provenance.get("tool") or adapter, + "capability": provenance.get("capability"), + "command_id": command.get("command_id"), + "workflow_id": command.get("workflow_id"), + "run_id": command.get("run_id"), }, ) diff --git a/app/services/asset_manifest.py b/app/services/asset_manifest.py index c3139517..ea742244 100644 --- a/app/services/asset_manifest.py +++ b/app/services/asset_manifest.py @@ -395,6 +395,7 @@ def adapt_legacy_sidecar( parameters=params, timing={ "created_at": legacy.get("created_at"), + "queued_at": legacy.get("queued_at"), "started_at": started_at, "completed_at": completed_at, "inference_ms": inference_ms, diff --git a/app/services/generation_provenance.py b/app/services/generation_provenance.py index 85adc22a..6af97619 100644 --- a/app/services/generation_provenance.py +++ b/app/services/generation_provenance.py @@ -14,6 +14,7 @@ INITIATORS = frozenset({"user", "wizard", "system", "unknown"}) +_SUBMISSION_COMMAND_FIELDS = ("command_id", "workflow_id", "run_id") class CommandContext(TypedDict, total=False): @@ -70,6 +71,37 @@ def resolve_generation_location( return {"workspace_id": None, "output_folder": None} +def normalize_submission_provenance(value: Any) -> GenerationProvenance: + """Validate the optional provenance attached to a generation request. + + This is attribution data, not an authorization boundary. Runtime-owned + identifiers (job/task/pipeline) and the physical output folder are added + by the backend and therefore cannot be supplied by the browser. + """ + raw = value if isinstance(value, Mapping) else {} + actor = _clean(raw.get("actor")) or "unknown" + if actor not in INITIATORS: + actor = "unknown" + capability = _clean(raw.get("capability")) + workspace_id = _clean(raw.get("workspace_id")) + command_raw = raw.get("command") if isinstance(raw.get("command"), Mapping) else {} + command: CommandContext = {} + for key in _SUBMISSION_COMMAND_FIELDS: + cleaned = _clean(command_raw.get(key)) + if cleaned: + command[key] = cleaned[:200] + result: GenerationProvenance = { + "actor": actor, + "tool": "studio", + "command": command, + } + if capability: + result["capability"] = capability[:200] + if workspace_id: + result["workspace_id"] = workspace_id[:200] + return result + + def provenance_from_manifest(manifest: Mapping[str, Any] | None) -> GenerationProvenance: """Project a canonical manifest onto initiator vs provider/model vs location.""" value = manifest if isinstance(manifest, Mapping) else {} @@ -101,5 +133,6 @@ def provenance_from_manifest(manifest: Mapping[str, Any] | None) -> GenerationPr __all__ = [ "CommandContext", "GenerationLocation", "GenerationProvenance", "INITIATORS", - "provenance_from_manifest", "resolve_generation_location", + "normalize_submission_provenance", "provenance_from_manifest", + "resolve_generation_location", ] diff --git a/app/services/model3d_service.py b/app/services/model3d_service.py index 8240e7b6..70ee8f3e 100644 --- a/app/services/model3d_service.py +++ b/app/services/model3d_service.py @@ -517,6 +517,31 @@ def _physical_output_folder(value: Any) -> str | None: return name or None +def _publish_model3d_result(job: dict[str, Any], output_path: str | Path, sidecar: dict[str, Any]) -> None: + provenance = job.get("provenance") if isinstance(job.get("provenance"), dict) else {} + command = provenance.get("command") if isinstance(provenance.get("command"), dict) else {} + payload = dict(sidecar) + payload.setdefault("job_id", job.get("job_id")) + payload.setdefault("task_id", job.get("task_id")) + payload.setdefault("root_task_id", job.get("root_task_id") or job.get("task_id")) + payload["created_at"] = job.get("created_at") or payload.get("created_at") or time.time() + payload["queued_at"] = job.get("created_at") or payload.get("queued_at") + payload["started_at"] = job.get("started_at") or payload.get("started_at") + payload["completed_at"] = job.get("finished_at") or time.time() + for key in ("command_id", "workflow_id", "run_id"): + if command.get(key): + payload.setdefault(key, command[key]) + publish_generation_sidecar_best_effort( + output_path, + payload, + workspace_id=provenance.get("workspace_id"), + output_folder=_physical_output_folder(job.get("workspace")), + tool=provenance.get("tool") or "model3d", + actor=provenance.get("actor"), + capability=provenance.get("capability"), + ) + + def _prune_finished_jobs_locked() -> None: """Drop old terminal jobs; callers must hold _lock.""" now = time.time() @@ -622,6 +647,7 @@ def _start_remote_job( "model_id": model_id, "provider": provider, "workspace": str(workspace or "default"), + "provenance": dict(body.get("provenance") or {}), "created_at": time.time(), "updated_at": time.time(), "request": { @@ -658,7 +684,10 @@ def cancelled() -> bool: services = _services() stem = f"{provider}-{job_id[:8]}" try: - _update_job(job_id, status="running", phase="running", progress=0.1, message=f"Calling {provider}") + _update_job( + job_id, status="running", phase="running", progress=0.1, + message=f"Calling {provider}", started_at=time.time(), + ) if cancelled(): _settle_cancelled_job(job_id) return @@ -700,6 +729,19 @@ def cancelled() -> bool: else: raise RuntimeError(f"Unknown 3D provider: {provider}") filename = result["filename"] + with _lock: + current_job = dict(_jobs.get(job_id) or job) + _publish_model3d_result(current_job, os.path.join(output_dir, filename), { + "generation_mode": "model3d", + "mode": "model3d", + "params": { + "prompt": request_data.get("prompt"), + "model_id": request_data.get("model"), + "model_type": request_data.get("model"), + "provider": provider, + "image_path": os.path.basename(str(request_data.get("image_path") or "")) or None, + }, + }) _update_job( job_id, status="completed", @@ -770,6 +812,7 @@ def start_job( "operation": request_data["operation"], "model_id": request_data["model"]["id"], "workspace": str(workspace or "default"), + "provenance": dict(body.get("provenance") or {}), "created_at": time.time(), "updated_at": time.time(), "request": request_data, @@ -907,6 +950,7 @@ def _spawn_worker_if_active( "phase": "starting", "message": message, "progress": 0.02, + "started_at": time.time(), "updated_at": time.time(), }) process = subprocess.Popen( @@ -932,6 +976,7 @@ def _run_job_serialized(job_id: str, output_dir: str) -> None: _update_job( job_id, status="running", phase="simulated_inference", progress=0.08, message="Simulating Hunyuan3D inference…", + started_at=time.time(), ) try: output = execution_mode.create_artifact( @@ -943,6 +988,20 @@ def _run_job_serialized(job_id: str, output_dir: str) -> None: cancelled=lambda: bool((_jobs.get(job_id) or {}).get("cancel_requested")), ) filename = os.path.basename(output) + with _lock: + current_job = dict(_jobs.get(job_id) or job) + _publish_model3d_result(current_job, output, { + "generation_mode": "model3d", + "mode": "model3d", + "simulated": True, + "execution_mode": "simulate", + "params": { + **request_data.get("settings", {}), + "model_id": (request_data.get("model") or {}).get("id"), + "model_type": (request_data.get("model") or {}).get("id"), + "provider": "hunyuan3d", + }, + }) _update_job( job_id, status="completed", phase="completed", progress=1.0, message="3D model ready · simulated artifact", filename=filename, @@ -1129,7 +1188,8 @@ def _watchdog() -> None: # The mesh is on disk. Sidecar/status failures must not delete it. generation_committed = True - publish_generation_sidecar_best_effort( + _publish_model3d_result( + current_job, output_path, { "generation_mode": "model3d", @@ -1149,8 +1209,6 @@ def _watchdog() -> None: "images": request_data["images"], }, }, - output_folder=_physical_output_folder(current_job.get("workspace")), - tool="model3d", ) _update_job( job_id, diff --git a/tests/test_execution_mode.py b/tests/test_execution_mode.py index ca05dfa4..70a637cd 100644 --- a/tests/test_execution_mode.py +++ b/tests/test_execution_mode.py @@ -13,7 +13,12 @@ import pytest from services import execution_mode -from services.asset_manifest import SCHEMA_NAME, publish_generation_sidecar, read_asset_manifest +from services.asset_manifest import ( + SCHEMA_NAME, + publish_generation_sidecar, + publish_generation_sidecar_best_effort, + read_asset_manifest, +) def _load_launch_function(name, namespace): @@ -219,14 +224,23 @@ def create_artifact(*_args, **_kwargs): "record_job_outputs": lambda _job, names: recorded.extend(names), "finish_job": lambda *_args, **_kwargs: True, "publish_generation_sidecar": publish_generation_sidecar, + "publish_generation_sidecar_best_effort": publish_generation_sidecar_best_effort, "os": os, "time": time, } + _load_launch_function("_publish_generation_sidecar_for_studio_job", namespace) worker = _load_launch_function("_run_simulated_generation", namespace) job = { "id": "sim-job-1", + "created_at": 1_700_000_000.0, + "started_at": 1_700_000_002.0, "out_dir": str(tmp_path), "workspace": "night-shift", + "provenance": { + "actor": "wizard", + "capability": "start_generation", + "command": {"command_id": "command-1"}, + }, "task_id": "task-1", "root_task_id": "task-1", "params": { @@ -246,8 +260,16 @@ def create_artifact(*_args, **_kwargs): assert raw["job_id"] == "sim-job-1" assert loaded is not None assert loaded["asset"]["kind"] == "video" - assert loaded["origin"]["workspace_id"] == "night-shift" + assert "workspace_id" not in loaded["origin"] + assert loaded["origin"]["output_folder"] == "night-shift" + assert loaded["origin"]["actor"] == "wizard" + assert loaded["origin"]["capability"] == "start_generation" + assert loaded["execution"]["command_id"] == "command-1" + assert loaded["generation"]["model"]["provider"] == "local" assert loaded["execution"]["mode"] == "simulate" + assert loaded["timing"]["queued_at"] == "2023-11-14T22:13:20Z" + assert loaded["timing"]["started_at"] == "2023-11-14T22:13:22Z" + assert loaded["timing"]["queue_ms"] == 2_000 assert loaded["technical"]["published_on_generate"] is True diff --git a/tests/test_generation_provenance_submission.py b/tests/test_generation_provenance_submission.py new file mode 100644 index 00000000..c6a7ca0e --- /dev/null +++ b/tests/test_generation_provenance_submission.py @@ -0,0 +1,43 @@ +import unittest + +from app.services.generation_provenance import normalize_submission_provenance + + +class GenerationSubmissionProvenanceTests(unittest.TestCase): + def test_normalizes_browser_owned_fields_without_accepting_runtime_ids(self): + value = normalize_submission_provenance({ + "actor": "wizard", + "tool": "spoofed", + "capability": "start_generation", + "workspace_id": "collection-7", + "output_folder": "not-browser-owned", + "command": { + "command_id": "command-1", + "workflow_id": "workflow-2", + "run_id": "run-3", + "task_id": "spoofed-task", + "job_id": "spoofed-job", + "pipeline_id": "spoofed-pipeline", + }, + }) + self.assertEqual(value["actor"], "wizard") + self.assertEqual(value["tool"], "studio") + self.assertEqual(value["capability"], "start_generation") + self.assertEqual(value["workspace_id"], "collection-7") + self.assertEqual(value["command"], { + "command_id": "command-1", + "workflow_id": "workflow-2", + "run_id": "run-3", + }) + self.assertNotIn("output_folder", value) + + def test_invalid_or_missing_actor_remains_unknown(self): + self.assertEqual(normalize_submission_provenance(None)["actor"], "unknown") + self.assertEqual( + normalize_submission_provenance({"actor": "administrator"})["actor"], + "unknown", + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_h3_preplan_job_contract.py b/tests/test_h3_preplan_job_contract.py index f40c155d..58ee6674 100644 --- a/tests/test_h3_preplan_job_contract.py +++ b/tests/test_h3_preplan_job_contract.py @@ -146,12 +146,19 @@ def task_context_scope(**_context): "services.task_manager", task_context_scope=task_context_scope, ) + provenance_module = _module( + "services.generation_provenance", + normalize_submission_provenance=lambda value: value or { + "actor": "unknown", "tool": "studio", "command": {}, + }, + ) services_package = _module( "services", h3_window_planner=planner_module, llm_service=llm_module, minimax_h3_duration=duration_module, task_manager=task_module, + generation_provenance=provenance_module, ) services_package.__path__ = [] monkeypatch.setitem(sys.modules, "services", services_package) @@ -171,6 +178,11 @@ def task_context_scope(**_context): duration_module, ) monkeypatch.setitem(sys.modules, "services.task_manager", task_module) + monkeypatch.setitem( + sys.modules, + "services.generation_provenance", + provenance_module, + ) def _base_body() -> dict: @@ -187,7 +199,12 @@ def _base_body() -> dict: } -def _harness(monkeypatch, tmp_path: Path, planner) -> tuple[dict, dict, float]: +def _harness( + monkeypatch, + tmp_path: Path, + planner, + publisher=None, +) -> tuple[dict, dict, float]: _DeferredThread.instances = [] _install_h3_fakes(monkeypatch, planner) events: list[tuple[str, str]] = [] @@ -201,10 +218,11 @@ def _harness(monkeypatch, tmp_path: Path, planner) -> tuple[dict, dict, float]: "minimax_h3_text_encoder_variants": {"qwen-test": {}}, } - def publish(job: dict) -> dict: - task_id = f"task-generation-{job['id']}" - events.append(("publish", job["id"])) - return {"id": task_id, "root_id": task_id} + if publisher is None: + def publisher(job: dict) -> dict: + task_id = f"task-generation-{job['id']}" + events.append(("publish", job["id"])) + return {"id": task_id, "root_id": task_id} def persist(job: dict) -> None: events.append(("persist", job["id"])) @@ -283,8 +301,11 @@ def acknowledge_cancel(job: dict, **updates) -> bool: "_normalize_video_prompt_type": lambda _body: None, "_normalize_image_prompt_type": lambda _body: None, "_get_active_workspace": lambda: "default", + "normalize_submission_provenance": lambda value: value or { + "actor": "unknown", "tool": "studio", "command": {}, + }, "_workspace_dir": lambda _workspace=None: str(tmp_path), - "_publish_generation_task": publish, + "_publish_generation_task": publisher, "_persist_generation_job": persist, "_remove_persisted_generation_job": lambda job: events.append(("remove", job["id"])), "_cancel_h3_idle_release": lambda: None, @@ -343,6 +364,27 @@ def slow_planner(*_args, **_kwargs): assert _DeferredThread.instances[0].target.__name__ == "_run_generation_with_preparation" +def test_submission_keeps_canonical_ids_when_activity_publication_fails( + tmp_path, monkeypatch, +): + def broken_publisher(_job: dict) -> dict: + raise RuntimeError("activity unavailable") + + namespace, result, _elapsed = _harness( + monkeypatch, + tmp_path, + lambda *_args, **_kwargs: {}, + publisher=broken_publisher, + ) + + job = namespace["_jobs"][result["job_id"]] + assert result["status"] == "queued" + assert result["task_id"] == f"task-generation-{result['job_id']}" + assert result["root_task_id"] == result["task_id"] + assert job["task_id"] == result["task_id"] + assert job["root_task_id"] == result["root_task_id"] + + def test_cancel_before_preplanning_never_enters_or_reserves_gpu( tmp_path, monkeypatch, ): diff --git a/tests/test_provenance_3d_director.py b/tests/test_provenance_3d_director.py index a88bb8dd..f0d6253f 100644 --- a/tests/test_provenance_3d_director.py +++ b/tests/test_provenance_3d_director.py @@ -155,6 +155,38 @@ def test_model3d_writer_passes_output_folder_not_workspace_collection(tmp_path, assert proven["model_id"] == "hunyuan3d-2-turbo" +def test_model3d_writer_keeps_wizard_command_provenance(tmp_path): + output = tmp_path / "wizard.glb" + output.write_bytes(b"mesh") + task_id = "canonical-model3d-wizard-job" + model3d_service._publish_model3d_result({ + "job_id": "wizard-job", + "task_id": task_id, + "root_task_id": task_id, + "workspace": "physical-folder", + "created_at": 1_700_000_000, + "started_at": 1_700_000_002, + "provenance": { + "actor": "wizard", + "tool": "studio", + "capability": "start_generation", + "command": {"command_id": "command-wizard-3d"}, + }, + }, output, { + "generation_mode": "model3d", + "params": {"model_type": "hunyuan3d-2-turbo", "provider": "hunyuan3d"}, + }) + + loaded = read_asset_manifest(output) + assert loaded["origin"]["tool"] == "studio" + assert loaded["origin"]["actor"] == "wizard" + assert loaded["origin"]["capability"] == "start_generation" + assert loaded["origin"]["output_folder"] == "physical-folder" + assert loaded["origin"].get("workspace_id") in (None, "") + assert loaded["execution"]["command_id"] == "command-wizard-3d" + assert loaded["execution"]["task_id"] == task_id + + def test_model3d_strips_absolute_output_folder_path(tmp_path, monkeypatch): job_id = "model3d-abs-folder" jobs_dir = tmp_path / "model3d-jobs" diff --git a/ui/e2e/live-specs/wizard-generation.spec.ts b/ui/e2e/live-specs/wizard-generation.spec.ts index 05cb4a5a..b2717daa 100644 --- a/ui/e2e/live-specs/wizard-generation.spec.ts +++ b/ui/e2e/live-specs/wizard-generation.spec.ts @@ -127,6 +127,7 @@ async function waitForTerminalRoot( return terminalStatus }, { timeout: 20 * 60_000, intervals: [500, 1_000, 2_000, 5_000] }).toMatch(/completed|failed|cancelled/) if (expectedStatus === 'completed') expect(terminalStatus).toBe('completed') + return taskId } async function attachEvidence(page: Page, request: APIRequestContext, testInfo: TestInfo, transcript: string) { @@ -162,7 +163,34 @@ test('wizard: Studio UI → canonical queue → generated video', async ({ page, ? 'Abre Studio → Video y rellena visiblemente el formulario con un plano de 5 segundos de un mago programador ante servidores. No lo generes.' : 'Abre Studio → Video, rellena visiblemente el formulario con un plano de 5 segundos de un mago programador ante servidores, y genéralo ahora. Decide tú los demás valores compatibles.') if (expectedMode !== 'plan') { - await waitForTerminalRoot(request, String(config.execution_workspace), before) + const workspace = String(config.execution_workspace) + const taskId = await waitForTerminalRoot(request, workspace, before, 'completed') + const tasks = await json(request, `/api/v1/tasks?status=all&workspace=${encodeURIComponent(workspace)}`) as { + tasks: Array<{ id: string; result_refs?: string[]; metadata?: Record }> + } + const task = tasks.tasks.find(item => item.id === taskId) + const wizardTrace = await page.evaluate(() => ( + window as Window & { __HOCUSPOCUS_WIZARD_TRACE__?: Array> } + ).__HOCUSPOCUS_WIZARD_TRACE__ || []) as Array<{ results?: Array<{ action?: { type?: string }; command?: { commandId?: string }; report?: { taskId?: string } }> }> + const generationResult = wizardTrace.flatMap(item => item.results || []) + .find(item => item.action?.type === 'start_generation') + const commandId = generationResult?.command?.commandId + expect(commandId).toBeTruthy() + expect(generationResult?.report?.taskId).toBe(taskId) + expect(task?.metadata?.actor).toBe('wizard') + expect(task?.metadata?.capability).toBe('start_generation') + expect(task?.metadata?.command_id).toBe(commandId) + const outputName = task?.result_refs?.[0] + expect(outputName).toBeTruthy() + const metadata = await json( + request, + `/api/v1/outputs/${encodeURIComponent(String(outputName))}/metadata?workspace=${encodeURIComponent(workspace)}`, + ) as { origin?: Record; execution?: Record } + expect(metadata.origin?.actor).toBe('wizard') + expect(metadata.origin?.capability).toBe('start_generation') + expect(metadata.origin?.output_folder).toBe(workspace) + expect(metadata.origin).not.toHaveProperty('workspace_id') + expect(metadata.execution?.command_id).toBe(commandId) } else { expect(await rootTaskIds(request, String(config.execution_workspace))).toEqual(before) } diff --git a/ui/src/api/model3d.ts b/ui/src/api/model3d.ts index 1d40e661..3ed466d9 100644 --- a/ui/src/api/model3d.ts +++ b/ui/src/api/model3d.ts @@ -42,6 +42,8 @@ export interface Hunyuan3DCapabilities { export interface Hunyuan3DJob { job_id: string + task_id?: string + root_task_id?: string operation?: 'generate' | 'retexture' status: 'queued' | 'waiting' | 'waiting_resource' | 'running' | 'cancelling' | 'completed' | 'failed' | 'cancelled' progress: number @@ -84,6 +86,7 @@ export async function startHunyuan3DJob(params: { reduce_face?: boolean target_face_num?: number mc_algo?: string + provenance?: Record }): Promise { const res = await fetch(`${BASE}/api/v1/model3d/generate`, { method: 'POST', diff --git a/ui/src/components/Sidebar/GenerateButton.tsx b/ui/src/components/Sidebar/GenerateButton.tsx index 7f9f7987..68f7149d 100644 --- a/ui/src/components/Sidebar/GenerateButton.tsx +++ b/ui/src/components/Sidebar/GenerateButton.tsx @@ -2,6 +2,7 @@ import { useState, useEffect } from 'react' import { Play, AlertTriangle } from 'lucide-react' import { useStore } from '../../stores/useStore' import { splitPromptSchedule } from '../../lib/promptScheduler' +import { newUserGenerationContext } from '../../features/studio/generationProvenance' export function GenerateButton() { const jobs = useStore(s => s.jobs) @@ -54,7 +55,7 @@ export function GenerateButton() { const handleClick = () => { if (blocked) return setCooldown(true) - startGeneration() + startGeneration(undefined, newUserGenerationContext()) setSidebarOpen(false) } diff --git a/ui/src/features/agent/agentActions.ts b/ui/src/features/agent/agentActions.ts index f5262966..5b329bdb 100644 --- a/ui/src/features/agent/agentActions.ts +++ b/ui/src/features/agent/agentActions.ts @@ -2906,7 +2906,10 @@ export async function executeAgentActions( const reusedTaskId = String(reused?.taskId || '') const reusedJob = reusedTaskId ? useStore.getState().jobs.find(job => ( - job.id === reusedTaskId || `task-generation-${job.id}` === reusedTaskId + job.id === reusedTaskId + || job.taskId === reusedTaskId + || job.rootTaskId === reusedTaskId + || `task-generation-${job.id}` === reusedTaskId )) : undefined const reusableTask = !reusedJob || !['failed', 'cancelled', 'canceled'].includes(reusedJob.status) diff --git a/ui/src/features/agent/applicationAdapters.ts b/ui/src/features/agent/applicationAdapters.ts index bf9c33f2..4c6fe893 100644 --- a/ui/src/features/agent/applicationAdapters.ts +++ b/ui/src/features/agent/applicationAdapters.ts @@ -40,6 +40,7 @@ import type { AgentTrackCharacterKitJobAction, AgentUpdateCharacterKitAction, } from './characterKitActions' +import type { GenerationSubmissionContext } from '../studio/generationProvenance' export interface AdapterOutcome { message: string @@ -67,10 +68,10 @@ export interface StudioAdapter { prepareImage(action: AgentPrepareImageAction): Promise prepareAudio(action: AgentPrepareAudioAction): Promise prepare3d(action: AgentPrepare3dAction): Promise - startGeneration(action: AgentStartGenerationAction): Promise + startGeneration(action: AgentStartGenerationAction, context?: GenerationSubmissionContext): Promise attachReferences(action: AgentAttachStudioReferencesAction): Promise configureLoras(action: AgentConfigureStudioLorasAction): Promise - queueSfxPack(action: AgentQueueSfxPackAction): Promise + queueSfxPack(action: AgentQueueSfxPackAction, context?: GenerationSubmissionContext): Promise } export interface StoryLabAdapter { @@ -314,9 +315,11 @@ export function createDefaultApplicationAdapters(): WizardApplicationAdapters { const { prepare3dForm } = await import('../studio/adapters') return presentStudioSliceResult(await prepare3dForm(action), '3D') }, - async startGeneration(action) { + async startGeneration(action, context) { const { startGeneration } = await import('../studio/adapters') - const result = await startGeneration() + const result = await startGeneration(context || { + actor: 'wizard', capability: action.type, + }) const presented = await presentStudioSliceResult(result, 'Studio generation') const taskId = result.taskIds[0] return { @@ -344,9 +347,11 @@ export function createDefaultApplicationAdapters(): WizardApplicationAdapters { const { configureLoras } = await import('../studio/adapters') return presentStudioSliceResult(await configureLoras(action), 'Image / Video') }, - async queueSfxPack(action) { + async queueSfxPack(action, context) { const { queueSfx } = await import('../studio/adapters') - return presentStudioSliceResult(await queueSfx(action), 'Audio → SFX') + return presentStudioSliceResult(await queueSfx(action, context || { + actor: 'wizard', capability: action.type, + }), 'Audio → SFX') }, } adapters.storyLab = { diff --git a/ui/src/features/agent/capabilityRegistry.ts b/ui/src/features/agent/capabilityRegistry.ts index 9d7c1c7a..ec2f1dce 100644 --- a/ui/src/features/agent/capabilityRegistry.ts +++ b/ui/src/features/agent/capabilityRegistry.ts @@ -43,6 +43,7 @@ import type { AgentApplyCharacterKitPresetAction, AgentAttachCharacterKitReferen import { registerStudioCapabilities } from './studioCapabilities' import { registerNavigationQueueCapabilities } from './navigationQueueCapabilities' import { registerEditorAuxCapabilities } from './editorAuxCapabilities' +import type { GenerationSubmissionContext } from '../studio/generationProvenance' export const AGENT_TABS = [ 'studio', 'director', 'productions', 'images', 'videos', 'audio', '3d', @@ -75,6 +76,7 @@ export interface CapabilityExecutionContext { adapters: WizardApplicationAdapters workspace?: string onStep?: (message: string) => void + generationContext?: GenerationSubmissionContext } export interface CapabilityDefinition { diff --git a/ui/src/features/agent/capabilityRunner.ts b/ui/src/features/agent/capabilityRunner.ts index a8ee3125..d0d870eb 100644 --- a/ui/src/features/agent/capabilityRunner.ts +++ b/ui/src/features/agent/capabilityRunner.ts @@ -88,9 +88,18 @@ export async function runRegisteredCapability( requireConfirmation(prepared, definition.confirmation === 'required') const executionCommandId = commandId() + const executionContext: CapabilityRunnerOptions = { + ...options, + generationContext: { + ...options.generationContext, + actor: 'wizard', + capability: prepared.type, + commandId: executionCommandId, + }, + } stage(options, 'execute', action.type) - const executed = await definition.execute(prepared, options) + const executed = await definition.execute(prepared, executionContext) stage(options, 'correlate', action.type) const target = definition.correlate(prepared, executed) @@ -117,7 +126,7 @@ export async function runRegisteredCapability( } stage(options, 'track', action.type) - const tracked = await definition.track(prepared, executed, options) + const tracked = await definition.track(prepared, executed, executionContext) stage(options, 'report', action.type) const message = definition.summarize(prepared, tracked) diff --git a/ui/src/features/agent/studioCapabilities.ts b/ui/src/features/agent/studioCapabilities.ts index ebb25f1d..a7214855 100644 --- a/ui/src/features/agent/studioCapabilities.ts +++ b/ui/src/features/agent/studioCapabilities.ts @@ -276,7 +276,7 @@ export function registerStudioCapabilities(register: typeof defineCapability): v resolve: sfxAction, validate(action) { return action.confirm === true && action.clips.length > 0 ? validType('queue_sfx_pack', action) : ['confirmed SFX clips are required'] }, async prepare(action) { return action }, - async execute(action, context) { return context.adapters.studio.queueSfxPack(action) }, + async execute(action, context) { return context.adapters.studio.queueSfxPack(action, context.generationContext) }, correlate(_action, outcome) { return outcome.target }, async track(_action, outcome) { return outcome }, report: { targetKind: 'studio_sfx_pack', successState: 'completed' }, summarize(_action, outcome) { return outcome.message }, presentation: commonPresentation(['audio-mode', 'sfx-pack', 'queue']), @@ -293,7 +293,7 @@ export function registerStudioCapabilities(register: typeof defineCapability): v resolve(raw) { return raw.type === 'start_generation' && raw.confirm === true ? { type: 'start_generation', confirm: true } : null }, validate(action) { return validType('start_generation', action) }, async prepare(action) { return action }, - async execute(action, context) { return context.adapters.studio.startGeneration(action) }, + async execute(action, context) { return context.adapters.studio.startGeneration(action, context.generationContext) }, correlate(_action, outcome) { return outcome.target }, async track(_action, outcome) { return outcome }, report: { targetKind: 'generation_task', successState: 'completed' }, summarize(_action, outcome) { return outcome.message }, presentation: commonPresentation(['generate', 'queue']), diff --git a/ui/src/features/studio/actions.ts b/ui/src/features/studio/actions.ts index af21b252..f896de18 100644 --- a/ui/src/features/studio/actions.ts +++ b/ui/src/features/studio/actions.ts @@ -11,6 +11,10 @@ import type { PrepareVideoCommand, QueueSfxPackCommand, } from './commands' +import { + generationProvenancePayload, + type GenerationSubmissionContext, +} from './generationProvenance' export type StudioSfxClip = { name: string @@ -114,7 +118,10 @@ export async function selectAudioModel( return selectedModel?.name || selected } -export async function queueSfxPack(action: QueueSfxPackCommand): Promise { +export async function queueSfxPack( + action: QueueSfxPackCommand, + context?: GenerationSubmissionContext, +): Promise { if (!action.confirm) throw new Error('Encolar el pack de SFX requiere confirm=true tras una petición explícita.') if (!action.clips.length) throw new Error('El pack de SFX no incluye clips.') openStudioAudio('sfx') @@ -124,7 +131,7 @@ export async function queueSfxPack(action: QueueSfxPackCommand): Promise !before.has(job)) if (!created) throw new Error(`HocusPocus no encoló el efecto ${clip.name}.`) if (created.status === 'failed') throw new Error(created.error || created.message || `Falló ${clip.name}.`) @@ -339,7 +346,7 @@ export async function prepareAudio(action: PrepareAudioCommand): Promise { +export async function startPreparedGeneration(context?: GenerationSubmissionContext): Promise { const state = useStore.getState() if (state.generationMode === 'model3d') { const { startHunyuan3DJob } = await import('../../api/client') @@ -350,29 +357,31 @@ export async function startPreparedGeneration(): Promise { workspace: state.activeWorkspace || 'default', preset: prepared3dPreset, seed: typeof state.params.seed === 'number' ? state.params.seed : 1234, + provenance: generationProvenancePayload(context), }) - if (!job.job_id) throw new Error('Hunyuan3D devolvió éxito sin jobId; no considero la generación encolada.') + if (!job.task_id) throw new Error('Hunyuan3D devolvió éxito sin taskId; no considero la generación encolada.') return studioResult( 'generation', 'Studio generation', - `He enviado el modelo 3D a Hunyuan3D (${job.job_id}). Aparecerá en la galería 3D al terminar.`, - { taskId: job.job_id }, + `He enviado el modelo 3D a Hunyuan3D (${job.task_id}). Aparecerá en la galería 3D al terminar.`, + { taskId: job.task_id }, ) } const before = useStore.getState().jobs const knownJobs = new Set(before) - await useStore.getState().startGeneration() + await useStore.getState().startGeneration(undefined, context) const created = useStore.getState().jobs.find(job => !knownJobs.has(job)) if (!created) throw new Error('HocusPocus no creó una tarea; revisa los requisitos del modelo y los campos visibles.') if (created.status === 'failed') throw new Error(created.error || created.message || 'La generación no pudo entrar en cola.') - if (!created.id) throw new Error('HocusPocus devolvió éxito sin taskId; no considero la generación encolada.') + const taskId = created.taskId + if (!taskId) throw new Error('HocusPocus devolvió éxito sin taskId; no considero la generación encolada.') const mode = useStore.getState().generationMode const kind = mode === 'image' ? 'imagen' : mode === 'audio' ? 'pista de audio' : 'vídeo' return studioResult( 'generation', 'Studio generation', - `He enviado la ${kind} a la cola (${created.id}).`, - { taskId: created.id }, + `He enviado la ${kind} a la cola (${taskId}).`, + { taskId }, ) } diff --git a/ui/src/features/studio/adapters.ts b/ui/src/features/studio/adapters.ts index 15005ccd..0c21be28 100644 --- a/ui/src/features/studio/adapters.ts +++ b/ui/src/features/studio/adapters.ts @@ -25,6 +25,7 @@ import type { QueueSfxPackCommand, QueueTaskCommand, } from './commands' +import type { GenerationSubmissionContext } from './generationProvenance' export async function inspect(command: InspectQueueCommand) { return inspectCanonicalQueue(command.scope) @@ -58,8 +59,8 @@ export async function prepare3dForm(command: Prepare3dCommand) { return prepare3d(command) } -export async function startGeneration() { - return startPreparedGeneration() +export async function startGeneration(context?: GenerationSubmissionContext) { + return startPreparedGeneration(context) } export async function attachReferences(command: AttachStudioReferencesCommand) { @@ -70,6 +71,6 @@ export async function configureLoras(command: ConfigureStudioLorasCommand) { return configureStudioLoras(command) } -export async function queueSfx(command: QueueSfxPackCommand) { - return queueSfxPack(command) +export async function queueSfx(command: QueueSfxPackCommand, context?: GenerationSubmissionContext) { + return queueSfxPack(command, context) } diff --git a/ui/src/features/studio/generationProvenance.ts b/ui/src/features/studio/generationProvenance.ts new file mode 100644 index 00000000..d7487e33 --- /dev/null +++ b/ui/src/features/studio/generationProvenance.ts @@ -0,0 +1,35 @@ +export type GenerationActor = 'user' | 'wizard' | 'system' | 'unknown' + +export interface GenerationSubmissionContext { + actor: GenerationActor + capability?: string + commandId?: string + workflowId?: string + runId?: string + workspaceCollectionId?: string +} + +export function generationProvenancePayload(context?: GenerationSubmissionContext) { + if (!context) return undefined + const command = { + ...(context.commandId ? { command_id: context.commandId } : {}), + ...(context.workflowId ? { workflow_id: context.workflowId } : {}), + ...(context.runId ? { run_id: context.runId } : {}), + } + return { + actor: context.actor, + tool: 'studio', + ...(context.capability ? { capability: context.capability } : {}), + ...(context.workspaceCollectionId ? { workspace_id: context.workspaceCollectionId } : {}), + command, + } +} + +export function newUserGenerationContext(): GenerationSubmissionContext { + return { + actor: 'user', + capability: 'start_generation', + commandId: globalThis.crypto?.randomUUID?.() + || `command-${Date.now()}-${Math.random().toString(36).slice(2)}`, + } +} diff --git a/ui/src/stores/useStore.ts b/ui/src/stores/useStore.ts index 0c3ff0f4..dc6b960a 100644 --- a/ui/src/stores/useStore.ts +++ b/ui/src/stores/useStore.ts @@ -24,6 +24,10 @@ import { splitStudioClipPrompts, } from '../features/studio/studioRestore' import { GALLERY_LIST_FILTERS, galleryListQuery } from '../lib/galleryListQuery' +import { + generationProvenancePayload, + type GenerationSubmissionContext, +} from '../features/studio/generationProvenance' const DASHBOARD_PIPELINE_PAGE_SIZE = 8 const CIVIT_DOWNLOAD_POLL_MS = 2000 @@ -1609,7 +1613,10 @@ interface AppState { isGenerating: boolean promptSchedulerEnabled: boolean setPromptSchedulerEnabled: (enabled: boolean) => void - startGeneration: (scheduledPrompt?: ScheduledPromptSubmission) => Promise + startGeneration: ( + scheduledPrompt?: ScheduledPromptSubmission, + submissionContext?: GenerationSubmissionContext, + ) => Promise stopGeneration: (jobId?: string) => void dismissJob: (jobId: string) => void reconnectJobs: () => Promise @@ -4522,7 +4529,7 @@ export const useStore = create((set, get) => { promptSchedulerEnabled: false, setPromptSchedulerEnabled: (enabled) => set({ promptSchedulerEnabled: enabled }), - startGeneration: async (scheduledPrompt) => { + startGeneration: async (scheduledPrompt, submissionContext) => { const initialState = get() // Studio Prompt Scheduler: each non-empty line becomes its own normal @@ -4543,7 +4550,7 @@ export const useStore = create((set, get) => { prompt: scheduledPrompts[index], position: index + 1, total: scheduledPrompts.length, - }) + }, submissionContext) } return } @@ -5162,6 +5169,8 @@ export const useStore = create((set, get) => { } const params: Record = { ...state.params, generation_mode: state.generationMode, workspace: state.activeWorkspace } + const provenance = generationProvenancePayload(submissionContext) + if (provenance) params.provenance = provenance if (scheduledPrompt) { params.prompt = scheduledPrompt.prompt params.repeat_generation = 1 @@ -5905,7 +5914,7 @@ export const useStore = create((set, get) => { }) try { - const { job_id, h3_window_plan } = await api.submitGeneration(params) + const { job_id, task_id, root_task_id, h3_window_plan } = await api.submitGeneration(params) if (h3_window_plan) { const planFps = state.modelOptions?.fps ?? 24 @@ -5923,6 +5932,8 @@ export const useStore = create((set, get) => { jobs: s.jobs.map(j => j === newJob ? { ...j, id: job_id, + taskId: task_id || undefined, + rootTaskId: root_task_id || task_id || undefined, status: 'queued', message: 'Queued...', h3WindowPlan: h3_window_plan ?? null, diff --git a/ui/src/types/index.ts b/ui/src/types/index.ts index 0c558047..a84bb4c8 100644 --- a/ui/src/types/index.ts +++ b/ui/src/types/index.ts @@ -299,6 +299,9 @@ export interface GenerationDetails { export interface GenerationJob { id: string + /** Canonical Activity identity; distinct from the backend polling job id. */ + taskId?: string + rootTaskId?: string status: 'queued' | 'waiting_resource' | 'running' | 'cancelling' | 'completed' | 'failed' | 'cancelled' progress: number step: number diff --git a/ui/tests/agentActions.test.mjs b/ui/tests/agentActions.test.mjs index 3f1c6e07..5fdd9fec 100644 --- a/ui/tests/agentActions.test.mjs +++ b/ui/tests/agentActions.test.mjs @@ -998,6 +998,7 @@ test('start_generation reports the real taskId and an identical repeat reuses it jobs: useStore.getState().jobs, } let generationCalls = 0 + const generationContexts = [] useStore.setState({ modelsLoaded: true, loadOutputs: async () => {}, @@ -1017,11 +1018,14 @@ test('start_generation reports the real taskId and an identical repeat reuses it params: { ...useStore.getState().params, model_type: 'flux-test', prompt: '' }, jobs: [], loadModelOptions: async () => {}, - startGeneration: async () => { + startGeneration: async (_scheduledPrompt, context) => { generationCalls += 1 + generationContexts.push(context) useStore.setState({ jobs: [{ id: `job-studio-${generationCalls}`, + taskId: `canonical-generation-job-studio-${generationCalls}`, + rootTaskId: `canonical-generation-job-studio-${generationCalls}`, status: 'queued', progress: 0, step: 0, @@ -1045,11 +1049,14 @@ test('start_generation reports the real taskId and an identical repeat reuses it assert.equal(first[0].ok, true) assert.equal(first[1].ok, true) assert.equal(first[1].report.state, 'queued') - assert.equal(first[1].report.taskId, 'job-studio-1') + assert.equal(first[1].report.taskId, 'canonical-generation-job-studio-1') assert.equal(generationCalls, 1) + assert.equal(generationContexts[0].actor, 'wizard') + assert.equal(generationContexts[0].capability, 'start_generation') + assert.equal(generationContexts[0].commandId, first[1].command.commandId) const second = await executeAgentActions([prepare, { type: 'start_generation', confirm: true }]) assert.match(second[1].message, /Reutilizo/) - assert.equal(second[1].report.taskId, 'job-studio-1') + assert.equal(second[1].report.taskId, 'canonical-generation-job-studio-1') assert.equal(generationCalls, 1) useStore.setState({ @@ -1060,7 +1067,7 @@ test('start_generation reports the real taskId and an identical repeat reuses it const retry = await executeAgentActions([prepare, { type: 'start_generation', confirm: true }]) assert.equal(retry[1].ok, true) assert.doesNotMatch(retry[1].message, /Reutilizo/) - assert.equal(retry[1].report.taskId, 'job-studio-2') + assert.equal(retry[1].report.taskId, 'canonical-generation-job-studio-2') assert.equal(generationCalls, 2) } finally { useStore.setState(original) diff --git a/ui/tests/agentContract.test.mjs b/ui/tests/agentContract.test.mjs index 6b4256ee..86d8e89d 100644 --- a/ui/tests/agentContract.test.mjs +++ b/ui/tests/agentContract.test.mjs @@ -152,6 +152,31 @@ test('common capability runner follows every stage and reports the verified adap assert.equal(result.commandResult.navigationTarget, undefined) }) +test('Wizard Studio generation receives the command context before execution', async () => { + const { resolveAndRunRegisteredCapability } = await import('../src/features/agent/capabilityRunner.ts') + let received + const result = await resolveAndRunRegisteredCapability('start_generation', { + type: 'start_generation', confirm: true, + }, { + workspace: 'physical-output-folder', + adapters: { + studio: { + async startGeneration(_action, context) { + received = context + return { + message: 'Queued', taskId: 'canonical-generation-demo', + target: { kind: 'generation_task', id: 'canonical-generation-demo', title: 'Generation' }, + } + }, + }, + }, + }) + assert.equal(received.actor, 'wizard') + assert.equal(received.capability, 'start_generation') + assert.equal(received.commandId, result.command.commandId) + assert.equal('workspaceCollectionId' in received, false) +}) + test('application adapters navigate and verify targets without rendering React', async () => { const { createDefaultApplicationAdapters } = await import('../src/features/agent/applicationAdapters.ts') const { useStore } = await import('../src/stores/useStore.ts')