Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
7f83808
Add queue monitoring metrics and heartbeats
barrfalk Oct 1, 2026
e5a3d5b
Update app.module.ts
barrfalk Oct 1, 2026
574381f
Add live queue monitoring and batch progress
barrfalk Oct 1, 2026
128a57a
Revamp monitoring view and load-test autobind
barrfalk Oct 2, 2026
879feb2
Prevent load-test auto-bind key hijacking
barrfalk Oct 2, 2026
d01afdb
Fix queue heartbeat and live refresh behavior
barrfalk Oct 2, 2026
a70e5e3
Use active-minute send rate in monitoring
barrfalk Oct 2, 2026
83a57c8
Refine Redis fragmentation signal in monitoring
barrfalk Oct 2, 2026
6624da5
Drain queue workers on shutdown
barrfalk Oct 4, 2026
697fa8a
Deduplicate GC Notify module imports
barrfalk Oct 4, 2026
110145c
Allow queue jobs to stall up to three times before failing
barrfalk Oct 4, 2026
e2daf87
Harden merge retries and batch recovery
barrfalk Oct 4, 2026
d48d48a
Add shutdown diagnostics for nonzero exits
barrfalk Oct 4, 2026
87888c8
Trigger PR deploy
barrfalk Oct 4, 2026
080b868
Make backend heap and memory request configurable
barrfalk Oct 4, 2026
7b66cd9
Merge remote-tracking branch 'origin/main' into CCP-6026
barrfalk Oct 5, 2026
5d51dd4
Unify delivery recovery and add CHES timeouts
barrfalk Oct 5, 2026
1d51d40
Add delivery reconciler monitoring section
barrfalk Oct 5, 2026
51360ea
Limit CHES email concurrency via Redis
barrfalk Oct 5, 2026
7442b08
Add provider circuit breakers and owed-send flow
barrfalk Oct 6, 2026
1678543
Merge branch 'main' into CCP-6026
barrfalk Oct 6, 2026
5cc02db
Merge branch 'main' into CCP-6026
barrfalk Oct 6, 2026
2dcfac3
Merge branch 'main' into CCP-6026
barrfalk Oct 6, 2026
48f837e
Fix queue shutdown drain timing
barrfalk Oct 6, 2026
b6d42bd
Move TEST/PROD deploy params to chart values
barrfalk Oct 6, 2026
adb955a
Update queue.module.shutdown.spec.ts
barrfalk Oct 6, 2026
ab51302
Potential fix for pull request finding 'Overwritten property'
barrfalk Oct 6, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 51 additions & 0 deletions api-gateway/templates/routes.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -628,6 +628,40 @@ routes:
request_buffering: true
response_buffering: true

# ============================================================================
# MONITORING ADMIN ROUTES (admin only)
# ============================================================================
# /api/v1/frontend/admin/monitoring/queues (queue, worker and Redis health)
- name: ${ROUTE_PREFIX}admin-monitoring-queues-frontend
tags: [${CONFIG_TAG}]
hosts:
- ${GATEWAY_HOSTNAME}
paths:
- ${PATH_PREFIX}/api/v1/frontend/admin/monitoring/queues
methods:
- GET
strip_path: true
https_redirect_status_code: 426
path_handling: v0
request_buffering: true
response_buffering: true

# /api/v1/frontend/admin/monitoring/queues/events (SSE change signals)
# Buffering off, as for notification-request-events: a buffered stream never reaches the page.
- name: ${ROUTE_PREFIX}admin-monitoring-queue-events-frontend
tags: [${CONFIG_TAG}]
hosts:
- ${GATEWAY_HOSTNAME}
paths:
- ${PATH_PREFIX}/api/v1/frontend/admin/monitoring/queues/events
methods:
- GET
strip_path: true
https_redirect_status_code: 426
path_handling: v0
request_buffering: false
response_buffering: false

# ============================================================================
# FEATURE FLAG ADMIN ROUTES (admin only)
# ============================================================================
Expand Down Expand Up @@ -1679,6 +1713,23 @@ plugins:
route: ${ROUTE_PREFIX}admin-users-list-frontend
config: *jwtConfig

# ============================================================================
# MONITORING ADMIN ROUTES PLUGINS
# ============================================================================
# Admin Monitoring Queues - Frontend Admin Token
- name: jwt-keycloak
tags: [${CONFIG_TAG}]
enabled: true
route: ${ROUTE_PREFIX}admin-monitoring-queues-frontend
config: *jwtConfig

# Admin Monitoring Queue Events (SSE) - Frontend Admin Token
- name: jwt-keycloak
tags: [${CONFIG_TAG}]
enabled: true
route: ${ROUTE_PREFIX}admin-monitoring-queue-events-frontend
config: *jwtConfig

# ============================================================================
# FEATURE FLAG ADMIN ROUTES PLUGINS
# ============================================================================
Expand Down
5 changes: 3 additions & 2 deletions backend/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ COPY --from=build /app/dist ./dist
EXPOSE 3000
HEALTHCHECK --interval=30s --timeout=3s CMD curl -f http://localhost:3000/api

# Nonroot user, limit heap size to 150 MB
# Nonroot user. The heap limit is set per environment by the Helm chart (NODE_OPTIONS from
# backend.heapMb), not baked in here, so it can be raised without rebuilding the image.
USER nonroot
CMD ["--max-old-space-size=150", "/app/dist/main"]
CMD ["/app/dist/main"]
18 changes: 17 additions & 1 deletion backend/src/adapters/adapters.module.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,16 @@ import { NodemailerEmailTransport } from '../../src/adapters/implementations/del
import { LogEmailTransport } from '../../src/adapters/implementations/delivery/email/log/log-email.adapter'
import { TwilioSmsTransport } from '../../src/adapters/implementations/delivery/sms/twilio/twilio-sms.adapter'
import { AcsSmsTransport } from '../../src/adapters/implementations/delivery/sms/acs/acs-sms.adapter'
import { SmsCircuitBreaker } from '../../src/adapters/implementations/delivery/sms/sms-circuit-breaker'

describe('AdaptersModule', () => {
it('forRoot returns dynamic module with all adapters and maps', () => {
const dynamic = AdaptersModule.forRoot()

expect(dynamic.module).toBe(AdaptersModule)
expect(dynamic.global).toBe(true)
expect(dynamic.providers).toHaveLength(11)
expect(dynamic.providers).toHaveLength(12)
expect(dynamic.providers).toContain(SmsCircuitBreaker)
expect(dynamic.exports).toContain(EMAIL_ADAPTER)
expect(dynamic.exports).toContain(SMS_ADAPTER)
expect(dynamic.exports).toContain(EMAIL_ADAPTER_MAP)
Expand All @@ -39,4 +41,18 @@ describe('AdaptersModule', () => {
expect(adapterProviders).toContain(TwilioSmsTransport)
expect(adapterProviders).toContain(AcsSmsTransport)
})

it('sends SMS through the circuit breaker, whichever provider is configured', () => {
const provider = (AdaptersModule.forRoot().providers ?? []).find(
(p) => (p as { provide?: unknown }).provide === SMS_ADAPTER,
) as { useFactory: (...args: unknown[]) => unknown }
const acs = { name: 'acs', send: vi.fn() }
const wrapped = { name: 'acs', send: vi.fn() }
const breaker = { wrap: vi.fn(() => wrapped) }

const adapter = provider.useFactory({ acs }, { get: () => 'acs' }, breaker)

expect(breaker.wrap).toHaveBeenCalledWith(acs)
expect(adapter).toBe(wrapped)
})
})
9 changes: 6 additions & 3 deletions backend/src/adapters/adapters.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { NodemailerEmailTransport } from './implementations/delivery/email/nodem
import { LogEmailTransport } from './implementations/delivery/email/log/log-email.adapter'
import { TwilioSmsTransport } from './implementations/delivery/sms/twilio/twilio-sms.adapter'
import { AcsSmsTransport } from './implementations/delivery/sms/acs/acs-sms.adapter'
import { SmsCircuitBreaker } from './implementations/delivery/sms/sms-circuit-breaker'
import {
EMAIL_ADAPTER,
EMAIL_ADAPTER_MAP,
Expand Down Expand Up @@ -37,6 +38,7 @@ export class AdaptersModule {
LogEmailTransport,
TwilioSmsTransport,
AcsSmsTransport,
SmsCircuitBreaker,
{
provide: EMAIL_ADAPTER_MAP,
useFactory: (
Expand Down Expand Up @@ -83,16 +85,17 @@ export class AdaptersModule {
useFactory: (
map: Record<string, ISmsTransport>,
configService: ConfigService,
breaker: SmsCircuitBreaker,
): ISmsTransport => {
const key = configService.get<string>('delivery.sms') ?? 'acs'
// As above: a `:passthrough` suffix names no ISmsTransport, so fall back for DI.
if (key?.includes(':passthrough')) {
return map['acs']
return breaker.wrap(map['acs'])
}
const provider = key?.includes(':') ? key.split(':')[0] : key
return map[provider] ?? map['acs']
return breaker.wrap(map[provider] ?? map['acs'])
},
inject: [SMS_ADAPTER_MAP, ConfigService],
inject: [SMS_ADAPTER_MAP, ConfigService, SmsCircuitBreaker],
},
InMemoryTemplateStore,
{ provide: SENDER_STORE, useClass: InMemorySenderStore },
Expand Down
75 changes: 75 additions & 0 deletions backend/src/adapters/delivery-errors.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
import { HttpException, HttpStatus } from '@nestjs/common'

/**
* The provider could not take the message right now - it is down, restarting, overloaded or
* rate limiting - and it accepted nothing. The message itself is fine: send it again later.
*
* Workers must not mark a recipient failed for this. A provider outage of a few seconds would
* otherwise fail every recipient that reached it in those seconds, and since such errors come
* back instantly, a large send would burn through its whole backlog. The work is left owed and
* retried; the circuit breaker in front of the provider stops calls while it recovers.
*
* Not used when the outcome is unknown - a timeout after the request was sent may have been
* delivered, and sending it again could deliver it twice.
*/
export class TransientDeliveryError extends HttpException {
readonly transient = true

constructor(message: string, status: number = HttpStatus.SERVICE_UNAVAILABLE) {
super(message, status)
}
}

export function isTransientDeliveryError(error: unknown): error is TransientDeliveryError {
return error instanceof TransientDeliveryError
}

/** HTTP statuses that mean "not now" rather than "never": the request was not acted on. */
export function isTransientHttpStatus(status: number): boolean {
return status === 408 || status === 429 || status >= 500
}

/**
* Network errors raised before a request reached the provider: nothing was sent. A connection
* reset or a body cut off mid-response is not here - the provider may already have acted.
*/
const NOT_SENT_NETWORK_CODES = new Set(['ECONNREFUSED', 'ENOTFOUND', 'EAI_AGAIN', 'EHOSTUNREACH'])

/**
* A network error raised before the request reached the provider. The code is on the error itself
* (Azure SDK, axios) or on its `cause` (fetch's `TypeError('fetch failed')`).
*/
export function isNotSentNetworkError(error: unknown): boolean {
const { code, cause } = (error ?? {}) as { code?: unknown; cause?: { code?: unknown } }
return [code, cause?.code].some(
(value) => typeof value === 'string' && NOT_SENT_NETWORK_CODES.has(value),
)
}

/**
* An SDK error from a provider, as a TransientDeliveryError when it means the provider took
* nothing: an HTTP status of 408, 429 or 5xx (`statusCode` on Azure's RestError, `status` on
* Twilio's RestException), or a network error before the request was sent. Null otherwise.
*/
export function transientFromProviderError(
error: unknown,
label: string,
): TransientDeliveryError | null {
if (isTransientDeliveryError(error)) return error
// Our own exceptions have a runtime `status` too; only a raw SDK error is classified here.
if (error instanceof HttpException) return null
const { statusCode, status, message } = (error ?? {}) as {
statusCode?: unknown
status?: unknown
message?: unknown
}
const httpStatus = typeof statusCode === 'number' ? statusCode : status
const text = typeof message === 'string' ? message : String(error)
if (typeof httpStatus === 'number' && isTransientHttpStatus(httpStatus)) {
return new TransientDeliveryError(`${label}: upstream ${httpStatus} - ${text}`, 502)
}
if (isNotSentNetworkError(error)) {
return new TransientDeliveryError(`${label} unreachable: ${text}`)
}
return null
}
8 changes: 8 additions & 0 deletions backend/src/adapters/delivery-keys.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
import { redisKey } from '../common/redis/redis-namespace'

/** CHES email requests in flight across all pods (RedisConcurrencyLimiter). */
export const CHES_IN_FLIGHT_KEY = redisKey('ches:in-flight')
/** The circuit breaker in front of CHES (RedisCircuitBreaker); read by admin monitoring. */
export const CHES_CIRCUIT_KEY = redisKey('ches:circuit')
/** The circuit breaker in front of the SMS provider; read by admin monitoring. */
export const SMS_CIRCUIT_KEY = redisKey('sms:circuit')
Loading
Loading