| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299 |
- """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/<gateway_id>/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/<gateway_id>/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/<gateway_id>/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/<task_id>/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/<gateway_id>/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/<gateway_id>/tasks/<task_id>/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/<gateway_id>/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/<gateway_id>/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/<gateway_id>/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/<gateway_id>/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),
- ))
|