| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246 |
- """HTTP boundary for SLO monitoring and evidence-bound data incidents."""
- from __future__ import annotations
- import uuid
- from contextlib import suppress
- from flask import g, jsonify, request
- from app import db
- from app.api.data_factory import bp
- from app.models.result import failed, success
- def get_data_observability_service():
- from app.core.events.data_observability import DataObservabilityService
- from app.core.events.data_observability_repository import (
- SqlAlchemyDataObservabilityRepository,
- )
- return DataObservabilityService(
- SqlAlchemyDataObservabilityRepository(db.session),
- commit=db.session.commit,
- rollback=db.session.rollback,
- )
- def get_production_operations_service():
- import time
- from app.core.events.production_operations import ProductionOperationsService
- from app.core.events.production_operations_repository import (
- SqlAlchemyProductionOperationsRepository,
- )
- return ProductionOperationsService(
- SqlAlchemyProductionOperationsRepository(db.session), now=time.time
- )
- def _actor_uid():
- identity = getattr(g, "current_user", {}) or {}
- return identity.get("id") or identity.get("sub")
- def _payload():
- body = request.get_data(cache=True)
- if len(body) > 4096:
- raise ValueError("request body exceeds 4096 bytes")
- value = request.get_json(silent=True)
- if not isinstance(value, dict):
- raise ValueError("request body must be an object")
- return value
- def _no_store(value, status):
- response = jsonify(success(value))
- response.headers["Cache-Control"] = "no-store"
- return response, status
- def _error(error):
- db.session.rollback()
- if isinstance(error, LookupError):
- return jsonify(failed(str(error), code=404)), 404
- if isinstance(error, ValueError):
- return jsonify(failed(str(error), code=400)), 400
- if isinstance(error, RuntimeError):
- return jsonify(failed(str(error), code=409)), 409
- return jsonify(failed("数据可观测与事故请求处理失败", code=500)), 500
- def _production_delivery_error(error, service):
- """Rollback first, then record only a redacted rejection receipt."""
- db.session.rollback()
- message = str(error)
- reason_code = (
- "conflict"
- if isinstance(error, RuntimeError)
- or message in {"idempotency conflict", "delivery queue hard limit reached"}
- else "invalid_request"
- if isinstance(error, ValueError)
- else "internal_error"
- )
- with suppress(Exception):
- service.audit_rejection_independently(
- actor_uid=_actor_uid(),
- reason_code=reason_code,
- correlation_id=str(uuid.uuid4()),
- )
- status = {"invalid_request": 400, "conflict": 409, "internal_error": 500}[reason_code]
- return jsonify(failed(reason_code, code=status)), status
- @bp.after_app_request
- def _no_store_production_delivery_responses(response):
- if request.path == "/api/datafactory/observability/operations/deliveries":
- response.headers["Cache-Control"] = "no-store"
- return response
- @bp.get("/observability/overview")
- def get_observability_overview():
- try:
- return jsonify(success(get_data_observability_service().overview())), 200
- except Exception as error:
- return _error(error)
- @bp.get("/observability/slos")
- def list_slo_policies():
- try:
- return jsonify(success(get_data_observability_service().list_slos())), 200
- except Exception as error:
- return _error(error)
- @bp.post("/observability/slos")
- def create_slo_policy():
- try:
- record = get_data_observability_service().create_slo(
- _payload(), actor_uid=_actor_uid()
- )
- return jsonify(success(record)), 201
- except Exception as error:
- return _error(error)
- @bp.post("/observability/collect")
- def collect_observability_signals():
- try:
- _payload()
- record = get_data_observability_service().collect(
- actor_uid=_actor_uid()
- )
- return jsonify(success(record)), 200
- except Exception as error:
- return _error(error)
- @bp.post("/observability/operations/deliveries")
- def enqueue_production_delivery():
- service = None
- try:
- payload = _payload()
- if set(payload) != {"incident_uid", "alert_uid", "channel", "summary"}:
- raise ValueError("production delivery request has unsupported fields")
- service = get_production_operations_service()
- record = service.enqueue_incident_delivery(
- payload, actor_uid=_actor_uid()
- )
- db.session.commit()
- return _no_store(record, 201)
- except Exception as error:
- return _production_delivery_error(
- error, service or get_production_operations_service()
- )
- @bp.get("/observability/alerts")
- def list_observability_alerts():
- try:
- return jsonify(
- success(get_data_observability_service().list_alerts())
- ), 200
- except Exception as error:
- return _error(error)
- @bp.post("/observability/alerts/<alert_uid>/suppress")
- def suppress_observability_alert(alert_uid):
- try:
- payload = _payload()
- record = get_data_observability_service().suppress_alert(
- alert_uid,
- until=payload.get("until"),
- reason=payload.get("reason"),
- actor_uid=_actor_uid(),
- )
- return jsonify(success(record)), 200
- except Exception as error:
- return _error(error)
- @bp.post("/observability/alerts/<alert_uid>/acknowledge-delivery")
- def acknowledge_observability_alert_delivery(alert_uid):
- try:
- payload = _payload()
- record = get_data_observability_service().acknowledge_delivery(
- alert_uid,
- channel=payload.get("channel"),
- receipt_id=payload.get("receipt_id"),
- actor_uid=_actor_uid(),
- )
- return jsonify(success(record)), 200
- except Exception as error:
- return _error(error)
- @bp.get("/observability/incidents")
- def list_data_incidents():
- try:
- return jsonify(
- success(get_data_observability_service().list_incidents())
- ), 200
- except Exception as error:
- return _error(error)
- @bp.get("/observability/incidents/<incident_uid>")
- def get_data_incident(incident_uid):
- try:
- return jsonify(
- success(
- get_data_observability_service().incident_detail(incident_uid)
- )
- ), 200
- except Exception as error:
- return _error(error)
- @bp.post("/observability/incidents/<incident_uid>/escalate")
- def escalate_data_incident(incident_uid):
- try:
- payload = _payload()
- record = get_data_observability_service().escalate_incident(
- incident_uid,
- reason=payload.get("reason"),
- actor_uid=_actor_uid(),
- )
- return jsonify(success(record)), 200
- except Exception as error:
- return _error(error)
- @bp.post("/observability/incidents/<incident_uid>/close")
- def close_data_incident(incident_uid):
- try:
- record = get_data_observability_service().close_incident(
- incident_uid,
- _payload(),
- actor_uid=_actor_uid(),
- )
- return jsonify(success(record)), 200
- except Exception as error:
- return _error(error)
|