File size: 18,735 Bytes
2aa8b3a 3c241d8 2aa8b3a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 | """JSON-facing adapter used by the hidden, named Gradio agent endpoints."""
from __future__ import annotations
import json
import math
from collections.abc import Mapping
from typing import Any
from distinct_protocol import AgentCapabilities, AgentSnapshot, canonical_json
from distinct_protocol.handshake import RESPONSE_AUTH_KEY, HandshakeError, SignedSnapshot
from .control_plane import ControlPlane
from .models import ControlPlaneError, JobView
MAX_AGENT_BODY_BYTES = 1_000_000
PATH_SYNC = "/agent/sync"
PATH_ACCEPT = "/agent/accept"
PATH_COMPLETE = "/agent/complete"
PATH_CATALOGUE = "/agent/catalogue"
def _json_object(raw: str, *, maximum: int = MAX_AGENT_BODY_BYTES) -> Mapping[str, Any]:
if not isinstance(raw, str) or len(raw.encode("utf-8")) > maximum:
raise ValueError("request body is missing or too large")
try:
value = json.loads(raw, parse_constant=lambda item: (_ for _ in ()).throw(ValueError(item)))
except (json.JSONDecodeError, ValueError) as exc:
raise ValueError("request body is not valid JSON") from exc
if not isinstance(value, dict):
raise ValueError("request body must be a JSON object")
return value
def _job_dict(view: JobView, *, include_spec: bool = True) -> dict:
result = {
"id": view.spec.id,
"status": view.status.value,
"accepted_by": view.accepted_by,
"queue_position": view.queue_position,
"updated_at": view.updated_at,
"expires_at": view.expires_at,
"result": view.result,
"error": view.error,
}
if include_spec:
result["spec"] = view.spec.to_dict()
return result
def _response(value: Mapping[str, Any]) -> str:
return canonical_json(value)
class AgentApi:
"""Validate, authenticate and translate calls from community workers.
Every response to an authenticated call is signed on the way out. Under an
agent-pull transport the agent dials an address a human typed, so the
server-to-agent direction is the one carrying prompts and tool invocations
onto somebody's personal machine, and it is the direction that most needs
proving. The proof binds the agent's own request nonce, so a captured
response cannot answer a later request.
"""
def __init__(self, control_plane: ControlPlane) -> None:
self.control_plane = control_plane
def _signed(
self,
agent_id: str,
path: str,
request_nonce: str,
body: Mapping[str, Any],
*,
authenticated: bool = True,
) -> str:
"""Attach the response proof, when the caller earned one.
A rejection issued before the caller authenticated is **never** signed,
and `authenticated=False` is how a caller says so. That is not only
about there being no agreed key: the presence of a proof is itself a
fact about the server, and one worth nothing to a legitimate worker and
everything to somebody probing.
**This used to sign whenever it could, and that was an oracle.** A
request naming an agent that exists could be signed even when its
signature was wrong, and a request naming one that does not could not.
So an unauthenticated caller learned which agent ids are real by
looking for the `auth` field, whatever the error text said. Default
agent ids are machine hostnames, so the ids worth trying are guessable.
The uniform refusal in `ControlPlane.verify_agent_request` closes the
message half of that; this closes the signature half.
"""
payload = dict(body)
proof = None
if authenticated:
try:
proof = self.control_plane.sign_response(agent_id, path, request_nonce, payload)
except (ControlPlaneError, HandshakeError, TypeError, ValueError):
# A malformed nonce is itself one of the reasons a request gets
# rejected, so the proof cannot always be built even for a
# caller who authenticated. Unsigned is right: a verifying
# agent refuses it, which is the correct reading of an
# unprovable answer.
proof = None
if proof is not None:
payload[RESPONSE_AUTH_KEY] = proof
return _response(payload)
def pair(
self,
pairing_code: str,
capabilities_json: str,
signing_public_key_hex: str = "",
access_code: str = "",
) -> str:
"""Register an agent and pin the key it will sign advertisements with.
``signing_public_key_hex`` is the agent's Ed25519 public key, 32 bytes
as hex. It arrives once, here, and the server keeps it for the life of
the pairing. An agent that omits it can still pair and still work, but
its advertisements cannot be verified and are therefore not honoured.
"""
try:
capabilities = AgentCapabilities.from_dict(_json_object(capabilities_json, maximum=65_536))
public_key = b""
if signing_public_key_hex:
if not isinstance(signing_public_key_hex, str):
raise ValueError("signing public key must be hexadecimal text")
try:
public_key = bytes.fromhex(signing_public_key_hex)
except ValueError as exc:
raise ValueError("signing public key is not valid hexadecimal") from exc
credential = self.control_plane.pair_agent(
pairing_code,
capabilities,
signing_public_key=public_key,
# The agent prints this when it starts. Only its digest is
# kept; see ControlPlane.pair_agent.
access_code=access_code or "",
)
return _response(
{
"ok": True,
"agent_id": credential.agent_id,
"secret": credential.secret,
"issued_at": credential.issued_at,
# One cadence. ``heartbeat_interval_s`` is retained as a
# legacy alias for agents built before the rename; both
# always carry the same value.
"poll_interval_s": self.control_plane.limits.poll_interval_seconds,
"heartbeat_interval_s": self.control_plane.limits.poll_interval_seconds,
"offline_after_s": self.control_plane.limits.offline_after_seconds,
}
)
except (ControlPlaneError, TypeError, ValueError) as exc:
return _response({"ok": False, "error": str(exc)})
def register(
self,
claim_code: str,
capabilities_json: str,
signing_public_key_hex: str = "",
access_code: str = "",
) -> str:
"""Register a worker that will be claimed from the website afterwards.
The mirror image of :meth:`pair`. There, the code came from the server
and the volunteer carried it to the terminal; here the worker invents
one and the volunteer carries it to the browser. Nothing is copied out
of the page, which is what makes the failure this replaced -- somebody
inventing a value where the pairing code went -- impossible rather than
merely better explained.
This is the only entry point on the server that needs no account. What
it creates is inert: see ``ControlPlane.register_unclaimed``.
"""
try:
capabilities = AgentCapabilities.from_dict(
_json_object(capabilities_json, maximum=65_536)
)
public_key = b""
if signing_public_key_hex:
if not isinstance(signing_public_key_hex, str):
raise ValueError("signing public key must be hexadecimal text")
try:
public_key = bytes.fromhex(signing_public_key_hex)
except ValueError as exc:
raise ValueError("signing public key is not valid hexadecimal") from exc
credential = self.control_plane.register_unclaimed(
capabilities,
claim_code,
signing_public_key=public_key,
access_code=access_code or "",
)
return _response(
{
"ok": True,
"agent_id": credential.agent_id,
"secret": credential.secret,
"issued_at": credential.issued_at,
"claim_state": "waiting",
"claim_ttl_s": self.control_plane.limits.claim_ttl_seconds,
"poll_interval_s": self.control_plane.limits.poll_interval_seconds,
"heartbeat_interval_s": self.control_plane.limits.poll_interval_seconds,
"offline_after_s": self.control_plane.limits.offline_after_seconds,
}
)
except (ControlPlaneError, TypeError, ValueError) as exc:
return _response({"ok": False, "error": str(exc)})
def claim_state(self, agent_id: str) -> str:
"""Whether this worker has been claimed yet.
Says nothing a caller could not already know: a worker asks about the
record it created, and the answer is one of three words.
"""
try:
return _response(
{"ok": True, "claim_state": self.control_plane.claim_state(str(agent_id or ""))}
)
except (ControlPlaneError, TypeError, ValueError) as exc:
return _response({"ok": False, "error": str(exc)})
def sync(
self,
agent_id: str,
timestamp: float,
nonce: str,
payload_json: str,
signature: str,
) -> str:
# False until `_verify` returns. Every rejection below is
# therefore unsigned unless the caller proved who they are,
# and the presence of a proof stops being a fact an
# unauthenticated caller can read.
authenticated = False
try:
self._verify(agent_id, timestamp, nonce, PATH_SYNC, payload_json, signature)
authenticated = True
payload = _json_object(payload_json)
snapshot = AgentSnapshot.from_dict(payload.get("snapshot") or {})
self.control_plane.record_snapshot(agent_id, snapshot)
self._apply_advertisement(agent_id, payload)
offers = self.control_plane.offers_for_agent(agent_id)
cancellations = self.control_plane.cancellation_ids_for_agent(agent_id)
return self._signed(
agent_id,
PATH_SYNC,
nonce,
{
"ok": True,
# The worker must inspect the complete immutable spec before
# accepting it into its own bounded queue. In particular,
# model and exact tool-version checks happen agent-side.
"offers": [_job_dict(job, include_spec=True) for job in offers],
"cancellations": list(cancellations),
},
)
except (ControlPlaneError, HandshakeError, TypeError, ValueError) as exc:
return self._signed(
agent_id,
PATH_SYNC,
nonce,
{"ok": False, "error": str(exc)},
authenticated=authenticated,
)
def _apply_advertisement(self, agent_id: str, payload: Mapping[str, Any]) -> None:
"""Narrow this agent's capability to what its operator approved.
The advertisement carries its own detached signature, so the server can
check that the agent really made this statement rather than inferring it
from the fact that the request authenticated. That distinction is what
lets the advertisement be stored, rendered and quoted later and still
mean something.
"""
raw = payload.get("signed_snapshot")
if raw is None:
return
if not isinstance(raw, Mapping):
raise ValueError("signed_snapshot must be an object")
envelope = SignedSnapshot.from_dict(raw)
if envelope.snapshot.agent_id != agent_id:
raise ValueError("signed snapshot agent id does not match credential")
if not self.control_plane.verify_signed_snapshot(agent_id, envelope):
raise ValueError("signed snapshot failed verification")
self.control_plane.record_advertised(agent_id, envelope.advertised)
def catalogue(
self,
agent_id: str,
timestamp: float,
nonce: str,
payload_json: str,
signature: str,
) -> str:
"""State what this server may ask for, signed and challenge-bound.
This is the pull half of the consent handshake: the agent fetches the
catalogue, shows its operator the intersection with what was offered,
and only an explicit approval bound to this catalogue's digest turns
into an advertised, enforceable grant.
"""
# False until `_verify` returns; see `_signed`.
authenticated = False
try:
self._verify(agent_id, timestamp, nonce, PATH_CATALOGUE, payload_json, signature)
authenticated = True
payload = _json_object(payload_json, maximum=4_096)
challenge = payload.get("challenge")
if not isinstance(challenge, str):
raise ValueError("challenge is required")
catalogue = self.control_plane.catalogue_for_agent(agent_id, challenge)
return self._signed(
agent_id,
PATH_CATALOGUE,
nonce,
{"ok": True, "catalogue": catalogue.to_dict()},
)
except (ControlPlaneError, HandshakeError, TypeError, ValueError) as exc:
return self._signed(
agent_id,
PATH_CATALOGUE,
nonce,
{"ok": False, "error": str(exc)},
authenticated=authenticated,
)
def accept(
self,
agent_id: str,
timestamp: float,
nonce: str,
payload_json: str,
signature: str,
) -> str:
# False until `_verify` returns. Every rejection below is
# therefore unsigned unless the caller proved who they are,
# and the presence of a proof stops being a fact an
# unauthenticated caller can read.
authenticated = False
try:
self._verify(agent_id, timestamp, nonce, PATH_ACCEPT, payload_json, signature)
authenticated = True
payload = _json_object(payload_json, maximum=16_384)
job_id = payload.get("job_id")
if not isinstance(job_id, str):
raise ValueError("job_id is required")
accepted = payload.get("accepted", True)
if not isinstance(accepted, bool):
raise ValueError("accepted must be true or false")
if accepted:
# The agent's queue is the only component that knows the real
# position; carry it through so the user sees a number now
# rather than after the next snapshot.
position = payload.get("queue_position")
if position is not None and (
isinstance(position, bool) or not isinstance(position, int)
):
raise ValueError("queue_position must be an integer or null")
job = self.control_plane.accept_offer(
agent_id, job_id, queue_position=position
)
else:
reason = payload.get("reason")
if not isinstance(reason, str):
raise ValueError("reason is required when an offer is rejected")
job = self.control_plane.reject_offer(agent_id, job_id, reason)
return self._signed(agent_id, PATH_ACCEPT, nonce, {"ok": True, "job": _job_dict(job)})
except (ControlPlaneError, TypeError, ValueError) as exc:
return self._signed(
agent_id,
PATH_ACCEPT,
nonce,
{"ok": False, "error": str(exc)},
authenticated=authenticated,
)
def complete(
self,
agent_id: str,
timestamp: float,
nonce: str,
payload_json: str,
signature: str,
) -> str:
# False until `_verify` returns. Every rejection below is
# therefore unsigned unless the caller proved who they are,
# and the presence of a proof stops being a fact an
# unauthenticated caller can read.
authenticated = False
try:
self._verify(agent_id, timestamp, nonce, PATH_COMPLETE, payload_json, signature)
authenticated = True
payload = _json_object(payload_json)
job_id = payload.get("job_id")
if not isinstance(job_id, str):
raise ValueError("job_id is required")
error = payload.get("error")
result = payload.get("result") or {}
if error is None:
job = self.control_plane.complete_job(agent_id, job_id, result)
elif isinstance(error, str):
job = self.control_plane.fail_job(agent_id, job_id, error, result=result)
else:
raise ValueError("error must be text or null")
return self._signed(
agent_id, PATH_COMPLETE, nonce, {"ok": True, "job": _job_dict(job, include_spec=False)}
)
except (ControlPlaneError, TypeError, ValueError) as exc:
return self._signed(
agent_id,
PATH_COMPLETE,
nonce,
{"ok": False, "error": str(exc)},
authenticated=authenticated,
)
def _verify(
self,
agent_id: str,
timestamp: float,
nonce: str,
path: str,
body: str,
signature: str,
) -> None:
if isinstance(timestamp, bool) or not isinstance(timestamp, int | float):
raise ValueError("timestamp must be numeric")
if not math.isfinite(float(timestamp)):
raise ValueError("timestamp must be finite")
self.control_plane.verify_agent_request(
agent_id,
timestamp,
nonce,
"POST",
path,
body,
signature,
)
|