diff --git a/docs/design/gimin-#338-observability-incident-drill.md b/docs/design/gimin-#338-observability-incident-drill.md new file mode 100644 index 00000000..d0763e05 --- /dev/null +++ b/docs/design/gimin-#338-observability-incident-drill.md @@ -0,0 +1,129 @@ +# 관측성 파이프라인 실제 장애 주입 검증 설계 + +## 배경 + +DocGrid는 Backend와 Embedding Provider 메트릭, 비동기 Queue Snapshot, Prometheus 경보 규칙과 +Alertmanager 라우팅을 제공한다. 기존 검증은 각 계층의 계약과 합성 시계열 전달을 빠르게 확인했지만, +실제 DB 상태와 Provider 프로세스에서 시작한 신호가 전체 경로를 지나가는지 한 번에 재현하지 않았다. + +또한 Queue Gauge는 Prometheus scrape callback에서 DB를 읽지 않고 전용 Scheduler가 만든 Snapshot을 +사용하도록 설계됐다. 단위 테스트만으로 이 경계를 확인하면 실제 HTTP scrape와 PostgreSQL 사이에 다른 +호출이 추가됐을 가능성을 배제하기 어렵다. + +## 목표와 성공 조건 + +다음 세 시나리오를 하나의 공통 Harness에서 독립적으로 실행한다. + +1. 실제 `/actuator/prometheus` 반복 호출 중 운영 Snapshot aggregate SQL이 증가하지 않는다. +2. 실제 문서 업로드로 생성한 Queue가 운영 규칙의 지속 시간 뒤 firing되고 Worker 재개 후 INDEXED와 + resolved 상태로 수렴한다. +3. 실제 BGE-M3 컨테이너 중단이 root-cause 경보로 전달되고 같은 배포의 종속 warning을 억제하며, + Provider 재시작 뒤 resolved가 전달된다. + +각 시나리오는 다음 공통 조건을 만족해야 PASS다. + +- 고유 Compose project, 동적 host port와 새 PostgreSQL volume을 사용한다. +- 외부 알림 시크릿 대신 실제 HTTP 요청을 받는 로컬 webhook receiver를 사용한다. +- timeout 안에 관측 상태가 나타나지 않으면 실패하고 진단 로그를 보존한다. +- 실행 종료 시 Backend 프로세스, 컨테이너, DB volume과 로컬 build 이미지를 정리한다. +- 환경·설정·UTC 시각·구간별 시간을 JSON으로 남긴다. + +## 구조 + +```text +실제 Backend ──/actuator/prometheus──┐ + ├─ Prometheus ─ Alertmanager ─ local webhook +실제 BGE-M3 ─────────/metrics────────┘ + │ + └─ stop/start 장애 주입 + +PostgreSQL ─ Queue row / pg_stat_statements + ▲ + └─ Observer Backend + Indexing Worker +``` + +공통 Compose에는 PostgreSQL 17 + pgvector, Redis, BGE-M3, Prometheus 3.5.5, Alertmanager 0.33.1과 +timestamp를 기록하는 webhook receiver가 포함된다. Backend는 동일한 실행 jar를 host의 서로 다른 동적 +포트에서 실행한다. Queue 실험의 Observer와 Worker는 같은 DB와 Local Storage를 공유한다. + +## 명령 인터페이스 + +이 작업은 제품 HTTP API를 추가하지 않는다. 저장소 루트에서 다음 명령을 제공한다. + +```bash +./monitoring/drills/run.sh scrape-load [--output-dir PATH] +./monitoring/drills/run.sh queue-recovery [--output-dir PATH] +./monitoring/drills/run.sh provider-outage [--output-dir PATH] +./monitoring/drills/run.sh all [--output-dir PATH] +``` + +`all`은 세 시나리오를 `scrape-load → queue-recovery → provider-outage` 순으로 실행하지만, 각 시나리오는 +자신의 Stack과 데이터를 새로 만든다. 중간 실패 시 뒤 시나리오는 실행하지 않고 실패 JSON과 해당 +Stack의 최근 로그를 남긴다. + +## 시나리오별 경계 + +### Scrape와 DB 부하 + +1. Snapshot 주기를 10분으로 설정해 최초 갱신 뒤 측정 창 안에 예약 갱신이 없도록 한다. +2. `pg_stat_statements`에서 세 aggregate SQL의 최초 호출 합계가 3인지 확인한다. +3. 통계를 reset하고 20개 Thread에서 Management endpoint를 300회 호출한다. +4. 모든 응답에 Queue Gauge가 존재하는지 확인하고 p50·p95·max를 계산한다. +5. 측정 뒤 aggregate SQL 호출 합계와 Snapshot 성공 Counter 변화가 모두 0인지 확인한다. + +이 방식은 Scheduler SQL과 scrape 유발 SQL을 구분한다. Snapshot Scheduler를 제거하거나 mocking하지 +않으며 실제 Backend와 PostgreSQL을 사용한다. + +### Queue 정체와 Worker 복구 + +1. Worker가 꺼진 Observer Backend에서 실제 ADMIN 로그인과 multipart 문서 업로드를 실행한다. +2. 이 문서의 Job만 10분 전 PENDING으로 옮겨 `oldest age > 300s`를 만든다. +3. Queue 실험에 섞일 수 있는 동일 업로드의 Outbox Event만 PROCESSED로 종결한다. +4. 실제 Gauge와 Prometheus target을 확인한 시점부터 운영 규칙의 `for: 1m`, `for: 5m`을 측정한다. +5. 두 firing webhook 뒤 Worker Backend를 시작한다. +6. Job `INDEXED`, 저장된 Embedding 양수, claimable Gauge 0과 두 resolved webhook을 확인한다. + +### Provider 장애와 억제 + +1. BGE-M3 readiness와 Prometheus `up=1`을 확인한 뒤 컨테이너를 stop한다. +2. `up=0`, pending, 운영 `for: 1m` 이후 firing과 critical webhook을 순서대로 기록한다. +3. 실제 root-cause alert가 Alertmanager에 존재할 때 동일 cluster/environment의 종속 warning 하나를 + Alertmanager API로 주입한다. +4. warning의 상태가 `suppressed`이고 warning용 `group_wait` 30초보다 긴 35초 동안 firing webhook이 + 0건인지 확인한다. +5. Provider를 start하고 Docker가 다시 할당한 동적 host port를 조회한다. +6. readiness, Prometheus alert 해제와 resolved webhook을 확인한다. + +Provider 중단과 root-cause 경보는 실제다. 종속 warning 주입은 inhibition 자체를 다른 Queue 대기시간과 +분리하기 위한 제한된 합성 입력이며 결과 문서에 이 경계를 명시한다. + +## 시간 설정 + +Prometheus의 scrape/evaluation 15초, 모든 운영 규칙의 `for`, Alertmanager `group_wait`은 제품 설정과 +같다. `group_interval`만 resolved 전달 실험을 반복 가능한 시간에 끝내기 위해 5분에서 30초로 줄인다. +따라서 firing 측정은 운영 경로를 그대로 나타내지만 resolved webhook 시간은 테스트 설정 결과다. + +## 보안과 데이터 경계 + +- 고정된 로컬 DB 암호와 JWT secret은 격리 Stack에서만 사용한다. +- ADMIN access token은 Python 프로세스 메모리 안에서만 사용하며 command argument와 로그에 쓰지 않는다. +- Slack·Discord·SMTP 자격 증명은 요구하지 않는다. +- 실험용 문서와 DB는 종료 시 삭제한다. +- BGE-M3 모델 cache만 외부 Docker volume에 유지해 반복 다운로드를 막는다. + +## 검증 자동화 + +기존 `monitoring/verify.sh`는 장시간 장애 주입을 CI마다 반복하지 않는다. 대신 다음을 빠르게 검사한다. + +- Drill Alertmanager 설정을 고정 버전 `amtool`로 파싱 +- Python 실행기 AST 파싱과 Shell 문법 +- 동적 mount를 포함한 Drill Compose 렌더링 +- 기존 Prometheus 규칙·rule test·Alertmanager 전달 E2E + +실제 장시간 실험은 명시적 명령으로 실행하고 결과 JSON을 `docs/test-results/evidence/issue-338/`에 +보존한다. + +Alertmanager 전달 E2E의 종속 warning은 root-cause가 Alertmanager에 먼저 등록된 뒤 firing되게 해, +동시 도착 순서가 inhibition 검증 결과를 바꾸지 않도록 한다. + +closes #338 diff --git a/docs/test-results/evidence/issue-338/provider-outage-webhooks.jsonl b/docs/test-results/evidence/issue-338/provider-outage-webhooks.jsonl new file mode 100644 index 00000000..a79dbe5a --- /dev/null +++ b/docs/test-results/evidence/issue-338/provider-outage-webhooks.jsonl @@ -0,0 +1,2 @@ +{"receivedAt":"2026-09-13T12:39:02.781656+00:00","payload":{"receiver":"local-webhook","status":"firing","alerts":[{"status":"firing","labels":{"alertname":"EmbeddingProviderDown","cluster":"docgrid-drill","environment":"local-drill","instance":"embedding-server:8000","job":"embedding-provider","service":"embedding-provider","severity":"critical"},"annotations":{"description":"Prometheus has failed to scrape the Embedding Provider for one minute.","summary":"Embedding Provider metric endpoint is unavailable"},"startsAt":"2026-09-13T12:38:52.76Z","endsAt":"0001-01-01T00:00:00Z","generatorURL":"http://34d39fe91415:9090/graph?g0.expr=up%7Bjob%3D%22embedding-provider%22%7D+%3D%3D+0&g0.tab=1","fingerprint":"dcc4408eec1a5b0a"}],"notification_reason":"first notification","groupLabels":{"alertname":"EmbeddingProviderDown","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-provider","severity":"critical"},"commonLabels":{"alertname":"EmbeddingProviderDown","cluster":"docgrid-drill","environment":"local-drill","instance":"embedding-server:8000","job":"embedding-provider","service":"embedding-provider","severity":"critical"},"commonAnnotations":{"description":"Prometheus has failed to scrape the Embedding Provider for one minute.","summary":"Embedding Provider metric endpoint is unavailable"},"externalURL":"http://d444d44ceae8:9093","version":"4","groupKey":"{}/{severity=\"critical\"}:{alertname=\"EmbeddingProviderDown\", cluster=\"docgrid-drill\", environment=\"local-drill\", service=\"embedding-provider\", severity=\"critical\"}","truncatedAlerts":0}} +{"receivedAt":"2026-09-13T12:40:32.789920+00:00","payload":{"receiver":"local-webhook","status":"resolved","alerts":[{"status":"resolved","labels":{"alertname":"EmbeddingProviderDown","cluster":"docgrid-drill","environment":"local-drill","instance":"embedding-server:8000","job":"embedding-provider","service":"embedding-provider","severity":"critical"},"annotations":{"description":"Prometheus has failed to scrape the Embedding Provider for one minute.","summary":"Embedding Provider metric endpoint is unavailable"},"startsAt":"2026-09-13T12:38:52.76Z","endsAt":"2026-09-13T12:40:07.76Z","generatorURL":"http://34d39fe91415:9090/graph?g0.expr=up%7Bjob%3D%22embedding-provider%22%7D+%3D%3D+0&g0.tab=1","fingerprint":"dcc4408eec1a5b0a"}],"notification_reason":"all alerts resolved","groupLabels":{"alertname":"EmbeddingProviderDown","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-provider","severity":"critical"},"commonLabels":{"alertname":"EmbeddingProviderDown","cluster":"docgrid-drill","environment":"local-drill","instance":"embedding-server:8000","job":"embedding-provider","service":"embedding-provider","severity":"critical"},"commonAnnotations":{"description":"Prometheus has failed to scrape the Embedding Provider for one minute.","summary":"Embedding Provider metric endpoint is unavailable"},"externalURL":"http://d444d44ceae8:9093","version":"4","groupKey":"{}/{severity=\"critical\"}:{alertname=\"EmbeddingProviderDown\", cluster=\"docgrid-drill\", environment=\"local-drill\", service=\"embedding-provider\", severity=\"critical\"}","truncatedAlerts":0}} diff --git a/docs/test-results/evidence/issue-338/provider-outage.json b/docs/test-results/evidence/issue-338/provider-outage.json new file mode 100644 index 00000000..402fa15d --- /dev/null +++ b/docs/test-results/evidence/issue-338/provider-outage.json @@ -0,0 +1,42 @@ +{ + "scenario": "provider-outage", + "status": "PASS", + "measuredAt": "2026-09-13T12:40:32.921836+00:00", + "gitCommit": "9f92f04c638bc7281911886a16c06a0bb244a3da", + "host": { + "system": "Darwin", + "machine": "arm64" + }, + "images": { + "postgres": "pgvector/pgvector:0.8.1-pg17", + "prometheus": "prom/prometheus:v3.5.5", + "alertmanager": "quay.io/prometheus/alertmanager:v0.33.1" + }, + "conditions": { + "scrapeInterval": "15s", + "evaluationInterval": "15s", + "productionRuleFor": "1m", + "criticalGroupWait": "10s", + "testGroupInterval": "30s", + "productionGroupInterval": "5m" + }, + "timestamps": { + "providerStoppedAt": "2026-09-13T12:37:23.423467+00:00", + "upZeroObservedAt": "2026-09-13T12:37:38.950811+00:00", + "alertPendingAt": "2026-09-13T12:37:53.042109+00:00", + "alertFiringAt": "2026-09-13T12:38:53.234353+00:00", + "firingWebhookReceivedAt": "2026-09-13T12:39:02.781656+00:00", + "recoveryStartedAt": "2026-09-13T12:39:38.312330+00:00", + "providerReadyAt": "2026-09-13T12:39:43.637544+00:00", + "resolvedWebhookReceivedAt": "2026-09-13T12:40:32.789920+00:00" + }, + "measurements": { + "stopToConditionSeconds": 15.527, + "stopToFiringSeconds": 89.811, + "stopToWebhookSeconds": 99.358, + "restartToReadySeconds": 5.325, + "restartToResolvedWebhookSeconds": 54.478, + "dependentWarningInhibited": true, + "dependentWarningWebhookDeliveries": 0 + } +} diff --git a/docs/test-results/evidence/issue-338/queue-recovery-webhooks.jsonl b/docs/test-results/evidence/issue-338/queue-recovery-webhooks.jsonl new file mode 100644 index 00000000..e6d36d74 --- /dev/null +++ b/docs/test-results/evidence/issue-338/queue-recovery-webhooks.jsonl @@ -0,0 +1,4 @@ +{"receivedAt":"2026-09-13T12:31:41.284908+00:00","payload":{"receiver":"local-webhook","status":"firing","alerts":[{"status":"firing","labels":{"alertname":"DocGridEmbeddingWorkersUnavailable","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"critical"},"annotations":{"description":"Claimable embedding jobs have been waiting without a live worker for one minute.","summary":"DocGrid embedding queue has no active workers"},"startsAt":"2026-09-13T12:31:31.265Z","endsAt":"0001-01-01T00:00:00Z","generatorURL":"http://9e014350d14a:9090/graph?g0.expr=max+by+%28cluster%2C+environment%29+%28docgrid_embedding_claimable_jobs%29+%3E+0+and+max+by+%28cluster%2C+environment%29+%28docgrid_embedding_active_workers%29+%3D%3D+0&g0.tab=1","fingerprint":"66839df19bb4ff97"}],"notification_reason":"first notification","groupLabels":{"alertname":"DocGridEmbeddingWorkersUnavailable","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"critical"},"commonLabels":{"alertname":"DocGridEmbeddingWorkersUnavailable","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"critical"},"commonAnnotations":{"description":"Claimable embedding jobs have been waiting without a live worker for one minute.","summary":"DocGrid embedding queue has no active workers"},"externalURL":"http://3cf7ddb745e3:9093","version":"4","groupKey":"{}/{severity=\"critical\"}:{alertname=\"DocGridEmbeddingWorkersUnavailable\", cluster=\"docgrid-drill\", environment=\"local-drill\", service=\"embedding-worker\", severity=\"critical\"}","truncatedAlerts":0}} +{"receivedAt":"2026-09-13T12:36:01.278503+00:00","payload":{"receiver":"local-webhook","status":"firing","alerts":[{"status":"firing","labels":{"alertname":"DocGridEmbeddingQueueStalled","cluster":"docgrid-drill","dependency":"embedding-provider","environment":"local-drill","service":"embedding-worker","severity":"warning"},"annotations":{"description":"The oldest claimable embedding job has remained eligible for more than five minutes.","summary":"DocGrid embedding queue is stalled"},"startsAt":"2026-09-13T12:35:31.265Z","endsAt":"0001-01-01T00:00:00Z","generatorURL":"http://9e014350d14a:9090/graph?g0.expr=max+by+%28cluster%2C+environment%29+%28docgrid_embedding_claimable_jobs%29+%3E+0+and+max+by+%28cluster%2C+environment%29+%28docgrid_embedding_oldest_claimable_age_seconds%29+%3E+300&g0.tab=1","fingerprint":"f4a15170e154fd71"}],"notification_reason":"first notification","groupLabels":{"alertname":"DocGridEmbeddingQueueStalled","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"warning"},"commonLabels":{"alertname":"DocGridEmbeddingQueueStalled","cluster":"docgrid-drill","dependency":"embedding-provider","environment":"local-drill","service":"embedding-worker","severity":"warning"},"commonAnnotations":{"description":"The oldest claimable embedding job has remained eligible for more than five minutes.","summary":"DocGrid embedding queue is stalled"},"externalURL":"http://3cf7ddb745e3:9093","version":"4","groupKey":"{}:{alertname=\"DocGridEmbeddingQueueStalled\", cluster=\"docgrid-drill\", environment=\"local-drill\", service=\"embedding-worker\", severity=\"warning\"}","truncatedAlerts":0}} +{"receivedAt":"2026-09-13T12:36:41.302049+00:00","payload":{"receiver":"local-webhook","status":"resolved","alerts":[{"status":"resolved","labels":{"alertname":"DocGridEmbeddingWorkersUnavailable","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"critical"},"annotations":{"description":"Claimable embedding jobs have been waiting without a live worker for one minute.","summary":"DocGrid embedding queue has no active workers"},"startsAt":"2026-09-13T12:31:31.265Z","endsAt":"2026-09-13T12:36:31.265Z","generatorURL":"http://9e014350d14a:9090/graph?g0.expr=max+by+%28cluster%2C+environment%29+%28docgrid_embedding_claimable_jobs%29+%3E+0+and+max+by+%28cluster%2C+environment%29+%28docgrid_embedding_active_workers%29+%3D%3D+0&g0.tab=1","fingerprint":"66839df19bb4ff97"}],"notification_reason":"all alerts resolved","groupLabels":{"alertname":"DocGridEmbeddingWorkersUnavailable","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"critical"},"commonLabels":{"alertname":"DocGridEmbeddingWorkersUnavailable","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"critical"},"commonAnnotations":{"description":"Claimable embedding jobs have been waiting without a live worker for one minute.","summary":"DocGrid embedding queue has no active workers"},"externalURL":"http://3cf7ddb745e3:9093","version":"4","groupKey":"{}/{severity=\"critical\"}:{alertname=\"DocGridEmbeddingWorkersUnavailable\", cluster=\"docgrid-drill\", environment=\"local-drill\", service=\"embedding-worker\", severity=\"critical\"}","truncatedAlerts":0}} +{"receivedAt":"2026-09-13T12:37:01.280511+00:00","payload":{"receiver":"local-webhook","status":"resolved","alerts":[{"status":"resolved","labels":{"alertname":"DocGridEmbeddingQueueStalled","cluster":"docgrid-drill","dependency":"embedding-provider","environment":"local-drill","service":"embedding-worker","severity":"warning"},"annotations":{"description":"The oldest claimable embedding job has remained eligible for more than five minutes.","summary":"DocGrid embedding queue is stalled"},"startsAt":"2026-09-13T12:35:31.265Z","endsAt":"2026-09-13T12:36:31.265Z","generatorURL":"http://9e014350d14a:9090/graph?g0.expr=max+by+%28cluster%2C+environment%29+%28docgrid_embedding_claimable_jobs%29+%3E+0+and+max+by+%28cluster%2C+environment%29+%28docgrid_embedding_oldest_claimable_age_seconds%29+%3E+300&g0.tab=1","fingerprint":"f4a15170e154fd71"}],"notification_reason":"all alerts resolved","groupLabels":{"alertname":"DocGridEmbeddingQueueStalled","cluster":"docgrid-drill","environment":"local-drill","service":"embedding-worker","severity":"warning"},"commonLabels":{"alertname":"DocGridEmbeddingQueueStalled","cluster":"docgrid-drill","dependency":"embedding-provider","environment":"local-drill","service":"embedding-worker","severity":"warning"},"commonAnnotations":{"description":"The oldest claimable embedding job has remained eligible for more than five minutes.","summary":"DocGrid embedding queue is stalled"},"externalURL":"http://3cf7ddb745e3:9093","version":"4","groupKey":"{}:{alertname=\"DocGridEmbeddingQueueStalled\", cluster=\"docgrid-drill\", environment=\"local-drill\", service=\"embedding-worker\", severity=\"warning\"}","truncatedAlerts":0}} diff --git a/docs/test-results/evidence/issue-338/queue-recovery.json b/docs/test-results/evidence/issue-338/queue-recovery.json new file mode 100644 index 00000000..648f912b --- /dev/null +++ b/docs/test-results/evidence/issue-338/queue-recovery.json @@ -0,0 +1,47 @@ +{ + "scenario": "queue-recovery", + "status": "PASS", + "measuredAt": "2026-09-13T12:37:01.517068+00:00", + "gitCommit": "9f92f04c638bc7281911886a16c06a0bb244a3da", + "host": { + "system": "Darwin", + "machine": "arm64" + }, + "images": { + "postgres": "pgvector/pgvector:0.8.1-pg17", + "prometheus": "prom/prometheus:v3.5.5", + "alertmanager": "quay.io/prometheus/alertmanager:v0.33.1" + }, + "conditions": { + "queueAgeAtStart": ">10m", + "scrapeInterval": "15s", + "evaluationInterval": "15s", + "workersUnavailableRuleFor": "1m", + "queueStalledRuleFor": "5m", + "criticalGroupWait": "10s", + "warningGroupWait": "30s", + "testGroupInterval": "30s", + "productionGroupInterval": "5m" + }, + "timestamps": { + "conditionObservedAt": "2026-09-13T12:30:20.850153+00:00", + "workersUnavailableFiringAt": "2026-09-13T12:31:33.077816+00:00", + "workersUnavailableWebhookAt": "2026-09-13T12:31:41.284908+00:00", + "queueStalledFiringAt": "2026-09-13T12:35:32.444496+00:00", + "queueStalledWebhookAt": "2026-09-13T12:36:01.278503+00:00", + "recoveryStartedAt": "2026-09-13T12:36:01.565641+00:00", + "jobIndexedAt": "2026-09-13T12:36:10.164030+00:00", + "workersUnavailableResolvedWebhookAt": "2026-09-13T12:36:41.302049+00:00", + "queueStalledResolvedWebhookAt": "2026-09-13T12:37:01.280511+00:00" + }, + "measurements": { + "conditionToWorkersUnavailableFiringSeconds": 72.228, + "conditionToWorkersUnavailableWebhookSeconds": 80.435, + "conditionToQueueStalledFiringSeconds": 311.594, + "conditionToQueueStalledWebhookSeconds": 340.428, + "workerStartToIndexedSeconds": 8.598, + "workerStartToAllResolvedWebhooksSeconds": 59.715, + "storedEmbeddings": 1, + "finalJobStatus": "INDEXED" + } +} diff --git a/docs/test-results/evidence/issue-338/scrape-load.json b/docs/test-results/evidence/issue-338/scrape-load.json new file mode 100644 index 00000000..77900cfe --- /dev/null +++ b/docs/test-results/evidence/issue-338/scrape-load.json @@ -0,0 +1,31 @@ +{ + "scenario": "scrape-load", + "status": "PASS", + "measuredAt": "2026-09-13T12:29:48.309670+00:00", + "gitCommit": "9f92f04c638bc7281911886a16c06a0bb244a3da", + "host": { + "system": "Darwin", + "machine": "arm64" + }, + "images": { + "postgres": "pgvector/pgvector:0.8.1-pg17", + "prometheus": "prom/prometheus:v3.5.5", + "alertmanager": "quay.io/prometheus/alertmanager:v0.33.1" + }, + "conditions": { + "requests": 300, + "concurrency": 20, + "snapshotInterval": "10m", + "endpoint": "/actuator/prometheus" + }, + "measurements": { + "initialSnapshotAggregateSqlCalls": 3, + "aggregateSqlCallsCausedByScrapes": 0, + "snapshotRefreshesDuringScrapes": 0, + "latencyMilliseconds": { + "p50": 8.404, + "p95": 41.746, + "max": 80.882 + } + } +} diff --git a/docs/test-results/evidence/issue-338/summary.json b/docs/test-results/evidence/issue-338/summary.json new file mode 100644 index 00000000..2800e060 --- /dev/null +++ b/docs/test-results/evidence/issue-338/summary.json @@ -0,0 +1,14 @@ +{ + "status": "PASS", + "measuredAt": "2026-09-13T12:40:35.000016+00:00", + "scenarios": [ + "scrape-load", + "queue-recovery", + "provider-outage" + ], + "resultFiles": [ + "scrape-load.json", + "queue-recovery.json", + "provider-outage.json" + ] +} diff --git a/docs/test-results/gimin-#338-observability-incident-drill.md b/docs/test-results/gimin-#338-observability-incident-drill.md new file mode 100644 index 00000000..4ab56ecc --- /dev/null +++ b/docs/test-results/gimin-#338-observability-incident-drill.md @@ -0,0 +1,251 @@ +# 관측성 파이프라인 실제 장애 주입 검증 결과 + +## 결론 + +2026-09-13에 커밋 `9f92f04c638bc7281911886a16c06a0bb244a3da`의 Harness로 세 시나리오를 +연속 실행했고 모두 PASS했다. + +| 검증 항목 | 결과 | +|---|---:| +| `/actuator/prometheus` 호출 | 300회, 동시성 20 | +| 최초 Snapshot aggregate SQL | 3회 | +| 300회 scrape가 추가한 aggregate SQL | **0회** | +| scrape 응답시간 | p50 8.404ms · p95 41.746ms · max 80.882ms | +| Worker 부재 condition → firing | 72.228초 | +| Worker 부재 condition → webhook | 80.435초 | +| Queue 정체 condition → firing | 311.594초 | +| Queue 정체 condition → webhook | 340.428초 | +| Worker 시작 → 실제 Job `INDEXED` | 8.598초 | +| Worker 시작 → 두 resolved webhook 완료 | 59.715초 | +| Provider stop → `up=0` | 15.527초 | +| Provider stop → firing | 89.811초 | +| Provider stop → webhook | 99.358초 | +| Provider restart → readiness | 5.325초 | +| Provider restart → resolved webhook | 54.478초 | +| 종속 warning 상태 | `suppressed` | +| 종속 warning firing 전달 | **0건** | +| Backend 회귀 테스트 | **1,171개 성공, 실패·건너뜀 0개** | +| Alertmanager routing E2E | firing 2초 · resolved 2초 | + +## 실행 환경 + +| 항목 | 값 | +|---|---| +| Host | macOS Darwin 25.6.0, arm64 | +| Java | Eclipse Temurin 17.0.18+8 | +| Docker Engine | 29.4.1 | +| Docker Compose | 5.1.3 | +| PostgreSQL | `pgvector/pgvector:0.8.1-pg17` | +| Prometheus | `prom/prometheus:v3.5.5` | +| Alertmanager | `quay.io/prometheus/alertmanager:v0.33.1` | +| BGE-M3 | 저장소 `backend/embedding-server/Dockerfile` build | +| cluster / environment | `docgrid-drill` / `local-drill` | + +모든 시나리오는 새 Compose project, 동적 loopback port와 새 PostgreSQL volume을 사용했다. BGE-M3 +모델 cache만 재사용했다. 최종 실행 뒤 `docker ps`와 drill 이름의 volume 목록이 비어 있는 것을 +확인했다. + +## 실행 명령 + +```bash +./monitoring/drills/run.sh all \ + --output-dir /private/tmp/docgrid-drill-final +``` + +실행기는 Backend `bootJar`를 한 번 만든 뒤 다음 순서로 시나리오를 실행했다. + +```text +scrape-load +→ Stack 삭제 +→ queue-recovery +→ Stack 삭제 +→ provider-outage +→ Stack 삭제 +``` + +최종 원본 결과: + +- [`scrape-load.json`](evidence/issue-338/scrape-load.json) +- [`queue-recovery.json`](evidence/issue-338/queue-recovery.json) +- [`provider-outage.json`](evidence/issue-338/provider-outage.json) +- [`summary.json`](evidence/issue-338/summary.json) +- [`queue-recovery-webhooks.jsonl`](evidence/issue-338/queue-recovery-webhooks.jsonl) +- [`provider-outage-webhooks.jsonl`](evidence/issue-338/provider-outage-webhooks.jsonl) + +## Scrape DB 부하 + +### 목적 + +Queue Gauge callback이 DB를 직접 조회하지 않고 마지막 메모리 Snapshot만 읽는다는 것을 실제 +Management HTTP 요청과 PostgreSQL 통계로 검증했다. + +### 절차 + +1. `shared_preload_libraries=pg_stat_statements`, `pg_stat_statements.track=all`로 PostgreSQL을 실행했다. +2. Worker·Sync Dispatcher를 끄고 Snapshot 주기를 10분으로 설정한 실제 Backend를 시작했다. +3. `docgrid_operational_snapshot_refresh_total{outcome="success"} >= 1`을 기다렸다. +4. 세 Snapshot SQL을 query text의 table과 alias로 식별해 호출 합계가 3임을 확인했다. +5. `pg_stat_statements_reset()`을 실행했다. +6. Python `ThreadPoolExecutor` 20개에서 `/actuator/prometheus`를 300회 호출했다. +7. 모든 응답에 `docgrid_embedding_claimable_jobs`가 존재하는지 확인했다. +8. Snapshot 성공 Counter와 세 SQL 호출 합계를 다시 읽었다. + +### 결과와 판정 + +| 항목 | 결과 | 성공 조건 | +|---|---:|---:| +| 성공 HTTP scrape | 300 | 300 | +| 최초 Snapshot SQL | 3 | 3 | +| scrape 구간 Snapshot 갱신 | 0 | 0 | +| scrape 구간 aggregate SQL | 0 | 0 | +| p50 | 8.404ms | 측정값 기록 | +| p95 | 41.746ms | 측정값 기록 | +| max | 80.882ms | 측정값 기록 | + +Snapshot 주기를 10분으로 둔 이유는 Scheduler 실행과 scrape 유발 SQL을 분리하기 위해서다. 이 결과는 +“DB Snapshot 갱신 비용이 없다”가 아니라 **“scrape 횟수가 Snapshot SQL 횟수를 증가시키지 않는다”**는 +경계를 증명한다. 갱신 한 번의 DB 비용은 별도로 확인한 aggregate SQL 3회다. + +## Queue 정체와 Worker 복구 + +### 목적 + +실제 DB Queue 상태가 Gauge와 운영 Prometheus 경보를 만들고, Worker가 돌아오면 실제 Embedding 저장과 +함께 경보가 해제되는지 검증했다. + +### 절차 + +1. 실제 Backend `/auth/login`에서 격리 DB의 seed ADMIN으로 로그인했다. +2. 실제 multipart `/api/documents` 요청으로 TXT 문서를 업로드해 PENDING Embedding Job을 생성했다. +3. Worker는 끈 채 해당 Job의 `created_at`만 10분 전으로 옮겼다. +4. 동일 업로드의 Outbox Event는 이 실험의 Sync 경보를 섞지 않도록 PROCESSED로 종결했다. +5. 실제 endpoint에서 claimable 1건, active Worker 0개, oldest age 300초 초과를 확인했다. +6. Prometheus가 Backend를 정상 scrape하고 같은 조건을 읽은 UTC 시각을 시작점으로 기록했다. +7. `DocGridEmbeddingWorkersUnavailable`과 `DocGridEmbeddingQueueStalled`의 firing·webhook을 기다렸다. +8. 같은 DB와 Local Storage를 보는 실제 Worker Backend를 시작했다. +9. Job `INDEXED`, 저장 Embedding 1건, claimable Gauge 0을 확인했다. +10. 두 Prometheus 경보 해제와 두 resolved webhook을 확인했다. + +### 결과와 판정 + +| 경계 | 실측 | 설정과의 관계 | +|---|---:|---| +| Worker 부재 firing | 72.228초 | `for: 1m` + 15초 평가 정렬 | +| Worker 부재 webhook | 80.435초 | firing + critical `group_wait: 10s` | +| Queue 정체 firing | 311.594초 | `for: 5m` + 15초 평가 정렬 | +| Queue 정체 webhook | 340.428초 | firing + warning `group_wait: 30s` | +| Worker → INDEXED | 8.598초 | 실제 BGE-M3 호출 포함 | +| Worker → resolved 완료 | 59.715초 | 15초 평가 + test `group_interval: 30s` | +| 최종 Job 상태 | INDEXED | INDEXED | +| 저장 Embedding | 1건 | 1건 이상 | + +각 firing 시간은 설정된 `for`보다 짧지 않았다. 최대 약 15초의 차이는 scrape/evaluation 주기의 시작점 +정렬에서 발생한다. webhook 시각도 critical과 warning의 서로 다른 group wait 순서를 따랐다. + +## Provider 장애와 경보 억제 + +### 목적 + +실제 BGE-M3 프로세스 장애가 Prometheus와 Alertmanager를 통과하는지, root-cause 경보가 같은 배포의 +파생 warning을 가리고 복구 뒤 resolved를 전달하는지 검증했다. + +### 절차 + +1. BGE-M3 `/health/ready`와 Prometheus `up{job="embedding-provider"}=1`을 확인했다. +2. `docker compose stop embedding-server`로 실제 Provider를 중단했다. +3. `up=0`, `EmbeddingProviderDown` pending·firing, firing webhook 시각을 각각 기록했다. +4. root-cause 경보가 활성화된 상태에서 같은 `cluster/environment`, + `severity="warning"`, `dependency="embedding-provider"`인 종속 경보를 Alertmanager API에 넣었다. +5. Alertmanager의 종속 경보 상태가 `suppressed`인지 확인했다. +6. warning `group_wait: 30s`보다 긴 35초 뒤에도 종속 firing webhook이 0건인지 확인했다. +7. 종속 입력을 종료하고 실제 Provider 컨테이너를 시작했다. +8. 재할당된 host port를 다시 읽어 readiness, Prometheus 해제와 resolved webhook을 확인했다. + +### 결과와 판정 + +| 경계 | 실측 | 성공 조건 | +|---|---:|---| +| Provider stop → `up=0` | 15.527초 | 60초 이내 | +| Provider stop → firing | 89.811초 | 운영 `for: 1m` 통과 | +| Provider stop → webhook | 99.358초 | firing 뒤 critical group wait 통과 | +| restart → readiness | 5.325초 | 900초 이내 | +| restart → resolved webhook | 54.478초 | 90초 이내(test 설정) | +| 종속 warning | suppressed | suppressed | +| 종속 firing 전달 | 0건 | 0건 | + +Provider 중단부터 firing까지의 89.811초에는 첫 실패 scrape 대기, 1분 `for`와 평가 정렬이 포함된다. +firing부터 webhook까지 약 9.5초는 critical `group_wait: 10s`와 일치한다. + +## 운영값과 해석 한계 + +- Prometheus 15초 scrape/evaluation, 모든 `for`, warning 30초·critical 10초 group wait은 운영값이다. +- resolved 실험의 Alertmanager `group_interval`은 운영 5분 대신 30초다. 따라서 54.478초와 59.715초는 + 테스트 수치이며 운영 resolved의 최악 시간으로 사용할 수 없다. +- Provider 실험의 root-cause는 실제 BGE-M3 중단에서 발생했다. 종속 warning 하나만 inhibition을 Queue + 대기시간과 분리하기 위해 Alertmanager API로 주입했다. +- HTTP 지연은 단일 arm64 개발 장비의 로컬 loopback 결과다. 절대 성능 목표나 운영 서버 SLA로 + 일반화하지 않고, 같은 조건의 회귀 비교 기준으로 사용한다. +- 이 실험은 실제 Slack·Discord 사업자 응답을 검증하지 않는다. Alertmanager가 실제 HTTP webhook + payload를 보냈고 로컬 receiver가 이를 수신한 경계까지 검증한다. +- 직접 비교할 구형 scrape-time DB 조회 구현이 존재하지 않으므로 임의의 before 수치를 만들지 않았다. + 대신 300회 scrape와 aggregate SQL 0회의 인과 경계를 PostgreSQL 통계로 확인했다. + +## 회귀 검증 + +장애 주입 Harness가 기존 Backend 동작이나 Alertmanager 전달 테스트를 깨지 않았는지 별도로 확인했다. +Backend 테스트는 기존 개발 DB를 사용하지 않고 Drill Compose의 PostgreSQL·Redis 서비스만 고유 +Compose project와 새 volume으로 실행했다. Compose가 할당한 loopback port를 `DB_PORT`, +`REDIS_PORT`로 전달했고 테스트 전용 `JWT_SECRET`을 사용했다. + +```bash +DRILL_TMP_DIR=/private/tmp/docgrid-backend-tests-338 \ + docker compose -p docgrid-backend-tests-338 \ + -f monitoring/drills/docker-compose.yml up -d --wait postgres redis + +DB_HOST=127.0.0.1 DB_PORT=<동적 PostgreSQL port> \ +DB_PASSWORD=drill_password \ +REDIS_HOST=127.0.0.1 REDIS_PORT=<동적 Redis port> \ +JWT_SECRET=<테스트 전용 값> \ + ./backend/gradlew -p backend test --no-daemon --rerun-tasks +``` + +JUnit XML 197개 suite를 합산한 결과는 다음과 같다. + +| 항목 | 결과 | +|---|---:| +| 테스트 | 1,171 | +| 성공 | 1,171 | +| 실패 | 0 | +| 오류 | 0 | +| 건너뜀 | 0 | +| JUnit suite 누적 실행시간 | 214.825초 | +| Gradle 전체 실행시간 | 4분 51초 | + +통합 테스트 종료 뒤 테스트 전용 Compose project와 PostgreSQL volume을 삭제했다. + +## 추가 검증 + +```bash +./monitoring/verify.sh --e2e +``` + +다음 항목이 함께 통과했다. + +- Prometheus 설정과 운영 규칙 19개 +- Prometheus rule test 3개 +- 기본·합성 E2E·Drill·채널 예제 6개를 포함한 Alertmanager 설정 9개 +- Drill Python AST와 Shell 문법 +- 기본·monitoring·Drill Compose 렌더링 +- 실제 Prometheus → Alertmanager → webhook grouped firing·resolved 전달: 각 2초 +- Provider root-cause가 먼저 활성화된 뒤 종속 warning이 webhook에 전달되지 않는 inhibition + +테스트 과정에서 Docker의 동적 host port가 Provider `stop/start` 뒤 바뀔 수 있음을 발견했다. 초기 +Harness는 중단 전 포트를 계속 조회해 readiness를 기다렸고, 재시작 직후 port mapping을 다시 읽도록 +수정했다. 수정 후 개별 Provider 실험과 세 시나리오 연속 실행이 모두 통과했다. + +추가 monitoring E2E에서는 root-cause와 종속 warning이 같은 평가 시각에 발생하면 Alertmanager 도착 +순서에 따라 warning이 먼저 전달될 수 있는 테스트 경합도 재현했다. 종속 warning에 테스트 전용 +`for: 2s`를 적용해 root-cause 등록 이후 inhibition을 판정하도록 만들었고, 수정 후 grouped firing, +inhibition, resolved 전달이 모두 통과했다. 운영 경보의 지속 시간이나 routing 설정은 바꾸지 않았다. + +closes #338 diff --git a/monitoring/alertmanager/tests/rules-firing.yml b/monitoring/alertmanager/tests/rules-firing.yml index 79743c0a..d325ba31 100644 --- a/monitoring/alertmanager/tests/rules-firing.yml +++ b/monitoring/alertmanager/tests/rules-firing.yml @@ -27,6 +27,8 @@ groups: - alert: DocGridAlertmanagerE2EDerived expr: vector(1) + # Root-cause를 Alertmanager에 먼저 등록해 동시 도착 순서에 따른 inhibition 경합을 제거한다. + for: 2s labels: cluster: docgrid-e2e dependency: embedding-provider diff --git a/monitoring/drills/README.md b/monitoring/drills/README.md new file mode 100644 index 00000000..ed33b8c4 --- /dev/null +++ b/monitoring/drills/README.md @@ -0,0 +1,92 @@ +# Observability incident drills + +DocGrid의 실제 Backend·PostgreSQL·BGE-M3와 운영 Prometheus 규칙을 연결해 관측성 경로를 +재현한다. 외부 Slack·Discord 주소 대신 격리된 webhook receiver를 사용하며 각 실행은 고유한 +Compose project와 PostgreSQL volume을 만들고 종료할 때 삭제한다. + +## 사전 조건 + +- Docker와 Docker Compose +- Java 17 +- Python 3.9 이상 +- 최초 BGE-M3 실행에 필요한 모델 다운로드 또는 기존 `docgrid_huggingface-cache` Docker volume + +실행기는 캐시 volume이 없으면 생성한다. 모델 파일은 재실행 비용을 줄이기 위해 실험 종료 후에도 +보존하지만, PostgreSQL·Prometheus·webhook 결과와 Backend 프로세스는 매번 정리한다. + +## 실행 + +저장소 루트에서 실행한다. + +```bash +# 300회 동시 scrape와 PostgreSQL 집계 Query 수 +./monitoring/drills/run.sh scrape-load + +# 실제 문서 Job의 Queue 경보와 Worker 복구 +./monitoring/drills/run.sh queue-recovery + +# 실제 BGE-M3 중단·재시작과 종속 경보 억제 +./monitoring/drills/run.sh provider-outage + +# 세 실험을 순서대로 각각 새 Stack에서 실행 +./monitoring/drills/run.sh all +``` + +기본 결과는 `backend/build/reports/observability-drill/<실행시각>/`에 생성된다. 검토할 경로를 +고정하려면 `--output-dir`을 지정한다. + +```bash +./monitoring/drills/run.sh all --output-dir /tmp/docgrid-observability-result +``` + +scrape 횟수와 동시성은 재현 환경에 맞게 조절할 수 있다. + +```bash +DRILL_SCRAPE_REQUESTS=1000 \ +DRILL_SCRAPE_CONCURRENCY=50 \ +./monitoring/drills/run.sh scrape-load +``` + +## 실험 경계 + +### `scrape-load` + +1. `pg_stat_statements`가 활성화된 실제 PostgreSQL과 Worker가 꺼진 Backend를 실행한다. +2. `OperationalMetricsSnapshotRefresher`의 최초 성공과 aggregate SQL 3회를 확인한다. +3. Snapshot 주기를 10분으로 고정하고 통계를 초기화한다. +4. `/actuator/prometheus`를 기본 20개 Thread에서 300회 호출한다. +5. p50·p95·max와 scrape 구간의 Snapshot 갱신 수, aggregate SQL 호출 수를 기록한다. + +Snapshot 갱신을 멈춘 것이 아니라 측정 창보다 긴 정상 설정값을 주입한다. 따라서 결과의 SQL 0회는 +scrape callback이 DB를 호출하지 않았다는 의미이며, 주기적 갱신 자체의 비용은 최초 SQL 3회로 별도 +확인한다. + +### `queue-recovery` + +1. Worker가 꺼진 실제 Backend의 로그인·multipart API로 TXT 문서를 업로드한다. +2. 생성된 Job만 10분 전 대기 상태로 옮기고, 이 실험과 무관한 Outbox Event는 종결한다. +3. 실제 Gauge에서 claimable 1건, active Worker 0개, oldest age 300초 초과를 확인한다. +4. 운영 규칙의 `for: 1m`, `for: 5m`을 그대로 사용해 Worker 부재와 Queue 정체 경보를 기다린다. +5. 같은 DB와 Local Storage를 보는 두 번째 Backend Worker를 시작한다. +6. 실제 BGE-M3 호출, Embedding 저장, Job `INDEXED`, Gauge 회복과 두 resolved webhook을 확인한다. + +### `provider-outage` + +1. BGE-M3 readiness와 Prometheus `up=1`을 먼저 확인한다. +2. 실제 Provider 컨테이너를 중단하고 운영 `EmbeddingProviderDown` 규칙의 pending·firing을 기다린다. +3. 실제 root-cause 경보가 Alertmanager에 있는 동안 동일 배포의 종속 warning 하나를 API로 주입한다. +4. 종속 warning의 `suppressed` 상태와 firing webhook 0건을 함께 확인한다. +5. Provider를 재시작하고 새 동적 host 포트, readiness, Prometheus 해제와 resolved webhook을 확인한다. + +종속 warning만 inhibition 경계를 분리하기 위해 직접 주입한다. Provider 중단·복구와 root-cause 경보는 +실제 BGE-M3와 운영 Prometheus 규칙에서 발생한다. + +## 운영 설정과 다른 값 + +Prometheus의 scrape 15초, evaluation 15초와 모든 `for` 시간, Alertmanager의 warning 30초·critical +10초 `group_wait`은 운영값을 그대로 사용한다. `group_interval`만 resolved 실험 시간을 제한하기 위해 +운영 5분에서 30초로 줄인다. JSON에는 두 값을 모두 기록하므로 테스트 복구 시간과 운영 최악 시간을 +구분할 수 있다. + +각 결과 JSON은 조건, UTC 시각, 구간별 초 단위 측정값과 검증된 최종 상태를 기록한다. 실패하면 같은 +결과 디렉터리에 traceback과 Backend·Compose·webhook 로그를 남기며 Stack은 동일하게 정리한다. diff --git a/monitoring/drills/alertmanager.yml b/monitoring/drills/alertmanager.yml new file mode 100644 index 00000000..d8726f08 --- /dev/null +++ b/monitoring/drills/alertmanager.yml @@ -0,0 +1,29 @@ +global: + resolve_timeout: 5m + +route: + receiver: local-webhook + group_by: [alertname, cluster, environment, severity, service] + group_wait: 30s + # 운영값 5m은 반복 가능한 복구 실험 시간을 과도하게 늘리므로 resolved 관측에서만 30s로 줄인다. + group_interval: 30s + repeat_interval: 4h + routes: + - receiver: local-webhook + matchers: [severity="critical"] + group_wait: 10s + repeat_interval: 1h + +receivers: + - name: local-webhook + webhook_configs: + - url: http://webhook-receiver:8080/alerts + send_resolved: true + +inhibit_rules: + - source_matchers: [severity="critical", service="embedding-provider"] + target_matchers: [severity="warning", service="embedding-provider"] + equal: [cluster, environment] + - source_matchers: [severity="critical", service="embedding-provider"] + target_matchers: [severity="warning", dependency="embedding-provider"] + equal: [cluster, environment] diff --git a/monitoring/drills/docker-compose.yml b/monitoring/drills/docker-compose.yml new file mode 100644 index 00000000..20ce095e --- /dev/null +++ b/monitoring/drills/docker-compose.yml @@ -0,0 +1,113 @@ +services: + postgres: + image: pgvector/pgvector:0.8.1-pg17 + command: + - postgres + - -c + - shared_preload_libraries=pg_stat_statements + - -c + - pg_stat_statements.track=all + environment: + POSTGRES_DB: app + POSTGRES_USER: app + POSTGRES_PASSWORD: drill_password + ports: + - "127.0.0.1::5432" + volumes: + - postgres-data:/var/lib/postgresql/data + - ../../docker/postgres/init/001-enable-vector.sql:/docker-entrypoint-initdb.d/001-enable-vector.sql:ro + - ./init/002-enable-pg-stat-statements.sql:/docker-entrypoint-initdb.d/002-enable-pg-stat-statements.sql:ro + healthcheck: + test: ["CMD-SHELL", "pg_isready -U app -d app"] + interval: 1s + timeout: 2s + retries: 30 + + redis: + image: redis:7-alpine + ports: + - "127.0.0.1::6379" + healthcheck: + test: ["CMD", "redis-cli", "ping"] + interval: 1s + timeout: 1s + retries: 30 + + embedding-server: + build: + context: ../../backend/embedding-server + dockerfile: Dockerfile + environment: + EMBEDDING_PROVIDER_MAX_CONCURRENCY: 1 + EMBEDDING_PROVIDER_MAX_QUEUE_SIZE: 1 + EMBEDDING_PROVIDER_QUEUE_WAIT_TIMEOUT_SECONDS: 15 + ports: + - "127.0.0.1::8000" + volumes: + - huggingface-cache:/root/.cache/huggingface + healthcheck: + test: ["CMD", "python3", "-c", "import urllib.request; urllib.request.urlopen('http://localhost:8000/health/ready')"] + interval: 5s + timeout: 5s + retries: 180 + start_period: 30s + + webhook-receiver: + image: python:3.13-alpine + environment: + EVENT_LOG: /results/events.jsonl + command: ["python", "/drill/webhook_receiver.py"] + volumes: + - ./webhook_receiver.py:/drill/webhook_receiver.py:ro + - ${DRILL_TMP_DIR}/webhook:/results + healthcheck: + test: ["CMD", "python", "-c", "import urllib.request; urllib.request.urlopen('http://localhost:8080/health')"] + interval: 1s + timeout: 1s + retries: 30 + + alertmanager: + image: quay.io/prometheus/alertmanager:v0.33.1 + command: + - "--config.file=/etc/alertmanager/alertmanager.yml" + - "--storage.path=/alertmanager" + - "--cluster.listen-address=" + ports: + - "127.0.0.1::9093" + volumes: + - ./alertmanager.yml:/etc/alertmanager/alertmanager.yml:ro + depends_on: + webhook-receiver: + condition: service_healthy + healthcheck: + test: ["CMD", "wget", "--spider", "--quiet", "http://localhost:9093/-/ready"] + interval: 1s + timeout: 1s + retries: 30 + + prometheus: + image: prom/prometheus:v3.5.5 + command: + - "--config.file=/etc/prometheus/prometheus.yml" + - "--storage.tsdb.path=/prometheus" + ports: + - "127.0.0.1::9090" + extra_hosts: + - "host.docker.internal:host-gateway" + volumes: + - ${DRILL_TMP_DIR}/prometheus.yml:/etc/prometheus/prometheus.yml:ro + - ../prometheus/rules:/etc/prometheus/rules:ro + depends_on: + alertmanager: + condition: service_healthy + healthcheck: + test: ["CMD", "wget", "--spider", "--quiet", "http://localhost:9090/-/ready"] + interval: 1s + timeout: 1s + retries: 30 + +volumes: + postgres-data: + huggingface-cache: + external: true + name: docgrid_huggingface-cache diff --git a/monitoring/drills/init/002-enable-pg-stat-statements.sql b/monitoring/drills/init/002-enable-pg-stat-statements.sql new file mode 100644 index 00000000..841ff0c4 --- /dev/null +++ b/monitoring/drills/init/002-enable-pg-stat-statements.sql @@ -0,0 +1 @@ +CREATE EXTENSION IF NOT EXISTS pg_stat_statements; diff --git a/monitoring/drills/run.sh b/monitoring/drills/run.sh new file mode 100755 index 00000000..0aedceb0 --- /dev/null +++ b/monitoring/drills/run.sh @@ -0,0 +1,5 @@ +#!/bin/sh +set -eu + +SCRIPT_DIR=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) +exec python3 "$SCRIPT_DIR/run_drill.py" "$@" diff --git a/monitoring/drills/run_drill.py b/monitoring/drills/run_drill.py new file mode 100755 index 00000000..f8f3100a --- /dev/null +++ b/monitoring/drills/run_drill.py @@ -0,0 +1,963 @@ +"""Run isolated, measurable incident drills against DocGrid's real monitoring path.""" + +import argparse +import concurrent.futures +import json +import math +import os +import platform +import shutil +import signal +import socket +import subprocess +import sys +import tempfile +import time +import traceback +import uuid +import urllib.error +import urllib.parse +import urllib.request +from datetime import datetime, timedelta, timezone +from pathlib import Path + + +DRILL_DIR = Path(__file__).resolve().parent +ROOT_DIR = DRILL_DIR.parents[1] +COMPOSE_FILE = DRILL_DIR / "docker-compose.yml" +CLUSTER = "docgrid-drill" +ENVIRONMENT = "local-drill" +SNAPSHOT_QUERY_CALLS_SQL = """ +SELECT COALESCE(SUM(calls), 0)::bigint +FROM pg_stat_statements +WHERE query NOT ILIKE '%pg_stat_statements%' + AND ( + (query ILIKE '%FROM embedding_jobs job%' AND query ILIKE '%AS claimable_jobs%') + OR (query ILIKE '%FROM rag_responses response%' AND query ILIKE '%AS oldest_processing_at%') + OR (query ILIKE '%FROM sync_outbox_events event%' AND query ILIKE '%AS claimable_events%') + ) +""" + + +def utc_now(): + """Return an RFC 3339 UTC timestamp suitable for result evidence and Alertmanager.""" + return datetime.now(timezone.utc) + + +def iso(timestamp=None): + """Serialize a UTC timestamp without losing subsecond measurement context.""" + return (timestamp or utc_now()).isoformat() + + +def free_port(): + """Reserve an available loopback port long enough to choose distinct Backend listeners.""" + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as listener: + listener.bind(("127.0.0.1", 0)) + return listener.getsockname()[1] + + +def run(command, *, cwd=ROOT_DIR, env=None, capture=True, check=True): + """Run one external command without a shell so measured values cannot become shell code.""" + return subprocess.run( + [str(part) for part in command], + cwd=cwd, + env=env, + check=check, + text=True, + capture_output=capture, + ) + + +def request(url, *, method="GET", payload=None, headers=None, timeout=10): + """Perform one local HTTP request and return status plus raw response bytes.""" + body = None if payload is None else json.dumps(payload).encode("utf-8") + request_headers = dict(headers or {}) + if body is not None: + request_headers["Content-Type"] = "application/json" + http_request = urllib.request.Request(url, data=body, headers=request_headers, method=method) + with urllib.request.urlopen(http_request, timeout=timeout) as response: + return response.status, response.read() + + +def request_json(url, *, method="GET", payload=None, headers=None, timeout=10): + """Perform one local HTTP request and decode a non-empty JSON body.""" + status, body = request(url, method=method, payload=payload, headers=headers, timeout=timeout) + return status, None if not body else json.loads(body) + + +def wait_for(description, predicate, timeout_seconds, interval_seconds=1): + """Poll an observable boundary with a deadline and periodic progress output.""" + deadline = time.monotonic() + timeout_seconds + next_progress = time.monotonic() + 30 + last_error = None + while time.monotonic() < deadline: + try: + value = predicate() + if value: + return value + except (AssertionError, OSError, ValueError, urllib.error.URLError) as error: + last_error = error + now = time.monotonic() + if now >= next_progress: + print(f"Waiting for {description}...", flush=True) + next_progress = now + 30 + time.sleep(interval_seconds) + detail = f": {last_error}" if last_error else "" + raise TimeoutError(f"Timed out after {timeout_seconds}s waiting for {description}{detail}") + + +def nearest_rank(values, percentile): + """Calculate an explicit nearest-rank percentile without external benchmark libraries.""" + ordered = sorted(values) + index = max(0, math.ceil(percentile * len(ordered)) - 1) + return ordered[index] + + +def build_backend(): + """Build the executable Backend jar once before scenarios that launch real processes.""" + print("Building Backend bootJar...", flush=True) + run(["./backend/gradlew", "-p", "backend", "bootJar", "--no-daemon"], capture=False) + candidates = sorted((ROOT_DIR / "backend" / "build" / "libs").glob("*.jar")) + candidates = [candidate for candidate in candidates if not candidate.name.endswith("-plain.jar")] + if len(candidates) != 1: + raise RuntimeError(f"Expected one executable jar, found: {candidates}") + return candidates[0] + + +class DrillEnvironment: + """Own one isolated Compose project, its Backend processes, ports, logs, and cleanup boundary.""" + + def __init__(self, scenario, output_dir): + self.scenario = scenario + self.output_dir = output_dir + self.output_dir.mkdir(parents=True, exist_ok=True) + self.log_dir = self.output_dir / "logs" / scenario + self.log_dir.mkdir(parents=True, exist_ok=True) + self.temp_dir = Path(tempfile.mkdtemp(prefix=f"docgrid-{scenario}-")) + (self.temp_dir / "webhook").mkdir() + (self.temp_dir / "webhook" / "events.jsonl").touch() + self.project_name = f"docgrid-drill-{scenario.replace('_', '-')}-{os.getpid()}" + self.compose_env = os.environ.copy() + self.compose_env["DRILL_TMP_DIR"] = str(self.temp_dir) + self.backend_processes = [] + self.backend_logs = [] + self.postgres_port = None + self.redis_port = None + self.embedding_port = None + self.prometheus_url = None + self.alertmanager_url = None + + def compose(self, *arguments, capture=True, check=True): + """Address only this scenario's Compose project.""" + return run( + ["docker", "compose", "-p", self.project_name, "-f", COMPOSE_FILE, *arguments], + cwd=DRILL_DIR, + env=self.compose_env, + capture=capture, + check=check, + ) + + def start_core(self, *, include_provider): + """Start fresh database, cache, local receiver, Alertmanager, and optionally BGE-M3.""" + services = ["postgres", "redis", "webhook-receiver", "alertmanager"] + if include_provider: + run(["docker", "volume", "create", "docgrid_huggingface-cache"]) + services.append("embedding-server") + print(f"Starting isolated services for {self.scenario}...", flush=True) + self.compose("up", "-d", "--wait", *services, capture=False) + self.postgres_port = self.service_port("postgres", 5432) + self.redis_port = self.service_port("redis", 6379) + self.alertmanager_url = f"http://127.0.0.1:{self.service_port('alertmanager', 9093)}" + if include_provider: + self.embedding_port = self.service_port("embedding-server", 8000) + + def service_port(self, service, container_port): + """Resolve a dynamic host port assigned to one service.""" + address = self.compose("port", service, str(container_port)).stdout.strip() + return int(address.rsplit(":", 1)[1]) + + def psql(self, sql): + """Execute one SQL statement inside this scenario's PostgreSQL container.""" + result = self.compose( + "exec", "-T", "postgres", "psql", "-v", "ON_ERROR_STOP=1", + "-U", "app", "-d", "app", "-Atc", sql, + ) + return result.stdout.strip() + + def start_backend(self, jar_path, name, *, worker_enabled, snapshot_interval): + """Launch a real Backend process against this scenario's database and shared local files.""" + api_port = free_port() + management_port = free_port() + log_path = self.log_dir / f"backend-{name}.log" + log_file = log_path.open("w", encoding="utf-8") + environment = os.environ.copy() + environment.update({ + "SPRING_PROFILES_ACTIVE": "local", + "SPRING_DEVTOOLS_RESTART_ENABLED": "false", + "DB_HOST": "127.0.0.1", + "DB_PORT": str(self.postgres_port), + "DB_NAME": "app", + "DB_USER": "app", + "DB_PASSWORD": "drill_password", + "DB_SCHEMA": "public", + "DB_SSLMODE": "disable", + "REDIS_HOST": "127.0.0.1", + "REDIS_PORT": str(self.redis_port), + "JWT_SECRET": "docgrid-observability-incident-drill-secret-key-2026", + "JWT_EXPIRATION": "3600", + "MINIO_ENDPOINT": "http://127.0.0.1:65534", + "MINIO_ACCESS_KEY": "unused", + "MINIO_SECRET_KEY": "unused", + "STORAGE_TYPE": "local", + "STORAGE_BUCKET": "docgrid-drill", + "STORAGE_LOCAL_ROOT": str(self.temp_dir / "storage"), + "EMBEDDING_SERVER_URL": ( + f"http://127.0.0.1:{self.embedding_port}" + if self.embedding_port else "http://127.0.0.1:65534" + ), + "INDEXING_WORKER_ENABLED": str(worker_enabled).lower(), + "INDEXING_WORKER_NAME": f"observability-drill-{name}", + "INDEXING_WORKER_POLLING_INTERVAL": "200ms", + "INDEXING_WORKER_IDLE_MAX_POLLING_INTERVAL": "1s", + "INDEXING_WORKER_HEARTBEAT_INTERVAL": "1s", + "INDEXING_WORKER_DEAD_THRESHOLD": "30s", + "INDEXING_WORKER_MAX_CONCURRENCY": "1", + "SYNC_DISPATCHER_ENABLED": "false", + "SYNC_RECONCILIATION_ENABLED": "false", + "MANAGEMENT_METRICS_SNAPSHOT_INTERVAL": snapshot_interval, + "SERVER_PORT": str(api_port), + "MANAGEMENT_PORT": str(management_port), + }) + process = subprocess.Popen( + ["java", "-jar", str(jar_path)], + cwd=ROOT_DIR, + env=environment, + stdout=log_file, + stderr=subprocess.STDOUT, + start_new_session=True, + ) + self.backend_processes.append(process) + self.backend_logs.append(log_file) + readiness_url = f"http://127.0.0.1:{management_port}/actuator/health/readiness" + try: + wait_for( + f"Backend {name} readiness", + lambda: request(readiness_url, timeout=2)[0] == 200, + 180, + ) + except Exception: + log_file.flush() + raise RuntimeError(f"Backend {name} did not become ready; see {log_path}") + return { + "name": name, + "apiPort": api_port, + "managementPort": management_port, + "metricsUrl": f"http://127.0.0.1:{management_port}/actuator/prometheus", + } + + def write_prometheus_config(self, backend_management_port=None): + """Render the production rule files with this scenario's concrete Backend target.""" + scrape_configs = [] + if backend_management_port is not None: + scrape_configs.append(f""" + - job_name: docgrid-backend + metrics_path: /actuator/prometheus + static_configs: + - targets: [\"host.docker.internal:{backend_management_port}\"] + labels: + cluster: {CLUSTER} + environment: {ENVIRONMENT} +""") + scrape_configs.append(""" + - job_name: embedding-provider + metrics_path: /metrics + static_configs: + - targets: ["embedding-server:8000"] + labels: + cluster: docgrid-drill + environment: local-drill +""") + config = """global: + scrape_interval: 15s + evaluation_interval: 15s + +rule_files: + - /etc/prometheus/rules/*.yml + +alerting: + alertmanagers: + - static_configs: + - targets: [\"alertmanager:9093\"] + +scrape_configs: +""" + "".join(scrape_configs) + config_path = self.temp_dir / "prometheus.yml" + config_path.write_text(config, encoding="utf-8") + shutil.copy2(config_path, self.log_dir / "prometheus.yml") + + def start_prometheus(self): + """Start Prometheus after its dynamic host target configuration has been rendered.""" + self.compose("up", "-d", "--wait", "prometheus", capture=False) + self.prometheus_url = f"http://127.0.0.1:{self.service_port('prometheus', 9090)}" + + def prometheus_query(self, expression): + """Return the vector result of one instant Prometheus query.""" + query = urllib.parse.urlencode({"query": expression}) + _, response = request_json(f"{self.prometheus_url}/api/v1/query?{query}") + if response["status"] != "success": + raise AssertionError(response) + return response["data"]["result"] + + def prometheus_value(self, expression): + """Return the first numeric instant-vector sample, or None when the vector is absent.""" + result = self.prometheus_query(expression) + return None if not result else float(result[0]["value"][1]) + + def prometheus_alert(self, alertname, expected_state=None): + """Find a production rule alert and optionally require its pending or firing state.""" + _, response = request_json(f"{self.prometheus_url}/api/v1/alerts") + for alert in response["data"]["alerts"]: + if alert["labels"].get("alertname") != alertname: + continue + if expected_state is None or alert.get("state") == expected_state: + return alert + return None + + def webhook_events(self): + """Read complete timestamped webhook deliveries while tolerating a concurrently written file.""" + path = self.temp_dir / "webhook" / "events.jsonl" + events = [] + for line in path.read_text(encoding="utf-8").splitlines(): + if line.strip(): + events.append(json.loads(line)) + return events + + def webhook_delivery(self, alertname, status): + """Find the first webhook containing one alert with the requested name and lifecycle status.""" + for event in self.webhook_events(): + for alert in event["payload"].get("alerts", []): + if (alert.get("labels", {}).get("alertname") == alertname + and alert.get("status") == status): + return event + return None + + def alertmanager_alert(self, alertname): + """Find one active Alertmanager alert, including its inhibition state.""" + _, alerts = request_json(f"{self.alertmanager_url}/api/v2/alerts") + return next( + (alert for alert in alerts if alert.get("labels", {}).get("alertname") == alertname), + None, + ) + + def post_alertmanager_alert(self, labels, *, starts_at, ends_at): + """Inject one explicitly labeled dependent alert to test inhibition against a real root cause.""" + payload = [{ + "labels": labels, + "annotations": {"summary": "Incident drill dependent warning"}, + "startsAt": iso(starts_at), + "endsAt": iso(ends_at), + "generatorURL": "https://github.com/DocGrid/docgrid/issues/338", + }] + request(f"{self.alertmanager_url}/api/v2/alerts", method="POST", payload=payload) + + def close(self): + """Stop child processes, persist bounded diagnostic evidence, and remove scenario state.""" + for process in reversed(self.backend_processes): + if process.poll() is None: + try: + os.killpg(process.pid, signal.SIGTERM) + process.wait(timeout=20) + except (ProcessLookupError, subprocess.TimeoutExpired): + try: + os.killpg(process.pid, signal.SIGKILL) + except ProcessLookupError: + pass + for log_file in self.backend_logs: + log_file.close() + + event_log = self.temp_dir / "webhook" / "events.jsonl" + if event_log.exists(): + shutil.copy2(event_log, self.log_dir / "webhook-events.jsonl") + compose_logs = self.compose("logs", "--no-color", "--tail", "300", check=False) + (self.log_dir / "compose.log").write_text(compose_logs.stdout + compose_logs.stderr, encoding="utf-8") + self.compose("down", "--volumes", "--remove-orphans", "--rmi", "local", check=False) + shutil.rmtree(self.temp_dir, ignore_errors=True) + + +def metric_value(metrics_text, name, labels=None): + """Extract one Prometheus text sample with an optional exact label subset.""" + required_labels = labels or {} + for line in metrics_text.splitlines(): + if not line.startswith(name): + continue + sample, raw_value = line.rsplit(" ", 1) + if any(f'{key}="{value}"' not in sample for key, value in required_labels.items()): + continue + return float(raw_value) + raise AssertionError(f"Metric {name} with labels {required_labels} was not exposed") + + +def fetch_metrics(backend): + """Read the real Spring Boot Prometheus endpoint as text.""" + _, body = request(backend["metricsUrl"], timeout=10) + return body.decode("utf-8") + + +def login_and_upload(backend, temp_dir): + """Create a real pending Embedding Job through authentication and multipart upload APIs.""" + _, login_response = request_json( + f"http://127.0.0.1:{backend['apiPort']}/auth/login", + method="POST", + payload={"email": "kcw130502@gmail.com", "password": "admin1234"}, + ) + access_token = login_response["data"]["accessToken"] + document_path = temp_dir / "queue-stall.txt" + document_path.write_text( + "A recovered indexing worker must claim this queued document and complete its embedding.", + encoding="utf-8", + ) + boundary = f"docgrid-drill-{uuid.uuid4().hex}" + parts = [] + + def append_field(name, value): + parts.extend([ + f"--{boundary}\r\n".encode(), + f'Content-Disposition: form-data; name="{name}"\r\n\r\n'.encode(), + str(value).encode(), + b"\r\n", + ]) + + parts.extend([ + f"--{boundary}\r\n".encode(), + b'Content-Disposition: form-data; name="file"; filename="queue-stall.txt"\r\n', + b"Content-Type: text/plain\r\n\r\n", + document_path.read_bytes(), + b"\r\n", + ]) + append_field("title", "Observability Queue Recovery Drill") + append_field("description", "Actual queue alert and worker recovery") + append_field("visibility", "PRIVATE") + parts.append(f"--{boundary}--\r\n".encode()) + upload_request = urllib.request.Request( + f"http://127.0.0.1:{backend['apiPort']}/api/documents", + data=b"".join(parts), + headers={ + "Authorization": f"Bearer {access_token}", + "Content-Type": f"multipart/form-data; boundary={boundary}", + }, + method="POST", + ) + with urllib.request.urlopen(upload_request, timeout=30) as response: + payload = json.loads(response.read()) + return payload["data"] + + +def base_result(scenario): + """Record reproducibility metadata shared by all measured scenarios.""" + commit = run(["git", "rev-parse", "HEAD"]).stdout.strip() + return { + "scenario": scenario, + "status": "PASS", + "measuredAt": iso(), + "gitCommit": commit, + "host": {"system": platform.system(), "machine": platform.machine()}, + "images": { + "postgres": "pgvector/pgvector:0.8.1-pg17", + "prometheus": "prom/prometheus:v3.5.5", + "alertmanager": "quay.io/prometheus/alertmanager:v0.33.1", + }, + } + + +def run_scrape_load(output_dir, jar_path): + """Measure endpoint latency and prove repeated scrapes do not execute operational SQL.""" + environment = DrillEnvironment("scrape-load", output_dir) + try: + environment.start_core(include_provider=False) + backend = environment.start_backend( + jar_path, "metrics", worker_enabled=False, snapshot_interval="10m" + ) + wait_for( + "the first operational snapshot", + lambda: metric_value( + fetch_metrics(backend), + "docgrid_operational_snapshot_refresh_total", + {"outcome": "success"}, + ) >= 1, + 30, + ) + + initial_snapshot_calls = int(environment.psql(SNAPSHOT_QUERY_CALLS_SQL)) + if initial_snapshot_calls != 3: + raise AssertionError(f"Expected three initial aggregate queries, got {initial_snapshot_calls}") + refresh_before = metric_value( + fetch_metrics(backend), + "docgrid_operational_snapshot_refresh_total", + {"outcome": "success"}, + ) + environment.psql("SELECT pg_stat_statements_reset()") + + request_count = int(os.environ.get("DRILL_SCRAPE_REQUESTS", "300")) + concurrency = int(os.environ.get("DRILL_SCRAPE_CONCURRENCY", "20")) + + def measured_scrape(_): + started = time.perf_counter() + metrics = fetch_metrics(backend) + metric_value(metrics, "docgrid_embedding_claimable_jobs") + return (time.perf_counter() - started) * 1000 + + print(f"Running {request_count} scrapes with concurrency {concurrency}...", flush=True) + with concurrent.futures.ThreadPoolExecutor(max_workers=concurrency) as executor: + latencies = list(executor.map(measured_scrape, range(request_count))) + + snapshot_calls_during_scrape = int(environment.psql(SNAPSHOT_QUERY_CALLS_SQL)) + refresh_after = metric_value( + fetch_metrics(backend), + "docgrid_operational_snapshot_refresh_total", + {"outcome": "success"}, + ) + if snapshot_calls_during_scrape != 0: + raise AssertionError( + f"Scrapes executed {snapshot_calls_during_scrape} operational aggregate queries" + ) + if refresh_after != refresh_before: + raise AssertionError("Snapshot scheduler refreshed during the isolated scrape measurement") + + result = base_result("scrape-load") + result["conditions"] = { + "requests": request_count, + "concurrency": concurrency, + "snapshotInterval": "10m", + "endpoint": "/actuator/prometheus", + } + result["measurements"] = { + "initialSnapshotAggregateSqlCalls": initial_snapshot_calls, + "aggregateSqlCallsCausedByScrapes": snapshot_calls_during_scrape, + "snapshotRefreshesDuringScrapes": int(refresh_after - refresh_before), + "latencyMilliseconds": { + "p50": round(nearest_rank(latencies, 0.50), 3), + "p95": round(nearest_rank(latencies, 0.95), 3), + "max": round(max(latencies), 3), + }, + } + return result + finally: + environment.close() + + +def run_provider_outage(output_dir): + """Stop and restore the real BGE-M3 container through production Prometheus alert rules.""" + environment = DrillEnvironment("provider-outage", output_dir) + try: + environment.start_core(include_provider=True) + environment.write_prometheus_config() + environment.start_prometheus() + wait_for( + "Prometheus to scrape the ready Embedding Provider", + lambda: environment.prometheus_value('up{job="embedding-provider"}') == 1, + 60, + ) + + print("Stopping the real Embedding Provider container...", flush=True) + stopped_at = utc_now() + environment.compose("stop", "embedding-server", capture=False) + condition_at = wait_for( + "Prometheus up=0 provider condition", + lambda: utc_now() if environment.prometheus_value( + 'up{job="embedding-provider"}' + ) == 0 else None, + 60, + ) + pending_at = wait_for( + "EmbeddingProviderDown pending state", + lambda: utc_now() if environment.prometheus_alert( + "EmbeddingProviderDown", "pending" + ) else None, + 60, + ) + firing_at = wait_for( + "EmbeddingProviderDown firing state", + lambda: utc_now() if environment.prometheus_alert( + "EmbeddingProviderDown", "firing" + ) else None, + 120, + 2, + ) + firing_delivery = wait_for( + "EmbeddingProviderDown firing webhook", + lambda: environment.webhook_delivery("EmbeddingProviderDown", "firing"), + 60, + ) + + # A real Provider root cause is combined with one explicit downstream warning to isolate inhibition. + derived_labels = { + "alertname": "DocGridEmbeddingQueueStalled", + "cluster": CLUSTER, + "environment": ENVIRONMENT, + "severity": "warning", + "service": "embedding-worker", + "dependency": "embedding-provider", + } + derived_started = utc_now() + environment.post_alertmanager_alert( + derived_labels, + starts_at=derived_started, + ends_at=derived_started + timedelta(minutes=10), + ) + inhibited = wait_for( + "dependent warning inhibition", + lambda: ( + alert if (alert := environment.alertmanager_alert( + "DocGridEmbeddingQueueStalled" + )) and alert.get("status", {}).get("state") == "suppressed" else None + ), + 30, + ) + time.sleep(35) + if environment.webhook_delivery("DocGridEmbeddingQueueStalled", "firing"): + raise AssertionError("The inhibited dependent warning reached the webhook receiver") + environment.post_alertmanager_alert( + derived_labels, + starts_at=derived_started, + ends_at=utc_now(), + ) + + print("Restarting the real Embedding Provider container...", flush=True) + recovery_started = utc_now() + environment.compose("start", "embedding-server", capture=False) + # Docker may assign a new ephemeral host port when a stopped container starts again. + environment.embedding_port = environment.service_port("embedding-server", 8000) + wait_for( + "Embedding Provider readiness after restart", + lambda: request( + f"http://127.0.0.1:{environment.embedding_port}/health/ready", + timeout=5, + )[0] == 200, + 900, + 5, + ) + provider_ready_at = utc_now() + wait_for( + "Prometheus to observe provider recovery", + lambda: environment.prometheus_value('up{job="embedding-provider"}') == 1, + 60, + ) + wait_for( + "EmbeddingProviderDown rule resolution", + lambda: not environment.prometheus_alert("EmbeddingProviderDown"), + 60, + ) + resolved_delivery = wait_for( + "EmbeddingProviderDown resolved webhook", + lambda: environment.webhook_delivery("EmbeddingProviderDown", "resolved"), + 90, + ) + + firing_received = datetime.fromisoformat(firing_delivery["receivedAt"]) + resolved_received = datetime.fromisoformat(resolved_delivery["receivedAt"]) + result = base_result("provider-outage") + result["conditions"] = { + "scrapeInterval": "15s", + "evaluationInterval": "15s", + "productionRuleFor": "1m", + "criticalGroupWait": "10s", + "testGroupInterval": "30s", + "productionGroupInterval": "5m", + } + result["timestamps"] = { + "providerStoppedAt": iso(stopped_at), + "upZeroObservedAt": iso(condition_at), + "alertPendingAt": iso(pending_at), + "alertFiringAt": iso(firing_at), + "firingWebhookReceivedAt": iso(firing_received), + "recoveryStartedAt": iso(recovery_started), + "providerReadyAt": iso(provider_ready_at), + "resolvedWebhookReceivedAt": iso(resolved_received), + } + result["measurements"] = { + "stopToConditionSeconds": round((condition_at - stopped_at).total_seconds(), 3), + "stopToFiringSeconds": round((firing_at - stopped_at).total_seconds(), 3), + "stopToWebhookSeconds": round((firing_received - stopped_at).total_seconds(), 3), + "restartToReadySeconds": round((provider_ready_at - recovery_started).total_seconds(), 3), + "restartToResolvedWebhookSeconds": round( + (resolved_received - recovery_started).total_seconds(), 3 + ), + "dependentWarningInhibited": bool(inhibited), + "dependentWarningWebhookDeliveries": sum( + 1 for event in environment.webhook_events() + for alert in event["payload"].get("alerts", []) + if alert.get("labels", {}).get("alertname") == "DocGridEmbeddingQueueStalled" + and alert.get("status") == "firing" + ), + } + return result + finally: + environment.close() + + +def run_queue_recovery(output_dir, jar_path): + """Create a real queued document, observe production alerts, then recover it with a Worker.""" + environment = DrillEnvironment("queue-recovery", output_dir) + try: + environment.start_core(include_provider=True) + observer = environment.start_backend( + jar_path, "observer", worker_enabled=False, snapshot_interval="2s" + ) + upload = login_and_upload(observer, environment.temp_dir) + job_id = int(upload["embeddingJobId"]) + + # Keep this experiment scoped to Embedding Queue alerts; Outbox behavior has its own rule suite. + environment.psql( + "UPDATE sync_outbox_events SET status='PROCESSED', processed_at=CURRENT_TIMESTAMP, " + "updated_at=CURRENT_TIMESTAMP WHERE event_id=(SELECT source_event_id FROM embedding_jobs " + f"WHERE id={job_id})" + ) + environment.psql( + "UPDATE embedding_jobs SET created_at=CURRENT_TIMESTAMP - INTERVAL '10 minutes', " + f"next_retry_at=NULL WHERE id={job_id}" + ) + wait_for( + "real claimable Queue gauges", + lambda: ( + metrics if ( + metric_value(metrics := fetch_metrics(observer), "docgrid_embedding_claimable_jobs") >= 1 + and metric_value(metrics, "docgrid_embedding_active_workers") == 0 + and metric_value( + metrics, "docgrid_embedding_oldest_claimable_age_seconds" + ) > 300 + ) else None + ), + 30, + ) + + environment.write_prometheus_config(observer["managementPort"]) + environment.start_prometheus() + wait_for( + "Prometheus to scrape the real Backend", + lambda: environment.prometheus_value('up{job="docgrid-backend"}') == 1, + 60, + ) + condition_at = wait_for( + "Prometheus Queue condition", + lambda: utc_now() if ( + (environment.prometheus_value( + "max(docgrid_embedding_claimable_jobs)" + ) or 0) >= 1 + and (environment.prometheus_value( + "max(docgrid_embedding_oldest_claimable_age_seconds)" + ) or 0) > 300 + ) else None, + 60, + ) + unavailable_firing_at = wait_for( + "DocGridEmbeddingWorkersUnavailable firing state", + lambda: utc_now() if environment.prometheus_alert( + "DocGridEmbeddingWorkersUnavailable", "firing" + ) else None, + 150, + 2, + ) + unavailable_delivery = wait_for( + "workers unavailable firing webhook", + lambda: environment.webhook_delivery( + "DocGridEmbeddingWorkersUnavailable", "firing" + ), + 60, + ) + stalled_firing_at = wait_for( + "DocGridEmbeddingQueueStalled production firing state", + lambda: utc_now() if environment.prometheus_alert( + "DocGridEmbeddingQueueStalled", "firing" + ) else None, + 390, + 5, + ) + stalled_delivery = wait_for( + "queue stalled firing webhook", + lambda: environment.webhook_delivery("DocGridEmbeddingQueueStalled", "firing"), + 75, + ) + + print("Starting a real indexing Worker to drain the Queue...", flush=True) + recovery_started = utc_now() + environment.start_backend( + jar_path, "worker", worker_enabled=True, snapshot_interval="2s" + ) + indexed_at = wait_for( + "the queued Embedding Job to reach INDEXED", + lambda: utc_now() if environment.psql( + f"SELECT status FROM embedding_jobs WHERE id={job_id}" + ) == "INDEXED" else None, + 300, + 2, + ) + embedding_count = int(environment.psql( + "SELECT COUNT(*) FROM embeddings WHERE document_version_id=" + f"{int(upload['documentVersionId'])}" + )) + if embedding_count <= 0: + raise AssertionError("Worker marked the Job INDEXED without stored embeddings") + wait_for( + "observer gauges to reflect Queue recovery", + lambda: ( + True if metric_value( + fetch_metrics(observer), "docgrid_embedding_claimable_jobs" + ) == 0 else False + ), + 30, + ) + wait_for( + "both Queue alerts to resolve in Prometheus", + lambda: ( + not environment.prometheus_alert("DocGridEmbeddingWorkersUnavailable") + and not environment.prometheus_alert("DocGridEmbeddingQueueStalled") + ), + 60, + ) + unavailable_resolved = wait_for( + "workers unavailable resolved webhook", + lambda: environment.webhook_delivery( + "DocGridEmbeddingWorkersUnavailable", "resolved" + ), + 90, + ) + stalled_resolved = wait_for( + "queue stalled resolved webhook", + lambda: environment.webhook_delivery("DocGridEmbeddingQueueStalled", "resolved"), + 90, + ) + + unavailable_received = datetime.fromisoformat(unavailable_delivery["receivedAt"]) + stalled_received = datetime.fromisoformat(stalled_delivery["receivedAt"]) + unavailable_resolved_at = datetime.fromisoformat(unavailable_resolved["receivedAt"]) + stalled_resolved_at = datetime.fromisoformat(stalled_resolved["receivedAt"]) + result = base_result("queue-recovery") + result["conditions"] = { + "queueAgeAtStart": ">10m", + "scrapeInterval": "15s", + "evaluationInterval": "15s", + "workersUnavailableRuleFor": "1m", + "queueStalledRuleFor": "5m", + "criticalGroupWait": "10s", + "warningGroupWait": "30s", + "testGroupInterval": "30s", + "productionGroupInterval": "5m", + } + result["timestamps"] = { + "conditionObservedAt": iso(condition_at), + "workersUnavailableFiringAt": iso(unavailable_firing_at), + "workersUnavailableWebhookAt": iso(unavailable_received), + "queueStalledFiringAt": iso(stalled_firing_at), + "queueStalledWebhookAt": iso(stalled_received), + "recoveryStartedAt": iso(recovery_started), + "jobIndexedAt": iso(indexed_at), + "workersUnavailableResolvedWebhookAt": iso(unavailable_resolved_at), + "queueStalledResolvedWebhookAt": iso(stalled_resolved_at), + } + result["measurements"] = { + "conditionToWorkersUnavailableFiringSeconds": round( + (unavailable_firing_at - condition_at).total_seconds(), 3 + ), + "conditionToWorkersUnavailableWebhookSeconds": round( + (unavailable_received - condition_at).total_seconds(), 3 + ), + "conditionToQueueStalledFiringSeconds": round( + (stalled_firing_at - condition_at).total_seconds(), 3 + ), + "conditionToQueueStalledWebhookSeconds": round( + (stalled_received - condition_at).total_seconds(), 3 + ), + "workerStartToIndexedSeconds": round( + (indexed_at - recovery_started).total_seconds(), 3 + ), + "workerStartToAllResolvedWebhooksSeconds": round(max( + (unavailable_resolved_at - recovery_started).total_seconds(), + (stalled_resolved_at - recovery_started).total_seconds(), + ), 3), + "storedEmbeddings": embedding_count, + "finalJobStatus": "INDEXED", + } + return result + finally: + environment.close() + + +def write_result(output_dir, result): + """Write stable, reviewable JSON evidence for one completed scenario.""" + path = output_dir / f"{result['scenario']}.json" + path.write_text(json.dumps(result, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") + print(f"{result['scenario']}: PASS -> {path}", flush=True) + return path + + +def parse_arguments(): + """Expose each experiment independently while retaining one sequential all command.""" + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument( + "scenario", + choices=["scrape-load", "queue-recovery", "provider-outage", "all"], + ) + parser.add_argument( + "--output-dir", + type=Path, + default=ROOT_DIR / "backend" / "build" / "reports" / "observability-drill" + / datetime.now().strftime("%Y%m%d-%H%M%S"), + ) + return parser.parse_args() + + +def main(): + """Build only when required, run selected isolated scenarios, and preserve failure diagnostics.""" + arguments = parse_arguments() + output_dir = arguments.output_dir.resolve() + output_dir.mkdir(parents=True, exist_ok=True) + scenarios = ( + ["scrape-load", "queue-recovery", "provider-outage"] + if arguments.scenario == "all" else [arguments.scenario] + ) + jar_path = build_backend() if any( + scenario in {"scrape-load", "queue-recovery"} for scenario in scenarios + ) else None + results = [] + try: + for scenario in scenarios: + if scenario == "scrape-load": + result = run_scrape_load(output_dir, jar_path) + elif scenario == "queue-recovery": + result = run_queue_recovery(output_dir, jar_path) + else: + result = run_provider_outage(output_dir) + write_result(output_dir, result) + results.append(result) + except Exception as error: + failure = { + "status": "FAIL", + "scenario": scenario, + "failedAt": iso(), + "error": str(error), + "traceback": traceback.format_exc(), + } + (output_dir / f"{scenario}-failure.json").write_text( + json.dumps(failure, ensure_ascii=False, indent=2) + "\n", + encoding="utf-8", + ) + print(f"{scenario}: FAIL -> {error}", file=sys.stderr, flush=True) + return 1 + + if len(results) > 1: + summary = { + "status": "PASS", + "measuredAt": iso(), + "scenarios": [result["scenario"] for result in results], + "resultFiles": [f"{result['scenario']}.json" for result in results], + } + (output_dir / "summary.json").write_text( + json.dumps(summary, ensure_ascii=False, indent=2) + "\n", + encoding="utf-8", + ) + print(f"Observability incident drill completed: {output_dir}", flush=True) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/monitoring/drills/webhook_receiver.py b/monitoring/drills/webhook_receiver.py new file mode 100644 index 00000000..48693360 --- /dev/null +++ b/monitoring/drills/webhook_receiver.py @@ -0,0 +1,47 @@ +"""Capture timestamped Alertmanager webhook payloads for local incident drills.""" + +import json +import os +from datetime import datetime, timezone +from http.server import BaseHTTPRequestHandler, HTTPServer +from pathlib import Path + + +EVENT_LOG = Path(os.environ["EVENT_LOG"]) + + +class WebhookHandler(BaseHTTPRequestHandler): + """Expose a health probe and append each alert delivery as one durable JSON line.""" + + def do_GET(self): + if self.path != "/health": + self.send_error(404) + return + self.send_response(200) + self.end_headers() + + def do_POST(self): + if self.path != "/alerts": + self.send_error(404) + return + + content_length = int(self.headers.get("Content-Length", "0")) + payload = json.loads(self.rfile.read(content_length)) + event = { + "receivedAt": datetime.now(timezone.utc).isoformat(), + "payload": payload, + } + with EVENT_LOG.open("a", encoding="utf-8") as output: + output.write(json.dumps(event, separators=(",", ":")) + "\n") + output.flush() + + self.send_response(200) + self.end_headers() + + def log_message(self, format, *args): + """Suppress access logs so failures retain the useful experiment output.""" + + +if __name__ == "__main__": + EVENT_LOG.parent.mkdir(parents=True, exist_ok=True) + HTTPServer(("0.0.0.0", 8080), WebhookHandler).serve_forever() diff --git a/monitoring/verify.sh b/monitoring/verify.sh index c5cdb27a..dc92aadd 100755 --- a/monitoring/verify.sh +++ b/monitoring/verify.sh @@ -48,6 +48,9 @@ docker run --rm --entrypoint=amtool \ docker run --rm --entrypoint=amtool \ -v "$ROOT_DIR/monitoring/alertmanager/tests/alertmanager.yml:/etc/alertmanager/alertmanager.yml:ro" \ "$ALERTMANAGER_IMAGE" check-config /etc/alertmanager/alertmanager.yml +docker run --rm --entrypoint=amtool \ + -v "$ROOT_DIR/monitoring/drills/alertmanager.yml:/etc/alertmanager/alertmanager.yml:ro" \ + "$ALERTMANAGER_IMAGE" check-config /etc/alertmanager/alertmanager.yml for config_file in "$ROOT_DIR"/monitoring/alertmanager/examples/*.yml; do docker run --rm --entrypoint=amtool \ -v "$config_file:/etc/alertmanager/alertmanager.yml:ro" \ @@ -55,7 +58,19 @@ for config_file in "$ROOT_DIR"/monitoring/alertmanager/examples/*.yml; do "$ALERTMANAGER_IMAGE" check-config /etc/alertmanager/alertmanager.yml done -# 4. 사용자가 실행할 두 Compose 형태가 모두 정상 렌더링되는지 확인한다. +# 4. 장시간 실험을 다시 실행하지 않고도 Drill 실행기 문법과 격리 Compose를 검증한다. +python3 - "$ROOT_DIR/monitoring/drills/run_drill.py" <<'PY' +import ast +import pathlib +import sys + +ast.parse(pathlib.Path(sys.argv[1]).read_text(encoding="utf-8")) +PY +sh -n "$ROOT_DIR/monitoring/drills/run.sh" +DRILL_TMP_DIR="$SECRET_DIR" docker compose \ + -f "$ROOT_DIR/monitoring/drills/docker-compose.yml" config --quiet + +# 5. 사용자가 실행할 두 Compose 형태가 모두 정상 렌더링되는지 확인한다. docker compose -f "$ROOT_DIR/docker-compose.yml" config --quiet docker compose -f "$ROOT_DIR/docker-compose.yml" --profile monitoring config --quiet