From 373ad979f4475dc82453783c58008c32f1b8ffb8 Mon Sep 17 00:00:00 2001 From: Your Name Date: Sun, 19 Jul 2026 00:14:43 +0800 Subject: [PATCH] fix(sre): bind Agent99 closure to typed verifier --- apps/api/src/api/v1/agents.py | 8 + .../services/agent99_completion_callback.py | 260 +++++++++++--- .../agent99_controlled_dispatch_ledger.py | 110 +++++- .../test_agent99_completion_callback_api.py | 320 ++++++++++++++++++ .../test_agent99_controlled_dispatch_p1.py | 152 ++++++++- 5 files changed, 799 insertions(+), 51 deletions(-) diff --git a/apps/api/src/api/v1/agents.py b/apps/api/src/api/v1/agents.py index 3b7b2f49f..0d8dcec57 100644 --- a/apps/api/src/api/v1/agents.py +++ b/apps/api/src/api/v1/agents.py @@ -1425,6 +1425,14 @@ async def post_agent99_completion_callback( status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail="agent99_completion_durable_readback_failed", ) + if receipt.get("ok") is not True: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=str( + receipt.get("status") + or "agent99_completion_closure_authority_pending" + ), + ) return receipt diff --git a/apps/api/src/services/agent99_completion_callback.py b/apps/api/src/services/agent99_completion_callback.py index 4a771e36e..6d336ed88 100644 --- a/apps/api/src/services/agent99_completion_callback.py +++ b/apps/api/src/services/agent99_completion_callback.py @@ -12,12 +12,20 @@ from src.models.agent99_completion import Agent99CompletionCallbackRequest from src.repositories.alert_operation_log_repository import ( get_alert_operation_log_repository, ) +from src.services.agent99_controlled_dispatch_ledger import ( + parse_agent99_dispatch_identity, + read_agent99_dispatch_receipt_for_run, +) from src.services.channel_hub import ( build_external_alert_provider_event_id, record_external_alert_event, ) _OPERATION_INCIDENT_ID_MAX_LENGTH = 30 +_AGENT99_CLOSURE_DOMAINS = frozenset({ + "control_plane_recovery", + "windows_vmware", +}) def _operation_incident_id(alert_id: str | None) -> str | None: @@ -40,31 +48,37 @@ async def _load_callback_readback( result = await db.execute( text( """ + WITH event_readback AS ( + SELECT + event_row.event_id::text AS event_id, + event_row.run_id::text AS run_id + FROM awooop_conversation_event event_row + WHERE event_row.project_id = :project_id + AND event_row.channel_type = 'internal' + AND event_row.provider_event_id = :provider_event_id + LIMIT 1 + ), + operation_counts AS ( + SELECT + count(*) AS operation_count, + count(*) FILTER ( + WHERE operation_row.event_type = 'EXECUTION_COMPLETED' + ) AS execution_count, + count(*) FILTER ( + WHERE operation_row.event_type = 'RESOLVED' + ) AS resolved_count + FROM alert_operation_log operation_row + WHERE operation_row.context ->> 'callback_id' = :callback_id + AND operation_row.context ->> 'project_id' = :project_id + ) SELECT - event_row.event_id::text AS event_id, - event_row.run_id::text AS run_id, - ( - SELECT count(*) - FROM alert_operation_log operation_row - WHERE operation_row.context ->> 'callback_id' = :callback_id - ) AS operation_count, - ( - SELECT count(*) - FROM alert_operation_log operation_row - WHERE operation_row.context ->> 'callback_id' = :callback_id - AND operation_row.event_type = 'EXECUTION_COMPLETED' - ) AS execution_count, - ( - SELECT count(*) - FROM alert_operation_log operation_row - WHERE operation_row.context ->> 'callback_id' = :callback_id - AND operation_row.event_type = 'RESOLVED' - ) AS resolved_count - FROM awooop_conversation_event event_row - WHERE event_row.project_id = :project_id - AND event_row.channel_type = 'internal' - AND event_row.provider_event_id = :provider_event_id - LIMIT 1 + event_readback.event_id, + event_readback.run_id, + operation_counts.operation_count, + operation_counts.execution_count, + operation_counts.resolved_count + FROM operation_counts + LEFT JOIN event_readback ON TRUE """ ), { @@ -77,7 +91,12 @@ async def _load_callback_readback( return dict(row) if row else None -def _operation_context(payload: Agent99CompletionCallbackRequest) -> dict[str, Any]: +def _operation_context( + payload: Agent99CompletionCallbackRequest, + *, + effective_outcome_state: str, + closure_authority: dict[str, Any], +) -> dict[str, Any]: return { "schema_version": payload.schema_version, "callback_id": payload.callback_id, @@ -90,7 +109,8 @@ def _operation_context(payload: Agent99CompletionCallbackRequest) -> dict[str, A "correlation_key": payload.correlation_key, "source": payload.source, "mode": payload.mode, - "outcome_state": payload.outcome_state, + "outcome_state": effective_outcome_state, + "claimed_outcome_state": payload.outcome_state, "controlled_apply": payload.controlled_apply, "transport_ok": payload.transport_ok, "verifier_name": payload.verifier_name, @@ -105,6 +125,17 @@ def _operation_context(payload: Agent99CompletionCallbackRequest) -> dict[str, A "telegram": payload.telegram.model_dump(mode="json"), "problem": payload.problem.model_dump(mode="json"), "occurred_at": payload.occurred_at.isoformat(), + "callback_authority_verified": bool( + closure_authority.get("verified") is True + ), + "closure_allowed": bool(closure_authority.get("verified") is True), + "closure_blocker": str( + closure_authority.get("blocker") or "unknown" + )[:120], + "typed_domain": str( + closure_authority.get("typed_domain") or "" + )[:120], + "route_id": str(closure_authority.get("route_id") or "")[:180], "stores_secret": False, "raw_log_stored": False, } @@ -117,12 +148,139 @@ def _owner_plane(alert_kind: str | None) -> str: return "awoooi" +def _blocked_closure_authority(blocker: str) -> dict[str, Any]: + return { + "required": True, + "verified": False, + "blocker": blocker, + "typed_domain": "", + "route_id": "", + } + + +async def _load_callback_closure_authority( + payload: Agent99CompletionCallbackRequest, +) -> dict[str, Any]: + """Bind a terminal callback to the typed dispatch and verifier ledger.""" + + receipt = await read_agent99_dispatch_receipt_for_run( + project_id=payload.project_id, + run_id=payload.run_id, + ) + if not isinstance(receipt, dict): + return _blocked_closure_authority("dispatch_run_not_found") + public_identity = receipt.get("identity") + if not isinstance(public_identity, dict): + return _blocked_closure_authority("dispatch_identity_missing") + try: + identity = parse_agent99_dispatch_identity(public_identity) + except ValueError: + return _blocked_closure_authority("dispatch_identity_invalid") + + expected_incident_ids = { + identity.incident_id, + _operation_incident_id(identity.incident_id), + } + if ( + identity.project_id != payload.project_id + or str(identity.run_id) != payload.run_id + or identity.trace_id != payload.trace_id + or identity.work_item_id != payload.work_item_id + or payload.alert_id not in expected_incident_ids + or str(receipt.get("run_id") or "") != payload.run_id + or str(receipt.get("trace_id") or "") != payload.trace_id + ): + return _blocked_closure_authority("dispatch_callback_identity_mismatch") + + scope = receipt.get("dispatch_scope") + if not isinstance(scope, dict): + return _blocked_closure_authority("dispatch_scope_missing") + typed_domain = str(scope.get("typed_domain") or "") + route_id = str(scope.get("route_id") or "") + expected_verifier = str(scope.get("verifier") or "") + if ( + scope.get("schema_version") != "agent99_dispatch_scope_v1" + or typed_domain not in _AGENT99_CLOSURE_DOMAINS + or scope.get("executor") != "Agent99" + or route_id != identity.route_id + or not expected_verifier + ): + return _blocked_closure_authority("typed_dispatch_scope_mismatch") + if ( + bool(scope.get("controlled_apply_requested") is True) + != payload.controlled_apply + or ( + payload.controlled_apply + and receipt.get("controlled_apply_authorized") is not True + ) + or receipt.get("dispatch_promoted") is not True + ): + return _blocked_closure_authority("controlled_dispatch_not_authorized") + if ( + payload.verifier_name != expected_verifier + or ( + str(scope.get("suggested_mode") or "") + and str(scope.get("suggested_mode")) != payload.mode + ) + ): + return _blocked_closure_authority("callback_route_projection_mismatch") + + verifier = receipt.get("verifier") + if not isinstance(verifier, dict): + return _blocked_closure_authority("independent_verifier_missing") + if ( + verifier.get("schema_version") + != "agent99_independent_verifier_receipt_v1" + or verifier.get("status") != "success" + or str(verifier.get("run_id") or "") != payload.run_id + or str(verifier.get("trace_id") or "") != payload.trace_id + or str(verifier.get("work_item_id") or "") != payload.work_item_id + or verifier.get("outcome_state") != "resolved" + or verifier.get("transport_ok") is not True + or verifier.get("verifier_passed") is not True + or verifier.get("source_event_resolved") is not True + or receipt.get("post_verifier_passed") is not True + or str(receipt.get("run_state") or "") not in {"waiting_tool", "completed"} + ): + return _blocked_closure_authority("independent_verifier_not_verified") + if ( + payload.transport_ok is not True + or payload.verifier_passed is not True + or payload.source_event_resolved is not True + ): + return _blocked_closure_authority("callback_verifier_projection_mismatch") + return { + "required": True, + "verified": True, + "blocker": "none", + "typed_domain": typed_domain, + "route_id": route_id, + } + + async def record_agent99_completion_callback( payload: Agent99CompletionCallbackRequest, ) -> dict[str, Any]: """Write one callback to the AWOOOI/AwoooP truth chain and read it back.""" - stage = payload.outcome_state + closure_authority = ( + await _load_callback_closure_authority(payload) + if payload.outcome_state == "resolved" + else { + "required": False, + "verified": False, + "blocker": "not_terminal_claim", + "typed_domain": "", + "route_id": "", + } + ) + terminal_resolved = bool( + payload.outcome_state == "resolved" + and closure_authority.get("verified") is True + ) + stage = "resolved" if terminal_resolved else ( + "verifying" if payload.outcome_state == "resolved" else payload.outcome_state + ) provider_event_id = build_external_alert_provider_event_id( "agent99", payload.callback_id, @@ -133,27 +291,32 @@ async def record_agent99_completion_callback( provider_event_id=provider_event_id, callback_id=payload.callback_id, ) - duplicate = before is not None - terminal_resolved = bool( - payload.outcome_state == "resolved" - and payload.verifier_passed - and payload.source_event_resolved + duplicate = bool( + before + and ( + before.get("event_id") + or int(before.get("operation_count") or 0) > 0 + ) ) severity = ( "info" - if payload.outcome_state == "resolved" + if stage == "resolved" else "critical" - if payload.outcome_state == "failed" + if stage == "failed" else "warning" ) - context = _operation_context(payload) + context = _operation_context( + payload, + effective_outcome_state=stage, + closure_authority=closure_authority, + ) operation_incident_id = _operation_incident_id(payload.alert_id) event_id = await record_external_alert_event( project_id=payload.project_id, provider="agent99", event_id=payload.callback_id, stage=stage, - title=f"Agent99 {payload.mode} {payload.outcome_state}", + title=f"Agent99 {payload.mode} {stage}", severity=severity, namespace="agent99", target_resource=payload.alert_service or payload.mode, @@ -162,10 +325,12 @@ async def record_agent99_completion_callback( labels={ "agent": "agent99", "mode": payload.mode, - "outcome_state": payload.outcome_state, + "outcome_state": stage, + "claimed_outcome_state": payload.outcome_state, "verifier_passed": str(payload.verifier_passed).lower(), "source_event_resolved": str(payload.source_event_resolved).lower(), "owner_plane": _owner_plane(payload.alert_kind), + "callback_authority_verified": str(terminal_resolved).lower(), }, annotations={ "run_id": payload.run_id, @@ -184,7 +349,7 @@ async def record_agent99_completion_callback( incident_id=operation_incident_id, actor="agent99_completion_callback", action_detail=f"{payload.mode}:{payload.outcome_state}", - success=payload.verifier_passed, + success=bool(payload.verifier_passed and stage != "verifying"), context=context, ) if terminal_resolved and int((before or {}).get("resolved_count") or 0) == 0: @@ -212,17 +377,30 @@ async def record_agent99_completion_callback( and execution_count > 0 and (not terminal_resolved or resolved_count > 0) ) + closure_pending = bool( + payload.outcome_state == "resolved" and not terminal_resolved + ) return { "schema_version": "agent99_completion_callback_receipt_v1", - "ok": durable_readback, - "status": "duplicate_verified" if duplicate else "recorded_verified", + "ok": bool(durable_readback and not closure_pending), + "status": ( + "closure_authority_pending" + if closure_pending + else "duplicate_verified" + if duplicate + else "recorded_verified" + ), "accepted": True, "duplicate": duplicate, "callback_id": payload.callback_id, "run_id": payload.run_id, "trace_id": payload.trace_id, "work_item_id": payload.work_item_id, - "outcome_state": payload.outcome_state, + "outcome_state": stage, + "claimed_outcome_state": payload.outcome_state, + "callback_authority_verified": terminal_resolved, + "closure_allowed": terminal_resolved, + "closure_blocker": str(closure_authority.get("blocker") or "unknown"), "conversation_event_id": str(event_id), "awooop_run_id": (readback or {}).get("run_id"), "operation_receipt_count": operation_count, diff --git a/apps/api/src/services/agent99_controlled_dispatch_ledger.py b/apps/api/src/services/agent99_controlled_dispatch_ledger.py index 27ce37e96..c4fe253af 100644 --- a/apps/api/src/services/agent99_controlled_dispatch_ledger.py +++ b/apps/api/src/services/agent99_controlled_dispatch_ledger.py @@ -1271,6 +1271,57 @@ class PostgresAgent99DispatchLedger: "runtime_closure_verified": False, } + dispatch_scope = ( + envelope.get("dispatch_scope") + if isinstance(envelope.get("dispatch_scope"), dict) + else {} + ) + expected_mode = str( + dispatch_scope.get("suggested_mode") or "" + ) + expected_verifier = str( + dispatch_scope.get("verifier") or "" + ) + requested_controlled_apply = bool( + dispatch_scope.get("controlled_apply_requested") is True + ) + supplied_controlled_apply = outcome_receipt.get( + "controlledApply" + ) + if ( + dispatch_scope.get("schema_version") + != "agent99_dispatch_scope_v1" + or dispatch_scope.get("route_id") != identity.route_id + or not str( + dispatch_scope.get("canonical_asset_id") or "" + ) + or not str(dispatch_scope.get("typed_domain") or "") + or dispatch_scope.get("executor") != "Agent99" + or not expected_mode + or not expected_verifier + or mode != expected_mode + or str(outcome.get("verifierName") or "") + != expected_verifier + or not isinstance(supplied_controlled_apply, bool) + or supplied_controlled_apply != requested_controlled_apply + ): + return { + "status": "verifier_dispatch_scope_mismatch_fail_closed", + "receipt_persisted": False, + "runtime_closure_verified": False, + } + if ( + requested_controlled_apply + and envelope.get("controlled_apply_authorized") is not True + ): + return { + "status": ( + "verifier_controlled_apply_not_authorized_fail_closed" + ), + "receipt_persisted": False, + "runtime_closure_verified": False, + } + if ( current.state in {"pending", "running"} or envelope.get("dispatch_promoted") is not True @@ -1282,8 +1333,10 @@ class PostgresAgent99DispatchLedger: "dispatch_accepted": True, "inbox_triggered": True, "dispatch_promoted": True, + # Outcome readback may prove delivery, but it cannot + # retroactively grant mutation authority. "controlled_apply_authorized": bool( - outcome_receipt.get("controlledApply") is True + envelope.get("controlled_apply_authorized") is True ), "dispatch_receipt": { "schema_version": "agent99_sre_dispatch_receipt_v1", @@ -2476,6 +2529,50 @@ class PostgresAgent99DispatchLedger: ) return None + async def read_for_run( + self, + *, + project_id: str, + run_id: str, + ) -> dict[str, Any] | None: + """Read one exact Agent99 dispatch run; malformed identities fail closed.""" + + try: + resolved_run_id = UUID(str(run_id or "").strip()) + except (TypeError, ValueError, AttributeError): + return None + try: + async with get_db_context(project_id or "awoooi") as db: + result = await db.execute( + select( + AwoooPRunState.run_id, + AwoooPRunState.state, + AwoooPRunState.trace_id, + AwoooPRunState.error_detail, + ).where( + AwoooPRunState.project_id == (project_id or "awoooi"), + AwoooPRunState.agent_id == AGENT99_DISPATCH_AGENT_ID, + AwoooPRunState.run_id == resolved_run_id, + ) + ) + row = result.one_or_none() + envelope = _parse_envelope(row.error_detail if row else None) + if not isinstance(envelope, dict) or row is None: + return None + return { + **envelope, + "run_state": str(row.state), + "run_id": str(row.run_id), + "trace_id": str(row.trace_id or ""), + } + except Exception as exc: + logger.warning( + "agent99_dispatch_exact_receipt_read_failed", + run_id=str(resolved_run_id), + error=str(exc), + ) + return None + async def list_reconcilable( self, *, @@ -2600,6 +2697,17 @@ async def read_agent99_dispatch_receipt( ) +async def read_agent99_dispatch_receipt_for_run( + *, + project_id: str, + run_id: str, +) -> dict[str, Any] | None: + return await _ledger.read_for_run( + project_id=project_id, + run_id=run_id, + ) + + async def list_agent99_reconcilable_dispatches( *, project_id: str = "awoooi", diff --git a/apps/api/tests/test_agent99_completion_callback_api.py b/apps/api/tests/test_agent99_completion_callback_api.py index ab1f2d2ba..dbe2d7f0c 100644 --- a/apps/api/tests/test_agent99_completion_callback_api.py +++ b/apps/api/tests/test_agent99_completion_callback_api.py @@ -17,6 +17,9 @@ from src.api.v1 import agents as agents_api from src.core.config import settings from src.models.agent99_completion import Agent99CompletionCallbackRequest from src.services import agent99_completion_callback as callback_service +from src.services.agent99_controlled_dispatch_ledger import ( + build_agent99_dispatch_identity, +) def payload() -> dict: @@ -132,6 +135,27 @@ def test_completion_callback_requires_durable_readback(monkeypatch) -> None: assert response.json()["detail"] == "agent99_completion_durable_readback_failed" +def test_completion_callback_keeps_unverified_closure_pending(monkeypatch) -> None: + monkeypatch.setattr(settings, "AGENT99_SRE_ALERT_RELAY_TOKEN", "expected") + + async def fake_record(_request): + return { + "ok": False, + "durable_readback": True, + "status": "closure_authority_pending", + } + + monkeypatch.setattr(agents_api, "record_agent99_completion_callback", fake_record) + response = app_client().post( + "/api/v1/agents/agent99/completion-callback", + headers={"X-Agent99-Completion-Token": "expected"}, + json=payload(), + ) + + assert response.status_code == 503 + assert response.json()["detail"] == "closure_authority_pending" + + def test_completion_callback_returns_same_trace_receipt(monkeypatch) -> None: monkeypatch.setattr(settings, "AGENT99_SRE_ALERT_RELAY_TOKEN", "expected") @@ -168,8 +192,301 @@ class FakeOperationRepository: return SimpleNamespace(id=f"operation-{len(self.events)}") +def authorize_terminal_callback(monkeypatch) -> None: + async def fake_authority(_payload): + return { + "required": True, + "verified": True, + "blocker": "none", + "typed_domain": "control_plane_recovery", + "route_id": "agent99:host_recovery:Recover", + } + + monkeypatch.setattr( + callback_service, + "_load_callback_closure_authority", + fake_authority, + ) + + +def typed_dispatch_callback_contract() -> tuple[ + Agent99CompletionCallbackRequest, + dict, +]: + identity = build_agent99_dispatch_identity( + project_id="awoooi", + incident_id="INC-20260711-001", + source_fingerprint="agent99-callback-source", + route_id="agent99_cold_start_recovery", + work_item_id="agent99-incident:INC-20260711-001", + ) + verifier_name = "recover_post_condition_v1" + request = Agent99CompletionCallbackRequest.model_validate({ + **payload(), + "run_id": str(identity.run_id), + "trace_id": identity.trace_id, + "work_item_id": identity.work_item_id, + "verifier_name": verifier_name, + }) + receipt = { + "schema_version": "agent99_controlled_dispatch_receipt_v1", + "status": "verifier_passed_learning_writeback_pending", + "identity": identity.public_dict(), + "run_id": str(identity.run_id), + "trace_id": identity.trace_id, + "run_state": "waiting_tool", + "dispatch_promoted": True, + "controlled_apply_authorized": True, + "post_verifier_passed": True, + "dispatch_scope": { + "schema_version": "agent99_dispatch_scope_v1", + "kind": "host_recovery", + "suggested_mode": "Recover", + "target_resource": "cold-start-gate", + "controlled_apply_requested": True, + "canonical_asset_id": "windows-vmware:host_99", + "typed_domain": "control_plane_recovery", + "executor": "Agent99", + "verifier": verifier_name, + "route_id": identity.route_id, + }, + "verifier": { + "schema_version": "agent99_independent_verifier_receipt_v1", + "status": "success", + "run_id": str(identity.run_id), + "trace_id": identity.trace_id, + "work_item_id": identity.work_item_id, + "outcome_state": "resolved", + "transport_ok": True, + "verifier_passed": True, + "source_event_resolved": True, + }, + } + return request, receipt + + +@pytest.mark.asyncio +async def test_callback_closure_authority_requires_exact_typed_ledger( + monkeypatch, +) -> None: + request, dispatch_receipt = typed_dispatch_callback_contract() + + async def fake_read(**_kwargs): + return dispatch_receipt + + monkeypatch.setattr( + callback_service, + "read_agent99_dispatch_receipt_for_run", + fake_read, + ) + + authority = await callback_service._load_callback_closure_authority(request) + + assert authority == { + "required": True, + "verified": True, + "blocker": "none", + "typed_domain": "control_plane_recovery", + "route_id": "agent99_cold_start_recovery", + } + + +@pytest.mark.asyncio +async def test_callback_closure_authority_rejects_cross_domain_scope( + monkeypatch, +) -> None: + request, dispatch_receipt = typed_dispatch_callback_contract() + dispatch_receipt["dispatch_scope"]["typed_domain"] = "docker_compose" + + async def fake_read(**_kwargs): + return dispatch_receipt + + monkeypatch.setattr( + callback_service, + "read_agent99_dispatch_receipt_for_run", + fake_read, + ) + + authority = await callback_service._load_callback_closure_authority(request) + + assert authority["verified"] is False + assert authority["blocker"] == "typed_dispatch_scope_mismatch" + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "field", + ["transport_ok", "verifier_passed", "source_event_resolved"], +) +async def test_callback_closure_authority_rejects_false_verifier_projection( + monkeypatch, + field: str, +) -> None: + request, dispatch_receipt = typed_dispatch_callback_contract() + request = request.model_copy(update={field: False}) + + async def fake_read(**_kwargs): + return dispatch_receipt + + monkeypatch.setattr( + callback_service, + "read_agent99_dispatch_receipt_for_run", + fake_read, + ) + + authority = await callback_service._load_callback_closure_authority(request) + + assert authority["verified"] is False + assert authority["blocker"] == "callback_verifier_projection_mismatch" + + +@pytest.mark.asyncio +async def test_completion_service_does_not_resolve_from_caller_booleans( + monkeypatch, +) -> None: + request = Agent99CompletionCallbackRequest.model_validate(payload()) + event_id = UUID("3c9e8ee7-fb21-5420-afae-ef928f75a58c") + repository = FakeOperationRepository() + readbacks = iter([ + None, + { + "event_id": str(event_id), + "run_id": "5e5b0080-e8c9-53a7-872b-0fc1ce80f11f", + "operation_count": 1, + "execution_count": 1, + "resolved_count": 0, + }, + ]) + + async def fake_authority(_payload): + return { + "required": True, + "verified": False, + "blocker": "dispatch_run_not_found", + "typed_domain": "", + "route_id": "", + } + + async def fake_readback(**_kwargs): + return next(readbacks) + + async def fake_record_event(**kwargs): + assert kwargs["stage"] == "verifying" + assert kwargs["severity"] == "warning" + return event_id + + monkeypatch.setattr( + callback_service, + "_load_callback_closure_authority", + fake_authority, + ) + monkeypatch.setattr(callback_service, "_load_callback_readback", fake_readback) + monkeypatch.setattr(callback_service, "record_external_alert_event", fake_record_event) + monkeypatch.setattr( + callback_service, + "get_alert_operation_log_repository", + lambda: repository, + ) + + receipt = await callback_service.record_agent99_completion_callback(request) + + assert receipt["ok"] is False + assert receipt["durable_readback"] is True + assert receipt["outcome_state"] == "verifying" + assert receipt["closure_allowed"] is False + assert receipt["closure_blocker"] == "dispatch_run_not_found" + assert [item[0] for item in repository.events] == ["EXECUTION_COMPLETED"] + assert repository.events[0][1]["success"] is False + + +@pytest.mark.asyncio +async def test_completion_service_retry_adds_only_resolved_receipt( + monkeypatch, +) -> None: + request = Agent99CompletionCallbackRequest.model_validate(payload()) + verifying_event_id = UUID("3c9e8ee7-fb21-5420-afae-ef928f75a58c") + resolved_event_id = UUID("5e5b0080-e8c9-53a7-872b-0fc1ce80f11f") + awooop_run_id = "8b877ed3-037f-5463-954b-f4b77a77d786" + repository = FakeOperationRepository() + authority_ready = {"value": False} + readbacks = iter([ + None, + { + "event_id": str(verifying_event_id), + "run_id": awooop_run_id, + "operation_count": 1, + "execution_count": 1, + "resolved_count": 0, + }, + { + "event_id": None, + "run_id": None, + "operation_count": 1, + "execution_count": 1, + "resolved_count": 0, + }, + { + "event_id": str(resolved_event_id), + "run_id": awooop_run_id, + "operation_count": 2, + "execution_count": 1, + "resolved_count": 1, + }, + ]) + event_calls: list[dict] = [] + + async def fake_authority(_payload): + verified = authority_ready["value"] + return { + "required": True, + "verified": verified, + "blocker": "none" if verified else "independent_verifier_missing", + "typed_domain": "control_plane_recovery" if verified else "", + "route_id": "agent99_cold_start_recovery" if verified else "", + } + + async def fake_readback(**_kwargs): + return next(readbacks) + + async def fake_record_event(**kwargs): + event_calls.append(kwargs) + return ( + resolved_event_id + if kwargs["stage"] == "resolved" + else verifying_event_id + ) + + monkeypatch.setattr( + callback_service, + "_load_callback_closure_authority", + fake_authority, + ) + monkeypatch.setattr(callback_service, "_load_callback_readback", fake_readback) + monkeypatch.setattr(callback_service, "record_external_alert_event", fake_record_event) + monkeypatch.setattr( + callback_service, + "get_alert_operation_log_repository", + lambda: repository, + ) + + pending = await callback_service.record_agent99_completion_callback(request) + authority_ready["value"] = True + resolved = await callback_service.record_agent99_completion_callback(request) + + assert pending["ok"] is False + assert resolved["ok"] is True + assert resolved["duplicate"] is True + assert [item[0] for item in repository.events] == [ + "EXECUTION_COMPLETED", + "RESOLVED", + ] + assert [call["stage"] for call in event_calls] == ["verifying", "resolved"] + assert [call["is_duplicate"] for call in event_calls] == [False, True] + + @pytest.mark.asyncio async def test_completion_service_writes_event_run_and_operation_receipts(monkeypatch) -> None: + authorize_terminal_callback(monkeypatch) request = Agent99CompletionCallbackRequest.model_validate(payload()) event_id = UUID("3c9e8ee7-fb21-5420-afae-ef928f75a58c") awooop_run_id = "5e5b0080-e8c9-53a7-872b-0fc1ce80f11f" @@ -221,6 +538,7 @@ async def test_completion_service_writes_event_run_and_operation_receipts(monkey async def test_completion_service_compacts_long_alert_id_for_operation_log( monkeypatch, ) -> None: + authorize_terminal_callback(monkeypatch) external_alert_id = "awoooi-agent99-ae36a454-c9eb-5f2d-9e35-be852b3eac9b" request = Agent99CompletionCallbackRequest.model_validate( {**payload(), "alert_id": external_alert_id} @@ -267,6 +585,7 @@ async def test_completion_service_compacts_long_alert_id_for_operation_log( @pytest.mark.asyncio async def test_completion_service_is_idempotent(monkeypatch) -> None: + authorize_terminal_callback(monkeypatch) request = Agent99CompletionCallbackRequest.model_validate(payload()) event_id = UUID("3c9e8ee7-fb21-5420-afae-ef928f75a58c") existing = { @@ -304,6 +623,7 @@ async def test_completion_service_is_idempotent(monkeypatch) -> None: @pytest.mark.asyncio async def test_completion_service_repairs_partial_terminal_receipt(monkeypatch) -> None: + authorize_terminal_callback(monkeypatch) request = Agent99CompletionCallbackRequest.model_validate(payload()) event_id = UUID("3c9e8ee7-fb21-5420-afae-ef928f75a58c") before = { diff --git a/apps/api/tests/test_agent99_controlled_dispatch_p1.py b/apps/api/tests/test_agent99_controlled_dispatch_p1.py index ca3e979b8..3b9104734 100644 --- a/apps/api/tests/test_agent99_controlled_dispatch_p1.py +++ b/apps/api/tests/test_agent99_controlled_dispatch_p1.py @@ -45,11 +45,13 @@ def _outcome(identity, *, transport_ok: bool = True) -> dict: return { "identity": identity.public_dict(), "controlledApply": True, + "mode": "Recover", "outcome": { "identity": identity.public_dict(), "schemaVersion": "agent99_outcome_contract_v1", "state": "resolved", "transportOk": transport_ok, + "verifierName": "recover_post_condition_v1", "verifierPassed": True, "sourceEventResolved": True, "verifiedAt": "2026-07-11T20:00:00+08:00", @@ -57,6 +59,30 @@ def _outcome(identity, *, transport_ok: bool = True) -> dict: } +def _promoted_dispatch(identity, *, authorized: bool = True) -> dict: + return build_agent99_dispatch_receipt_envelope( + identity=identity, + dispatch_receipt={ + "kind": "host_recovery", + "suggested_mode": "Recover", + "target_resource": "cold-start-gate", + "controlled_apply_requested": True, + "accepted": True, + "inbox_triggered": True, + "queue_accepted": True, + "dispatch_identity_matched": True, + "delivery_certainty": "delivered", + "dispatch_scope": { + "canonical_asset_id": "windows-vmware:host_99", + "typed_domain": "control_plane_recovery", + "executor": "Agent99", + "verifier": "recover_post_condition_v1", + }, + }, + controlled_apply_authorized=authorized, + ) + + def test_pending_inbox_never_promotes_or_authorizes() -> None: receipt = build_agent99_dispatch_receipt_envelope( identity=_identity(), @@ -91,6 +117,45 @@ def test_complete_identity_is_required_and_recomputed() -> None: parse_agent99_dispatch_identity(mismatched) +@pytest.mark.asyncio +async def test_dispatch_receipt_read_is_exact_run_only(monkeypatch) -> None: + identity = _identity() + ledger = PostgresAgent99DispatchLedger() + assert await ledger.read_for_run( + project_id=identity.project_id, + run_id="free-form-shadow-run", + ) is None + + envelope = _promoted_dispatch(identity) + + class ReadDB: + async def execute(self, _statement): + return _ScalarResult( + row=SimpleNamespace( + run_id=identity.run_id, + state="waiting_tool", + trace_id=identity.trace_id, + error_detail=json.dumps(envelope), + ) + ) + + monkeypatch.setattr( + ledger_module, + "get_db_context", + lambda _project_id: _Context(ReadDB()), + ) + + receipt = await ledger.read_for_run( + project_id=identity.project_id, + run_id=str(identity.run_id), + ) + + assert receipt is not None + assert receipt["run_id"] == str(identity.run_id) + assert receipt["trace_id"] == identity.trace_id + assert receipt["run_state"] == "waiting_tool" + + def test_approval_identity_cannot_alias_run_or_idempotency() -> None: common = { "project_id": "awoooi", @@ -335,15 +400,7 @@ async def test_verifier_requires_transport_and_all_evidence(monkeypatch) -> None assert missing["status"] == "verifier_evidence_missing_fail_closed" assert missing["receipt_persisted"] is False - promoted = build_agent99_dispatch_receipt_envelope( - identity=identity, - dispatch_receipt={ - "accepted": True, - "inbox_triggered": True, - "delivery_certainty": "delivered", - }, - controlled_apply_authorized=True, - ) + promoted = _promoted_dispatch(identity) class VerifierDB: def __init__(self) -> None: @@ -397,6 +454,83 @@ async def test_verifier_requires_transport_and_all_evidence(monkeypatch) -> None assert passed["runtime_closure_verified"] is False +@pytest.mark.asyncio +@pytest.mark.parametrize("mismatch", ["verifier", "mode", "controlled_apply"]) +async def test_verifier_rejects_dispatch_scope_mismatch( + monkeypatch, + mismatch: str, +) -> None: + identity = _identity() + promoted = _promoted_dispatch(identity) + outcome = _outcome(identity) + if mismatch == "verifier": + outcome["outcome"]["verifierName"] = "foreign-verifier" + elif mismatch == "mode": + outcome["mode"] = "Status" + else: + outcome["controlledApply"] = False + + class ScopeDB: + async def execute(self, _statement): + return _ScalarResult( + row=SimpleNamespace( + state="waiting_tool", + error_detail=json.dumps(promoted), + ) + ) + + monkeypatch.setattr( + ledger_module, + "get_db_context", + lambda _project_id: _Context(ScopeDB()), + ) + + result = await PostgresAgent99DispatchLedger().record_verifier( + identity=identity, + outcome_receipt=outcome, + evidence_refs=_evidence_refs(), + ) + + assert result["status"] == "verifier_dispatch_scope_mismatch_fail_closed" + assert result["receipt_persisted"] is False + assert result["runtime_closure_verified"] is False + + +@pytest.mark.asyncio +async def test_verifier_never_upgrades_controlled_apply_from_outcome( + monkeypatch, +) -> None: + identity = _identity() + promoted = _promoted_dispatch(identity, authorized=False) + + class ScopeDB: + async def execute(self, _statement): + return _ScalarResult( + row=SimpleNamespace( + state="waiting_tool", + error_detail=json.dumps(promoted), + ) + ) + + monkeypatch.setattr( + ledger_module, + "get_db_context", + lambda _project_id: _Context(ScopeDB()), + ) + + result = await PostgresAgent99DispatchLedger().record_verifier( + identity=identity, + outcome_receipt=_outcome(identity), + evidence_refs=_evidence_refs(), + ) + + assert result["status"] == ( + "verifier_controlled_apply_not_authorized_fail_closed" + ) + assert result["receipt_persisted"] is False + assert result["runtime_closure_verified"] is False + + @pytest.mark.asyncio async def test_learning_writeback_is_the_only_same_run_terminal_step( monkeypatch,