"""Bounded HTTP boundary for human and machine edge-gateway operations.""" from __future__ import annotations import hmac import ipaddress import logging import re from datetime import UTC, datetime from urllib.parse import unquote from cryptography import x509 from cryptography.hazmat.primitives import hashes from flask import current_app, g, jsonify, request from app.api.data_source import bp from app.core.edge_gateway.contracts import EdgeContractError, preflight_json from app.core.edge_gateway.service import ( EdgeGatewayAuthenticationError, EdgeGatewayError, EdgeGatewayPayloadTooLargeError, EdgeGatewayService, EdgeGatewayValidationError, ) from app.models.result import failed, success logger = logging.getLogger(__name__) MAX_EDGE_REQUEST_BYTES = 262_144 MAX_EDGE_REQUEST_NODES = 5_000 MAX_EDGE_RESPONSE_BYTES = 1_048_576 def _service() -> EdgeGatewayService: from app import db from app.core.edge_gateway.repository import EdgeGatewayRepository return EdgeGatewayService( EdgeGatewayRepository(db.session), signing_private_key=current_app.config.get( "EDGE_GATEWAY_SIGNING_PRIVATE_KEY" ), signing_private_key_file=current_app.config.get( "EDGE_GATEWAY_SIGNING_PRIVATE_KEY_FILE" ), signing_public_key=current_app.config.get( "EDGE_GATEWAY_SIGNING_PUBLIC_KEY" ), signing_key_id=current_app.config.get("EDGE_GATEWAY_SIGNING_KEY_ID"), production=( str(current_app.config.get("FLASK_ENV", "")).lower() == "production" ), ) def _actor_uid(): identity = getattr(g, "current_user", {}) or {} return identity.get("id") or identity.get("sub") def _payload() -> dict: if request.content_length is not None and request.content_length > MAX_EDGE_REQUEST_BYTES: raise EdgeGatewayPayloadTooLargeError("edge request is too large") raw = request.get_data(cache=True) if len(raw) > MAX_EDGE_REQUEST_BYTES: raise EdgeGatewayPayloadTooLargeError("edge request is too large") value = request.get_json(silent=True) if value is None and not raw: return {} if not isinstance(value, dict): raise EdgeGatewayValidationError("edge request body is invalid") try: preflight_json( value, max_depth=32, max_nodes=MAX_EDGE_REQUEST_NODES, max_string_bytes=131_072, max_total_bytes=MAX_EDGE_REQUEST_BYTES, ) except EdgeContractError as exc: raise EdgeGatewayPayloadTooLargeError("edge request exceeds safe limits") from exc return value def _tls_certificate_sha256() -> str: configured = current_app.config.get("EDGE_MTLS_TRUSTED_PROXY_IPS", ()) if isinstance(configured, str): configured = tuple( item.strip() for item in configured.split(",") if item.strip() ) if not isinstance(configured, (tuple, list, set, frozenset)) or not configured: raise EdgeGatewayAuthenticationError( "trusted mTLS proxy is not configured" ) try: trusted = {str(ipaddress.ip_address(item)) for item in configured} remote = str(ipaddress.ip_address(request.remote_addr or "")) except ValueError as exc: raise EdgeGatewayAuthenticationError( "request did not originate from a trusted mTLS proxy" ) from exc if remote not in trusted: raise EdgeGatewayAuthenticationError( "request did not originate from a trusted mTLS proxy" ) if request.headers.get("X-DataOps-Edge-Client-Verify") != "SUCCESS": raise EdgeGatewayAuthenticationError( "verified TLS client certificate is required" ) escaped_certificate = request.headers.get("X-DataOps-Edge-Client-Cert", "") if not escaped_certificate or len(escaped_certificate.encode()) > 32_768: raise EdgeGatewayAuthenticationError( "verified TLS client certificate is required" ) try: certificate = x509.load_pem_x509_certificate( unquote(escaped_certificate).encode() ) except ValueError as exc: raise EdgeGatewayAuthenticationError( "verified TLS client certificate is invalid" ) from exc now = datetime.now(UTC) if ( certificate.not_valid_before_utc > now or certificate.not_valid_after_utc <= now ): raise EdgeGatewayAuthenticationError( "verified TLS client certificate is not currently valid" ) proof = certificate.fingerprint(hashes.SHA256()).hex() header = request.headers.get("X-Edge-Certificate-SHA256") if header and ( not re.fullmatch(r"[0-9a-f]{64}", header) or not hmac.compare_digest(header, proof) ): raise EdgeGatewayAuthenticationError( "TLS client certificate proof does not match" ) return proof def _machine_auth(gateway_id: str, payload: dict) -> dict: return { "credential": request.headers.get("X-Edge-Credential", ""), "certificate_sha256": _tls_certificate_sha256(), "gateway_id": gateway_id, "environment": payload.get("environment"), "network_zone": payload.get("network_zone"), "generation": payload.get("generation"), } def _respond(call, created: bool = False): try: result = call() response = jsonify(success(result)) if len(response.get_data()) > MAX_EDGE_RESPONSE_BYTES: response = jsonify(failed( "边缘网关响应超过安全上限", code=500, error={"code": "EDGE_GATEWAY_RESPONSE_LIMIT"}, )) response.headers["Cache-Control"] = "no-store" return response, 500 response.headers["Cache-Control"] = "no-store" return response, 201 if created else 200 except EdgeGatewayError as exc: logger.warning("edge gateway request rejected: code=%s", exc.code) response = jsonify(failed("边缘网关请求被拒绝", code=exc.status_code, error={"code": exc.code})) response.headers["Cache-Control"] = "no-store" return response, exc.status_code except Exception: logger.exception("edge gateway request failed") response = jsonify(failed("边缘网关操作失败", code=500, error={"code": "EDGE_GATEWAY_INTERNAL"})) response.headers["Cache-Control"] = "no-store" return response, 500 @bp.errorhandler(EdgeGatewayError) def edge_gateway_error(error): response = jsonify(failed("边缘网关请求被拒绝", code=error.status_code, error={"code": error.code})) response.headers["Cache-Control"] = "no-store" return response, error.status_code @bp.route("/edge/enrollments", methods=["POST"]) def edge_create_enrollment(): payload = _payload() return _respond(lambda: _service().create_enrollment( gateway_name=payload.get("gateway_name"), environment=payload.get("environment"), network_zone=payload.get("network_zone"), policy_digest=payload.get("policy_digest"), allowed_control_hosts=payload.get("allowed_control_hosts"), allowed_proxy_hosts=payload.get("allowed_proxy_hosts", []), expected_certificate_sha256=payload.get("expected_certificate_sha256"), ttl_seconds=payload.get("ttl_seconds", 600), actor_uid=_actor_uid(), ), created=True) @bp.route("/edge/register", methods=["POST"]) def edge_register(): payload = _payload() return _respond(lambda: _service().register( enrollment_token=request.headers.get("X-Edge-Enrollment", ""), certificate_sha256=_tls_certificate_sha256(), gateway_id=payload.get("gateway_id"), environment=payload.get("environment"), network_zone=payload.get("network_zone"), policy_digest=payload.get("policy_digest"), allowed_control_hosts=payload.get("allowed_control_hosts"), allowed_proxy_hosts=payload.get("allowed_proxy_hosts", []), version=payload.get("version"), ), created=True) @bp.route("/edge/gateways", methods=["GET"]) def edge_list_gateways(): return _respond(lambda: {"gateways": _service().list_gateways( limit=request.args.get("limit", 50), offset=request.args.get("offset", 0), )}) @bp.route("/edge/gateways//heartbeat", methods=["POST"]) def edge_heartbeat(gateway_id): payload = _payload() return _respond(lambda: _service().heartbeat( **_machine_auth(gateway_id, payload), version=payload.get("version"), safe_summary=payload.get("safe_summary", {}), )) @bp.route("/edge/gateways//rotate", methods=["POST"]) def edge_rotate(gateway_id): payload = _payload() return _respond(lambda: _service().rotate( gateway_id, certificate_sha256=payload.get("certificate_sha256"), actor_uid=_actor_uid(), request_id=payload.get("request_id"), )) @bp.route("/edge/gateways//revoke", methods=["POST"]) def edge_revoke(gateway_id): return _respond(lambda: {"revoked": _service().revoke(gateway_id, actor_uid=_actor_uid())}) @bp.route("/edge/tasks", methods=["POST"]) def edge_issue_task(): return _respond(lambda: _service().issue_task(_payload(), actor_uid=_actor_uid()), created=True) @bp.route("/edge/tasks//cancel", methods=["POST"]) def edge_cancel_task(task_id): return _respond(lambda: {"cancel_requested": _service().cancel_task(task_id, actor_uid=_actor_uid())}) @bp.route("/edge/gateways//tasks/pull", methods=["POST"]) def edge_pull_task(gateway_id): payload = _payload() return _respond(lambda: _service().pull_task(**_machine_auth(gateway_id, payload))) @bp.route("/edge/gateways//tasks//outcome", methods=["POST"]) def edge_task_outcome(gateway_id, task_id): payload = _payload() return _respond(lambda: _service().task_outcome( task_id, outcome=payload.get("outcome"), lease_token=payload.get("lease_token"), safe_summary=payload.get("safe_summary", {}), **_machine_auth(gateway_id, payload), )) @bp.route("/edge/gateways//events", methods=["POST"]) def edge_accept_event(gateway_id): payload = _payload() return _respond(lambda: _service().accept_event( event=payload.get("event"), lease_token=payload.get("lease_token"), **_machine_auth(gateway_id, payload), )) @bp.route("/edge/gateways//reconcile", methods=["POST"]) def edge_reconcile(gateway_id): payload = _payload() return _respond(lambda: _service().reconcile( limit=payload.get("limit", 50), cancel_cursor=payload.get("cancel_cursor"), release_cursor=payload.get("release_cursor"), **_machine_auth(gateway_id, payload), )) @bp.route("/edge/gateways//releases", methods=["POST"]) def edge_offer_release(gateway_id): payload = _payload() return _respond(lambda: _service().offer_release( gateway_id, version=payload.get("version"), artifact_digest=payload.get("artifact_digest"), artifact_name=payload.get("artifact_name"), deadline_at=payload.get("deadline_at"), rollback_version=payload.get("rollback_version"), request_id=payload.get("request_id"), actor_uid=_actor_uid(), ), created=True) @bp.route("/edge/gateways//releases/ack", methods=["POST"]) def edge_ack_release(gateway_id): payload = _payload() return _respond(lambda: _service().acknowledge_release( release_id=payload.get("release_id"), outcome=payload.get("outcome"), safe_summary=payload.get("safe_summary", {}), **_machine_auth(gateway_id, payload), ))