edge_routes.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299
  1. """Bounded HTTP boundary for human and machine edge-gateway operations."""
  2. from __future__ import annotations
  3. import hmac
  4. import ipaddress
  5. import logging
  6. import re
  7. from datetime import UTC, datetime
  8. from urllib.parse import unquote
  9. from cryptography import x509
  10. from cryptography.hazmat.primitives import hashes
  11. from flask import current_app, g, jsonify, request
  12. from app.api.data_source import bp
  13. from app.core.edge_gateway.contracts import EdgeContractError, preflight_json
  14. from app.core.edge_gateway.service import (
  15. EdgeGatewayAuthenticationError,
  16. EdgeGatewayError,
  17. EdgeGatewayPayloadTooLargeError,
  18. EdgeGatewayService,
  19. EdgeGatewayValidationError,
  20. )
  21. from app.models.result import failed, success
  22. logger = logging.getLogger(__name__)
  23. MAX_EDGE_REQUEST_BYTES = 262_144
  24. MAX_EDGE_REQUEST_NODES = 5_000
  25. MAX_EDGE_RESPONSE_BYTES = 1_048_576
  26. def _service() -> EdgeGatewayService:
  27. from app import db
  28. from app.core.edge_gateway.repository import EdgeGatewayRepository
  29. return EdgeGatewayService(
  30. EdgeGatewayRepository(db.session),
  31. signing_private_key=current_app.config.get(
  32. "EDGE_GATEWAY_SIGNING_PRIVATE_KEY"
  33. ),
  34. signing_private_key_file=current_app.config.get(
  35. "EDGE_GATEWAY_SIGNING_PRIVATE_KEY_FILE"
  36. ),
  37. signing_public_key=current_app.config.get(
  38. "EDGE_GATEWAY_SIGNING_PUBLIC_KEY"
  39. ),
  40. signing_key_id=current_app.config.get("EDGE_GATEWAY_SIGNING_KEY_ID"),
  41. production=(
  42. str(current_app.config.get("FLASK_ENV", "")).lower() == "production"
  43. ),
  44. )
  45. def _actor_uid():
  46. identity = getattr(g, "current_user", {}) or {}
  47. return identity.get("id") or identity.get("sub")
  48. def _payload() -> dict:
  49. if request.content_length is not None and request.content_length > MAX_EDGE_REQUEST_BYTES:
  50. raise EdgeGatewayPayloadTooLargeError("edge request is too large")
  51. raw = request.get_data(cache=True)
  52. if len(raw) > MAX_EDGE_REQUEST_BYTES:
  53. raise EdgeGatewayPayloadTooLargeError("edge request is too large")
  54. value = request.get_json(silent=True)
  55. if value is None and not raw:
  56. return {}
  57. if not isinstance(value, dict):
  58. raise EdgeGatewayValidationError("edge request body is invalid")
  59. try:
  60. preflight_json(
  61. value, max_depth=32, max_nodes=MAX_EDGE_REQUEST_NODES,
  62. max_string_bytes=131_072, max_total_bytes=MAX_EDGE_REQUEST_BYTES,
  63. )
  64. except EdgeContractError as exc:
  65. raise EdgeGatewayPayloadTooLargeError("edge request exceeds safe limits") from exc
  66. return value
  67. def _tls_certificate_sha256() -> str:
  68. configured = current_app.config.get("EDGE_MTLS_TRUSTED_PROXY_IPS", ())
  69. if isinstance(configured, str):
  70. configured = tuple(
  71. item.strip() for item in configured.split(",") if item.strip()
  72. )
  73. if not isinstance(configured, (tuple, list, set, frozenset)) or not configured:
  74. raise EdgeGatewayAuthenticationError(
  75. "trusted mTLS proxy is not configured"
  76. )
  77. try:
  78. trusted = {str(ipaddress.ip_address(item)) for item in configured}
  79. remote = str(ipaddress.ip_address(request.remote_addr or ""))
  80. except ValueError as exc:
  81. raise EdgeGatewayAuthenticationError(
  82. "request did not originate from a trusted mTLS proxy"
  83. ) from exc
  84. if remote not in trusted:
  85. raise EdgeGatewayAuthenticationError(
  86. "request did not originate from a trusted mTLS proxy"
  87. )
  88. if request.headers.get("X-DataOps-Edge-Client-Verify") != "SUCCESS":
  89. raise EdgeGatewayAuthenticationError(
  90. "verified TLS client certificate is required"
  91. )
  92. escaped_certificate = request.headers.get("X-DataOps-Edge-Client-Cert", "")
  93. if not escaped_certificate or len(escaped_certificate.encode()) > 32_768:
  94. raise EdgeGatewayAuthenticationError(
  95. "verified TLS client certificate is required"
  96. )
  97. try:
  98. certificate = x509.load_pem_x509_certificate(
  99. unquote(escaped_certificate).encode()
  100. )
  101. except ValueError as exc:
  102. raise EdgeGatewayAuthenticationError(
  103. "verified TLS client certificate is invalid"
  104. ) from exc
  105. now = datetime.now(UTC)
  106. if (
  107. certificate.not_valid_before_utc > now
  108. or certificate.not_valid_after_utc <= now
  109. ):
  110. raise EdgeGatewayAuthenticationError(
  111. "verified TLS client certificate is not currently valid"
  112. )
  113. proof = certificate.fingerprint(hashes.SHA256()).hex()
  114. header = request.headers.get("X-Edge-Certificate-SHA256")
  115. if header and (
  116. not re.fullmatch(r"[0-9a-f]{64}", header)
  117. or not hmac.compare_digest(header, proof)
  118. ):
  119. raise EdgeGatewayAuthenticationError(
  120. "TLS client certificate proof does not match"
  121. )
  122. return proof
  123. def _machine_auth(gateway_id: str, payload: dict) -> dict:
  124. return {
  125. "credential": request.headers.get("X-Edge-Credential", ""),
  126. "certificate_sha256": _tls_certificate_sha256(),
  127. "gateway_id": gateway_id,
  128. "environment": payload.get("environment"),
  129. "network_zone": payload.get("network_zone"),
  130. "generation": payload.get("generation"),
  131. }
  132. def _respond(call, created: bool = False):
  133. try:
  134. result = call()
  135. response = jsonify(success(result))
  136. if len(response.get_data()) > MAX_EDGE_RESPONSE_BYTES:
  137. response = jsonify(failed(
  138. "边缘网关响应超过安全上限", code=500,
  139. error={"code": "EDGE_GATEWAY_RESPONSE_LIMIT"},
  140. ))
  141. response.headers["Cache-Control"] = "no-store"
  142. return response, 500
  143. response.headers["Cache-Control"] = "no-store"
  144. return response, 201 if created else 200
  145. except EdgeGatewayError as exc:
  146. logger.warning("edge gateway request rejected: code=%s", exc.code)
  147. response = jsonify(failed("边缘网关请求被拒绝", code=exc.status_code, error={"code": exc.code}))
  148. response.headers["Cache-Control"] = "no-store"
  149. return response, exc.status_code
  150. except Exception:
  151. logger.exception("edge gateway request failed")
  152. response = jsonify(failed("边缘网关操作失败", code=500, error={"code": "EDGE_GATEWAY_INTERNAL"}))
  153. response.headers["Cache-Control"] = "no-store"
  154. return response, 500
  155. @bp.errorhandler(EdgeGatewayError)
  156. def edge_gateway_error(error):
  157. response = jsonify(failed("边缘网关请求被拒绝", code=error.status_code, error={"code": error.code}))
  158. response.headers["Cache-Control"] = "no-store"
  159. return response, error.status_code
  160. @bp.route("/edge/enrollments", methods=["POST"])
  161. def edge_create_enrollment():
  162. payload = _payload()
  163. return _respond(lambda: _service().create_enrollment(
  164. gateway_name=payload.get("gateway_name"), environment=payload.get("environment"),
  165. network_zone=payload.get("network_zone"), policy_digest=payload.get("policy_digest"),
  166. allowed_control_hosts=payload.get("allowed_control_hosts"),
  167. allowed_proxy_hosts=payload.get("allowed_proxy_hosts", []),
  168. expected_certificate_sha256=payload.get("expected_certificate_sha256"),
  169. ttl_seconds=payload.get("ttl_seconds", 600), actor_uid=_actor_uid(),
  170. ), created=True)
  171. @bp.route("/edge/register", methods=["POST"])
  172. def edge_register():
  173. payload = _payload()
  174. return _respond(lambda: _service().register(
  175. enrollment_token=request.headers.get("X-Edge-Enrollment", ""),
  176. certificate_sha256=_tls_certificate_sha256(),
  177. gateway_id=payload.get("gateway_id"), environment=payload.get("environment"),
  178. network_zone=payload.get("network_zone"), policy_digest=payload.get("policy_digest"),
  179. allowed_control_hosts=payload.get("allowed_control_hosts"),
  180. allowed_proxy_hosts=payload.get("allowed_proxy_hosts", []), version=payload.get("version"),
  181. ), created=True)
  182. @bp.route("/edge/gateways", methods=["GET"])
  183. def edge_list_gateways():
  184. return _respond(lambda: {"gateways": _service().list_gateways(
  185. limit=request.args.get("limit", 50), offset=request.args.get("offset", 0),
  186. )})
  187. @bp.route("/edge/gateways/<gateway_id>/heartbeat", methods=["POST"])
  188. def edge_heartbeat(gateway_id):
  189. payload = _payload()
  190. return _respond(lambda: _service().heartbeat(
  191. **_machine_auth(gateway_id, payload), version=payload.get("version"),
  192. safe_summary=payload.get("safe_summary", {}),
  193. ))
  194. @bp.route("/edge/gateways/<gateway_id>/rotate", methods=["POST"])
  195. def edge_rotate(gateway_id):
  196. payload = _payload()
  197. return _respond(lambda: _service().rotate(
  198. gateway_id, certificate_sha256=payload.get("certificate_sha256"), actor_uid=_actor_uid(),
  199. request_id=payload.get("request_id"),
  200. ))
  201. @bp.route("/edge/gateways/<gateway_id>/revoke", methods=["POST"])
  202. def edge_revoke(gateway_id):
  203. return _respond(lambda: {"revoked": _service().revoke(gateway_id, actor_uid=_actor_uid())})
  204. @bp.route("/edge/tasks", methods=["POST"])
  205. def edge_issue_task():
  206. return _respond(lambda: _service().issue_task(_payload(), actor_uid=_actor_uid()), created=True)
  207. @bp.route("/edge/tasks/<task_id>/cancel", methods=["POST"])
  208. def edge_cancel_task(task_id):
  209. return _respond(lambda: {"cancel_requested": _service().cancel_task(task_id, actor_uid=_actor_uid())})
  210. @bp.route("/edge/gateways/<gateway_id>/tasks/pull", methods=["POST"])
  211. def edge_pull_task(gateway_id):
  212. payload = _payload()
  213. return _respond(lambda: _service().pull_task(**_machine_auth(gateway_id, payload)))
  214. @bp.route("/edge/gateways/<gateway_id>/tasks/<task_id>/outcome", methods=["POST"])
  215. def edge_task_outcome(gateway_id, task_id):
  216. payload = _payload()
  217. return _respond(lambda: _service().task_outcome(
  218. task_id, outcome=payload.get("outcome"), lease_token=payload.get("lease_token"),
  219. safe_summary=payload.get("safe_summary", {}), **_machine_auth(gateway_id, payload),
  220. ))
  221. @bp.route("/edge/gateways/<gateway_id>/events", methods=["POST"])
  222. def edge_accept_event(gateway_id):
  223. payload = _payload()
  224. return _respond(lambda: _service().accept_event(
  225. event=payload.get("event"), lease_token=payload.get("lease_token"),
  226. **_machine_auth(gateway_id, payload),
  227. ))
  228. @bp.route("/edge/gateways/<gateway_id>/reconcile", methods=["POST"])
  229. def edge_reconcile(gateway_id):
  230. payload = _payload()
  231. return _respond(lambda: _service().reconcile(
  232. limit=payload.get("limit", 50),
  233. cancel_cursor=payload.get("cancel_cursor"),
  234. release_cursor=payload.get("release_cursor"),
  235. **_machine_auth(gateway_id, payload),
  236. ))
  237. @bp.route("/edge/gateways/<gateway_id>/releases", methods=["POST"])
  238. def edge_offer_release(gateway_id):
  239. payload = _payload()
  240. return _respond(lambda: _service().offer_release(
  241. gateway_id, version=payload.get("version"), artifact_digest=payload.get("artifact_digest"),
  242. artifact_name=payload.get("artifact_name"), deadline_at=payload.get("deadline_at"),
  243. rollback_version=payload.get("rollback_version"), request_id=payload.get("request_id"), actor_uid=_actor_uid(),
  244. ), created=True)
  245. @bp.route("/edge/gateways/<gateway_id>/releases/ack", methods=["POST"])
  246. def edge_ack_release(gateway_id):
  247. payload = _payload()
  248. return _respond(lambda: _service().acknowledge_release(
  249. release_id=payload.get("release_id"), outcome=payload.get("outcome"),
  250. safe_summary=payload.get("safe_summary", {}), **_machine_auth(gateway_id, payload),
  251. ))