Skip to content

feat: CCP-6026-Redis Monitoring - #253

Merged
barrfalk merged 27 commits into
mainfrom
CCP-6026
Oct 6, 2026
Merged

barrfalk merged 27 commits into
mainfrom
CCP-6026

Conversation

@barrfalk

@barrfalk barrfalk commented Oct 1, 2026 •

Copy link
Copy Markdown
Collaborator

Description

CCP-6026. This PR adds an admin Monitoring page for queues, workers and Redis, and makes bulk delivery recover from crashes, restarts and lost jobs without losing or duplicating messages.

Monitoring page (/admin/monitoring, NOTIFY_ADMIN SSO role only)

  • Overview tab:
    • summary cards for the last hour: delivery time median and p95, messages sent/failed, sustained sending rate, failure %, requests received;
    • Right now: pending messages per channel and an estimated clear time;
    • progress of batches in flight, per recipient.
  • Failures tab: recent failed jobs with tenant names. Email addresses and phone numbers in provider error messages are masked.
  • System tab:
    • per-queue backlog, rates and health reasons;
    • worker pods from Redis heartbeats, including a "Shutting down" badge while a pod drains;
    • Delivery recovery: the reconciler's last pass, retried / re-queued / given-up counts for the last hour, and recent recoveries;
    • Redis memory, eviction policy and connections.
  • Live updates: over SSE (GET /api/v1/frontend/admin/monitoring/queues/events), driven by Bull's Redis pub/sub and throttled to one refresh every 2s. Falls back to polling every 15s if the stream drops. Warnings from every tab are shown in a banner at the top.
  • Backend: GET /api/v1/frontend/admin/monitoring/queues (NotifyAdminGuard), plus two new Kong routes in routes.yaml. Every read goes to Redis or Postgres, so any pod gives the same answer.
  • Bull metrics are turned on (QUEUE_METRICS_MAX_DATA_POINTS), and a per-minute "added" counter supplies the fill rate.

Delivery robustness

  • Recipients live in Postgres, not in job payloads.
    • A merge batch job carries only shared content. Each recipient's params go in the new notification_request_detail.params column (V73).
    • Batch size defaults to 25 (BATCH_SIZE).
  • Retries are idempotent.
    • A batch retry skips recipients already sent.
    • A re-run of ingestion writes no rows twice and skips channels that have already finished.
  • Graceful shutdown:
    • On SIGTERM a pod stops taking jobs and drains the ones it holds, for up to QUEUE_SHUTDOWN_DRAIN_MS (105s, from the chart). terminationGracePeriodSeconds is 120.
    • The worker heartbeat reports draining.
  • maxStalledCount is 3 (QUEUE_MAX_STALLED_COUNT), so a pod restart no longer fails a long batch outright.
  • DeliveryReconcilerService replaces PendingNotificationRetryService.
    • One single-flight pass every 60s, behind a Redis lock. It compares what Postgres says is owed with what the queues hold.

    • A live job is left alone; a failed job is retried; a missing job is rebuilt. After 5 attempts, whatever is still owed is marked failed and the request is settled.

    • Five kinds of stuck work:

      Kind Found by
      pending request PENDING > 30s (Redis was down at acceptance)
      ingestion request QUEUED > 10 min and never fanned out
      scheduled request SCHEDULED > 10 min past its send time (its delayed job was lost or failed)
      batch merge batch recipients in flight > 10 min
      delivery plain-send recipients in flight > 10 min
    • Each pass is recorded in Redis for the Monitoring page.

  • CHES calls time out after CHES_TIMEOUT_MS (120s) instead of hanging a worker slot.
  • A provider outage no longer fails recipients (email and SMS). Errors are classified at the adapter: a 5xx, 429, 408, or a network error before the request was sent is a TransientDeliveryError. The email worker then pauses the batch, leaving unsent recipients owed for Bull's retry and then the reconciler, rather than marking each one failed. SMS works the same way: ACS and Twilio SDK errors are classified, and recipients a provider could not take are flagged transient and left owed. Previously an 11-second CHES restart failed about 1,500 recipients in a 2,000-recipient test. An email timeout still fails just that recipient, because CHES may have queued it and a resend could duplicate it.
  • A shared circuit breaker sits in front of CHES and of the SMS provider (Redis, every pod). 5 transient failures or timeouts in 10s open it: sends wait instead of failing, and after 30s one probe send decides whether to resume. A local replay of a 5-second CHES outage over 60 sends hit CHES 5 times during the outage and failed no one; all 60 were sent once it recovered.
  • Fixes found along the way:
    • A batch that exhausted its retries marked every recipient of the whole send failed, including ones already sent. It now fails only its own unsent recipients.
    • CHES 429 was treated as a permanent 400.
    • A CHES token could expire while a send waited for a slot. The token is now fetched inside the slot, and a 401 is retried once with a new token.
  • The Monitoring page gains Sends in progress (one row per request, with batches rolled up and a time-left estimate) and a Delivery providers panel (circuit state for CHES and ACS, and CHES requests in flight). The "arriving faster" and "oldest job waited" warnings no longer fire on a normal large send: they need sustained growth, or a queue that isn't moving.
  • CHES calls are capped across all pods (CHES_MAX_CONCURRENT_REQUESTS, 5) by RedisConcurrencyLimiter. CHES accepts about one email a second however many requests arrive together. Without a cap, a large merge piles requests up inside CHES until some pass the timeout: with a 30s timeout, a 400-recipient test failed one recipient this way. The cap moves that wait to our side, where nothing times out, and sends just as many emails a minute. If Redis is unreachable, sends go ahead without the cap.

Fixes found along the way

  • Every queue worker ran twice in each pod. GcNotifyModule.forRoot imported QueueModule/NotifyModule through forwardRef, which created a second instance. The e2e spec now asserts each module is created once.
  • Refreshing any /admin/* page returned 403. WAF rule 1004 in frontend/coraza.conf no longer blocks /admin. That URL is the SPA's admin area; access is enforced by the API. routes/waf-routes.spec.ts checks every route against rules 1004 and 1007.
  • Load-test bind:
  • Pending sweep never ran its re-queued jobs: it added named 'process' jobs that no handler picked up. It also sent scheduled sends immediately. Both are fixed by the reconciler.
  • Heap: the image no longer hard-codes --max-old-space-size=150. The chart sets NODE_OPTIONS from backend.heapMb (512). The memory request is 200Mi by default and 384Mi in TEST/PROD (merge.yml).
  • Shutdown diagnostics: installShutdownDiagnostics logs the signal, the exit code and the last error on any nonzero exit.

Type of change

  • Bug fix (non-breaking change which fixes an issue)
  • New feature (non-breaking change which adds functionality)
  • Breaking change (fix or feature that would cause existing functionality to not work as expected)
  • This change requires a documentation update
  • Documentation update

Rollout notes

  1. Deploy TEST and PROD with no bulk send in flight.
    • Those environments use a rolling update. During one, an old pod can pick up a merge batch queued by a new pod.
    • The new batch job no longer carries its recipients, so the old pod would not find them.
    • DEV and PR environments use Recreate and are unaffected.
  2. Migrations are additive:
    • V72: a partial index on notification_request_detail(last_attempt_at) for sent/failed rows, used by the monitoring counts;
    • V73: a nullable params jsonb column on notification_request_detail.
  3. Env vars are already set in app-api-secrets in f6bc3f-dev, -test and -prod. Nothing to do before merge.
    • BATCH_SIZE=25, QUEUE_MAX_STALLED_COUNT=3, QUEUE_METRICS_MAX_DATA_POINTS=1440, CHES_TIMEOUT_MS=120000, CHES_MAX_CONCURRENT_REQUESTS=5
    • DELIVERY_RECONCILE_INTERVAL_MS=60000, DELIVERY_RECONCILE_STALE_MS=600000, DELIVERY_RECONCILE_MAX_ATTEMPTS=5, DELIVERY_RECONCILE_PENDING_STALE_MS=30000
    • MONITORING_* thresholds
    • The old BATCH_RECONCILE_* keys have been removed.
  4. Chart-derived values (not secrets): NODE_OPTIONS, QUEUE_SHUTDOWN_DRAIN_MS, terminationGracePeriodSeconds: 120, backend.memoryRequest.
  5. Gateway: the two new monitoring routes in api-gateway/templates/routes.yaml go out with the usual gateway publish.

How Has This Been Tested?

  • New unit tests
  • New integrated tests
  • New component tests
  • New end-to-end tests
  • New user flow tests
  • No new tests are required
  • Manual tests (description below)
  • Updated existing tests

Automated checks (all green locally):

  • Backend npm run test:unit:cov: 1946 passed. This includes test/app.e2e-spec.ts against local Postgres and Redis.
  • Backend npm run lint and npm run build pass.
  • Frontend npx vitest run: 665 passed. Frontend npm run lint and npm run build pass.

Against real Postgres and Redis:

  • Every reconciler kind: a stuck job is retried, re-queued or given up as expected, and the pass is recorded and read back.
  • A scheduled send whose job was lost is re-queued. One whose job is still delayed, one not yet due, and one still within the 10-minute grace are left alone.
  • The delivery-time and sending-rate SQL.
  • An end-to-end merge of 60 recipients.

In the PR environment (OpenShift):

  • Deleted the backend pod in the middle of a bulk send. The pod drained ("Draining queue workers (up to 105s)") and the batches in progress finished. The new pod sent the remaining recipients and the request reached COMPLETED.
  • Sends left pending by an earlier build were picked up and delivered by the reconciler.

Checklist

  • I have read the CONTRIBUTING doc
  • I have performed a self-review of my own code
  • I have commented my code, particularly in hard-to-understand areas
  • I have made corresponding changes to the documentation
  • My changes generate no new warnings
  • I have added tests that prove my fix is effective or that my feature works
  • New and existing unit tests pass locally with my changes
  • Any dependent changes have already been accepted and merged

Further comments

Why recipients moved out of job payloads. A merge job used to carry its full recipient list in Redis, so a large send was held twice: in Redis and in Postgres. Retries also had no record of who had already been sent to. With Postgres as the record of what is owed, every recovery path (Bull retry, stalled job, reconciler) can be repeated safely. Delivery is at-least-once, and workers are idempotent.

Known limits, not addressed here:

  • Drain window vs. batch length. At 25 recipients a merge batch takes about 2 minutes, longer than the 105s drain window. A pod that stops mid-batch leaves its last few recipients to the next pod. That is verified to work, but those recipients go out later. Lowering BATCH_SIZE or lengthening the grace period in TEST/PROD would close the gap.
  • Rescheduling (cancelOrRescheduleNotification) updates delayed_send_time but doesn't move the delayed Bull job. The send still goes out at the original time. This predates this PR; the new scheduled finder neither causes nor fixes it.
  • Shutdown diagnostics stay in place. The cluster pod that died about 1s into shutdown hasn't recurred since the heap change and main's Loki transport change. They're kept so a recurrence is diagnosable.

AGENTS.md and CLAUDE.md are git-ignored, so their updates (the env-var rule and the reconciler notes) aren't in this diff.


Thanks for the PR!

Deployments, as required, will be available below:

Please create PRs in draft mode. Mark as ready to enable:

After merge, new images are deployed in:

🤖 Generated with Claude Code


Thanks for the PR!

Deployments, as required, will be available below:

Please create PRs in draft mode. Mark as ready to enable:

After merge, new images are deployed in:

Adds backend foundations for queue monitoring: new monitoring thresholds and DTO schemas, Bull throughput helpers, and worker heartbeat publishing to Redis so admin views can report live worker state across pods. Queue creation now enables Bull per-minute metrics, tracks added jobs with TTL-based Redis counters, and exposes a configurable history window (`QUEUE_METRICS_MAX_DATA_POINTS`, default 1440). QueueModule now starts/stops heartbeat publishing and registers tracked queues with their concurrency. Includes unit coverage to ensure queue metrics are enabled and default retention is one day.
@barrfalk
barrfalk marked this pull request as draft October 1, 2026 22:19
barrfalk and others added 18 commits October 1, 2026 16:28
Adds admin queue monitoring for recipient-level backlog and service health, plus live SSE change signals when queue state moves. The backend now aggregates per-channel message throughput, tracks stalled active jobs, exposes active merge-batch progress from Bull jobs, and includes a V72 index to support the hourly stats query. The frontend dashboard shows message backlog, active batches, and drill-down metrics, with matching specs covering the new monitoring views and progress reporting. No public API contract changes; this is admin-only monitoring.
Adds an overview layer to queue monitoring with last-hour totals, request counts, and delivery-time percentiles from backend data, and updates the admin UI to use tabbed Overview/Failures/System views with summary cards and URL-persisted tab state. It also fixes load-test API key autobind by introducing a dedicated guard so the route works with the global JwtGuard while still enforcing the feature flag and gateway-header checks. Finally, Coraza WAF path rules were adjusted to allow SPA /admin routes, with a route-vs-WAF regression test to prevent refresh/direct-link 403s.
Harden `autoBindApiKeyForLoadTest` so it no longer rebinds an existing credential from another tenant to the load-test tenant. The service now logs and throws a `ConflictException` when a key is already owned by a different tenant, while still allowing idempotent reuse of keys already bound to the load-test tenant. Added focused unit tests covering new bind, idempotent behavior, cross-tenant rejection, and environment guardrails in test/prod namespaces.
Update worker heartbeat tracking to use per-job IDs instead of a simple counter so local stalled-job `failed` events for jobs from other pods do not corrupt active/failed metrics. Added backend tests for unknown finished jobs and cross-pod stalled-job scenarios.

In Queue Monitoring, hide the manual Refresh button while the live stream is connected (it already auto-refetches) and show it again when the stream drops. Added a frontend test to cover this visibility toggle.
Replaces the monitoring overview’s `sentPerMinute` with `sendingRate` (`perMinute` + `activeMinutes`) on both backend and frontend DTOs/interfaces. The new backend calculation averages only minutes with sends and ignores the current partial minute unless it is the only active minute, preventing idle time from diluting throughput; UI copy and specs were updated to show the new metric and cover null/aggregation behavior.
Improve the Redis stats UI to avoid misleading fragmentation ratios on small instances. Fragmentation is now hidden with a “Not meaningful at low memory use” hint below 100MB used memory, and high fragmentation (>1.5) is explicitly flagged with guidance. Added/updated QueueMonitoring tests to cover normal, low-memory, and high-fragmentation scenarios.
Add graceful queue shutdown so workers mark themselves as draining, finish in-flight jobs within the pod grace period, and then remove their heartbeat instead of appearing stalled. Expose the new draining state through backend monitoring and the admin UI. Also make retried email merge batches skip recipients already sent by an earlier attempt to avoid duplicate delivery during pod replacement or shutdown.
Remove unnecessary `forwardRef()` wrappers from `GcNotifyModule.forRoot()` when importing `NotifyModule` and `QueueModule`, preventing duplicate static module instantiation. Add an e2e guard test that scans the Nest `ModulesContainer` (excluding `TypeOrmModule` dynamic `forFeature` entries) to ensure non-TypeORM modules are only created once, protecting against double worker/heartbeat startup and related false stalled-job monitoring.
Bull's default maxStalledCount of 1 fails a job the second time its pod
dies mid-run, e.g. a pod deleted and then restarted during one bulk send.
The two merge batches in PR-253 ended in the failed set that way, leaving
80 recipients pending. Retries are safe now that merge batches skip
recipients already sent. Overridable with QUEUE_MAX_STALLED_COUNT.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Refactors mail-merge queue payload handling so jobs carry shared content/params only, while per-recipient params are persisted on `notification_request_detail` (with migration `V73`). Adds merge reconstruction helpers and updates pending-notification retry to rebuild proper merge ingestion payloads and enqueue unnamed jobs so workers can process them.

Introduces `BatchReconcilerService` to detect stale merge batches, retry/requeue missing or failed batch jobs from stored request data, and fail remaining recipients after a capped number of recovery attempts to prevent batches from staying in-flight indefinitely. Includes new unit coverage for merge batch builders, batch reconciliation behavior, and pending retry merge rebuilding.
Introduce a shutdown diagnostics helper that logs synchronous stderr markers when SIGTERM is received and when the process exits with a non-zero code. It now records the last app error from `StructuredLoggerService` so failed shutdowns include likely cause/context, wires diagnostics into app bootstrap, and adds unit tests covering SIGTERM, non-zero exit reporting, and clean exits.
Moves the Node heap limit out of the backend image and into Helm so environments can tune it without rebuilding. The deployment now sets NODE_OPTIONS from `backend.heapMb` (default 512), and backend memory requests are templated via `backend.memoryRequest` (default 200Mi). TEST and PROD deploy params in `merge.yml` now override memory request to 384Mi while keeping autoscaling settings unchanged.
Replace the separate pending-notification and batch recovery services with a single DeliveryReconciler that restores stuck pending, ingestion, batch, and plain-delivery work from Postgres. It rebuilds scheduled sends correctly, avoids duplicating detail rows or re-sending finished channels on ingestion reruns, and adds CHES request deadlines so hung token or email calls fail fast with clear 504 errors.
Adds Redis-backed reconcile activity tracking (last pass, action counts, recent actions) and exposes it through the admin monitoring API with new reconciler DTOs and status evaluation. Updates the admin Monitoring UI to show a Delivery recovery panel, include reconciler warnings in the top banner, and display recent recoveries. Also extends delivery reconciliation to handle overdue scheduled sends and adds comprehensive backend/frontend test coverage for the new behavior.
Add a Redis-backed concurrency limiter and apply it to CHES email POST requests so in-flight CHES calls are capped across pods. The limiter uses expiring leases, releases slots on success/failure, and fails open if Redis is unavailable.

Update CHES config defaults by increasing `CHES_TIMEOUT_MS` to 120000 and adding `CHES_MAX_CONCURRENT_REQUESTS` (default 5, `0` disables limiting). Extend tests to cover limiter behavior and verify CHES adapter uses the shared limit only when enabled.
Introduces shared transient delivery error handling and a Redis-backed circuit breaker, then applies it to CHES and SMS so provider outages pause sends instead of failing recipients immediately. SMS adapter wiring now always runs through a breaker, provider/HTTP/network failures are classified consistently, and delivery workers keep outage-affected recipients owed (including merge batches) while only failing unsent rows when retries are exhausted for non-transient errors. Monitoring was expanded to show provider circuit state and sends-in-progress rollups (replacing active batch rows), with backend/frontend DTO, service, UI, and test updates to match.
@barrfalk
barrfalk marked this pull request as ready for review October 6, 2026 16:47
Comment thread backend/src/queue/workers/sms-delivery.worker.spec.ts Fixed
barrfalk and others added 4 commits October 6, 2026 12:02
This change renames the shutdown timeout env var to QUEUE_SHUTDOWN_DRAIN_MS and tightens the graceful shutdown sequence. Queue workers are paused and closed in stages (ingestion → delivery → webhooks) with a bounded drain window before app shutdown, while timed-out jobs are logged as re-run elsewhere. This avoids partial merge batches being left mid-flight during deploys or scale-downs.
Removes long TEST and PROD `params` overrides from `merge.yml` so rollout strategy, autoscaling, PDB behavior, and memory sizing come from `values-test.yaml` and `values-prod.yaml` directly. Both env value files now explicitly set `backend.memoryRequest` to `384Mi` and update quota/headroom comments to match, keeping deployment shape and resource planning reviewable in git.
Co-authored-by: Copilot Autofix powered by AI <223894421+github-code-quality[bot]@users.noreply.github.com>
@revanth-banala
revanth-banala self-requested a review October 6, 2026 19:33
@barrfalk
barrfalk merged commit 80c257d into main Oct 6, 2026
26 checks passed
@barrfalk
barrfalk deleted the CCP-6026 branch October 6, 2026 19:35
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants