diff --git a/.env.example b/.env.example index 6ad2e05e..e3c2b266 100644 --- a/.env.example +++ b/.env.example @@ -60,6 +60,7 @@ MINIO_SECRET_KEY=minioadmin1234 # 선택형 monitoring Profile: docker compose --profile monitoring up -d # Backend 운영 엔드포인트 전용 포트입니다. 외부 인터넷에 직접 공개하지 마세요. # MANAGEMENT_PORT=8081 +# MANAGEMENT_METRICS_SNAPSHOT_INTERVAL=15s # PROMETHEUS_PORT=9090 # PROMETHEUS_RETENTION=7d # PROMETHEUS_MEMORY_LIMIT=512m diff --git a/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetrics.java b/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetrics.java new file mode 100644 index 00000000..c614f2a1 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetrics.java @@ -0,0 +1,128 @@ +package com.opensource.docgrid.global.observability; + +import java.time.Clock; +import java.time.Duration; +import java.time.Instant; +import java.time.LocalDateTime; +import java.util.concurrent.atomic.AtomicReference; + +import org.springframework.stereotype.Component; + +import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.MeterRegistry; + +/** + * 마지막 정상 DB Snapshot을 메모리에 보관하고 이를 Micrometer Gauge로 노출한다. + * + *
Gauge callback은 AtomicReference와 Clock만 읽는다. DB 갱신 실패 시 이전 값을 유지하고 Snapshot
+ * 나이와 갱신 실패 Counter를 올려, 0으로 덮인 값이 정상 상태로 오인되는 것을 막는다.
+ */
+@Component
+public class OperationalMetrics {
+
+ private final Clock clock;
+ private final AtomicReference DB 집계 결과와 Prometheus scrape 사이의 경계다. 오래된 항목이 없으면 해당 시각은 {@code null}이며,
+ * 모든 개수는 음수가 될 수 없다.
+ */
+public record OperationalMetricsSnapshot(
+ long embeddingClaimableJobs,
+ long embeddingProcessingJobs,
+ LocalDateTime embeddingOldestClaimableAt,
+ long embeddingActiveWorkers,
+ long ragProcessingJobs,
+ LocalDateTime ragOldestProcessingAt,
+ long syncClaimableEvents,
+ long syncProcessingEvents,
+ LocalDateTime syncOldestClaimableAt
+) {
+
+ /** 집계 Query의 잘못된 결과가 Gauge에 게시되지 않도록 불변 조건을 확인한다. */
+ public OperationalMetricsSnapshot {
+ if (embeddingClaimableJobs < 0 || embeddingProcessingJobs < 0 || embeddingActiveWorkers < 0
+ || ragProcessingJobs < 0 || syncClaimableEvents < 0 || syncProcessingEvents < 0) {
+ throw new IllegalArgumentException("운영 상태 개수는 음수일 수 없습니다.");
+ }
+ requireTimestampWhenPresent(embeddingOldestClaimableAt, embeddingClaimableJobs, "Embedding");
+ requireTimestampWhenPresent(ragOldestProcessingAt, ragProcessingJobs, "RAG");
+ requireTimestampWhenPresent(syncOldestClaimableAt, syncClaimableEvents, "Sync Outbox");
+ }
+
+ private static void requireTimestampWhenPresent(LocalDateTime oldestAt, long count, String name) {
+ if ((count == 0) != Objects.isNull(oldestAt)) {
+ throw new IllegalArgumentException(name + " 개수와 가장 오래된 시각이 일치하지 않습니다.");
+ }
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotQueryService.java b/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotQueryService.java
new file mode 100644
index 00000000..c2dc3fc9
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotQueryService.java
@@ -0,0 +1,174 @@
+package com.opensource.docgrid.global.observability;
+
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.time.Duration;
+import java.time.LocalDateTime;
+
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 비동기 Pipeline의 현재 운영 상태를 행 로딩 없이 목적별 aggregate query로 읽는다.
+ *
+ * 호출자는 이 짧은 읽기 Transaction의 결과만 메모리에 보관한다. Prometheus scrape 경로에서는 이
+ * 서비스를 호출하지 않아, 감시 트래픽이 DB 부하로 이어지지 않게 한다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional(readOnly = true, timeout = 2)
+public class OperationalMetricsSnapshotQueryService {
+
+ private static final String EMBEDDING_SNAPSHOT_SQL = """
+ SELECT
+ (
+ SELECT COUNT(*)
+ FROM embedding_jobs job
+ WHERE job.status = 'PENDING'
+ AND (job.next_retry_at IS NULL OR job.next_retry_at <= ?)
+ ) AS claimable_jobs,
+ (
+ SELECT COUNT(*)
+ FROM embedding_jobs job
+ WHERE job.status = 'PROCESSING'
+ ) AS processing_jobs,
+ (
+ SELECT MIN(COALESCE(job.next_retry_at, job.created_at))
+ FROM embedding_jobs job
+ WHERE job.status = 'PENDING'
+ AND (job.next_retry_at IS NULL OR job.next_retry_at <= ?)
+ ) AS oldest_claimable_at,
+ (
+ SELECT COUNT(*)
+ FROM worker_nodes worker
+ WHERE worker.status IN ('ACTIVE', 'IDLE')
+ AND worker.last_heartbeat_at > ?
+ ) AS active_workers
+ """;
+
+ private static final String RAG_SNAPSHOT_SQL = """
+ SELECT
+ (
+ SELECT COUNT(*)
+ FROM rag_responses response
+ WHERE response.status = 'PROCESSING'
+ ) AS processing_jobs,
+ (
+ SELECT MIN(response.created_at)
+ FROM rag_responses response
+ WHERE response.status = 'PROCESSING'
+ ) AS oldest_processing_at
+ """;
+
+ private static final String SYNC_SNAPSHOT_SQL = """
+ SELECT
+ (
+ SELECT COUNT(*)
+ FROM sync_outbox_events event
+ WHERE event.status = 'PENDING'
+ AND event.available_at <= ?
+ ) AS claimable_events,
+ (
+ SELECT COUNT(*)
+ FROM sync_outbox_events event
+ WHERE event.status = 'PROCESSING'
+ ) AS processing_events,
+ (
+ SELECT MIN(event.available_at)
+ FROM sync_outbox_events event
+ WHERE event.status = 'PENDING'
+ AND event.available_at <= ?
+ ) AS oldest_claimable_at
+ """;
+
+ private final JdbcTemplate jdbcTemplate;
+
+ @Value("${indexing.worker.dead-threshold:30s}")
+ private Duration workerDeadThreshold;
+
+ /**
+ * 동일한 관측 시각을 세 Queue 집계에 적용해 경계 시각의 포함 여부가 서로 어긋나지 않게 한다.
+ */
+ public OperationalMetricsSnapshot load(LocalDateTime observedAt) {
+ // 1. Worker의 실효 상태는 관리자 조회와 같은 Heartbeat 만료 기준으로 계산한다.
+ LocalDateTime heartbeatDeadline = observedAt.minus(workerDeadThreshold);
+
+ // 2. Queue별 집계는 각각 한 SQL로 끝내 Entity 수에 비례하는 조회를 만들지 않는다.
+ EmbeddingSnapshot embedding = jdbcTemplate.queryForObject(
+ EMBEDDING_SNAPSHOT_SQL,
+ this::mapEmbeddingSnapshot,
+ observedAt,
+ observedAt,
+ heartbeatDeadline
+ );
+ RagSnapshot rag = jdbcTemplate.queryForObject(RAG_SNAPSHOT_SQL, this::mapRagSnapshot);
+ SyncSnapshot sync = jdbcTemplate.queryForObject(
+ SYNC_SNAPSHOT_SQL,
+ this::mapSyncSnapshot,
+ observedAt,
+ observedAt
+ );
+
+ // 3. 세 Query가 모두 성공한 경우에만 하나의 원자적 Snapshot 후보를 반환한다.
+ return new OperationalMetricsSnapshot(
+ embedding.claimableJobs(),
+ embedding.processingJobs(),
+ embedding.oldestClaimableAt(),
+ embedding.activeWorkers(),
+ rag.processingJobs(),
+ rag.oldestProcessingAt(),
+ sync.claimableEvents(),
+ sync.processingEvents(),
+ sync.oldestClaimableAt()
+ );
+ }
+
+ private EmbeddingSnapshot mapEmbeddingSnapshot(ResultSet resultSet, int rowNumber) throws SQLException {
+ return new EmbeddingSnapshot(
+ resultSet.getLong("claimable_jobs"),
+ resultSet.getLong("processing_jobs"),
+ resultSet.getObject("oldest_claimable_at", LocalDateTime.class),
+ resultSet.getLong("active_workers")
+ );
+ }
+
+ private RagSnapshot mapRagSnapshot(ResultSet resultSet, int rowNumber) throws SQLException {
+ return new RagSnapshot(
+ resultSet.getLong("processing_jobs"),
+ resultSet.getObject("oldest_processing_at", LocalDateTime.class)
+ );
+ }
+
+ private SyncSnapshot mapSyncSnapshot(ResultSet resultSet, int rowNumber) throws SQLException {
+ return new SyncSnapshot(
+ resultSet.getLong("claimable_events"),
+ resultSet.getLong("processing_events"),
+ resultSet.getObject("oldest_claimable_at", LocalDateTime.class)
+ );
+ }
+
+ /** Embedding Queue와 유효 Worker를 한 SQL에서 읽은 내부 집계 값이다. */
+ private record EmbeddingSnapshot(
+ long claimableJobs,
+ long processingJobs,
+ LocalDateTime oldestClaimableAt,
+ long activeWorkers
+ ) {
+ }
+
+ /** RAG PROCESSING Queue를 한 SQL에서 읽은 내부 집계 값이다. */
+ private record RagSnapshot(long processingJobs, LocalDateTime oldestProcessingAt) {
+ }
+
+ /** Sync Outbox Queue를 한 SQL에서 읽은 내부 집계 값이다. */
+ private record SyncSnapshot(
+ long claimableEvents,
+ long processingEvents,
+ LocalDateTime oldestClaimableAt
+ ) {
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotRefresher.java b/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotRefresher.java
new file mode 100644
index 00000000..fe05adbe
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotRefresher.java
@@ -0,0 +1,91 @@
+package com.opensource.docgrid.global.observability;
+
+import java.time.Clock;
+import java.time.Duration;
+import java.time.LocalDateTime;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.event.EventListener;
+import org.springframework.stereotype.Component;
+
+import jakarta.annotation.PreDestroy;
+import lombok.extern.slf4j.Slf4j;
+
+/**
+ * 전용 단일 Thread에서 운영 DB Snapshot을 주기적으로 갱신한다.
+ *
+ * 기존 Worker·RAG·Outbox Scheduler와 실행 자원을 공유하지 않으므로 긴 Queue 작업이 관측 상태 갱신을
+ * 지연시키지 않는다. 각 실행은 실패를 격리해 다음 주기에도 계속 시도한다.
+ */
+@Component
+@Slf4j
+public class OperationalMetricsSnapshotRefresher {
+
+ private final OperationalMetricsSnapshotQueryService queryService;
+ private final OperationalMetrics operationalMetrics;
+ private final Clock clock;
+ private final long refreshIntervalMillis;
+ private final AtomicBoolean started = new AtomicBoolean();
+
+ private ScheduledExecutorService executor;
+
+ /** 양수인 갱신 주기를 밀리초 단위로 고정해 Scheduler의 busy loop를 막는다. */
+ public OperationalMetricsSnapshotRefresher(
+ OperationalMetricsSnapshotQueryService queryService,
+ OperationalMetrics operationalMetrics,
+ Clock clock,
+ @Value("${management.metrics.docgrid.snapshot-interval:15s}") Duration refreshInterval
+ ) {
+ if (refreshInterval == null || refreshInterval.isZero() || refreshInterval.isNegative()) {
+ throw new IllegalArgumentException("운영 Metrics Snapshot 갱신 주기는 0보다 커야 합니다.");
+ }
+ this.queryService = queryService;
+ this.operationalMetrics = operationalMetrics;
+ this.clock = clock;
+ this.refreshIntervalMillis = Math.max(1L, refreshInterval.toMillis());
+ }
+
+ /** 애플리케이션 준비 후 즉시 첫 Snapshot을 읽고 고정 지연 방식으로 다음 갱신을 예약한다. */
+ @EventListener(ApplicationReadyEvent.class)
+ public void start() {
+ if (!started.compareAndSet(false, true)) {
+ return;
+ }
+
+ executor = Executors.newSingleThreadScheduledExecutor(task -> {
+ Thread thread = new Thread(task, "operational-metrics-snapshot");
+ thread.setDaemon(true);
+ return thread;
+ });
+ executor.scheduleWithFixedDelay(this::refreshSafely, 0L, refreshIntervalMillis, TimeUnit.MILLISECONDS);
+ }
+
+ /** 테스트와 Scheduler가 같은 실패 격리 경로를 사용하도록 한 번의 갱신을 명시적으로 실행한다. */
+ void refreshSafely() {
+ try {
+ // 1. 모든 집계가 같은 now를 사용하도록 관측 시각을 한 번만 만든다.
+ LocalDateTime observedAt = LocalDateTime.now(clock);
+ OperationalMetricsSnapshot snapshot = queryService.load(observedAt);
+
+ // 2. Query가 모두 성공한 경우에만 마지막 정상 Snapshot과 성공 Counter를 갱신한다.
+ operationalMetrics.recordRefreshSuccess(snapshot, clock.instant());
+ } catch (RuntimeException exception) {
+ // 3. 실패는 이전 Snapshot을 보존한 채 Counter와 로그에만 남겨 다음 주기에서 재시도한다.
+ operationalMetrics.recordRefreshFailure();
+ log.warn("[OBSERVABILITY] 운영 Metrics Snapshot 갱신 실패", exception);
+ }
+ }
+
+ /** 애플리케이션 종료 시 전용 Thread가 남지 않도록 새 실행을 중단한다. */
+ @PreDestroy
+ public void stop() {
+ if (executor != null) {
+ executor.shutdownNow();
+ }
+ }
+}
diff --git a/backend/src/main/resources/application.yml b/backend/src/main/resources/application.yml
index d7bfd473..9bc7d0ab 100644
--- a/backend/src/main/resources/application.yml
+++ b/backend/src/main/resources/application.yml
@@ -128,6 +128,10 @@ management:
readiness:
# Redis는 인증 경로에서 fail-open이므로 트래픽 수신 가능 여부는 애플리케이션과 DB로 판단한다.
include: readinessState,db
+ metrics:
+ docgrid:
+ # Prometheus scrape는 메모리 Snapshot만 읽고, DB 집계는 이 독립 주기로 제한한다.
+ snapshot-interval: ${MANAGEMENT_METRICS_SNAPSHOT_INTERVAL:15s}
jwt:
secret: ${JWT_SECRET}
diff --git a/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotQueryServiceTest.java b/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotQueryServiceTest.java
new file mode 100644
index 00000000..0be4394e
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotQueryServiceTest.java
@@ -0,0 +1,209 @@
+package com.opensource.docgrid.global.observability;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.time.LocalDateTime;
+import java.util.UUID;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.autoconfigure.jdbc.AutoConfigureTestDatabase;
+import org.springframework.boot.test.autoconfigure.orm.jpa.DataJpaTest;
+import org.springframework.context.annotation.Import;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.test.annotation.DirtiesContext;
+import org.springframework.test.context.ActiveProfiles;
+import org.springframework.test.context.DynamicPropertyRegistry;
+import org.springframework.test.context.DynamicPropertySource;
+
+/**
+ * 운영 상태 aggregate SQL의 claim 가능 시각, 상태 수와 가장 오래된 시각을 실제 PostgreSQL에서 검증한다.
+ */
+@DataJpaTest
+@ActiveProfiles("test")
+@AutoConfigureTestDatabase(replace = AutoConfigureTestDatabase.Replace.NONE)
+@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@Import(OperationalMetricsSnapshotQueryService.class)
+@DisplayName("운영 상태 Snapshot Query 테스트")
+class OperationalMetricsSnapshotQueryServiceTest {
+
+ private static final String TEST_SCHEMA = "docgrid_operational_metrics_snapshot_query_test";
+ private static final LocalDateTime NOW = LocalDateTime.of(2026, 9, 13, 19, 0);
+
+ @Autowired private JdbcTemplate jdbcTemplate;
+ @Autowired private OperationalMetricsSnapshotQueryService queryService;
+
+ private Long versionId;
+ private Long embeddingModelId;
+ private Long queryId;
+
+ @DynamicPropertySource
+ static void configureDatabase(DynamicPropertyRegistry registry) {
+ registry.add("TEST_DB_SCHEMA", () -> TEST_SCHEMA);
+ registry.add("jwt.secret", () -> "docgrid-operational-metrics-query-test-secret-key-2026");
+ registry.add("indexing.worker.dead-threshold", () -> "30s");
+ }
+
+ @BeforeEach
+ void setUp() {
+ jdbcTemplate.execute("""
+ TRUNCATE TABLE
+ embedding_jobs,
+ worker_nodes,
+ rag_responses,
+ search_queries,
+ search_conversations,
+ sync_outbox_events,
+ document_versions,
+ documents,
+ users
+ RESTART IDENTITY CASCADE
+ """);
+
+ String suffix = UUID.randomUUID().toString();
+ Long userId = insertUser(suffix);
+ Long documentId = insertDocument(userId);
+ versionId = insertVersion(documentId, userId);
+ embeddingModelId = jdbcTemplate.queryForObject("""
+ SELECT id FROM embedding_models WHERE is_active = TRUE AND is_searchable = TRUE
+ """, Long.class);
+ Long conversationId = insertConversation(userId);
+ queryId = insertQuery(userId, conversationId);
+ }
+
+ @AfterAll
+ void dropSchema() {
+ jdbcTemplate.execute("DROP SCHEMA IF EXISTS " + TEST_SCHEMA + " CASCADE");
+ }
+
+ @Test
+ @DisplayName("미래 backoff를 제외하고 Queue별 수·나이와 유효 Worker를 집계한다")
+ void load_excludesFutureBackoffAndCountsOperationalState() {
+ insertEmbeddingJob("PENDING", NOW.minusMinutes(1), null);
+ insertEmbeddingJob("PENDING", NOW.minusDays(2), NOW.minusMinutes(2));
+ insertEmbeddingJob("PENDING", NOW.minusDays(3), NOW.plusMinutes(5));
+ insertEmbeddingJob("PROCESSING", NOW.minusMinutes(3), null);
+ insertWorker("fresh-active", "ACTIVE", NOW.minusSeconds(10));
+ insertWorker("expired-idle", "IDLE", NOW.minusSeconds(30));
+ insertWorker("fresh-stopped", "STOPPED", NOW.minusSeconds(5));
+
+ insertRagResponse("PROCESSING", NOW.minusMinutes(4));
+ insertRagResponse("SUCCESS", NOW.minusHours(1));
+
+ insertSyncEvent("PENDING", NOW.minusMinutes(3));
+ insertSyncEvent("PENDING", NOW.plusMinutes(3));
+ insertSyncEvent("PROCESSING", NOW.minusMinutes(5));
+
+ OperationalMetricsSnapshot snapshot = queryService.load(NOW);
+
+ assertThat(snapshot.embeddingClaimableJobs()).isEqualTo(2);
+ assertThat(snapshot.embeddingProcessingJobs()).isEqualTo(1);
+ assertThat(snapshot.embeddingOldestClaimableAt()).isEqualTo(NOW.minusMinutes(2));
+ assertThat(snapshot.embeddingActiveWorkers()).isEqualTo(1);
+ assertThat(snapshot.ragProcessingJobs()).isEqualTo(1);
+ assertThat(snapshot.ragOldestProcessingAt()).isEqualTo(NOW.minusMinutes(4));
+ assertThat(snapshot.syncClaimableEvents()).isEqualTo(1);
+ assertThat(snapshot.syncProcessingEvents()).isEqualTo(1);
+ assertThat(snapshot.syncOldestClaimableAt()).isEqualTo(NOW.minusMinutes(3));
+ }
+
+ @Test
+ @DisplayName("대상 행이 없으면 개수는 0이고 가장 오래된 시각은 null이다")
+ void load_returnsZerosAndNullAgesWhenQueuesAreEmpty() {
+ OperationalMetricsSnapshot snapshot = queryService.load(NOW);
+
+ assertThat(snapshot.embeddingClaimableJobs()).isZero();
+ assertThat(snapshot.embeddingProcessingJobs()).isZero();
+ assertThat(snapshot.embeddingOldestClaimableAt()).isNull();
+ assertThat(snapshot.embeddingActiveWorkers()).isZero();
+ assertThat(snapshot.ragProcessingJobs()).isZero();
+ assertThat(snapshot.ragOldestProcessingAt()).isNull();
+ assertThat(snapshot.syncClaimableEvents()).isZero();
+ assertThat(snapshot.syncProcessingEvents()).isZero();
+ assertThat(snapshot.syncOldestClaimableAt()).isNull();
+ }
+
+ private Long insertUser(String suffix) {
+ return jdbcTemplate.queryForObject("""
+ INSERT INTO users (email, password_hash, name, status, created_at, updated_at)
+ VALUES (?, 'password-hash', 'Operational Metrics User', 'ACTIVE', ?, ?)
+ RETURNING id
+ """, Long.class, "operational-metrics-" + suffix + "@example.com", NOW, NOW);
+ }
+
+ private Long insertDocument(Long userId) {
+ return jdbcTemplate.queryForObject("""
+ INSERT INTO documents (
+ owner_user_id, title, document_type, source_type, status, visibility, created_at, updated_at
+ ) VALUES (?, 'Operational Metrics Document', 'TXT', 'UPLOAD', 'INDEXING', 'PRIVATE', ?, ?)
+ RETURNING id
+ """, Long.class, userId, NOW, NOW);
+ }
+
+ private Long insertVersion(Long documentId, Long userId) {
+ return jdbcTemplate.queryForObject("""
+ INSERT INTO document_versions (
+ document_id, version_no, title_snapshot, content_type, status, created_by, created_at, updated_at
+ ) VALUES (?, 1, 'Operational Metrics Version', 'text/plain', 'EMBEDDING', ?, ?, ?)
+ RETURNING id
+ """, Long.class, documentId, userId, NOW, NOW);
+ }
+
+ private Long insertConversation(Long userId) {
+ return jdbcTemplate.queryForObject("""
+ INSERT INTO search_conversations (user_id, title, last_message_at, created_at, updated_at)
+ VALUES (?, 'Operational Metrics Conversation', ?, ?, ?)
+ RETURNING id
+ """, Long.class, userId, NOW, NOW, NOW);
+ }
+
+ private Long insertQuery(Long userId, Long conversationId) {
+ return jdbcTemplate.queryForObject("""
+ INSERT INTO search_queries (
+ user_id, conversation_id, query_text, search_type, top_k, status, created_at, updated_at
+ ) VALUES (?, ?, 'Operational metrics query', 'VECTOR', 5, 'SUCCESS', ?, ?)
+ RETURNING id
+ """, Long.class, userId, conversationId, NOW, NOW);
+ }
+
+ private void insertEmbeddingJob(String status, LocalDateTime createdAt, LocalDateTime nextRetryAt) {
+ jdbcTemplate.update("""
+ INSERT INTO embedding_jobs (
+ document_version_id, embedding_model_id, status, priority, retry_count,
+ max_retry_count, next_retry_at, created_at, updated_at
+ ) VALUES (?, ?, ?, 0, 0, 3, ?, ?, ?)
+ """, versionId, embeddingModelId, status, nextRetryAt, createdAt, createdAt);
+ }
+
+ private void insertWorker(String instanceId, String status, LocalDateTime heartbeatAt) {
+ jdbcTemplate.update("""
+ INSERT INTO worker_nodes (
+ worker_name, instance_id, host_name, ip_address, status,
+ last_heartbeat_at, started_at, created_at, updated_at
+ ) VALUES ('indexing-worker', ?, 'test-host', '127.0.0.1', ?, ?, ?, ?, ?)
+ """, instanceId, status, heartbeatAt, NOW.minusHours(1), NOW.minusHours(1), NOW);
+ }
+
+ private void insertRagResponse(String status, LocalDateTime createdAt) {
+ jdbcTemplate.update("""
+ INSERT INTO rag_responses (query_id, answer_text, status, created_at, updated_at)
+ VALUES (?, '', ?, ?, ?)
+ """, queryId, status, createdAt, createdAt);
+ }
+
+ private void insertSyncEvent(String status, LocalDateTime availableAt) {
+ String eventId = UUID.randomUUID().toString();
+ jdbcTemplate.update("""
+ INSERT INTO sync_outbox_events (
+ event_id, idempotency_key, aggregate_type, aggregate_id, event_type,
+ status, available_at, occurred_at, retry_count, max_retry_count, created_at, updated_at
+ ) VALUES (?::uuid, ?, 'DOCUMENT', 1, 'DOCUMENT_CREATED', ?, ?, ?, 0, 3, ?, ?)
+ """, eventId, "operational-metrics-" + eventId, status, availableAt,
+ availableAt.minusMinutes(1), availableAt.minusMinutes(1), availableAt.minusMinutes(1));
+ }
+}
diff --git a/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotRefresherTest.java b/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotRefresherTest.java
new file mode 100644
index 00000000..c341b400
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsSnapshotRefresherTest.java
@@ -0,0 +1,66 @@
+package com.opensource.docgrid.global.observability;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.BDDMockito.given;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.reset;
+import static org.mockito.Mockito.verifyNoInteractions;
+
+import java.time.Clock;
+import java.time.Duration;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneOffset;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
+
+/**
+ * 운영 Snapshot 갱신 성공·실패가 Metric 저장소 경계에서 격리되는지 검증한다.
+ */
+@DisplayName("운영 상태 Snapshot Refresher 테스트")
+class OperationalMetricsSnapshotRefresherTest {
+
+ private static final Instant NOW = Instant.parse("2026-09-13T10:00:00Z");
+ private static final Clock CLOCK = Clock.fixed(NOW, ZoneOffset.UTC);
+
+ private final OperationalMetricsSnapshotQueryService queryService =
+ mock(OperationalMetricsSnapshotQueryService.class);
+ private final SimpleMeterRegistry meterRegistry = new SimpleMeterRegistry();
+ private final OperationalMetrics metrics = new OperationalMetrics(meterRegistry, CLOCK);
+ private final OperationalMetricsSnapshotRefresher refresher =
+ new OperationalMetricsSnapshotRefresher(queryService, metrics, CLOCK, Duration.ofSeconds(15));
+
+ @Test
+ @DisplayName("세 집계가 성공하면 같은 결과를 게시하고 scrape는 QueryService를 호출하지 않는다")
+ void refreshesSnapshotAndServesScrapesFromMemory() {
+ LocalDateTime observedAt = LocalDateTime.now(CLOCK);
+ OperationalMetricsSnapshot snapshot = new OperationalMetricsSnapshot(
+ 2, 1, observedAt.minusSeconds(30), 1, 0, null, 0, 0, null
+ );
+ given(queryService.load(observedAt)).willReturn(snapshot);
+
+ refresher.refreshSafely();
+ reset(queryService);
+
+ assertThat(meterRegistry.get("docgrid.embedding.claimable.jobs").gauge().value()).isEqualTo(2.0);
+ assertThat(meterRegistry.get("docgrid.embedding.claimable.jobs").gauge().value()).isEqualTo(2.0);
+ verifyNoInteractions(queryService);
+ }
+
+ @Test
+ @DisplayName("집계 예외를 밖으로 전파하지 않고 실패 Counter만 증가시킨다")
+ void recordsFailureAndKeepsSchedulerAlive() {
+ LocalDateTime observedAt = LocalDateTime.now(CLOCK);
+ given(queryService.load(observedAt)).willThrow(new IllegalStateException("database unavailable"));
+
+ refresher.refreshSafely();
+
+ assertThat(meterRegistry.get("docgrid.operational.snapshot.refresh")
+ .tag("outcome", "failed").counter().count()).isEqualTo(1.0);
+ assertThat(meterRegistry.get("docgrid.operational.snapshot.age.seconds").gauge().value())
+ .isInfinite();
+ }
+}
diff --git a/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsTest.java b/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsTest.java
new file mode 100644
index 00000000..61a6792b
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/global/observability/OperationalMetricsTest.java
@@ -0,0 +1,90 @@
+package com.opensource.docgrid.global.observability;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import java.time.Clock;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+import io.micrometer.core.instrument.simple.SimpleMeterRegistry;
+
+/**
+ * 운영 Snapshot이 DB 접근 없이 Gauge로 변환되고 실패 때 마지막 정상 값을 보존하는지 검증한다.
+ */
+@DisplayName("운영 상태 Gauge 테스트")
+class OperationalMetricsTest {
+
+ private static final ZoneId ZONE_ID = ZoneId.of("Asia/Seoul");
+ private static final Instant NOW = Instant.parse("2026-09-13T10:00:00Z");
+
+ private final Clock clock = Clock.fixed(NOW, ZONE_ID);
+ private final SimpleMeterRegistry meterRegistry = new SimpleMeterRegistry();
+ private final OperationalMetrics metrics = new OperationalMetrics(meterRegistry, clock);
+
+ @Test
+ @DisplayName("마지막 정상 Snapshot의 Queue 수와 현재 기준 나이를 Gauge로 제공한다")
+ void exposesLastSuccessfulSnapshotAsGauges() {
+ LocalDateTime now = LocalDateTime.now(clock);
+ metrics.recordRefreshSuccess(new OperationalMetricsSnapshot(
+ 3,
+ 2,
+ now.minusSeconds(75),
+ 1,
+ 4,
+ now.minusSeconds(120),
+ 5,
+ 1,
+ now.minusSeconds(30)
+ ), NOW.minusSeconds(10));
+
+ assertThat(gauge("docgrid.embedding.claimable.jobs")).isEqualTo(3.0);
+ assertThat(gauge("docgrid.embedding.processing.jobs")).isEqualTo(2.0);
+ assertThat(gauge("docgrid.embedding.oldest.claimable.age.seconds")).isEqualTo(75.0);
+ assertThat(gauge("docgrid.embedding.active.workers")).isEqualTo(1.0);
+ assertThat(gauge("docgrid.rag.processing.jobs")).isEqualTo(4.0);
+ assertThat(gauge("docgrid.rag.oldest.processing.age.seconds")).isEqualTo(120.0);
+ assertThat(gauge("docgrid.sync.outbox.claimable.events")).isEqualTo(5.0);
+ assertThat(gauge("docgrid.sync.outbox.processing.events")).isEqualTo(1.0);
+ assertThat(gauge("docgrid.sync.outbox.oldest.claimable.age.seconds")).isEqualTo(30.0);
+ assertThat(gauge("docgrid.operational.snapshot.age.seconds")).isEqualTo(10.0);
+ assertThat(counter("success")).isEqualTo(1.0);
+ }
+
+ @Test
+ @DisplayName("갱신 실패는 마지막 정상 Gauge를 0으로 덮어쓰지 않는다")
+ void retainsLastGoodSnapshotAfterRefreshFailure() {
+ LocalDateTime now = LocalDateTime.now(clock);
+ metrics.recordRefreshSuccess(new OperationalMetricsSnapshot(
+ 2, 0, now.minusSeconds(15), 1, 0, null, 0, 0, null
+ ), NOW);
+
+ metrics.recordRefreshFailure();
+
+ assertThat(gauge("docgrid.embedding.claimable.jobs")).isEqualTo(2.0);
+ assertThat(gauge("docgrid.embedding.oldest.claimable.age.seconds")).isEqualTo(15.0);
+ assertThat(counter("success")).isEqualTo(1.0);
+ assertThat(counter("failed")).isEqualTo(1.0);
+ }
+
+ @Test
+ @DisplayName("첫 정상 갱신 전에는 Snapshot 나이를 무한대로 노출한다")
+ void exposesInfiniteAgeBeforeFirstSuccessfulRefresh() {
+ assertThat(gauge("docgrid.operational.snapshot.age.seconds")).isPositive().isInfinite();
+ assertThat(gauge("docgrid.embedding.claimable.jobs")).isZero();
+ }
+
+ private double gauge(String name) {
+ return meterRegistry.get(name).gauge().value();
+ }
+
+ private double counter(String outcome) {
+ return meterRegistry.get("docgrid.operational.snapshot.refresh")
+ .tag("outcome", outcome)
+ .counter()
+ .count();
+ }
+}
diff --git a/docs/design/gimin-#334-operational-queue-gauges.md b/docs/design/gimin-#334-operational-queue-gauges.md
new file mode 100644
index 00000000..54cbf22b
--- /dev/null
+++ b/docs/design/gimin-#334-operational-queue-gauges.md
@@ -0,0 +1,87 @@
+# 비동기 Queue 운영 상태 Gauge 설계 (#334)
+
+closes #334
+
+## 문제
+
+상태 전이 Counter는 최근 성공·재시도·실패 증가를 보여주지만 현재 Queue가 줄고 있는지는 보여주지
+않는다. DB의 전체 `PENDING` 수를 그대로 노출하면 Retry backoff로 아직 실행할 수 없는 항목도 backlog로
+계산된다. 또한 Gauge가 scrape 순간 Repository를 조회하면 Prometheus 수집 빈도와 Backend 인스턴스 수에
+비례해 DB 부하가 증가한다.
+
+## 수집 경계
+
+```text
+전용 단일 Thread (기본 15초)
+ → 짧은 read-only Transaction
+ → Embedding + 유효 Worker aggregate SQL 1회
+ → RAG aggregate SQL 1회
+ → Sync Outbox aggregate SQL 1회
+ → 세 Query가 모두 성공하면 AtomicReference를 한 번에 교체
+
+Prometheus scrape
+ → Micrometer Gauge
+ → AtomicReference와 Clock만 읽음
+ → DB 조회 0회
+```
+
+집계 Transaction은 2초 timeout을 사용한다. 한 Query라도 실패하면 새 결과를 버리고 마지막 정상
+Snapshot을 유지한다. 이를 0으로 바꾸면 실제 backlog가 사라진 것처럼 보이므로 실패 Counter와 마지막
+정상 갱신 이후 시간으로 별도 감시한다.
+
+## Queue별 의미
+
+### Embedding
+
+- claim 가능: `status = PENDING`이며 `next_retry_at IS NULL OR next_retry_at <= now`
+- oldest 기준: 최초 실행은 `created_at`, Retry는 `next_retry_at`
+- 처리 중: `status = PROCESSING`
+- 유효 Worker: 상태가 `ACTIVE` 또는 `IDLE`이고 `last_heartbeat_at > now - dead-threshold`
+
+오래전에 생성된 Job이 backoff 종료 직후 즉시 오래된 backlog로 판정되지 않도록 Retry Job의 대기 나이는
+`next_retry_at`부터 계산한다.
+
+### RAG
+
+- 처리 중: `status = PROCESSING`
+- oldest 기준: `created_at`
+
+RAG는 현재 별도 대기 상태 없이 PROCESSING 생성 후 단일 Worker가 처리하므로 PROCESSING 자체가 Queue다.
+
+### Sync Outbox
+
+- claim 가능: `status = PENDING`이며 `available_at <= now`
+- oldest 기준: `available_at`
+- 처리 중: `status = PROCESSING`
+
+## Metric 계약
+
+| Metric | Type | 의미 |
+|---|---|---|
+| `docgrid_embedding_claimable_jobs` | Gauge | 현재 claim 가능한 Embedding Job 수 |
+| `docgrid_embedding_processing_jobs` | Gauge | 처리 중인 Embedding Job 수 |
+| `docgrid_embedding_oldest_claimable_age_seconds` | Gauge | 가장 오래된 claim 가능 Job의 대기 시간 |
+| `docgrid_embedding_active_workers` | Gauge | Heartbeat가 유효한 ACTIVE·IDLE Worker 수 |
+| `docgrid_rag_processing_jobs` | Gauge | PROCESSING RAG 응답 수 |
+| `docgrid_rag_oldest_processing_age_seconds` | Gauge | 가장 오래된 PROCESSING 응답의 나이 |
+| `docgrid_sync_outbox_claimable_events` | Gauge | 현재 claim 가능한 Outbox Event 수 |
+| `docgrid_sync_outbox_processing_events` | Gauge | 처리 중인 Outbox Event 수 |
+| `docgrid_sync_outbox_oldest_claimable_age_seconds` | Gauge | 가장 오래된 claim 가능 Event의 대기 시간 |
+| `docgrid_operational_snapshot_age_seconds` | Gauge | 마지막 정상 Snapshot 갱신 이후 시간 |
+| `docgrid_operational_snapshot_refresh_total{outcome}` | Counter | `success`, `failed` 갱신 결과 |
+
+동적 식별자나 오류 메시지 label은 없다. 여러 Backend가 같은 DB를 읽으면 동일한 Gauge가 인스턴스별로
+노출되므로 Prometheus 규칙은 `sum` 대신 `max by (cluster, environment)`를 사용한다.
+
+## 경보
+
+| 경보 | 조건 | 지속 시간 |
+|---|---|---:|
+| `DocGridEmbeddingWorkersUnavailable` | claim 가능 Job > 0, 유효 Worker = 0 | 1분 |
+| `DocGridEmbeddingQueueStalled` | oldest claim 가능 Job > 300초 | 5분 |
+| `DocGridRagQueueStalled` | oldest PROCESSING 응답 > 150초 | 30초 |
+| `DocGridSyncOutboxQueueStalled` | oldest claim 가능 Event > 120초 | 2분 |
+| `DocGridOperationalSnapshotStale` | 마지막 정상 Snapshot > 60초 | 1분 |
+
+`snapshot_age`는 첫 정상 갱신 전 `+Inf`다. Backend가 올라왔지만 DB 집계를 한 번도 완료하지 못한 경우도
+수집 정상으로 오인하지 않기 위해서다.
diff --git a/docs/test-results/gimin-#334-operational-queue-gauges.md b/docs/test-results/gimin-#334-operational-queue-gauges.md
new file mode 100644
index 00000000..3d2c6f69
--- /dev/null
+++ b/docs/test-results/gimin-#334-operational-queue-gauges.md
@@ -0,0 +1,74 @@
+# 비동기 Queue 운영 상태 Gauge 검증 결과 (#334)
+
+검증일: 2026-09-13
+
+## Java 대상 테스트
+
+```bash
+DB_PORT=55433 \
+JWT_SECRET=docgrid-operational-metrics-test-secret-key-2026-with-at-least-32-bytes \
+./gradlew test \
+ --tests 'com.opensource.docgrid.global.observability.OperationalMetricsTest' \
+ --tests 'com.opensource.docgrid.global.observability.OperationalMetricsSnapshotRefresherTest' \
+ --tests 'com.opensource.docgrid.global.observability.OperationalMetricsSnapshotQueryServiceTest'
+```
+
+결과: **6 tests, 실패 0**
+
+검증 범위:
+
+- Embedding과 Sync Outbox의 미래 Retry·available 시각 제외
+- Retry Job의 오래된 `created_at` 대신 실행 가능 시각을 oldest 기준으로 사용
+- PROCESSING 수와 Queue가 비었을 때의 0·null 경계
+- Worker 상태와 Heartbeat 만료 경계를 함께 적용
+- 세 집계 성공 후 Atomic Snapshot 게시
+- 집계 실패 시 마지막 정상 Gauge 보존과 실패 Counter 증가
+- Gauge 반복 조회 시 QueryService 호출 0회
+- 첫 정상 갱신 전 Snapshot 나이 `+Inf`
+
+## Prometheus 규칙 테스트
+
+```bash
+docker run --rm --entrypoint=promtool \
+ -v "$PWD/monitoring/prometheus:/etc/prometheus:ro" \
+ prom/prometheus:v3.5.5 \
+ test rules /etc/prometheus/tests/docgrid-operational-alerts.test.yml
+```
+
+결과: **SUCCESS**
+
+검증 범위:
+
+- 같은 DB Gauge를 두 Backend가 노출해도 cluster당 경보 1개 생성
+- claim 가능 Job이 있고 Worker가 없을 때만 1분 후 경보
+- Embedding·RAG·Outbox 임계값과 `for` 지속 시간
+- 정상 RAG 처리 시간은 경보 제외
+- Snapshot stale 지속 시간
+
+## 전체 회귀 테스트
+
+```bash
+DB_PORT=55433 \
+JWT_SECRET=docgrid-operational-metrics-test-secret-key-2026-with-at-least-32-bytes \
+./gradlew test
+```
+
+결과: **1171 tests, 실패 0, errors 0, skipped 0**
+
+추가 검증:
+
+- `promtool check config`: 성공, rule file 3개·총 19개 규칙
+- `promtool check rules`: 성공
+- Backend·Pipeline·Operational 전체 rule suite: **SUCCESS**
+- 기본 Compose와 monitoring profile `docker compose config --quiet`: 성공
+
+## 실제 Management endpoint
+
+격리 Schema에서 Backend를 `SERVER_PORT=18080`, `MANAGEMENT_PORT=18081`로 실행한 뒤
+`GET /actuator/prometheus`를 확인했다.
+
+- Queue·Worker Gauge 9개 노출
+- `docgrid_operational_snapshot_age_seconds` 노출
+- `docgrid_operational_snapshot_refresh_total{outcome="success|failed"}` 노출
+- 빈 Queue의 수와 oldest age는 0
+- 15초 주기 갱신 성공 Counter 증가, 실패 Counter 0
diff --git a/monitoring/prometheus/README.md b/monitoring/prometheus/README.md
index 7e31978f..65a83020 100644
--- a/monitoring/prometheus/README.md
+++ b/monitoring/prometheus/README.md
@@ -76,11 +76,44 @@ HTTP 오류율에서는 Streamable HTTP 특성이 다른 `/mcp`를 제외한다.
| `DocGridRagProviderFallbackSpike` | 10분간 provider fallback 3회 이상 | 1분 |
| `DocGridRagTimeoutSweepSpike` | 10분간 timeout 강제 종료 3회 이상 | 1분 |
| `DocGridSyncOutboxTerminalFailure` | 15분간 새로운 최종 실패 1회 이상 | 즉시 |
+| `DocGridEmbeddingWorkersUnavailable` | claim 가능 Job이 있지만 유효 Worker가 없음 | 1분 |
+| `DocGridEmbeddingQueueStalled` | 가장 오래된 claim 가능 Job 나이가 5분 초과 | 5분 |
+| `DocGridRagQueueStalled` | 가장 오래된 PROCESSING 응답 나이가 150초 초과 | 30초 |
+| `DocGridSyncOutboxQueueStalled` | 가장 오래된 claim 가능 Event 나이가 2분 초과 | 2분 |
+| `DocGridOperationalSnapshotStale` | DB 운영 Snapshot을 1분 넘게 갱신하지 못함 | 1분 |
Counter는 DB 상태 전이를 수행한 Transaction이 커밋된 뒤에만 증가한다. Embedding 실패의
`failure_type`은 고정 enum이며 `retryable` label로 사용자 문서 오류와 운영 장애를 구분한다.
Job ID, 오류 메시지와 사용자 입력은 label에 포함하지 않는다.
+### 현재 Queue 상태 Gauge
+
+Backend는 기본 15초마다 별도 단일 Thread에서 Queue별 aggregate query를 실행하고 마지막 정상 결과를
+메모리에 보관한다. `/actuator/prometheus`의 Gauge callback은 이 메모리만 읽으므로 scrape 횟수가 DB
+조회 횟수를 늘리지 않는다. 주기는 `MANAGEMENT_METRICS_SNAPSHOT_INTERVAL`로 조정할 수 있다.
+
+`PENDING` 전체 개수가 아니라 현재 시각에 실제 claim 가능한 항목만 backlog에 포함한다. Embedding
+Retry는 `next_retry_at`, Sync Outbox는 `available_at`이 미래이면 제외한다. 가장 오래된 항목의 나이도
+최초 생성 시각과 Retry 실행 가능 시각을 구분해 계산한다.
+
+| Metric | 의미 |
+|---|---|
+| `docgrid_embedding_claimable_jobs` | 지금 claim 가능한 Embedding Job 수 |
+| `docgrid_embedding_processing_jobs` | 처리 중인 Embedding Job 수 |
+| `docgrid_embedding_oldest_claimable_age_seconds` | 가장 오래된 claim 가능 Job의 대기 시간 |
+| `docgrid_embedding_active_workers` | Heartbeat가 만료되지 않은 ACTIVE·IDLE Worker 수 |
+| `docgrid_rag_processing_jobs` | PROCESSING RAG 응답 수 |
+| `docgrid_rag_oldest_processing_age_seconds` | 가장 오래된 PROCESSING 응답의 나이 |
+| `docgrid_sync_outbox_claimable_events` | 지금 claim 가능한 Sync Outbox Event 수 |
+| `docgrid_sync_outbox_processing_events` | 처리 중인 Sync Outbox Event 수 |
+| `docgrid_sync_outbox_oldest_claimable_age_seconds` | 가장 오래된 claim 가능 Event의 대기 시간 |
+| `docgrid_operational_snapshot_age_seconds` | 마지막 정상 DB Snapshot 이후 경과 시간 |
+| `docgrid_operational_snapshot_refresh_total{outcome}` | Snapshot 갱신 성공·실패 횟수 |
+
+여러 Backend가 같은 DB를 수집하면 동일 Gauge가 인스턴스 수만큼 노출된다. 번들 규칙은 이 값을
+합산하지 않고 `max by (cluster, environment)`로 평가해 backlog를 중복 계산하지 않는다. 갱신 실패는
+마지막 정상 값을 유지하며, 첫 성공 전과 장시간 실패는 `snapshot_age` 경보로 드러난다.
+
## 설정 검증
로컬에 promtool을 설치하지 않아도 고정된 Prometheus 이미지로 검사할 수 있다.
@@ -109,5 +142,10 @@ docker run --rm --entrypoint=promtool \
prom/prometheus:v3.5.5 \
test rules /etc/prometheus/tests/docgrid-pipeline-alerts.test.yml
+docker run --rm --entrypoint=promtool \
+ -v "$PWD/monitoring/prometheus:/etc/prometheus:ro" \
+ prom/prometheus:v3.5.5 \
+ test rules /etc/prometheus/tests/docgrid-operational-alerts.test.yml
+
docker compose config
```
diff --git a/monitoring/prometheus/rules/docgrid-pipeline-alerts.yml b/monitoring/prometheus/rules/docgrid-pipeline-alerts.yml
index 5eff92a7..2769582d 100644
--- a/monitoring/prometheus/rules/docgrid-pipeline-alerts.yml
+++ b/monitoring/prometheus/rules/docgrid-pipeline-alerts.yml
@@ -74,3 +74,84 @@ groups:
annotations:
summary: A DocGrid Sync Outbox event reached terminal failure
description: At least one Sync Outbox event exhausted retries during the last fifteen minutes.
+
+ - alert: DocGridEmbeddingWorkersUnavailable
+ expr: |
+ max by (cluster, environment) (
+ docgrid_embedding_claimable_jobs
+ ) > 0
+ and
+ max by (cluster, environment) (
+ docgrid_embedding_active_workers
+ ) == 0
+ for: 1m
+ labels:
+ severity: critical
+ service: embedding-worker
+ annotations:
+ summary: DocGrid embedding queue has no active workers
+ description: Claimable embedding jobs have been waiting without a live worker for one minute.
+
+ - alert: DocGridEmbeddingQueueStalled
+ expr: |
+ max by (cluster, environment) (
+ docgrid_embedding_claimable_jobs
+ ) > 0
+ and
+ max by (cluster, environment) (
+ docgrid_embedding_oldest_claimable_age_seconds
+ ) > 300
+ for: 5m
+ labels:
+ severity: warning
+ service: embedding-worker
+ annotations:
+ summary: DocGrid embedding queue is stalled
+ description: The oldest claimable embedding job has remained eligible for more than five minutes.
+
+ - alert: DocGridRagQueueStalled
+ expr: |
+ max by (cluster, environment) (
+ docgrid_rag_processing_jobs
+ ) > 0
+ and
+ max by (cluster, environment) (
+ docgrid_rag_oldest_processing_age_seconds
+ ) > 150
+ for: 30s
+ labels:
+ severity: warning
+ service: rag-worker
+ annotations:
+ summary: DocGrid RAG queue is stalled
+ description: The oldest RAG response has remained in PROCESSING for more than 150 seconds.
+
+ - alert: DocGridSyncOutboxQueueStalled
+ expr: |
+ max by (cluster, environment) (
+ docgrid_sync_outbox_claimable_events
+ ) > 0
+ and
+ max by (cluster, environment) (
+ docgrid_sync_outbox_oldest_claimable_age_seconds
+ ) > 120
+ for: 2m
+ labels:
+ severity: warning
+ service: sync-outbox
+ annotations:
+ summary: DocGrid Sync Outbox queue is stalled
+ description: The oldest claimable Sync Outbox event has remained eligible for more than two minutes.
+
+ - alert: DocGridOperationalSnapshotStale
+ expr: |
+ max by (cluster, environment) (
+ docgrid_operational_snapshot_age_seconds
+ ) > 60
+ for: 1m
+ labels:
+ severity: critical
+ service: docgrid-backend
+ annotations:
+ summary: DocGrid operational metrics snapshot is stale
+ description: The Backend has not completed an operational DB snapshot refresh for more than one minute.
diff --git a/monitoring/prometheus/tests/docgrid-operational-alerts.test.yml b/monitoring/prometheus/tests/docgrid-operational-alerts.test.yml
new file mode 100644
index 00000000..4fe846f4
--- /dev/null
+++ b/monitoring/prometheus/tests/docgrid-operational-alerts.test.yml
@@ -0,0 +1,129 @@
+rule_files:
+ - /etc/prometheus/rules/docgrid-pipeline-alerts.yml
+
+evaluation_interval: 15s
+
+tests:
+ - name: claimable embedding work without workers alerts once per cluster
+ interval: 15s
+ input_series:
+ - series: 'docgrid_embedding_claimable_jobs{instance="backend-a:8081",environment="production",cluster="docgrid-production"}'
+ values: '4x8'
+ - series: 'docgrid_embedding_claimable_jobs{instance="backend-b:8081",environment="production",cluster="docgrid-production"}'
+ values: '4x8'
+ - series: 'docgrid_embedding_active_workers{instance="backend-a:8081",environment="production",cluster="docgrid-production"}'
+ values: '0x8'
+ - series: 'docgrid_embedding_active_workers{instance="backend-b:8081",environment="production",cluster="docgrid-production"}'
+ values: '0x8'
+ alert_rule_test:
+ - eval_time: 45s
+ alertname: DocGridEmbeddingWorkersUnavailable
+ exp_alerts: []
+ - eval_time: 1m
+ alertname: DocGridEmbeddingWorkersUnavailable
+ exp_alerts:
+ - exp_labels:
+ alertname: DocGridEmbeddingWorkersUnavailable
+ cluster: docgrid-production
+ environment: production
+ service: embedding-worker
+ severity: critical
+ exp_annotations:
+ summary: DocGrid embedding queue has no active workers
+ description: Claimable embedding jobs have been waiting without a live worker for one minute.
+
+ - name: duplicate backend gauges produce one embedding stall alert
+ interval: 15s
+ input_series:
+ - series: 'docgrid_embedding_claimable_jobs{instance="backend-a:8081",environment="production",cluster="docgrid-production"}'
+ values: '3x24'
+ - series: 'docgrid_embedding_claimable_jobs{instance="backend-b:8081",environment="production",cluster="docgrid-production"}'
+ values: '3x24'
+ - series: 'docgrid_embedding_oldest_claimable_age_seconds{instance="backend-a:8081",environment="production",cluster="docgrid-production"}'
+ values: '360x24'
+ - series: 'docgrid_embedding_oldest_claimable_age_seconds{instance="backend-b:8081",environment="production",cluster="docgrid-production"}'
+ values: '360x24'
+ alert_rule_test:
+ - eval_time: 5m
+ alertname: DocGridEmbeddingQueueStalled
+ exp_alerts:
+ - exp_labels:
+ alertname: DocGridEmbeddingQueueStalled
+ cluster: docgrid-production
+ environment: production
+ service: embedding-worker
+ severity: warning
+ exp_annotations:
+ summary: DocGrid embedding queue is stalled
+ description: The oldest claimable embedding job has remained eligible for more than five minutes.
+
+ - name: rag threshold distinguishes normal processing from a stalled queue
+ interval: 15s
+ input_series:
+ - series: 'docgrid_rag_processing_jobs{instance="backend-a:8081",environment="healthy",cluster="docgrid-healthy"}'
+ values: '1x8'
+ - series: 'docgrid_rag_oldest_processing_age_seconds{instance="backend-a:8081",environment="healthy",cluster="docgrid-healthy"}'
+ values: '149x8'
+ - series: 'docgrid_rag_processing_jobs{instance="backend-b:8081",environment="production",cluster="docgrid-production"}'
+ values: '1x8'
+ - series: 'docgrid_rag_oldest_processing_age_seconds{instance="backend-b:8081",environment="production",cluster="docgrid-production"}'
+ values: '180x8'
+ alert_rule_test:
+ - eval_time: 30s
+ alertname: DocGridRagQueueStalled
+ exp_alerts:
+ - exp_labels:
+ alertname: DocGridRagQueueStalled
+ cluster: docgrid-production
+ environment: production
+ service: rag-worker
+ severity: warning
+ exp_annotations:
+ summary: DocGrid RAG queue is stalled
+ description: The oldest RAG response has remained in PROCESSING for more than 150 seconds.
+
+ - name: sync outbox age waits for two minutes
+ interval: 15s
+ input_series:
+ - series: 'docgrid_sync_outbox_claimable_events{instance="backend:8081",environment="production",cluster="docgrid-production"}'
+ values: '2x12'
+ - series: 'docgrid_sync_outbox_oldest_claimable_age_seconds{instance="backend:8081",environment="production",cluster="docgrid-production"}'
+ values: '180x12'
+ alert_rule_test:
+ - eval_time: 1m45s
+ alertname: DocGridSyncOutboxQueueStalled
+ exp_alerts: []
+ - eval_time: 2m
+ alertname: DocGridSyncOutboxQueueStalled
+ exp_alerts:
+ - exp_labels:
+ alertname: DocGridSyncOutboxQueueStalled
+ cluster: docgrid-production
+ environment: production
+ service: sync-outbox
+ severity: warning
+ exp_annotations:
+ summary: DocGrid Sync Outbox queue is stalled
+ description: The oldest claimable Sync Outbox event has remained eligible for more than two minutes.
+
+ - name: stale operational snapshot alerts after one minute
+ interval: 15s
+ input_series:
+ - series: 'docgrid_operational_snapshot_age_seconds{instance="backend:8081",environment="production",cluster="docgrid-production"}'
+ values: '75x8'
+ alert_rule_test:
+ - eval_time: 45s
+ alertname: DocGridOperationalSnapshotStale
+ exp_alerts: []
+ - eval_time: 1m
+ alertname: DocGridOperationalSnapshotStale
+ exp_alerts:
+ - exp_labels:
+ alertname: DocGridOperationalSnapshotStale
+ cluster: docgrid-production
+ environment: production
+ service: docgrid-backend
+ severity: critical
+ exp_annotations:
+ summary: DocGrid operational metrics snapshot is stale
+ description: The Backend has not completed an operational DB snapshot refresh for more than one minute.