| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700 |
- """Fail-closed application service for enterprise edge gateways."""
- from __future__ import annotations
- import hashlib
- import hmac
- import re
- import secrets
- import stat
- import uuid
- from collections.abc import Mapping, Sequence
- from contextlib import suppress
- from datetime import UTC, datetime, timedelta
- from pathlib import Path
- from cryptography.hazmat.primitives import serialization
- from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
- from app.core.edge_gateway.contracts import (
- EdgeContractError,
- EdgeEventContract,
- EdgeTaskContract,
- SignedTaskEnvelope,
- canonical_json_bytes,
- canonical_sha256,
- canonical_timestamp,
- )
- from app.core.edge_gateway.policy import EdgeEgressPolicy, EdgePolicyError
- _SHA256 = re.compile(r"^[0-9a-f]{64}$")
- _VERSION = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._+-]{0,79}$")
- _RELEASE_VERSION = re.compile(
- r"^(0|[1-9][0-9]{0,9})\.(0|[1-9][0-9]{0,9})\.(0|[1-9][0-9]{0,9})$"
- )
- _HOST = re.compile(r"^(?=.{1,253}$)[a-z0-9](?:[a-z0-9.-]*[a-z0-9])?$")
- class EdgeGatewayError(Exception):
- status_code = 400
- code = "EDGE_GATEWAY_ERROR"
- class EdgeGatewayValidationError(EdgeGatewayError):
- code = "EDGE_GATEWAY_INVALID"
- class EdgeGatewayAuthenticationError(EdgeGatewayError):
- status_code = 401
- code = "EDGE_GATEWAY_AUTHENTICATION_FAILED"
- class EdgeGatewayConflictError(EdgeGatewayError):
- status_code = 409
- code = "EDGE_GATEWAY_CONFLICT"
- class EdgeGatewayPayloadTooLargeError(EdgeGatewayValidationError):
- status_code = 413
- code = "EDGE_GATEWAY_PAYLOAD_TOO_LARGE"
- class EdgeGatewayRateLimitError(EdgeGatewayError):
- status_code = 429
- code = "EDGE_GATEWAY_RATE_LIMITED"
- class EdgeGatewayNotFoundError(EdgeGatewayError):
- status_code = 404
- code = "EDGE_GATEWAY_NOT_FOUND"
- class EdgeGatewayConfigurationError(EdgeGatewayError):
- status_code = 503
- code = "EDGE_GATEWAY_SIGNING_UNAVAILABLE"
- def _digest(secret: str) -> str:
- return hashlib.sha256(secret.encode("utf-8")).hexdigest()
- def _uuid(value: object, label: str) -> str:
- try:
- return str(uuid.UUID(str(value)))
- except (ValueError, TypeError, AttributeError) as exc:
- raise EdgeGatewayValidationError(f"{label} is invalid") from exc
- def _sha(value: object, label: str) -> str:
- if not isinstance(value, str) or not _SHA256.fullmatch(value):
- raise EdgeGatewayValidationError(f"{label} is invalid")
- return value
- def _text(value: object, label: str, maximum: int) -> str:
- if not isinstance(value, str) or not value or value.strip() != value or len(value.encode()) > maximum or "\x00" in value:
- raise EdgeGatewayValidationError(f"{label} is invalid")
- return value
- def _version(value: object) -> str:
- value = _text(value, "version", 80)
- if not _VERSION.fullmatch(value):
- raise EdgeGatewayValidationError("version is invalid")
- return value
- def _release_version(value: object, label: str) -> str:
- value = _text(value, label, 80)
- if (
- not _RELEASE_VERSION.fullmatch(value)
- or any(int(part) > 2_147_483_647 for part in value.split("."))
- ):
- raise EdgeGatewayValidationError(f"{label} is invalid")
- return value
- def _release_version_tuple(value: str) -> tuple[int, ...]:
- return tuple(int(item) for item in value.split("."))
- def _artifact_name(value: object) -> str:
- value = _text(value, "artifact_name", 255)
- if (
- value in {".", ".."}
- or not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._+-]{0,254}", value)
- or "/" in value
- or "\\" in value
- ):
- raise EdgeGatewayValidationError("artifact_name is invalid")
- return value
- def _future_deadline(value: object, *, now: datetime) -> str:
- try:
- canonical = canonical_timestamp(value, "deadline_at")
- deadline = datetime.fromisoformat(canonical.replace("Z", "+00:00"))
- except (EdgeContractError, TypeError, ValueError) as exc:
- raise EdgeGatewayValidationError("deadline_at is invalid") from exc
- if deadline <= now or deadline > now + timedelta(days=30):
- raise EdgeGatewayValidationError("deadline_at is invalid")
- return canonical
- def _signing_private_key(value: object) -> Ed25519PrivateKey | None:
- if value is None or value == "":
- return None
- if isinstance(value, Ed25519PrivateKey):
- return value
- if isinstance(value, str):
- encoded = value.encode("utf-8")
- if re.fullmatch(r"[0-9a-fA-F]{64}", value):
- encoded = bytes.fromhex(value)
- elif isinstance(value, bytes):
- encoded = value
- else:
- raise EdgeGatewayConfigurationError("edge signing authority is invalid")
- try:
- if isinstance(encoded, bytes) and len(encoded) == 32:
- return Ed25519PrivateKey.from_private_bytes(encoded)
- key = serialization.load_pem_private_key(encoded, password=None)
- except (TypeError, ValueError) as exc:
- raise EdgeGatewayConfigurationError("edge signing authority is invalid") from exc
- if not isinstance(key, Ed25519PrivateKey):
- raise EdgeGatewayConfigurationError("edge signing authority is invalid")
- return key
- def _hosts(values: object, label: str, *, required: bool) -> list[str]:
- if not isinstance(values, Sequence) or isinstance(values, (str, bytes)):
- raise EdgeGatewayValidationError(f"{label} is invalid")
- result = []
- for item in values:
- if not isinstance(item, str) or item != item.lower() or not _HOST.fullmatch(item) or ".." in item:
- raise EdgeGatewayValidationError(f"{label} is invalid")
- result.append(item)
- if required and not result:
- raise EdgeGatewayValidationError(f"{label} is required")
- if len(result) != len(set(result)) or len(result) > 32:
- raise EdgeGatewayValidationError(f"{label} is invalid")
- return sorted(result)
- def _environment(value: object) -> str:
- if value not in {"development", "staging", "production"}:
- raise EdgeGatewayValidationError("environment is invalid")
- return str(value)
- def _page(limit: object = 50, offset: object = 0) -> tuple[int, int]:
- try:
- limit, offset = int(limit), int(offset)
- except (TypeError, ValueError) as exc:
- raise EdgeGatewayValidationError("pagination is invalid") from exc
- if not 1 <= limit <= 100 or not 0 <= offset <= 10_000:
- raise EdgeGatewayValidationError("pagination is invalid")
- return limit, offset
- def _safe_summary(value: object, classification: str = "diagnostic_summary") -> dict:
- if not isinstance(value, Mapping):
- raise EdgeGatewayValidationError("safe_summary is invalid")
- candidate = dict(value)
- try:
- EdgeEgressPolicy(allowed_control_hosts={"validation.invalid"}).validate_approved_payload(classification, candidate)
- except (EdgePolicyError, ValueError) as exc:
- raise EdgeGatewayValidationError("safe_summary contains unapproved content") from exc
- return candidate
- class EdgeGatewayService:
- def __init__(
- self, repository, *, signing_private_key=None,
- signing_private_key_file=None, signing_key_provider=None,
- signing_public_key=None, signing_key_id=None, production=False, clock=None,
- ):
- self.repository = repository
- sources = sum(
- item not in {None, ""}
- for item in (signing_private_key, signing_private_key_file)
- ) + int(signing_key_provider is not None)
- if sources > 1:
- raise EdgeGatewayConfigurationError("edge signing authority source is ambiguous")
- if production and signing_private_key not in {None, ""}:
- raise EdgeGatewayConfigurationError(
- "production inline edge signing private key is forbidden"
- )
- key_material = signing_private_key
- if signing_private_key_file not in {None, ""}:
- if not isinstance(signing_private_key_file, str):
- raise EdgeGatewayConfigurationError("edge signing key file is invalid")
- path = Path(signing_private_key_file)
- try:
- if (
- not path.is_absolute()
- or not path.is_file()
- or path.is_symlink()
- or stat.S_IMODE(path.stat().st_mode) != 0o600
- ):
- raise EdgeGatewayConfigurationError(
- "edge signing key file permissions must be 0600"
- )
- key_material = path.read_bytes()
- except OSError as exc:
- raise EdgeGatewayConfigurationError(
- "edge signing key file is unavailable"
- ) from exc
- elif signing_key_provider is not None:
- if not callable(signing_key_provider):
- raise EdgeGatewayConfigurationError("edge signing key provider is invalid")
- try:
- key_material = signing_key_provider()
- except Exception as exc:
- raise EdgeGatewayConfigurationError(
- "edge signing key provider is unavailable"
- ) from exc
- self._signing_private_key = _signing_private_key(key_material)
- if self._signing_private_key is not None and signing_public_key not in {None, ""}:
- if (
- not isinstance(signing_public_key, str)
- or not re.fullmatch(r"[0-9a-f]{64}", signing_public_key)
- ):
- raise EdgeGatewayConfigurationError("edge signing public key is invalid")
- actual_public = self._signing_private_key.public_key().public_bytes(
- serialization.Encoding.Raw,
- serialization.PublicFormat.Raw,
- ).hex()
- if not hmac.compare_digest(actual_public, signing_public_key):
- raise EdgeGatewayConfigurationError(
- "edge signing public key does not match private key"
- )
- elif production and self._signing_private_key is not None:
- raise EdgeGatewayConfigurationError(
- "production edge signing public key is required"
- )
- self._signing_key_id = signing_key_id
- self._clock = clock or (lambda: datetime.now(UTC))
- def _require_signer(self) -> tuple[Ed25519PrivateKey, str]:
- key_id = self._signing_key_id
- if (
- self._signing_private_key is None
- or not isinstance(key_id, str)
- or not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._:-]{0,127}", key_id)
- ):
- self.repository.rollback()
- raise EdgeGatewayConfigurationError(
- "edge signing authority is not configured"
- )
- return self._signing_private_key, key_id
- def _signed_task_envelope(self, contract: EdgeTaskContract) -> dict:
- private_key, key_id = self._require_signer()
- issued_at = canonical_timestamp(
- self._clock().astimezone(UTC).isoformat().replace("+00:00", "Z")
- )
- unsigned = {
- "task": contract.to_mapping(),
- "authority_key_id": key_id,
- "signature_algorithm": "Ed25519",
- "contract_digest": contract.digest,
- "gateway_id": contract.gateway_id,
- "environment": contract.environment,
- "network_zone": contract.network_zone,
- "policy_digest": contract.policy_digest,
- "purpose": contract.purpose,
- "issued_at": issued_at,
- "expires_at": contract.deadline_at,
- }
- try:
- signed_bytes = SignedTaskEnvelope.canonical_unsigned_bytes(unsigned)
- except EdgeContractError as exc:
- self.repository.rollback()
- raise EdgeGatewayValidationError(
- "edge task authority validity is invalid"
- ) from exc
- return {**unsigned, "signature": private_key.sign(signed_bytes).hex()}
- def _signed_release_manifest(self, unsigned: Mapping[str, object]) -> dict:
- private_key, key_id = self._require_signer()
- values = dict(unsigned)
- values["signature_algorithm"] = "Ed25519"
- values["key_id"] = key_id
- digest = canonical_sha256(values)
- return {
- **values,
- "manifest_digest": digest,
- "signature": private_key.sign(canonical_json_bytes(values)).hex(),
- }
- def _failure(
- self, gateway_id, event_type: str, reason_code: str, *, trusted: bool = False,
- ) -> int:
- count = self.repository.audit_failure(
- gateway_id, event_type, reason_code, trusted=trusted,
- )
- if count >= 10:
- raise EdgeGatewayRateLimitError("edge failure rate limit exceeded")
- return count
- def create_enrollment(
- self, *, gateway_name, environment, network_zone, policy_digest,
- allowed_control_hosts, allowed_proxy_hosts, ttl_seconds, actor_uid,
- expected_certificate_sha256,
- ) -> dict:
- ttl = int(ttl_seconds)
- if ttl < 60 or ttl > 86_400:
- raise EdgeGatewayValidationError("enrollment TTL is invalid")
- gateway_id, uid = str(uuid.uuid4()), str(uuid.uuid4())
- token = "dope_" + secrets.token_urlsafe(32)
- values = {
- "uid": uid, "gateway_id": gateway_id,
- "gateway_name": _text(gateway_name, "gateway_name", 200),
- "environment": _environment(environment),
- "network_zone": _text(network_zone, "network_zone", 120),
- "policy_digest": _sha(policy_digest, "policy_digest"),
- "expected_certificate_sha256": _sha(expected_certificate_sha256, "expected_certificate_sha256"),
- "control_hosts": _hosts(allowed_control_hosts, "allowed_control_hosts", required=True),
- "proxy_hosts": _hosts(allowed_proxy_hosts, "allowed_proxy_hosts", required=False),
- "token_hash": _digest(token), "ttl": ttl, "actor_uid": _uuid(actor_uid, "actor_uid"),
- }
- self.repository.create_enrollment(values)
- return {"enrollment_id": uid, "gateway_id": gateway_id, "enrollment_token": token, "expires_in": ttl, "returned_once": True}
- def register(
- self, *, enrollment_token, certificate_sha256, gateway_id, environment,
- network_zone, policy_digest, allowed_control_hosts, allowed_proxy_hosts, version,
- ) -> dict:
- gateway_id = _uuid(gateway_id, "gateway_id")
- if not isinstance(enrollment_token, str) or not enrollment_token.startswith("dope_") or len(enrollment_token) > 128:
- self._failure(gateway_id, "registration_rejected", "enrollment_invalid")
- raise EdgeGatewayAuthenticationError("edge enrollment is invalid or expired")
- try:
- certificate = _sha(certificate_sha256, "certificate_sha256")
- supplied = {
- "gateway_id": gateway_id, "environment": _environment(environment),
- "network_zone": _text(network_zone, "network_zone", 120),
- "policy_digest": _sha(policy_digest, "policy_digest"),
- "allowed_control_hosts": _hosts(allowed_control_hosts, "allowed_control_hosts", required=True),
- "allowed_proxy_hosts": _hosts(allowed_proxy_hosts, "allowed_proxy_hosts", required=False),
- }
- except EdgeGatewayValidationError as exc:
- self._failure(gateway_id, "registration_rejected", "binding_invalid")
- raise EdgeGatewayAuthenticationError("edge enrollment binding was rejected") from exc
- enrollment = self.repository.lock_enrollment(_digest(enrollment_token))
- now = datetime.now(UTC)
- if enrollment is None or enrollment["status"] != "pending" or enrollment["expires_at"] <= now:
- self._failure(supplied["gateway_id"], "registration_rejected", "enrollment_invalid")
- raise EdgeGatewayAuthenticationError("edge enrollment is invalid or expired")
- expected = {key: enrollment[key] for key in supplied}
- expected["allowed_control_hosts"] = sorted(expected["allowed_control_hosts"])
- expected["allowed_proxy_hosts"] = sorted(expected["allowed_proxy_hosts"])
- if any(not hmac.compare_digest(str(expected[key]), str(supplied[key])) for key in supplied) or not hmac.compare_digest(enrollment["expected_certificate_sha256"], certificate):
- self._failure(supplied["gateway_id"], "registration_rejected", "binding_mismatch")
- raise EdgeGatewayAuthenticationError("edge enrollment binding was rejected")
- credential = "dopg_" + secrets.token_urlsafe(40)
- try:
- result = self.repository.complete_registration(enrollment, _digest(credential), certificate, _version(version))
- except Exception:
- self.repository.rollback()
- raise
- return {**result, "credential": credential, "credential_returned_once": True, "certificate_sha256": certificate}
- def authenticate(self, *, credential, certificate_sha256, gateway_id, environment, network_zone, generation, lock=False) -> dict:
- try:
- gateway_id = _uuid(gateway_id, "gateway_id")
- generation = int(generation)
- supplied = (
- gateway_id, _environment(environment),
- _text(network_zone, "network_zone", 120), generation,
- _sha(certificate_sha256, "certificate_sha256"),
- )
- except EdgeGatewayValidationError as exc:
- with suppress(EdgeGatewayValidationError):
- self._failure(_uuid(gateway_id, "gateway_id"), "authentication_rejected", "binding_invalid")
- raise EdgeGatewayAuthenticationError("edge credential binding was rejected") from exc
- except (TypeError, ValueError) as exc:
- self._failure(gateway_id, "authentication_rejected", "binding_invalid")
- raise EdgeGatewayAuthenticationError("edge credential binding was rejected") from exc
- if not isinstance(credential, str) or not credential.startswith("dopg_") or len(credential) > 160:
- self._failure(gateway_id, "authentication_rejected", "credential_invalid")
- raise EdgeGatewayAuthenticationError("edge credential binding was rejected")
- row = self.repository.authenticate(_digest(credential), lock=lock)
- if row is None:
- self._failure(supplied[0], "authentication_rejected", "credential_invalid")
- raise EdgeGatewayAuthenticationError("edge credential binding was rejected")
- expected = (row["gateway_id"], row["environment"], row["network_zone"], row["generation"], row["certificate_sha256"])
- if any(not hmac.compare_digest(str(a), str(b)) for a, b in zip(supplied, expected, strict=True)):
- self._failure(row["gateway_id"], "authentication_rejected", "binding_mismatch")
- raise EdgeGatewayAuthenticationError("edge credential binding was rejected")
- return dict(row)
- def list_gateways(self, *, limit=50, offset=0) -> list[dict]:
- limit, offset = _page(limit, offset)
- return self.repository.list_gateways(limit=limit, offset=offset)
- def heartbeat(self, *, version, safe_summary, **auth) -> dict:
- version = _version(version)
- safe_summary = _safe_summary(safe_summary, "health_summary")
- identity = self.authenticate(**auth, lock=True)
- result = self.repository.heartbeat(identity["gateway_id"], version, safe_summary)
- result["status"] = "online"
- return result
- def rotate(self, gateway_id, *, certificate_sha256, request_id, actor_uid) -> dict:
- request_id = _text(request_id, "request_id", 255)
- certificate_sha256 = _sha(certificate_sha256, "certificate_sha256")
- request_digest = canonical_sha256({"gateway_id": str(gateway_id), "certificate_sha256": certificate_sha256})
- token = "dopg_" + secrets.token_urlsafe(40)
- gateway_id = _uuid(gateway_id, "gateway_id")
- result = self.repository.rotate(gateway_id, _digest(token), certificate_sha256, _uuid(actor_uid, "actor_uid"), request_id, request_digest)
- if result is None:
- raise EdgeGatewayNotFoundError("edge gateway was not found")
- if result.get("conflict"):
- raise EdgeGatewayConflictError("rotation request conflicts with an existing request")
- if result.get("replayed"):
- return {"gateway_id": gateway_id, "generation": result["generation"], "credential": None, "credential_returned_once": False, "replayed": True}
- return {**result, "gateway_id": gateway_id, "credential": token, "credential_returned_once": True}
- def revoke(self, gateway_id, *, actor_uid) -> bool:
- return self.repository.revoke(_uuid(gateway_id, "gateway_id"), _uuid(actor_uid, "actor_uid"))
- def issue_task(self, value: Mapping[str, object], *, actor_uid) -> dict:
- try:
- contract = EdgeTaskContract.from_mapping(value)
- except EdgeContractError as exc:
- raise EdgeGatewayValidationError("edge task contract is invalid") from exc
- if _uuid(contract.gateway_id, "gateway_id") != contract.gateway_id:
- raise EdgeGatewayValidationError("gateway_id is invalid")
- gateway = self.repository.gateway_binding(contract.gateway_id)
- if gateway is None or gateway["status"] not in {"active", "offline"}:
- self.repository.rollback()
- raise EdgeGatewayNotFoundError("edge gateway was not found")
- if any(
- getattr(contract, field) != gateway[field]
- for field in ("environment", "network_zone", "policy_digest")
- ):
- self._failure(
- contract.gateway_id, "binding_rejected", "task_binding_mismatch",
- trusted=True,
- )
- raise EdgeGatewayValidationError("edge task gateway binding is invalid")
- authority = self._signed_task_envelope(contract)
- row = self.repository.insert_task(
- contract.to_mapping(), contract.digest, authority,
- _uuid(actor_uid, "actor_uid"),
- )
- if not row:
- raise EdgeGatewayNotFoundError("edge gateway was not found")
- if row["contract_digest"] != contract.digest:
- self._failure(
- contract.gateway_id, "idempotency_conflict", "task_contract_conflict",
- trusted=True,
- )
- raise EdgeGatewayConflictError("task idempotency key conflicts with an existing contract")
- return {"task_id": row["task_id"], "contract_digest": row["contract_digest"]}
- def pull_task(self, **auth) -> dict:
- identity = self.authenticate(**auth, lock=True)
- token = "dopl_" + secrets.token_urlsafe(32)
- row = self.repository.pull_task(identity["gateway_id"], _digest(token))
- if row is None:
- return {"task": None}
- envelope = dict(row["authority"])
- if envelope.get("task") != row["contract"]:
- self.repository.rollback()
- raise EdgeGatewayConfigurationError("stored task authority is invalid")
- return {
- "task": row["contract"],
- "signed_task_envelope": envelope,
- "lease_token": token,
- "lease_expires_at": row["lease_expires_at"].isoformat(),
- }
- def cancel_task(self, task_id, *, actor_uid) -> bool:
- return self.repository.cancel_task(_text(task_id, "task_id", 255), _uuid(actor_uid, "actor_uid"))
- def task_outcome(self, task_id, *, outcome, lease_token, safe_summary, **auth) -> dict:
- if outcome not in {"completed", "failed", "cancelled"}:
- raise EdgeGatewayValidationError("task outcome is invalid")
- if not isinstance(lease_token, str) or not lease_token.startswith("dopl_") or len(lease_token) > 128:
- raise EdgeGatewayAuthenticationError("edge task lease binding was rejected")
- summary = _safe_summary(safe_summary)
- identity = self.authenticate(**auth, lock=True)
- result = self.repository.task_outcome(
- identity["gateway_id"], _text(task_id, "task_id", 255),
- _digest(lease_token), outcome, summary,
- )
- if result is None:
- self._failure(
- identity["gateway_id"], "binding_rejected", "task_outcome_rejected",
- trusted=True,
- )
- raise EdgeGatewayConflictError("task outcome was rejected")
- return result
- def reconcile(
- self, *, limit=50, cancel_cursor=None, release_cursor=None, **auth
- ) -> dict:
- limit, _ = _page(limit, 0)
- if cancel_cursor is not None:
- cancel_cursor = _text(cancel_cursor, "cancel_cursor", 255)
- if release_cursor is not None:
- release_cursor = _uuid(release_cursor, "release_cursor")
- identity = self.authenticate(**auth, lock=True)
- cancellations = self.repository.cancelled_tasks(
- identity["gateway_id"], limit=limit, cursor=cancel_cursor
- )
- releases = self.repository.unfinished_releases(
- identity["gateway_id"], limit=limit, cursor=release_cursor
- )
- result = {
- "cancelled_task_ids": cancellations["items"],
- "cancel_next_cursor": cancellations["next_cursor"],
- "release_offers": releases["items"],
- "release_next_cursor": releases["next_cursor"],
- "release_baseline": self.repository.latest_release_baseline(
- identity["gateway_id"]
- ),
- }
- self.repository.commit()
- return result
- def accept_event(self, *, event, lease_token=None, **auth) -> dict:
- identity = self.authenticate(**auth, lock=True)
- if not isinstance(event, Mapping):
- self.repository.rollback()
- raise EdgeGatewayValidationError("edge event contract is invalid")
- policy = EdgeEgressPolicy(allowed_control_hosts=set(identity["allowed_control_hosts"]), allowed_proxy_hosts=set(identity["allowed_proxy_hosts"]))
- try:
- approved = policy.approve_event(event)
- except EdgePolicyError as exc:
- self._failure(identity["gateway_id"], "policy_rejected", "egress_policy_rejected", trusted=True)
- raise EdgeGatewayValidationError("edge event is not approved for control-plane egress") from exc
- try:
- contract = EdgeEventContract.from_mapping(approved)
- except EdgeContractError as exc:
- self._failure(identity["gateway_id"], "binding_rejected", "event_contract_invalid", trusted=True)
- raise EdgeGatewayValidationError("edge event contract is invalid") from exc
- if contract.gateway_id != identity["gateway_id"] or contract.environment != identity["environment"] or contract.network_zone != identity["network_zone"] or contract.policy_digest != identity["policy_digest"]:
- self._failure(identity["gateway_id"], "binding_rejected", "event_gateway_mismatch", trusted=True)
- raise EdgeGatewayAuthenticationError("edge event binding was rejected")
- supplied_digest = contract.digest
- existing = self.repository.event_by_id(identity["gateway_id"], contract.event_id)
- if existing is not None:
- if not hmac.compare_digest(existing["contract_digest"], supplied_digest):
- self._failure(identity["gateway_id"], "idempotency_conflict", "event_contract_conflict", trusted=True)
- raise EdgeGatewayConflictError("event identifier conflicts with an existing contract")
- self.repository.rollback()
- return dict(existing["ack"])
- if not isinstance(lease_token, str) or not lease_token.startswith("dopl_") or len(lease_token) > 128:
- self._failure(identity["gateway_id"], "binding_rejected", "lease_invalid", trusted=True)
- raise EdgeGatewayAuthenticationError("edge task lease binding was rejected")
- remote_lease_digest = _digest(lease_token)
- task_contract = self.repository.task_lease_binding(
- contract.task_id, contract.gateway_id, remote_lease_digest
- )
- if task_contract is None:
- existing = self.repository.event_by_id(identity["gateway_id"], contract.event_id)
- if existing is not None and hmac.compare_digest(existing["contract_digest"], supplied_digest):
- self.repository.rollback()
- return dict(existing["ack"])
- self._failure(identity["gateway_id"], "binding_rejected", "lease_mismatch", trusted=True)
- raise EdgeGatewayAuthenticationError("edge task lease binding was rejected")
- if any(contract.to_mapping()[field] != task_contract[field] for field in (
- "gateway_id", "environment", "network_zone", "purpose", "policy_digest",
- "attempt",
- )):
- self._failure(identity["gateway_id"], "binding_rejected", "event_task_mismatch", trusted=True)
- raise EdgeGatewayValidationError("edge event task binding is invalid")
- ack = {
- "event_id": contract.event_id,
- "event_digest": contract.digest,
- "remote_lease_digest": remote_lease_digest,
- "received_at": datetime.now(UTC).isoformat().replace("+00:00", "Z"),
- "status": "accepted",
- }
- accepted = self.repository.insert_event(contract.to_mapping(), contract.digest, ack)
- if accepted is None:
- self._failure(identity["gateway_id"], "idempotency_conflict", "event_contract_conflict", trusted=True)
- raise EdgeGatewayConflictError("event identifier conflicts with an existing contract")
- return accepted
- def offer_release(
- self, gateway_id, *, version, artifact_digest, rollback_version,
- request_id, actor_uid, artifact_name=None, deadline_at=None,
- ) -> dict:
- gateway_id = _uuid(gateway_id, "gateway_id")
- version = _release_version(version, "version")
- artifact_digest = _sha(artifact_digest, "artifact_digest")
- rollback_version = _release_version(rollback_version, "rollback_version")
- if _release_version_tuple(version) <= _release_version_tuple(rollback_version):
- raise EdgeGatewayValidationError(
- "release version must be newer than rollback_version"
- )
- request_id = _text(request_id, "request_id", 255)
- now = self._clock().astimezone(UTC)
- supplied_deadline = deadline_at
- artifact_name = _artifact_name(
- artifact_name or f"edge-agent-{version}.bin"
- )
- deadline_at = _future_deadline(
- deadline_at
- or (now + timedelta(hours=24)).isoformat().replace("+00:00", "Z"),
- now=now,
- )
- unsigned = {
- "release_id": str(uuid.uuid4()), "version": version,
- "rollback_version": rollback_version,
- "artifact_digest": artifact_digest, "artifact_name": artifact_name,
- "deadline_at": deadline_at, "status": "offered",
- }
- manifest = self._signed_release_manifest(unsigned)
- request_digest = canonical_sha256({
- "gateway_id": gateway_id, "version": version,
- "artifact_digest": artifact_digest, "rollback_version": rollback_version,
- "artifact_name": artifact_name,
- "deadline_at": deadline_at if supplied_deadline is not None else None,
- })
- result = self.repository.offer_release(
- gateway_id, manifest, _uuid(actor_uid, "actor_uid"), request_id,
- request_digest,
- )
- if result.get("conflict"):
- raise EdgeGatewayConflictError("release request conflicts with an existing request")
- if not result:
- raise EdgeGatewayNotFoundError("edge gateway was not found")
- return result
- def acknowledge_release(self, *, release_id, outcome, safe_summary, **auth) -> dict:
- statuses = {"accepted": "accepted", "installed": "installed", "failed": "failed", "rollback": "rolled_back"}
- if outcome not in statuses:
- raise EdgeGatewayValidationError("release outcome is invalid")
- release_id = _uuid(release_id, "release_id")
- safe_summary = _safe_summary(safe_summary)
- identity = self.authenticate(**auth, lock=True)
- current = self.repository.release_for_ack(identity["gateway_id"], release_id)
- if current is None:
- self.repository.rollback()
- self._failure(
- identity["gateway_id"], "release_rejected",
- "release_transition_invalid", trusted=True,
- )
- raise EdgeGatewayConflictError("release acknowledgement was rejected")
- stored = dict(current["signed_manifest"])
- unsigned = {
- key: value for key, value in stored.items()
- if key not in {
- "manifest_digest", "signature", "signature_algorithm", "key_id"
- }
- }
- unsigned["status"] = statuses[outcome]
- manifest = self._signed_release_manifest(unsigned)
- row = self.repository.acknowledge_release(
- identity["gateway_id"], release_id, statuses[outcome], safe_summary,
- manifest,
- )
- if row is None:
- self._failure(identity["gateway_id"], "release_rejected", "release_transition_invalid", trusted=True)
- raise EdgeGatewayConflictError("release acknowledgement was rejected")
- return row
|