observability_routes.py 7.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246
  1. """HTTP boundary for SLO monitoring and evidence-bound data incidents."""
  2. from __future__ import annotations
  3. import uuid
  4. from contextlib import suppress
  5. from flask import g, jsonify, request
  6. from app import db
  7. from app.api.data_factory import bp
  8. from app.models.result import failed, success
  9. def get_data_observability_service():
  10. from app.core.events.data_observability import DataObservabilityService
  11. from app.core.events.data_observability_repository import (
  12. SqlAlchemyDataObservabilityRepository,
  13. )
  14. return DataObservabilityService(
  15. SqlAlchemyDataObservabilityRepository(db.session),
  16. commit=db.session.commit,
  17. rollback=db.session.rollback,
  18. )
  19. def get_production_operations_service():
  20. import time
  21. from app.core.events.production_operations import ProductionOperationsService
  22. from app.core.events.production_operations_repository import (
  23. SqlAlchemyProductionOperationsRepository,
  24. )
  25. return ProductionOperationsService(
  26. SqlAlchemyProductionOperationsRepository(db.session), now=time.time
  27. )
  28. def _actor_uid():
  29. identity = getattr(g, "current_user", {}) or {}
  30. return identity.get("id") or identity.get("sub")
  31. def _payload():
  32. body = request.get_data(cache=True)
  33. if len(body) > 4096:
  34. raise ValueError("request body exceeds 4096 bytes")
  35. value = request.get_json(silent=True)
  36. if not isinstance(value, dict):
  37. raise ValueError("request body must be an object")
  38. return value
  39. def _no_store(value, status):
  40. response = jsonify(success(value))
  41. response.headers["Cache-Control"] = "no-store"
  42. return response, status
  43. def _error(error):
  44. db.session.rollback()
  45. if isinstance(error, LookupError):
  46. return jsonify(failed(str(error), code=404)), 404
  47. if isinstance(error, ValueError):
  48. return jsonify(failed(str(error), code=400)), 400
  49. if isinstance(error, RuntimeError):
  50. return jsonify(failed(str(error), code=409)), 409
  51. return jsonify(failed("数据可观测与事故请求处理失败", code=500)), 500
  52. def _production_delivery_error(error, service):
  53. """Rollback first, then record only a redacted rejection receipt."""
  54. db.session.rollback()
  55. message = str(error)
  56. reason_code = (
  57. "conflict"
  58. if isinstance(error, RuntimeError)
  59. or message in {"idempotency conflict", "delivery queue hard limit reached"}
  60. else "invalid_request"
  61. if isinstance(error, ValueError)
  62. else "internal_error"
  63. )
  64. with suppress(Exception):
  65. service.audit_rejection_independently(
  66. actor_uid=_actor_uid(),
  67. reason_code=reason_code,
  68. correlation_id=str(uuid.uuid4()),
  69. )
  70. status = {"invalid_request": 400, "conflict": 409, "internal_error": 500}[reason_code]
  71. return jsonify(failed(reason_code, code=status)), status
  72. @bp.after_app_request
  73. def _no_store_production_delivery_responses(response):
  74. if request.path == "/api/datafactory/observability/operations/deliveries":
  75. response.headers["Cache-Control"] = "no-store"
  76. return response
  77. @bp.get("/observability/overview")
  78. def get_observability_overview():
  79. try:
  80. return jsonify(success(get_data_observability_service().overview())), 200
  81. except Exception as error:
  82. return _error(error)
  83. @bp.get("/observability/slos")
  84. def list_slo_policies():
  85. try:
  86. return jsonify(success(get_data_observability_service().list_slos())), 200
  87. except Exception as error:
  88. return _error(error)
  89. @bp.post("/observability/slos")
  90. def create_slo_policy():
  91. try:
  92. record = get_data_observability_service().create_slo(
  93. _payload(), actor_uid=_actor_uid()
  94. )
  95. return jsonify(success(record)), 201
  96. except Exception as error:
  97. return _error(error)
  98. @bp.post("/observability/collect")
  99. def collect_observability_signals():
  100. try:
  101. _payload()
  102. record = get_data_observability_service().collect(
  103. actor_uid=_actor_uid()
  104. )
  105. return jsonify(success(record)), 200
  106. except Exception as error:
  107. return _error(error)
  108. @bp.post("/observability/operations/deliveries")
  109. def enqueue_production_delivery():
  110. service = None
  111. try:
  112. payload = _payload()
  113. if set(payload) != {"incident_uid", "alert_uid", "channel", "summary"}:
  114. raise ValueError("production delivery request has unsupported fields")
  115. service = get_production_operations_service()
  116. record = service.enqueue_incident_delivery(
  117. payload, actor_uid=_actor_uid()
  118. )
  119. db.session.commit()
  120. return _no_store(record, 201)
  121. except Exception as error:
  122. return _production_delivery_error(
  123. error, service or get_production_operations_service()
  124. )
  125. @bp.get("/observability/alerts")
  126. def list_observability_alerts():
  127. try:
  128. return jsonify(
  129. success(get_data_observability_service().list_alerts())
  130. ), 200
  131. except Exception as error:
  132. return _error(error)
  133. @bp.post("/observability/alerts/<alert_uid>/suppress")
  134. def suppress_observability_alert(alert_uid):
  135. try:
  136. payload = _payload()
  137. record = get_data_observability_service().suppress_alert(
  138. alert_uid,
  139. until=payload.get("until"),
  140. reason=payload.get("reason"),
  141. actor_uid=_actor_uid(),
  142. )
  143. return jsonify(success(record)), 200
  144. except Exception as error:
  145. return _error(error)
  146. @bp.post("/observability/alerts/<alert_uid>/acknowledge-delivery")
  147. def acknowledge_observability_alert_delivery(alert_uid):
  148. try:
  149. payload = _payload()
  150. record = get_data_observability_service().acknowledge_delivery(
  151. alert_uid,
  152. channel=payload.get("channel"),
  153. receipt_id=payload.get("receipt_id"),
  154. actor_uid=_actor_uid(),
  155. )
  156. return jsonify(success(record)), 200
  157. except Exception as error:
  158. return _error(error)
  159. @bp.get("/observability/incidents")
  160. def list_data_incidents():
  161. try:
  162. return jsonify(
  163. success(get_data_observability_service().list_incidents())
  164. ), 200
  165. except Exception as error:
  166. return _error(error)
  167. @bp.get("/observability/incidents/<incident_uid>")
  168. def get_data_incident(incident_uid):
  169. try:
  170. return jsonify(
  171. success(
  172. get_data_observability_service().incident_detail(incident_uid)
  173. )
  174. ), 200
  175. except Exception as error:
  176. return _error(error)
  177. @bp.post("/observability/incidents/<incident_uid>/escalate")
  178. def escalate_data_incident(incident_uid):
  179. try:
  180. payload = _payload()
  181. record = get_data_observability_service().escalate_incident(
  182. incident_uid,
  183. reason=payload.get("reason"),
  184. actor_uid=_actor_uid(),
  185. )
  186. return jsonify(success(record)), 200
  187. except Exception as error:
  188. return _error(error)
  189. @bp.post("/observability/incidents/<incident_uid>/close")
  190. def close_data_incident(incident_uid):
  191. try:
  192. record = get_data_observability_service().close_incident(
  193. incident_uid,
  194. _payload(),
  195. actor_uid=_actor_uid(),
  196. )
  197. return jsonify(success(record)), 200
  198. except Exception as error:
  199. return _error(error)