diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 6f95c4d..9ad7958 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -67,20 +67,62 @@ jobs: raise time.sleep(0.5) - request = urllib.request.Request( - "http://127.0.0.1:8000/v1/agent-specs/validate", - data=json.dumps( - { - "agent_id": "ci-smoke-agent", - "version": "1.0.0", - "display_name": "CI Smoke Agent", - "description": "Validates the packaged service contract.", - "entrypoint": "https://agents.example.test/ci-smoke", - } - ).encode(), - headers={"Content-Type": "application/json"}, - method="POST", + def call(path, method="GET", payload=None): + request = urllib.request.Request( + f"http://127.0.0.1:8000{path}", + data=None if payload is None else json.dumps(payload).encode(), + headers={"Content-Type": "application/json"}, + method=method, + ) + with urllib.request.urlopen(request, timeout=2) as response: + return json.load(response) + + specification = { + "agent_id": "ci-smoke-agent", + "version": "1.0.0", + "display_name": "CI Smoke Agent", + "description": "Validates the packaged service contract.", + "entrypoint": "https://agents.example.test/ci-smoke", + } + assert call("/v1/agent-specs/validate", "POST", specification)["valid"] is True + registered = call( + "/v1/agents", + "POST", + {"spec": specification, "actor": "ci@example.test"}, + ) + assert registered["revision"] == 1 + activated = call( + "/v1/agents/ci-smoke-agent/status", + "PATCH", + { + "status": "active", + "expected_revision": 1, + "actor": "ci@example.test", + "reason": "Container readiness checks passed.", + }, + ) + assert activated["status"] == "active" + + approval = call( + "/v1/approvals", + "POST", + { + "agent_id": "ci-smoke-agent", + "action": "deployment.promote", + "risk": "high", + "actor": "ci-smoke-agent", + "reason": "Exercise the packaged governance loop.", + }, + ) + decided = call( + f"/v1/approvals/{approval['request_id']}/decision", + "POST", + { + "decision": "approve", + "actor": "ci-reviewer@example.test", + "reason": "Container smoke evidence passed.", + }, ) - with urllib.request.urlopen(request, timeout=2) as response: - assert json.load(response)["valid"] is True + assert decided["status"] == "approved" + assert call("/v1/audit-events")["count"] == 4 PY diff --git a/CHANGELOG.md b/CHANGELOG.md index ea8dba7..72b4338 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,3 +10,5 @@ for public contracts once they are declared stable. - Initial FastAPI service with liveness, readiness, and `AgentSpec` validation. - Unit, smoke, lint, type, dependency audit, and container build automation. - Risk-based review and delivery policy. +- Agent registration and lifecycle status APIs with optimistic revision checks. +- Human approval queue with single-decision enforcement and append-only audit events. diff --git a/README.md b/README.md index 4e19457..529fcbc 100644 --- a/README.md +++ b/README.md @@ -42,6 +42,10 @@ The API is then available at `http://127.0.0.1:8000`. Important endpoints: - `GET /health/live` - `GET /health/ready` - `POST /v1/agent-specs/validate` +- `POST /v1/agents` and `PATCH /v1/agents/{agent_id}/status` +- `GET` and `POST /v1/approvals` +- `POST /v1/approvals/{request_id}/decision` +- `GET /v1/audit-events` - `GET /docs` Container execution: @@ -59,8 +63,13 @@ long-running checks run after merge and on a schedule. See ## Project status -The repository is in foundation stage. Public API compatibility starts with the `v1` schema; -runtime, storage, and workflow adapters are not yet production-ready. +The first governance loop is available: register an agent, activate or pause it with optimistic +revision checks, request and decide human approval, and inspect the resulting audit events. +Public API compatibility starts with the `v1` schema. + +The bundled store is intentionally in-memory and intended for development and evaluation. Data +does not survive a process restart and must not be treated as a production system of record. +PostgreSQL persistence, authenticated actor identity, and durable workflows remain planned. ## License diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index d8b9088..681339f 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -15,16 +15,24 @@ Existing Agent v Control Plane API |- AgentSpec validation + |- agent lifecycle and optimistic revision checks + |- human approval queue and append-only audit events |- trace and replay (planned) - |- policy and approval (planned) |- evaluation gates (planned) `- version promotion (planned) ``` -The current code implements the API shell and the first versioned contract. PostgreSQL becomes +The current code implements the API shell, the first versioned contract, and an in-memory +governance loop. The storage protocol is owned by the control plane so PostgreSQL can replace +the development adapter without leaking database types into the public API. PostgreSQL becomes the source of truth when persistence is introduced. Vector databases remain derived indexes, not authoritative stores. +State changes use an expected revision to reject stale writers. Only active agents can request +approval. Approval requests are single-decision records: an approved or rejected request cannot +be overwritten. Audit events are append-only within the store and returned newest first. +Authentication and durable audit retention are required before production use. + ## Adapter policy Temporal, Mem0, DSPy, LangSmith, and other providers must sit behind owned interfaces. A vendor diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index d43d0e2..169ac7e 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -4,16 +4,17 @@ Roadmap items advance only when tied to a validated user problem and an acceptan ## Foundation -- Versioned AgentSpec and event contracts. -- API health, readiness, and failure conventions. -- Pull request governance and automated quality gates. +- [x] Versioned AgentSpec and event contracts. +- [x] API health, readiness, and failure conventions. +- [x] Pull request governance and automated quality gates. ## Reliability gateway - Framework-neutral trace ingestion. - Run replay and failure classification. - Tool-call schema validation and risk policy. -- Human approval queue and immutable audit record. +- [x] In-memory human approval queue and append-only audit contract. +- PostgreSQL-backed approval and immutable audit persistence. - Offline evaluation datasets and version promotion gates. ## Durable operations diff --git a/src/agent_control_plane/api.py b/src/agent_control_plane/api.py index 8f077b5..e2e33d8 100644 --- a/src/agent_control_plane/api.py +++ b/src/agent_control_plane/api.py @@ -1,19 +1,50 @@ -"""HTTP surface for the initial control-plane contract.""" +"""HTTP surface for the control-plane contract.""" -from fastapi import FastAPI +from typing import Annotated, NoReturn +from uuid import UUID + +from fastapi import FastAPI, HTTPException, Query, status from agent_control_plane import __version__ from agent_control_plane.models import ( + AgentRecord, + AgentRegistrationRequest, AgentSpec, AgentSpecValidationResponse, + AgentStatusUpdate, + ApprovalDecisionRequest, + ApprovalQueueResponse, + ApprovalRecord, + ApprovalRequestCreate, + ApprovalStatus, + AuditEventPage, HealthResponse, HealthStatus, ) +from agent_control_plane.store import ( + AgentAlreadyExistsError, + AgentNotActiveError, + AgentNotFoundError, + ApprovalAlreadyDecidedError, + ApprovalNotFoundError, + ControlPlaneStore, + InMemoryControlPlaneStore, + InvalidStatusTransitionError, + RevisionConflictError, +) SERVICE_NAME = "agent-control-plane" -def create_app() -> FastAPI: +def _raise_http_error(status_code: int, code: str, error: Exception) -> NoReturn: + raise HTTPException( + status_code=status_code, + detail={"code": code, "message": str(error)}, + ) from error + + +def create_app(store: ControlPlaneStore | None = None) -> FastAPI: + control_plane = store if store is not None else InMemoryControlPlaneStore() application = FastAPI( title="Agent Control Plane", description="Reliability and governance APIs for production AI agents.", @@ -41,6 +72,98 @@ async def validate_agent_spec(spec: AgentSpec) -> AgentSpecValidationResponse: schema_version=spec.schema_version, ) + @application.post( + "/v1/agents", + response_model=AgentRecord, + status_code=status.HTTP_201_CREATED, + tags=["agents"], + ) + async def register_agent(request: AgentRegistrationRequest) -> AgentRecord: + try: + return control_plane.register_agent(request) + except AgentAlreadyExistsError as error: + _raise_http_error(status.HTTP_409_CONFLICT, "agent_already_exists", error) + + @application.get("/v1/agents", response_model=list[AgentRecord], tags=["agents"]) + async def list_agents() -> tuple[AgentRecord, ...]: + return control_plane.list_agents() + + @application.get("/v1/agents/{agent_id}", response_model=AgentRecord, tags=["agents"]) + async def get_agent(agent_id: str) -> AgentRecord: + try: + return control_plane.get_agent(agent_id) + except AgentNotFoundError as error: + _raise_http_error(status.HTTP_404_NOT_FOUND, "agent_not_found", error) + + @application.patch( + "/v1/agents/{agent_id}/status", + response_model=AgentRecord, + tags=["agents"], + ) + async def update_agent_status(agent_id: str, update: AgentStatusUpdate) -> AgentRecord: + try: + return control_plane.update_agent_status(agent_id, update) + except AgentNotFoundError as error: + _raise_http_error(status.HTTP_404_NOT_FOUND, "agent_not_found", error) + except RevisionConflictError as error: + _raise_http_error(status.HTTP_409_CONFLICT, "revision_conflict", error) + except InvalidStatusTransitionError as error: + _raise_http_error(status.HTTP_409_CONFLICT, "invalid_status_transition", error) + + @application.post( + "/v1/approvals", + response_model=ApprovalRecord, + status_code=status.HTTP_201_CREATED, + tags=["approvals"], + ) + async def create_approval(request: ApprovalRequestCreate) -> ApprovalRecord: + try: + return control_plane.create_approval(request) + except AgentNotFoundError as error: + _raise_http_error(status.HTTP_404_NOT_FOUND, "agent_not_found", error) + except AgentNotActiveError as error: + _raise_http_error(status.HTTP_409_CONFLICT, "agent_not_active", error) + + @application.get( + "/v1/approvals/{request_id}", response_model=ApprovalRecord, tags=["approvals"] + ) + async def get_approval(request_id: UUID) -> ApprovalRecord: + try: + return control_plane.get_approval(request_id) + except ApprovalNotFoundError as error: + _raise_http_error(status.HTTP_404_NOT_FOUND, "approval_not_found", error) + + @application.get("/v1/approvals", response_model=ApprovalQueueResponse, tags=["approvals"]) + async def list_approvals( + approval_status: Annotated[ApprovalStatus | None, Query(alias="status")] = None, + agent_id: str | None = None, + ) -> ApprovalQueueResponse: + items = control_plane.list_approvals(status=approval_status, agent_id=agent_id) + return ApprovalQueueResponse(items=items, count=len(items)) + + @application.post( + "/v1/approvals/{request_id}/decision", + response_model=ApprovalRecord, + tags=["approvals"], + ) + async def decide_approval( + request_id: UUID, decision: ApprovalDecisionRequest + ) -> ApprovalRecord: + try: + return control_plane.decide_approval(request_id, decision) + except ApprovalNotFoundError as error: + _raise_http_error(status.HTTP_404_NOT_FOUND, "approval_not_found", error) + except ApprovalAlreadyDecidedError as error: + _raise_http_error(status.HTTP_409_CONFLICT, "approval_already_decided", error) + + @application.get("/v1/audit-events", response_model=AuditEventPage, tags=["audit"]) + async def list_audit_events( + agent_id: str | None = None, + limit: Annotated[int, Query(ge=1, le=500)] = 100, + ) -> AuditEventPage: + items = control_plane.list_audit_events(agent_id=agent_id, limit=limit) + return AuditEventPage(items=items, count=len(items)) + return application diff --git a/src/agent_control_plane/models.py b/src/agent_control_plane/models.py index f56db3c..5fece61 100644 --- a/src/agent_control_plane/models.py +++ b/src/agent_control_plane/models.py @@ -1,9 +1,16 @@ -"""Versioned public contracts for agents and service health.""" +"""Versioned public contracts for agents, approvals, audit, and service health.""" +from datetime import datetime from enum import StrEnum +from typing import Annotated +from uuid import UUID from pydantic import BaseModel, ConfigDict, Field, field_validator +AGENT_ID_PATTERN = r"^[a-z][a-z0-9-]{2,62}$" +ACTION_PATTERN = r"^[a-z][a-z0-9._-]{2,127}$" +ActorId = Annotated[str, Field(min_length=1, max_length=128, pattern=r"^[^\s]+$")] + class HealthStatus(StrEnum): OK = "ok" @@ -23,13 +30,38 @@ class RiskLevel(StrEnum): HIGH = "high" +class AgentRuntimeStatus(StrEnum): + REGISTERED = "registered" + ACTIVE = "active" + PAUSED = "paused" + + +class ApprovalStatus(StrEnum): + PENDING = "pending" + APPROVED = "approved" + REJECTED = "rejected" + + +class ApprovalDecision(StrEnum): + APPROVE = "approve" + REJECT = "reject" + + +class AuditEventType(StrEnum): + AGENT_REGISTERED = "agent.registered" + AGENT_STATUS_CHANGED = "agent.status_changed" + APPROVAL_REQUESTED = "approval.requested" + APPROVAL_APPROVED = "approval.approved" + APPROVAL_REJECTED = "approval.rejected" + + class AgentSpec(BaseModel): """Minimal, framework-neutral contract used to register an agent.""" model_config = ConfigDict(extra="forbid", frozen=True) schema_version: str = Field(default="v1", pattern=r"^v[1-9][0-9]*$") - agent_id: str = Field(pattern=r"^[a-z][a-z0-9-]{2,62}$") + agent_id: str = Field(pattern=AGENT_ID_PATTERN) version: str = Field(pattern=r"^[0-9]+\.[0-9]+\.[0-9]+$") display_name: str = Field(min_length=1, max_length=100) description: str = Field(min_length=1, max_length=500) @@ -54,3 +86,89 @@ class AgentSpecValidationResponse(BaseModel): valid: bool agent_id: str schema_version: str + + +class AgentRegistrationRequest(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + spec: AgentSpec + actor: ActorId + + +class AgentRecord(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + spec: AgentSpec + status: AgentRuntimeStatus + revision: int = Field(ge=1) + registered_at: datetime + updated_at: datetime + + +class AgentStatusUpdate(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + status: AgentRuntimeStatus + expected_revision: int = Field(ge=1) + actor: ActorId + reason: str = Field(min_length=1, max_length=500) + + +class ApprovalRequestCreate(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + agent_id: str = Field(pattern=AGENT_ID_PATTERN) + action: str = Field(pattern=ACTION_PATTERN) + risk: RiskLevel + actor: ActorId + reason: str = Field(min_length=1, max_length=500) + + +class ApprovalRecord(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + request_id: UUID + agent_id: str + action: str + risk: RiskLevel + status: ApprovalStatus + requested_by: str + request_reason: str + created_at: datetime + decided_at: datetime | None = None + decided_by: str | None = None + decision_reason: str | None = None + + +class ApprovalDecisionRequest(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + decision: ApprovalDecision + actor: ActorId + reason: str = Field(min_length=1, max_length=500) + + +class ApprovalQueueResponse(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + items: tuple[ApprovalRecord, ...] + count: int = Field(ge=0) + + +class AuditEvent(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + event_id: UUID + event_type: AuditEventType + agent_id: str + actor: str + occurred_at: datetime + resource_id: str + summary: str + + +class AuditEventPage(BaseModel): + model_config = ConfigDict(extra="forbid", frozen=True) + + items: tuple[AuditEvent, ...] + count: int = Field(ge=0) diff --git a/src/agent_control_plane/store.py b/src/agent_control_plane/store.py new file mode 100644 index 0000000..49f623d --- /dev/null +++ b/src/agent_control_plane/store.py @@ -0,0 +1,299 @@ +"""Storage boundary and the development in-memory control-plane adapter.""" + +from collections.abc import Callable +from datetime import UTC, datetime +from threading import RLock +from typing import Protocol +from uuid import UUID, uuid4 + +from agent_control_plane.models import ( + AgentRecord, + AgentRegistrationRequest, + AgentRuntimeStatus, + AgentStatusUpdate, + ApprovalDecision, + ApprovalDecisionRequest, + ApprovalRecord, + ApprovalRequestCreate, + ApprovalStatus, + AuditEvent, + AuditEventType, +) + + +class StoreError(Exception): + """Base error for expected control-plane storage failures.""" + + +class AgentAlreadyExistsError(StoreError): + pass + + +class AgentNotFoundError(StoreError): + pass + + +class AgentNotActiveError(StoreError): + pass + + +class RevisionConflictError(StoreError): + pass + + +class InvalidStatusTransitionError(StoreError): + pass + + +class ApprovalNotFoundError(StoreError): + pass + + +class ApprovalAlreadyDecidedError(StoreError): + pass + + +class ControlPlaneStore(Protocol): + def register_agent(self, request: AgentRegistrationRequest) -> AgentRecord: ... + + def get_agent(self, agent_id: str) -> AgentRecord: ... + + def list_agents(self) -> tuple[AgentRecord, ...]: ... + + def update_agent_status(self, agent_id: str, update: AgentStatusUpdate) -> AgentRecord: ... + + def create_approval(self, request: ApprovalRequestCreate) -> ApprovalRecord: ... + + def get_approval(self, request_id: UUID) -> ApprovalRecord: ... + + def list_approvals( + self, status: ApprovalStatus | None = None, agent_id: str | None = None + ) -> tuple[ApprovalRecord, ...]: ... + + def decide_approval( + self, request_id: UUID, decision: ApprovalDecisionRequest + ) -> ApprovalRecord: ... + + def list_audit_events( + self, agent_id: str | None = None, limit: int = 100 + ) -> tuple[AuditEvent, ...]: ... + + +def utc_now() -> datetime: + return datetime.now(UTC) + + +class InMemoryControlPlaneStore: + """Concurrency-safe adapter for development and single-process evaluation.""" + + _allowed_status_transitions = { + AgentRuntimeStatus.REGISTERED: {AgentRuntimeStatus.ACTIVE, AgentRuntimeStatus.PAUSED}, + AgentRuntimeStatus.ACTIVE: {AgentRuntimeStatus.PAUSED}, + AgentRuntimeStatus.PAUSED: {AgentRuntimeStatus.ACTIVE}, + } + + def __init__( + self, + *, + clock: Callable[[], datetime] = utc_now, + id_factory: Callable[[], UUID] = uuid4, + ) -> None: + self._clock = clock + self._id_factory = id_factory + self._agents: dict[str, AgentRecord] = {} + self._approvals: dict[UUID, ApprovalRecord] = {} + self._audit_events: list[AuditEvent] = [] + self._lock = RLock() + + def register_agent(self, request: AgentRegistrationRequest) -> AgentRecord: + with self._lock: + agent_id = request.spec.agent_id + if agent_id in self._agents: + raise AgentAlreadyExistsError(f"agent '{agent_id}' is already registered") + + timestamp = self._clock() + record = AgentRecord( + spec=request.spec, + status=AgentRuntimeStatus.REGISTERED, + revision=1, + registered_at=timestamp, + updated_at=timestamp, + ) + self._agents[agent_id] = record + self._append_event( + event_type=AuditEventType.AGENT_REGISTERED, + agent_id=agent_id, + actor=request.actor, + occurred_at=timestamp, + resource_id=agent_id, + summary=f"Registered agent version {request.spec.version}", + ) + return record + + def get_agent(self, agent_id: str) -> AgentRecord: + with self._lock: + try: + return self._agents[agent_id] + except KeyError as error: + raise AgentNotFoundError(f"agent '{agent_id}' was not found") from error + + def list_agents(self) -> tuple[AgentRecord, ...]: + with self._lock: + return tuple(self._agents[agent_id] for agent_id in sorted(self._agents)) + + def update_agent_status(self, agent_id: str, update: AgentStatusUpdate) -> AgentRecord: + with self._lock: + current = self.get_agent(agent_id) + if current.revision != update.expected_revision: + raise RevisionConflictError( + f"expected revision {update.expected_revision}, current revision is " + f"{current.revision}" + ) + if update.status not in self._allowed_status_transitions[current.status]: + raise InvalidStatusTransitionError( + f"cannot change agent status from '{current.status}' to '{update.status}'" + ) + + timestamp = self._clock() + updated = current.model_copy( + update={ + "status": update.status, + "revision": current.revision + 1, + "updated_at": timestamp, + } + ) + self._agents[agent_id] = updated + self._append_event( + event_type=AuditEventType.AGENT_STATUS_CHANGED, + agent_id=agent_id, + actor=update.actor, + occurred_at=timestamp, + resource_id=agent_id, + summary=f"Changed status from {current.status} to {update.status}: {update.reason}", + ) + return updated + + def create_approval(self, request: ApprovalRequestCreate) -> ApprovalRecord: + with self._lock: + agent = self.get_agent(request.agent_id) + if agent.status is not AgentRuntimeStatus.ACTIVE: + raise AgentNotActiveError( + f"agent '{request.agent_id}' must be active to request approval" + ) + timestamp = self._clock() + request_id = self._id_factory() + record = ApprovalRecord( + request_id=request_id, + agent_id=request.agent_id, + action=request.action, + risk=request.risk, + status=ApprovalStatus.PENDING, + requested_by=request.actor, + request_reason=request.reason, + created_at=timestamp, + ) + self._approvals[request_id] = record + self._append_event( + event_type=AuditEventType.APPROVAL_REQUESTED, + agent_id=request.agent_id, + actor=request.actor, + occurred_at=timestamp, + resource_id=str(request_id), + summary=f"Requested {request.risk} approval for {request.action}", + ) + return record + + def get_approval(self, request_id: UUID) -> ApprovalRecord: + with self._lock: + try: + return self._approvals[request_id] + except KeyError as error: + raise ApprovalNotFoundError( + f"approval request '{request_id}' was not found" + ) from error + + def list_approvals( + self, status: ApprovalStatus | None = None, agent_id: str | None = None + ) -> tuple[ApprovalRecord, ...]: + with self._lock: + records = ( + record + for record in self._approvals.values() + if (status is None or record.status == status) + and (agent_id is None or record.agent_id == agent_id) + ) + return tuple( + sorted(records, key=lambda record: (record.created_at, str(record.request_id))) + ) + + def decide_approval( + self, request_id: UUID, decision: ApprovalDecisionRequest + ) -> ApprovalRecord: + with self._lock: + current = self.get_approval(request_id) + if current.status is not ApprovalStatus.PENDING: + raise ApprovalAlreadyDecidedError( + f"approval request '{request_id}' has already been decided" + ) + + timestamp = self._clock() + status = ( + ApprovalStatus.APPROVED + if decision.decision is ApprovalDecision.APPROVE + else ApprovalStatus.REJECTED + ) + updated = current.model_copy( + update={ + "status": status, + "decided_at": timestamp, + "decided_by": decision.actor, + "decision_reason": decision.reason, + } + ) + self._approvals[request_id] = updated + self._append_event( + event_type=( + AuditEventType.APPROVAL_APPROVED + if status is ApprovalStatus.APPROVED + else AuditEventType.APPROVAL_REJECTED + ), + agent_id=current.agent_id, + actor=decision.actor, + occurred_at=timestamp, + resource_id=str(request_id), + summary=f"{status.value.capitalize()} {current.action}: {decision.reason}", + ) + return updated + + def list_audit_events( + self, agent_id: str | None = None, limit: int = 100 + ) -> tuple[AuditEvent, ...]: + with self._lock: + matching = ( + event + for event in reversed(self._audit_events) + if agent_id is None or event.agent_id == agent_id + ) + return tuple(event for _, event in zip(range(limit), matching, strict=False)) + + def _append_event( + self, + *, + event_type: AuditEventType, + agent_id: str, + actor: str, + occurred_at: datetime, + resource_id: str, + summary: str, + ) -> None: + self._audit_events.append( + AuditEvent( + event_id=self._id_factory(), + event_type=event_type, + agent_id=agent_id, + actor=actor, + occurred_at=occurred_at, + resource_id=resource_id, + summary=summary, + ) + ) diff --git a/tests/smoke/test_api_smoke.py b/tests/smoke/test_api_smoke.py index 72cc826..95f15b2 100644 --- a/tests/smoke/test_api_smoke.py +++ b/tests/smoke/test_api_smoke.py @@ -3,7 +3,7 @@ import httpx import pytest -from agent_control_plane.api import app +from agent_control_plane.api import create_app @pytest.fixture @@ -13,7 +13,7 @@ def anyio_backend() -> str: @pytest.fixture async def client() -> AsyncIterator[httpx.AsyncClient]: - transport = httpx.ASGITransport(app=app) + transport = httpx.ASGITransport(app=create_app()) async with httpx.AsyncClient(transport=transport, base_url="http://test") as test_client: yield test_client @@ -62,3 +62,199 @@ async def test_invalid_contract_fails_closed(client: httpx.AsyncClient) -> None: ) assert response.status_code == 422 + + +def registration_payload() -> dict[str, object]: + return { + "spec": { + "agent_id": "support-agent", + "version": "1.0.0", + "display_name": "Support Agent", + "description": "Handles support requests with governed actions.", + "entrypoint": "https://agents.example.test/support", + "capabilities": ["ticket.read", "ticket.refund"], + }, + "actor": "operator@example.test", + } + + +@pytest.mark.smoke +@pytest.mark.anyio +async def test_registration_approval_and_audit_loop(client: httpx.AsyncClient) -> None: + registration = await client.post("/v1/agents", json=registration_payload()) + assert registration.status_code == 201 + assert registration.json()["status"] == "registered" + assert registration.json()["revision"] == 1 + assert (await client.get("/v1/agents/support-agent")).status_code == 200 + agents = await client.get("/v1/agents") + assert [item["spec"]["agent_id"] for item in agents.json()] == ["support-agent"] + + activation = await client.patch( + "/v1/agents/support-agent/status", + json={ + "status": "active", + "expected_revision": 1, + "actor": "operator@example.test", + "reason": "Readiness checks passed.", + }, + ) + assert activation.status_code == 200 + assert activation.json()["revision"] == 2 + + approval = await client.post( + "/v1/approvals", + json={ + "agent_id": "support-agent", + "action": "ticket.refund", + "risk": "high", + "actor": "support-agent", + "reason": "Refund exceeds the automatic threshold.", + }, + ) + assert approval.status_code == 201 + request_id = approval.json()["request_id"] + assert (await client.get(f"/v1/approvals/{request_id}")).status_code == 200 + + queue = await client.get("/v1/approvals", params={"status": "pending"}) + assert queue.json()["count"] == 1 + + decision = await client.post( + f"/v1/approvals/{request_id}/decision", + json={ + "decision": "approve", + "actor": "reviewer@example.test", + "reason": "Evidence verified.", + }, + ) + assert decision.status_code == 200 + assert decision.json()["status"] == "approved" + + audit = await client.get("/v1/audit-events", params={"agent_id": "support-agent", "limit": 4}) + assert audit.status_code == 200 + assert [item["event_type"] for item in audit.json()["items"]] == [ + "approval.approved", + "approval.requested", + "agent.status_changed", + "agent.registered", + ] + + +@pytest.mark.smoke +@pytest.mark.anyio +async def test_governance_conflicts_fail_closed(client: httpx.AsyncClient) -> None: + assert (await client.post("/v1/agents", json=registration_payload())).status_code == 201 + + duplicate = await client.post("/v1/agents", json=registration_payload()) + assert duplicate.status_code == 409 + assert duplicate.json()["detail"]["code"] == "agent_already_exists" + + missing_agent = await client.get("/v1/agents/missing-agent") + assert missing_agent.status_code == 404 + assert missing_agent.json()["detail"]["code"] == "agent_not_found" + + stale = await client.patch( + "/v1/agents/support-agent/status", + json={ + "status": "active", + "expected_revision": 2, + "actor": "operator@example.test", + "reason": "Stale client state.", + }, + ) + assert stale.status_code == 409 + assert stale.json()["detail"]["code"] == "revision_conflict" + + inactive_approval = await client.post( + "/v1/approvals", + json={ + "agent_id": "support-agent", + "action": "ticket.refund", + "risk": "high", + "actor": "support-agent", + "reason": "The agent has not been activated.", + }, + ) + assert inactive_approval.status_code == 409 + assert inactive_approval.json()["detail"]["code"] == "agent_not_active" + + missing_status = await client.patch( + "/v1/agents/missing-agent/status", + json={ + "status": "active", + "expected_revision": 1, + "actor": "operator@example.test", + "reason": "Unknown agent.", + }, + ) + assert missing_status.status_code == 404 + assert missing_status.json()["detail"]["code"] == "agent_not_found" + + activated = await client.patch( + "/v1/agents/support-agent/status", + json={ + "status": "active", + "expected_revision": 1, + "actor": "operator@example.test", + "reason": "Readiness checks passed.", + }, + ) + assert activated.status_code == 200 + + invalid_transition = await client.patch( + "/v1/agents/support-agent/status", + json={ + "status": "active", + "expected_revision": 2, + "actor": "operator@example.test", + "reason": "No-op transition.", + }, + ) + assert invalid_transition.status_code == 409 + assert invalid_transition.json()["detail"]["code"] == "invalid_status_transition" + + unknown_agent = await client.post( + "/v1/approvals", + json={ + "agent_id": "missing-agent", + "action": "ticket.refund", + "risk": "high", + "actor": "missing-agent", + "reason": "This agent is not registered.", + }, + ) + assert unknown_agent.status_code == 404 + assert unknown_agent.json()["detail"]["code"] == "agent_not_found" + + missing_request_id = "00000000-0000-0000-0000-000000000000" + missing_approval = await client.get(f"/v1/approvals/{missing_request_id}") + assert missing_approval.status_code == 404 + assert missing_approval.json()["detail"]["code"] == "approval_not_found" + + approval = await client.post( + "/v1/approvals", + json={ + "agent_id": "support-agent", + "action": "ticket.refund", + "risk": "high", + "actor": "support-agent", + "reason": "Refund exceeds the automatic threshold.", + }, + ) + request_id = approval.json()["request_id"] + decision_payload = { + "decision": "reject", + "actor": "reviewer@example.test", + "reason": "Evidence is missing.", + } + assert ( + await client.post(f"/v1/approvals/{request_id}/decision", json=decision_payload) + ).status_code == 200 + repeated = await client.post(f"/v1/approvals/{request_id}/decision", json=decision_payload) + assert repeated.status_code == 409 + assert repeated.json()["detail"]["code"] == "approval_already_decided" + + missing_decision = await client.post( + f"/v1/approvals/{missing_request_id}/decision", json=decision_payload + ) + assert missing_decision.status_code == 404 + assert missing_decision.json()["detail"]["code"] == "approval_not_found" diff --git a/tests/unit/test_store.py b/tests/unit/test_store.py new file mode 100644 index 0000000..1ab2e66 --- /dev/null +++ b/tests/unit/test_store.py @@ -0,0 +1,222 @@ +from datetime import UTC, datetime +from uuid import UUID + +import pytest + +from agent_control_plane.models import ( + AgentRegistrationRequest, + AgentRuntimeStatus, + AgentSpec, + AgentStatusUpdate, + ApprovalDecision, + ApprovalDecisionRequest, + ApprovalRequestCreate, + ApprovalStatus, + AuditEventType, + RiskLevel, +) +from agent_control_plane.store import ( + AgentAlreadyExistsError, + AgentNotActiveError, + AgentNotFoundError, + ApprovalAlreadyDecidedError, + ApprovalNotFoundError, + InMemoryControlPlaneStore, + InvalidStatusTransitionError, + RevisionConflictError, +) + +NOW = datetime(2026, 8, 11, 12, 0, tzinfo=UTC) + + +def build_store() -> InMemoryControlPlaneStore: + return InMemoryControlPlaneStore(clock=lambda: NOW) + + +def registration(agent_id: str = "support-agent") -> AgentRegistrationRequest: + return AgentRegistrationRequest( + spec=AgentSpec( + agent_id=agent_id, + version="1.0.0", + display_name="Support Agent", + description="Handles support requests with governed actions.", + entrypoint="https://agents.example.test/support", + capabilities=("ticket.read", "ticket.reply"), + ), + actor="operator@example.test", + ) + + +def register(store: InMemoryControlPlaneStore, agent_id: str = "support-agent") -> None: + store.register_agent(registration(agent_id)) + + +def activate(store: InMemoryControlPlaneStore, agent_id: str = "support-agent") -> None: + store.update_agent_status( + agent_id, + AgentStatusUpdate( + status=AgentRuntimeStatus.ACTIVE, + expected_revision=1, + actor="operator@example.test", + reason="Readiness checks passed.", + ), + ) + + +def approval_request(agent_id: str = "support-agent") -> ApprovalRequestCreate: + return ApprovalRequestCreate( + agent_id=agent_id, + action="ticket.refund", + risk=RiskLevel.HIGH, + actor="support-agent", + reason="Refund exceeds the automatic approval threshold.", + ) + + +def test_register_agent_creates_initial_state_and_audit_event() -> None: + store = build_store() + + record = store.register_agent(registration()) + + assert record.status is AgentRuntimeStatus.REGISTERED + assert record.revision == 1 + assert record.registered_at == NOW + assert store.get_agent("support-agent") == record + event = store.list_audit_events()[0] + assert event.event_type is AuditEventType.AGENT_REGISTERED + assert event.actor == "operator@example.test" + + +def test_duplicate_agent_registration_is_rejected() -> None: + store = build_store() + register(store) + + with pytest.raises(AgentAlreadyExistsError): + register(store) + + +def test_agents_are_listed_in_stable_id_order() -> None: + store = build_store() + register(store, "zeta-agent") + register(store, "alpha-agent") + + records = store.list_agents() + + assert [record.spec.agent_id for record in records] == ["alpha-agent", "zeta-agent"] + + +def test_status_update_uses_optimistic_revision_and_creates_audit_event() -> None: + store = build_store() + register(store) + + updated = store.update_agent_status( + "support-agent", + AgentStatusUpdate( + status=AgentRuntimeStatus.ACTIVE, + expected_revision=1, + actor="operator@example.test", + reason="Production readiness checks passed.", + ), + ) + + assert updated.status is AgentRuntimeStatus.ACTIVE + assert updated.revision == 2 + assert store.list_audit_events()[0].event_type is AuditEventType.AGENT_STATUS_CHANGED + + +def test_stale_revision_and_invalid_transition_are_rejected() -> None: + store = build_store() + register(store) + + with pytest.raises(RevisionConflictError): + store.update_agent_status( + "support-agent", + AgentStatusUpdate( + status=AgentRuntimeStatus.ACTIVE, + expected_revision=2, + actor="operator@example.test", + reason="Uses a stale copy.", + ), + ) + + with pytest.raises(InvalidStatusTransitionError): + store.update_agent_status( + "support-agent", + AgentStatusUpdate( + status=AgentRuntimeStatus.REGISTERED, + expected_revision=1, + actor="operator@example.test", + reason="No-op transitions are not allowed.", + ), + ) + + +def test_approval_lifecycle_is_single_decision_and_audited() -> None: + store = build_store() + register(store) + activate(store) + pending = store.create_approval(approval_request()) + + approved = store.decide_approval( + pending.request_id, + ApprovalDecisionRequest( + decision=ApprovalDecision.APPROVE, + actor="reviewer@example.test", + reason="Customer identity and refund evidence verified.", + ), + ) + + assert approved.status is ApprovalStatus.APPROVED + assert approved.decided_by == "reviewer@example.test" + assert store.list_approvals(status=ApprovalStatus.PENDING) == () + assert [event.event_type for event in store.list_audit_events()] == [ + AuditEventType.APPROVAL_APPROVED, + AuditEventType.APPROVAL_REQUESTED, + AuditEventType.AGENT_STATUS_CHANGED, + AuditEventType.AGENT_REGISTERED, + ] + with pytest.raises(ApprovalAlreadyDecidedError): + store.decide_approval( + pending.request_id, + ApprovalDecisionRequest( + decision=ApprovalDecision.REJECT, + actor="second-reviewer@example.test", + reason="A second decision must not overwrite the first.", + ), + ) + + +def test_approval_requires_a_registered_agent() -> None: + store = build_store() + + with pytest.raises(AgentNotFoundError): + store.create_approval(approval_request("missing-agent")) + + +def test_approval_requires_an_active_agent() -> None: + store = build_store() + register(store) + + with pytest.raises(AgentNotActiveError): + store.create_approval(approval_request()) + + +def test_rejected_approval_is_audited_and_unknown_request_is_rejected() -> None: + store = build_store() + register(store) + activate(store) + pending = store.create_approval(approval_request()) + + rejected = store.decide_approval( + pending.request_id, + ApprovalDecisionRequest( + decision=ApprovalDecision.REJECT, + actor="reviewer@example.test", + reason="Required evidence is missing.", + ), + ) + + assert rejected.status is ApprovalStatus.REJECTED + assert store.list_audit_events(limit=1)[0].event_type is AuditEventType.APPROVAL_REJECTED + with pytest.raises(ApprovalNotFoundError): + store.get_approval(UUID(int=0))