|
@@ -205,6 +205,7 @@ async def run_real_e2e(
|
|
|
phase_two_root = await _wait_for_root(
|
|
phase_two_root = await _wait_for_root(
|
|
|
client,
|
|
client,
|
|
|
script_build_id,
|
|
script_build_id,
|
|
|
|
|
+ legacy_state=host.composition.mission_service.legacy_state,
|
|
|
expected_status="blocked",
|
|
expected_status="blocked",
|
|
|
expected_reason=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
|
|
expected_reason=PHASE_TWO_CANDIDATE_PORTFOLIO_READY,
|
|
|
timeout_seconds=timeout_seconds,
|
|
timeout_seconds=timeout_seconds,
|
|
@@ -243,6 +244,7 @@ async def run_real_e2e(
|
|
|
completed_root = await _wait_for_root(
|
|
completed_root = await _wait_for_root(
|
|
|
client,
|
|
client,
|
|
|
script_build_id,
|
|
script_build_id,
|
|
|
|
|
+ legacy_state=host.composition.mission_service.legacy_state,
|
|
|
expected_status="completed",
|
|
expected_status="completed",
|
|
|
expected_reason=None,
|
|
expected_reason=None,
|
|
|
timeout_seconds=timeout_seconds,
|
|
timeout_seconds=timeout_seconds,
|
|
@@ -278,9 +280,7 @@ async def run_real_e2e(
|
|
|
)
|
|
)
|
|
|
final_payload = _successful_json(finalize, "finalize")
|
|
final_payload = _successful_json(finalize, "finalize")
|
|
|
if not final_payload.get("committed"):
|
|
if not final_payload.get("committed"):
|
|
|
- raise RuntimeError(
|
|
|
|
|
- "final publication did not report a committed transaction"
|
|
|
|
|
- )
|
|
|
|
|
|
|
+ raise RuntimeError("final publication did not report a committed transaction")
|
|
|
|
|
|
|
|
publication_response = await client.get(
|
|
publication_response = await client.get(
|
|
|
f"/api/pattern/script_builds/{script_build_id}/publication",
|
|
f"/api/pattern/script_builds/{script_build_id}/publication",
|
|
@@ -427,6 +427,7 @@ async def _wait_for_root(
|
|
|
client: httpx.AsyncClient,
|
|
client: httpx.AsyncClient,
|
|
|
script_build_id: int,
|
|
script_build_id: int,
|
|
|
*,
|
|
*,
|
|
|
|
|
+ legacy_state: Any | None = None,
|
|
|
expected_status: str,
|
|
expected_status: str,
|
|
|
expected_reason: str | None,
|
|
expected_reason: str | None,
|
|
|
timeout_seconds: float,
|
|
timeout_seconds: float,
|
|
@@ -460,6 +461,12 @@ async def _wait_for_root(
|
|
|
last_state = state
|
|
last_state = state
|
|
|
if state[0] in {"failed", "cancelled"}:
|
|
if state[0] in {"failed", "cancelled"}:
|
|
|
raise RuntimeError(f"root task entered terminal failure state: {state}")
|
|
raise RuntimeError(f"root task entered terminal failure state: {state}")
|
|
|
|
|
+ if legacy_state is not None:
|
|
|
|
|
+ build_status = await legacy_state.get_status(script_build_id)
|
|
|
|
|
+ if build_status in {BuildStatus.FAILED, BuildStatus.STOPPED}:
|
|
|
|
|
+ raise RuntimeError(
|
|
|
|
|
+ f"script build entered terminal status={build_status.value}; root={state}"
|
|
|
|
|
+ )
|
|
|
if state[0] == expected_status and (expected_reason is None or state[1] == expected_reason):
|
|
if state[0] == expected_status and (expected_reason is None or state[1] == expected_reason):
|
|
|
return root
|
|
return root
|
|
|
await asyncio.sleep(poll_seconds)
|
|
await asyncio.sleep(poll_seconds)
|
|
@@ -486,9 +493,7 @@ async def _wait_for_build_status(
|
|
|
if last in expected_statuses:
|
|
if last in expected_statuses:
|
|
|
return last
|
|
return last
|
|
|
if last in {BuildStatus.FAILED, BuildStatus.STOPPING, BuildStatus.STOPPED}:
|
|
if last in {BuildStatus.FAILED, BuildStatus.STOPPING, BuildStatus.STOPPED}:
|
|
|
- raise RuntimeError(
|
|
|
|
|
- f"build entered {last.value} while waiting for {expected_values}"
|
|
|
|
|
- )
|
|
|
|
|
|
|
+ raise RuntimeError(f"build entered {last.value} while waiting for {expected_values}")
|
|
|
await asyncio.sleep(poll_seconds)
|
|
await asyncio.sleep(poll_seconds)
|
|
|
raise TimeoutError(
|
|
raise TimeoutError(
|
|
|
f"timed out waiting for build status={expected_values}; "
|
|
f"timed out waiting for build status={expected_values}; "
|
|
@@ -592,12 +597,8 @@ async def _collect_diagnostics(host: Any, script_build_id: int) -> dict[str, Any
|
|
|
preset_usage: dict[str, dict[str, int | float]] = {}
|
|
preset_usage: dict[str, dict[str, int | float]] = {}
|
|
|
domain_errors: list[dict[str, str]] = []
|
|
domain_errors: list[dict[str, str]] = []
|
|
|
trace_presets = {binding.root_trace_id: "script_planner"}
|
|
trace_presets = {binding.root_trace_id: "script_planner"}
|
|
|
- trace_presets.update(
|
|
|
|
|
- {item.worker_trace_id: item.worker_preset for item in attempts}
|
|
|
|
|
- )
|
|
|
|
|
- trace_presets.update(
|
|
|
|
|
- {item.validator_trace_id: item.validator_preset for item in validations}
|
|
|
|
|
- )
|
|
|
|
|
|
|
+ trace_presets.update({item.worker_trace_id: item.worker_preset for item in attempts})
|
|
|
|
|
+ trace_presets.update({item.validator_trace_id: item.validator_preset for item in validations})
|
|
|
traces_with_usage: set[str] = set()
|
|
traces_with_usage: set[str] = set()
|
|
|
usage_reader = getattr(
|
|
usage_reader = getattr(
|
|
|
host.composition.mission_service.runner.trace_store, "get_model_usage", None
|
|
host.composition.mission_service.runner.trace_store, "get_model_usage", None
|
|
@@ -607,9 +608,7 @@ async def _collect_diagnostics(host: Any, script_build_id: int) -> dict[str, Any
|
|
|
model_usage = await usage_reader(trace_id)
|
|
model_usage = await usage_reader(trace_id)
|
|
|
models = model_usage.get("models", []) if isinstance(model_usage, dict) else []
|
|
models = model_usage.get("models", []) if isinstance(model_usage, dict) else []
|
|
|
agent_models = [
|
|
agent_models = [
|
|
|
- item
|
|
|
|
|
- for item in models
|
|
|
|
|
- if isinstance(item, dict) and item.get("source") == "agent"
|
|
|
|
|
|
|
+ item for item in models if isinstance(item, dict) and item.get("source") == "agent"
|
|
|
]
|
|
]
|
|
|
if agent_models:
|
|
if agent_models:
|
|
|
traces_with_usage.add(trace_id)
|
|
traces_with_usage.add(trace_id)
|