Skip to content

Commit bc701a3

Browse files
committed
refactor: isolate operator HTTP routes
1 parent 1588f16 commit bc701a3

2 files changed

Lines changed: 303 additions & 209 deletions

File tree

src/agentnet/operations/http.py

Lines changed: 292 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,292 @@
1+
"""Authenticated operator, authority, incident, and version HTTP routes."""
2+
3+
from __future__ import annotations
4+
5+
from collections.abc import Awaitable, Callable, Mapping
6+
7+
from pydantic import BaseModel, ConfigDict, Field
8+
from starlette.requests import Request
9+
from starlette.responses import JSONResponse, Response
10+
from starlette.routing import Route
11+
12+
from agentnet.authorization.evidence import IssuanceAuthority, SignedAuthorityCommand
13+
from agentnet.core.app import CommunicationCore
14+
from agentnet.errors import AuthenticationError, AuthorizationError, ValidationError
15+
from agentnet.identity.actors import TrustedTransportContext, VerifiedActor
16+
from agentnet.operations.authority_inspection import DenialExplanationQuery
17+
from agentnet.operations.incident import IncidentModeChange
18+
19+
20+
BodyAndActor = Callable[
21+
[Request, CommunicationCore],
22+
Awaitable[tuple[bytes, VerifiedActor]],
23+
]
24+
AuthorityIssuer = Callable[..., IssuanceAuthority]
25+
26+
27+
class VersionRolloutBeginBody(BaseModel):
28+
model_config = ConfigDict(extra="forbid")
29+
30+
from_protocol_version: str = Field(
31+
pattern=r"^(0|[1-9][0-9]*)\.(0|[1-9][0-9]*)$"
32+
)
33+
to_protocol_version: str = Field(
34+
pattern=r"^(0|[1-9][0-9]*)\.(0|[1-9][0-9]*)$"
35+
)
36+
from_schema_version: int = Field(ge=1)
37+
to_schema_version: int = Field(ge=1)
38+
compatibility_deadline: int = Field(ge=1)
39+
40+
41+
class VersionRolloutAdvanceBody(BaseModel):
42+
model_config = ConfigDict(extra="forbid")
43+
44+
expected_phase: str = Field(
45+
pattern=r"^(expanded|migrated_backfilled|verified)$"
46+
)
47+
target_phase: str = Field(
48+
pattern=r"^(migrated_backfilled|verified|contracted)$"
49+
)
50+
verification_digest: str | None = Field(
51+
default=None,
52+
pattern=r"^[a-f0-9]{64}$",
53+
)
54+
55+
56+
class VersionRolloutRollbackBody(BaseModel):
57+
model_config = ConfigDict(extra="forbid")
58+
59+
verification_digest: str = Field(pattern=r"^[a-f0-9]{64}$")
60+
61+
62+
class VersionReplayBody(BaseModel):
63+
model_config = ConfigDict(extra="forbid")
64+
65+
peer_namespace: str = Field(pattern=r"^[a-z][a-z0-9_.-]{0,127}$")
66+
limit: int = Field(default=100, ge=1, le=1000)
67+
68+
69+
class IncidentModeChangeBody(BaseModel):
70+
model_config = ConfigDict(extra="forbid", strict=True)
71+
72+
change: IncidentModeChange
73+
command: SignedAuthorityCommand
74+
75+
76+
def _trusted_transport(request: Request) -> TrustedTransportContext:
77+
"""Return only transport state installed by proof authentication."""
78+
79+
transport = request.scope.get("agentnet.trusted_transport")
80+
if not isinstance(transport, TrustedTransportContext):
81+
raise AuthenticationError("verified transport context is unavailable")
82+
return transport
83+
84+
85+
def create_operator_routes(
86+
core: CommunicationCore,
87+
body_and_actor: BodyAndActor,
88+
issue_authority: AuthorityIssuer,
89+
response_headers: Mapping[str, str],
90+
) -> list[Route]:
91+
"""Mount only content-free operator and authority routes."""
92+
93+
async def operator_status(request: Request) -> Response:
94+
_body, actor = await body_and_actor(request, core)
95+
core._require(
96+
actor=actor,
97+
action="operator.status.read",
98+
resource="operator:self",
99+
)
100+
readiness = core.readiness()
101+
try:
102+
telemetry = {"available": True, **core.telemetry.operational_snapshot()}
103+
except Exception:
104+
telemetry = {
105+
"available": False,
106+
"counters": {},
107+
"latency_buckets": {},
108+
"gauges": {},
109+
}
110+
try:
111+
admission_controls = {
112+
"available": True,
113+
**core.quotas.content_free_status(),
114+
}
115+
except Exception:
116+
admission_controls = {"available": False}
117+
try:
118+
versioning = {
119+
"available": True,
120+
**core.versioning.content_free_status(),
121+
}
122+
except Exception:
123+
versioning = {"available": False}
124+
return JSONResponse(
125+
{
126+
"status": "ready" if readiness["ready"] else "degraded",
127+
"profile": readiness["profile"],
128+
"acceptance_fact": readiness["acceptance_fact"],
129+
"storage_ready": bool(readiness["storage"].get("ready")),
130+
"artifacts_ready": bool(readiness["artifacts"].get("ready")),
131+
"audit_valid": bool(readiness["audit"].get("valid")),
132+
"deployment_binding_ready": bool(
133+
readiness["deployment_binding"].get("ready")
134+
),
135+
"a2a_ready": bool(readiness["a2a_schema"].get("ready")),
136+
"scanner_trust_ready": bool(
137+
readiness["scanner_trust"].get("ready")
138+
),
139+
"telemetry": telemetry,
140+
"admission_controls": admission_controls,
141+
"versioning": versioning,
142+
}
143+
)
144+
145+
async def authority_inventory(request: Request) -> Response:
146+
if request.scope.get("query_string", b""):
147+
raise ValidationError(
148+
"authority inventory does not accept caller-selected scope"
149+
)
150+
await body_and_actor(request, core)
151+
inventory = core.authority_inspection.authority_inventory(
152+
transport=_trusted_transport(request),
153+
)
154+
return JSONResponse(
155+
{"authority": inventory.model_dump(mode="json")},
156+
headers=response_headers,
157+
)
158+
159+
async def explain_denial(request: Request) -> Response:
160+
if request.scope.get("query_string", b""):
161+
raise ValidationError(
162+
"denial explanation does not accept caller-selected scope"
163+
)
164+
await body_and_actor(request, core)
165+
query = DenialExplanationQuery(
166+
decision_id=request.path_params["decision_id"]
167+
)
168+
explanation = core.authority_inspection.explain_denial(
169+
transport=_trusted_transport(request),
170+
query=query,
171+
)
172+
return JSONResponse(
173+
{"explanation": explanation.model_dump(mode="json")},
174+
headers=response_headers,
175+
)
176+
177+
async def incident_status(request: Request) -> Response:
178+
_body, actor = await body_and_actor(request, core)
179+
resource = f"operator-domain:{core.config.domain_id}"
180+
core._require(
181+
actor=actor,
182+
action="operator.incident.read",
183+
resource=resource,
184+
)
185+
return JSONResponse(
186+
{
187+
"incident": core.incidents.state(
188+
core.config.domain_id
189+
).model_dump(mode="json")
190+
},
191+
headers=response_headers,
192+
)
193+
194+
async def set_incident_mode(request: Request) -> Response:
195+
body, actor = await body_and_actor(request, core)
196+
parsed = IncidentModeChangeBody.model_validate_json(body, strict=True)
197+
if parsed.change.domain_id != core.config.domain_id:
198+
raise AuthorizationError(
199+
"incident change crossed the authenticated domain"
200+
)
201+
resource, _exact_request = core.incidents.authority_binding(parsed.change)
202+
authority = issue_authority(
203+
core,
204+
actor=actor,
205+
action=core.incidents.ACTION,
206+
resource=resource,
207+
request={"request_digest": parsed.command.request_digest},
208+
)
209+
state = core.incidents.set_mode(
210+
parsed.change,
211+
authority=authority,
212+
command=parsed.command,
213+
)
214+
return JSONResponse(
215+
{"incident": state.model_dump(mode="json")},
216+
headers=response_headers,
217+
)
218+
219+
async def begin_version_rollout(request: Request) -> Response:
220+
body, actor = await body_and_actor(request, core)
221+
parsed = VersionRolloutBeginBody.model_validate_json(body)
222+
return JSONResponse(
223+
core.begin_version_rollout(actor=actor, **parsed.model_dump()),
224+
status_code=201,
225+
)
226+
227+
async def advance_version_rollout(request: Request) -> Response:
228+
body, actor = await body_and_actor(request, core)
229+
parsed = VersionRolloutAdvanceBody.model_validate_json(body)
230+
return JSONResponse(
231+
core.advance_version_rollout(
232+
actor=actor,
233+
rollout_id=request.path_params["rollout_id"],
234+
**parsed.model_dump(),
235+
)
236+
)
237+
238+
async def rollback_version_rollout(request: Request) -> Response:
239+
body, actor = await body_and_actor(request, core)
240+
parsed = VersionRolloutRollbackBody.model_validate_json(body)
241+
return JSONResponse(
242+
core.rollback_version_rollout(
243+
actor=actor,
244+
rollout_id=request.path_params["rollout_id"],
245+
verification_digest=parsed.verification_digest,
246+
)
247+
)
248+
249+
async def replay_version_events(request: Request) -> Response:
250+
body, actor = await body_and_actor(request, core)
251+
parsed = VersionReplayBody.model_validate_json(body)
252+
return JSONResponse(
253+
core.replay_unsupported_events(
254+
actor=actor,
255+
**parsed.model_dump(),
256+
)
257+
)
258+
259+
return [
260+
Route("/v1/operator/status", operator_status, methods=["GET"]),
261+
Route("/v1/authority", authority_inventory, methods=["GET"]),
262+
Route(
263+
"/v1/authority/denials/{decision_id}",
264+
explain_denial,
265+
methods=["GET"],
266+
),
267+
Route("/v1/operator/incident", incident_status, methods=["GET"]),
268+
Route("/v1/operator/incident", set_incident_mode, methods=["POST"]),
269+
Route(
270+
"/v1/operator/version-rollouts",
271+
begin_version_rollout,
272+
methods=["POST"],
273+
),
274+
Route(
275+
"/v1/operator/version-rollouts/{rollout_id}/advance",
276+
advance_version_rollout,
277+
methods=["POST"],
278+
),
279+
Route(
280+
"/v1/operator/version-rollouts/{rollout_id}/rollback",
281+
rollback_version_rollout,
282+
methods=["POST"],
283+
),
284+
Route(
285+
"/v1/operator/version-replay",
286+
replay_version_events,
287+
methods=["POST"],
288+
),
289+
]
290+
291+
292+
__all__ = ["create_operator_routes"]

0 commit comments

Comments
 (0)