"""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//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//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/") 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//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//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)