From d8ba617bcc22f538c12bb0980442e63f7f79ee87 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Sat, 1 Aug 2026 11:01:45 +0300 Subject: [PATCH 1/5] feat(replication): detect dead subscribers with app_versions canary Slot lag alone missed dead Google apply workers. Check subscription health on the replica and compare app_versions counts (cached 5m). Co-authored-by: Cursor --- src/pages/admin/dashboard/replication.vue | 184 ++++++++- .../functions/_backend/public/replication.ts | 349 +++++++++++++++++- tests/replication-data-canary.unit.test.ts | 90 +++++ 3 files changed, 620 insertions(+), 3 deletions(-) create mode 100644 tests/replication-data-canary.unit.test.ts diff --git a/src/pages/admin/dashboard/replication.vue b/src/pages/admin/dashboard/replication.vue index ed2ba408d4..dc16b7c9ed 100644 --- a/src/pages/admin/dashboard/replication.vue +++ b/src/pages/admin/dashboard/replication.vue @@ -32,6 +32,36 @@ interface ReplicationSlotLag { reasons: string[] } +interface SubscriptionHealth { + status: 'ok' | 'ko' | 'skipped' + checked_at: string + threshold_seconds: number + subscriptions: Array<{ + subname: string + enabled: boolean + has_apply_worker: boolean + apply_lag_seconds: number | null + last_msg_receipt_time: string | null + status: 'ok' | 'ko' + reasons: string[] + }> + reasons: string[] +} + +interface DataCanary { + status: 'ok' | 'ko' | 'skipped' + table: string + primary_count: number | null + replica_count: number | null + diff: number | null + diff_percent: number | null + threshold_percent: number + checked_at: string + expires_at: string + cached: boolean + reasons: string[] +} + interface ReplicationStatusResponse { status: 'ok' | 'ko' threshold_seconds: number @@ -44,6 +74,8 @@ interface ReplicationStatusResponse { max_lag_minutes?: number | null max_lag_slot: string | null slots: ReplicationSlotLag[] + subscription?: SubscriptionHealth + data_canary?: DataCanary error?: string message?: string error_message?: string @@ -100,6 +132,46 @@ const checkedAt = computed(() => { return formatLocalDateTime(data.value.checked_at) }) +const subscriptionStatus = computed(() => data.value?.subscription?.status?.toUpperCase() ?? '-') +const subscriptionColor = computed(() => { + const status = data.value?.subscription?.status + if (status === 'ok') + return 'text-emerald-500' + if (status === 'ko') + return 'text-rose-500' + return 'text-slate-500' +}) +const subscriptionSubtitle = computed(() => { + const sub = data.value?.subscription + if (!sub) + return 'No subscription data' + if (sub.reasons.length) + return sub.reasons.join(', ') + const healthy = sub.subscriptions.find(item => item.status === 'ok') + return healthy?.subname ?? `${sub.subscriptions.length} subscription(s)` +}) + +const canaryStatus = computed(() => data.value?.data_canary?.status?.toUpperCase() ?? '-') +const canaryColor = computed(() => { + const status = data.value?.data_canary?.status + if (status === 'ok') + return 'text-emerald-500' + if (status === 'ko') + return 'text-rose-500' + return 'text-slate-500' +}) +const canarySubtitle = computed(() => { + const canary = data.value?.data_canary + if (!canary) + return 'No canary data' + if (canary.reasons.length) + return canary.reasons.join(', ') + const primary = canary.primary_count == null ? '-' : formatNumberValue(canary.primary_count) + const replica = canary.replica_count == null ? '-' : formatNumberValue(canary.replica_count) + const cacheLabel = canary.cached ? 'cached' : 'fresh' + return `${canary.table}: ${primary} / ${replica} (${cacheLabel})` +}) + async function loadReplicationStatus() { isLoading.value = true errorMessage.value = null @@ -176,7 +248,7 @@ displayStore.defaultBack = '/dashboard' {{ t('replication') }}

- Logical replication slot lag monitoring + Slot lag, subscription workers, and app_versions count canary

@@ -206,13 +278,25 @@ displayStore.defaultBack = '/dashboard'
-
+
+ +
+
+
+
+

+ Subscription workers +

+

+ Live apply worker check on the read replica +

+
+
+ {{ data.subscription?.reasons?.join(', ') || 'No subscription rows' }} +
+
+ + + + + + + + + + + + + + + + + +
+ Name + + Worker + + Lag (s) + + Status +
+ {{ sub.subname }} + + {{ sub.has_apply_worker ? 'Yes' : 'No' }} + + {{ sub.apply_lag_seconds == null ? '-' : formatNumberValue(sub.apply_lag_seconds, { maximumFractionDigits: 1 }) }} + + + {{ sub.status.toUpperCase() }} + +
+
+
+ +
+
+

+ app_versions canary +

+

+ Count compare cached at most every 5 minutes +

+
+
+
+ Primary + {{ data.data_canary?.primary_count == null ? '-' : formatNumberValue(data.data_canary.primary_count) }} +
+
+ Replica + {{ data.data_canary?.replica_count == null ? '-' : formatNumberValue(data.data_canary.replica_count) }} +
+
+ Diff + + {{ data.data_canary?.diff == null ? '-' : formatNumberValue(data.data_canary.diff) }} + + +
+
+ Cache + {{ data.data_canary?.cached ? 'cached' : 'fresh' }} +
+
+ Notes + {{ data.data_canary?.reasons?.length ? data.data_canary.reasons.join(', ') : '-' }} +
+
+
+
+
diff --git a/supabase/functions/_backend/public/replication.ts b/supabase/functions/_backend/public/replication.ts index 8aa9bc8085..dcf2820e5d 100644 --- a/supabase/functions/_backend/public/replication.ts +++ b/supabase/functions/_backend/public/replication.ts @@ -1,4 +1,6 @@ +import type { Context } from 'hono' import { sql } from 'drizzle-orm' +import { CacheHelper } from '../utils/cache.ts' import { honoFactory, useCors } from '../utils/hono.ts' import { cloudlogErr } from '../utils/logging.ts' import { closeClient, getDrizzleClient, getPgClient, logPgError } from '../utils/pg.ts' @@ -6,8 +8,14 @@ import { validatePlatformAdminOrApiSecret } from '../utils/platform_admin_access const DEFAULT_THRESHOLD_SECONDS = 180 const DEFAULT_THRESHOLD_BYTES = 16 * 1024 * 1024 +const DATA_CANARY_TABLE = 'app_versions' +const DATA_CANARY_TTL_SECONDS = 300 +const DATA_CANARY_TTL_MS = DATA_CANARY_TTL_SECONDS * 1000 +const DATA_CANARY_THRESHOLD_PERCENT = 0.01 +const DATA_CANARY_CACHE_TIMEOUT_MS = 250 type SlotStatus = 'ok' | 'ko' +type CheckStatus = 'ok' | 'ko' | 'skipped' type ReplicationQueryMode = 'wal_stats' | 'replication_stats' | 'slots_only' interface ReplicationSlotLag { @@ -25,6 +33,51 @@ interface ReplicationSlotLag { reasons: string[] } +export interface SubscriptionWorkerRow { + subname: string + subenabled: boolean + has_apply_worker: boolean + apply_lag_seconds: number | null + last_msg_receipt_time: string | null +} + +export interface SubscriptionHealthResult { + status: CheckStatus + checked_at: string + threshold_seconds: number + subscriptions: Array<{ + subname: string + enabled: boolean + has_apply_worker: boolean + apply_lag_seconds: number | null + last_msg_receipt_time: string | null + status: SlotStatus + reasons: string[] + }> + reasons: string[] +} + +export interface DataCanaryResult { + status: CheckStatus + table: typeof DATA_CANARY_TABLE + primary_count: number | null + replica_count: number | null + diff: number | null + diff_percent: number | null + threshold_percent: number + checked_at: string + expires_at: string + cached: boolean + reasons: string[] +} + +interface DataCanaryCacheEntry extends Omit { + expiresAt: number +} + +const dataCanaryMemoryCache = new Map() +const dataCanaryInflight = new Map>() + function toNumber(value: unknown): number | null { if (value === null || value === undefined) return null @@ -34,6 +87,104 @@ function toNumber(value: unknown): number | null { return num } +function isReplicaDatabaseSource(source: string): boolean { + return source.startsWith('HYPERDRIVE_CAPGO_READ') || source === 'local_read_replica' +} + +export function evaluateAppVersionsCanary( + primaryCount: number, + replicaCount: number, + thresholdPercent = DATA_CANARY_THRESHOLD_PERCENT, +): Pick { + const diff = Math.abs(primaryCount - replicaCount) + const baseline = Math.max(primaryCount, replicaCount, 0) + const diffPercent = baseline === 0 ? 0 : diff / baseline + const reasons: string[] = [] + + if (primaryCount > 0 && replicaCount === 0) + reasons.push('replica_empty') + else if (diffPercent > thresholdPercent) + reasons.push('count_mismatch') + + return { + status: reasons.length > 0 ? 'ko' : 'ok', + diff, + diff_percent: Number(diffPercent.toFixed(6)), + reasons, + } +} + +export function evaluateSubscriptionHealth( + rows: SubscriptionWorkerRow[], + thresholdSeconds = DEFAULT_THRESHOLD_SECONDS, + checkedAt = new Date().toISOString(), +): SubscriptionHealthResult { + if (rows.length === 0) { + return { + status: 'ko', + checked_at: checkedAt, + threshold_seconds: thresholdSeconds, + subscriptions: [], + reasons: ['no_subscription'], + } + } + + const subscriptions = rows.map((row) => { + const reasons: string[] = [] + if (!row.subenabled) + reasons.push('subscription_disabled') + if (!row.has_apply_worker) + reasons.push('no_apply_worker') + if (row.apply_lag_seconds !== null && row.apply_lag_seconds > thresholdSeconds) + reasons.push('apply_lag_threshold_exceeded') + + return { + subname: row.subname, + enabled: row.subenabled, + has_apply_worker: row.has_apply_worker, + apply_lag_seconds: row.apply_lag_seconds, + last_msg_receipt_time: row.last_msg_receipt_time, + status: (reasons.length > 0 ? 'ko' : 'ok') as SlotStatus, + reasons, + } + }) + + const hasHealthy = subscriptions.some(sub => sub.status === 'ok') + const reasons = hasHealthy + ? [] + : [...new Set(subscriptions.flatMap(sub => sub.reasons))] + + return { + status: hasHealthy ? 'ok' : 'ko', + checked_at: checkedAt, + threshold_seconds: thresholdSeconds, + subscriptions, + reasons: reasons.length > 0 ? reasons : (hasHealthy ? [] : ['subscription_unhealthy']), + } +} + +function getFreshDataCanaryMemoryEntry(cacheKey: string, now = Date.now()): DataCanaryResult | null { + const cached = dataCanaryMemoryCache.get(cacheKey) + if (!cached) + return null + if (cached.expiresAt <= now) { + dataCanaryMemoryCache.delete(cacheKey) + return null + } + const { expiresAt: _expiresAt, ...payload } = cached + return { ...payload, cached: true } +} + +function setDataCanaryMemoryEntry(cacheKey: string, entry: DataCanaryCacheEntry) { + dataCanaryMemoryCache.set(cacheKey, entry) +} + +/** Test helper: clear in-process canary cache between unit tests. */ +export function clearDataCanaryCacheForTests() { + dataCanaryMemoryCache.clear() + dataCanaryInflight.clear() +} + function buildReplicationQuery(mode: ReplicationQueryMode) { const slotsCte = sql` WITH slots AS ( @@ -139,6 +290,161 @@ async function executeReplicationQuery( throw lastError } +async function querySubscriptionHealth( + drizzleClient: ReturnType, + thresholdSeconds: number, +): Promise { + const checkedAt = new Date().toISOString() + const result = await drizzleClient.execute(sql` + SELECT + s.subname, + s.subenabled, + COALESCE(bool_or(ss.pid IS NOT NULL AND ss.last_msg_receipt_time IS NOT NULL), false) AS has_apply_worker, + MAX(EXTRACT(EPOCH FROM (now() - ss.last_msg_receipt_time))) + FILTER (WHERE ss.last_msg_receipt_time IS NOT NULL) AS apply_lag_seconds, + MAX(ss.last_msg_receipt_time) AS last_msg_receipt_time + FROM pg_subscription s + LEFT JOIN pg_stat_subscription ss ON ss.subname = s.subname + GROUP BY s.subname, s.subenabled + ORDER BY s.subname + `) + + const rows: SubscriptionWorkerRow[] = (result.rows as any[]).map(row => ({ + subname: String(row.subname), + subenabled: Boolean(row.subenabled), + has_apply_worker: Boolean(row.has_apply_worker), + apply_lag_seconds: toNumber(row.apply_lag_seconds), + last_msg_receipt_time: row.last_msg_receipt_time + ? new Date(row.last_msg_receipt_time).toISOString() + : null, + })) + + return evaluateSubscriptionHealth(rows, thresholdSeconds, checkedAt) +} + +async function countAppVersions(drizzleClient: ReturnType): Promise { + const result = await drizzleClient.execute(sql` + SELECT COUNT(*)::bigint AS count + FROM public.app_versions + `) + const count = toNumber((result.rows as any[])[0]?.count) + if (count === null) + throw new Error('app_versions count missing') + return count +} + +async function queryDataCanary( + primaryDrizzle: ReturnType, + replicaDrizzle: ReturnType, +): Promise & { expiresAt: number }> { + const checkedAt = new Date().toISOString() + const expiresAt = Date.now() + DATA_CANARY_TTL_MS + const [primaryCount, replicaCount] = await Promise.all([ + countAppVersions(primaryDrizzle), + countAppVersions(replicaDrizzle), + ]) + const evaluation = evaluateAppVersionsCanary(primaryCount, replicaCount) + + return { + status: evaluation.status, + table: DATA_CANARY_TABLE, + primary_count: primaryCount, + replica_count: replicaCount, + diff: evaluation.diff, + diff_percent: evaluation.diff_percent, + threshold_percent: DATA_CANARY_THRESHOLD_PERCENT, + checked_at: checkedAt, + expiresAt, + reasons: evaluation.reasons, + } +} + +async function getCachedDataCanary( + c: Context, + primaryDrizzle: ReturnType, + replicaDrizzle: ReturnType, + replicaSource: string, +): Promise { + const cacheKey = `app_versions:${replicaSource}` + const memoryEntry = getFreshDataCanaryMemoryEntry(cacheKey) + if (memoryEntry) + return memoryEntry + + const existingQuery = dataCanaryInflight.get(cacheKey) + if (existingQuery) + return existingQuery + + const cacheHelper = new CacheHelper(c) + const cacheRequest = cacheHelper.buildRequest('/cache/replication-data-canary', { source: cacheKey }) + const cachedEntry = await cacheHelper.matchJson(cacheRequest, { + timeoutMs: DATA_CANARY_CACHE_TIMEOUT_MS, + }) + + if (cachedEntry && cachedEntry.expiresAt > Date.now()) { + setDataCanaryMemoryEntry(cacheKey, cachedEntry) + const { expiresAt: _expiresAt, ...payload } = cachedEntry + return { ...payload, cached: true } + } + + const existingQueryAfterCache = dataCanaryInflight.get(cacheKey) + if (existingQueryAfterCache) + return existingQueryAfterCache + + const query = queryDataCanary(primaryDrizzle, replicaDrizzle) + .then(async (result) => { + const cacheEntry: DataCanaryCacheEntry = { + status: result.status, + table: result.table, + primary_count: result.primary_count, + replica_count: result.replica_count, + diff: result.diff, + diff_percent: result.diff_percent, + threshold_percent: result.threshold_percent, + checked_at: result.checked_at, + expires_at: new Date(result.expiresAt).toISOString(), + expiresAt: result.expiresAt, + reasons: result.reasons, + } + setDataCanaryMemoryEntry(cacheKey, cacheEntry) + await cacheHelper.putJson(cacheRequest, cacheEntry, DATA_CANARY_TTL_SECONDS) + const { expiresAt: _expiresAt, ...payload } = cacheEntry + return { ...payload, cached: false } + }) + .finally(() => { + dataCanaryInflight.delete(cacheKey) + }) + + dataCanaryInflight.set(cacheKey, query) + return query +} + +function skippedSubscription(reason: string): SubscriptionHealthResult { + return { + status: 'skipped', + checked_at: new Date().toISOString(), + threshold_seconds: DEFAULT_THRESHOLD_SECONDS, + subscriptions: [], + reasons: [reason], + } +} + +function skippedDataCanary(reason: string): DataCanaryResult { + const checkedAt = new Date().toISOString() + return { + status: 'skipped', + table: DATA_CANARY_TABLE, + primary_count: null, + replica_count: null, + diff: null, + diff_percent: null, + threshold_percent: DATA_CANARY_THRESHOLD_PERCENT, + checked_at: checkedAt, + expires_at: checkedAt, + cached: false, + reasons: [reason], + } +} + export const app = honoFactory.createApp() app.use('*', useCors) @@ -154,6 +460,7 @@ app.get('/', async (c) => { const pgClient = getPgClient(c, false) const drizzleClient = getDrizzleClient(pgClient) + let replicaPgClient: ReturnType | null = null try { const { rows, mode } = await executeReplicationQuery({ requestId: c.get('requestId') }, drizzleClient) @@ -212,7 +519,41 @@ app.get('/', async (c) => { return acc }, { slot: null, lag: null }) - const overallStatus: SlotStatus = slots.length === 0 || slots.some(slot => slot.status === 'ko') ? 'ko' : 'ok' + let subscription = skippedSubscription('no_replica_connection') + let dataCanary = skippedDataCanary('no_replica_connection') + + try { + replicaPgClient = getPgClient(c, true) + const replicaSource = String(c.get('databaseSource') ?? c.res.headers.get('X-Database-Source') ?? '') + if (!isReplicaDatabaseSource(replicaSource)) { + subscription = skippedSubscription('no_replica_connection') + dataCanary = skippedDataCanary('no_replica_connection') + } + else { + const replicaDrizzle = getDrizzleClient(replicaPgClient) + subscription = await querySubscriptionHealth(replicaDrizzle, thresholdSeconds) + dataCanary = await getCachedDataCanary(c, drizzleClient, replicaDrizzle, replicaSource) + } + } + catch (error) { + cloudlogErr({ requestId: c.get('requestId'), message: 'replication_replica_check_failed', error }) + subscription = { + ...skippedSubscription('replica_check_failed'), + status: 'ko', + } + dataCanary = { + ...skippedDataCanary('replica_check_failed'), + status: 'ko', + } + } + + const slotStatus: SlotStatus = slots.length === 0 || slots.some(slot => slot.status === 'ko') ? 'ko' : 'ok' + const failingChecks = [ + slotStatus === 'ko', + subscription.status === 'ko', + dataCanary.status === 'ko', + ] + const overallStatus: SlotStatus = failingChecks.some(Boolean) ? 'ko' : 'ok' const response = { status: overallStatus, @@ -228,6 +569,8 @@ app.get('/', async (c) => { max_lag_minutes: maxLagSlot.lag !== null ? Number((maxLagSlot.lag / 60).toFixed(2)) : null, max_lag_slot: maxLagSlot.slot, slots, + subscription, + data_canary: dataCanary, } return c.json(response, overallStatus === 'ok' ? 200 : 503) @@ -250,9 +593,13 @@ app.get('/', async (c) => { max_lag_minutes: null, max_lag_slot: null, slots: [], + subscription: skippedSubscription('replication_lag_error'), + data_canary: skippedDataCanary('replication_lag_error'), }, 500) } finally { await closeClient(c, pgClient) + if (replicaPgClient) + await closeClient(c, replicaPgClient) } }) diff --git a/tests/replication-data-canary.unit.test.ts b/tests/replication-data-canary.unit.test.ts new file mode 100644 index 0000000000..220a9a4420 --- /dev/null +++ b/tests/replication-data-canary.unit.test.ts @@ -0,0 +1,90 @@ +import { afterEach, describe, expect, it, vi } from 'vitest' +import { + clearDataCanaryCacheForTests, + evaluateAppVersionsCanary, + evaluateSubscriptionHealth, +} from '../supabase/functions/_backend/public/replication.ts' + +describe('replication data canary evaluation', () => { + afterEach(() => { + clearDataCanaryCacheForTests() + }) + + it('accepts similar app_versions counts within 1%', () => { + expect(evaluateAppVersionsCanary(1000, 995)).toMatchObject({ + status: 'ok', + diff: 5, + reasons: [], + }) + }) + + it('rejects empty replica when primary has rows', () => { + expect(evaluateAppVersionsCanary(120, 0)).toMatchObject({ + status: 'ko', + reasons: ['replica_empty'], + }) + }) + + it('rejects count drift above threshold', () => { + expect(evaluateAppVersionsCanary(1000, 900)).toMatchObject({ + status: 'ko', + reasons: ['count_mismatch'], + }) + }) + + it('treats one healthy subscription as ok even with a disabled sibling', () => { + const result = evaluateSubscriptionHealth([ + { + subname: 'capgo_google_eu_2', + subenabled: false, + has_apply_worker: false, + apply_lag_seconds: null, + last_msg_receipt_time: null, + }, + { + subname: 'capgo_google_eu_2_sub', + subenabled: true, + has_apply_worker: true, + apply_lag_seconds: 2, + last_msg_receipt_time: '2026-08-01T07:00:00.000Z', + }, + ]) + + expect(result.status).toBe('ok') + expect(result.subscriptions.find(s => s.subname === 'capgo_google_eu_2_sub')?.status).toBe('ok') + }) + + it('marks enabled subscription without apply worker as ko', () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-08-01T07:00:00Z')) + + const result = evaluateSubscriptionHealth([ + { + subname: 'capgo_google_eu_2_sub', + subenabled: true, + has_apply_worker: false, + apply_lag_seconds: null, + last_msg_receipt_time: null, + }, + ]) + + expect(result.status).toBe('ko') + expect(result.reasons).toContain('no_apply_worker') + vi.useRealTimers() + }) + + it('marks apply lag above threshold as ko', () => { + const result = evaluateSubscriptionHealth([ + { + subname: 'capgo_google_eu_2_sub', + subenabled: true, + has_apply_worker: true, + apply_lag_seconds: 400, + last_msg_receipt_time: '2026-08-01T06:50:00.000Z', + }, + ], 180) + + expect(result.status).toBe('ko') + expect(result.reasons).toContain('apply_lag_threshold_exceeded') + }) +}) From 894594c78393ce5298d383d3732118afbd9ac246 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Sat, 1 Aug 2026 11:02:15 +0300 Subject: [PATCH 2/5] docs: add replication data canary PR screenshot Co-authored-by: Cursor --- .../pr-screenshots/replication-data-canary.webp | Bin 0 -> 21966 bytes 1 file changed, 0 insertions(+), 0 deletions(-) create mode 100644 docs/pr-screenshots/replication-data-canary.webp diff --git a/docs/pr-screenshots/replication-data-canary.webp b/docs/pr-screenshots/replication-data-canary.webp new file mode 100644 index 0000000000000000000000000000000000000000..7654f72c9e16a47e1e3772e8d2d43db335facad8 GIT binary patch literal 21966 zcmeFYbC4%ZyDi$Bwr$%sr)}G|G0k7w=CnO++qP}nw(flIxA(sLM%*27&fn)$Ra8_& zRaRC$&sr;2Wo0Qziir)W0s*Ou3M;57a1yrtJw{vs&H|nP_SLr-XmZNu=CmD z78LmNZTD#OSo;9<>Hb*#4S1KmxqkO~PrStY{?iF|h{c&1Jfb}cm3+97h ze|ZJs4&nIwxV!u79uW8a{atX8^8k4DR_1^G3jErB{q`d`CRi5O`}Y5qe9w5O{?PxN zzRx)a*nJlHdVFzw9lutO+t&b&1eOKPzC*qRUjN?Wn&$iaGsjQAwfp<~)SJO~5Agk6 zfmjvxHu8J-Q@7!(Y7L-nD^T8)v$!sQc3J%BHviUc`M-W5n@UBRRFGWHpe)WR3;9zw z9*Q%re3Hq;lcT!U)p?Fq(E2Q-LmaiOhF-+SM8@HNtZvH+RXi1QuDT)-*y*W>1+vL) z%aV~>_Vz8u40+42ivKdR5~<|!mZx7xZ-X*b|&inU`+;hP8(KF4-rnP}ba{KJu&`Cw9^-ZUS5)+8etQBU5& zmZXMbZt}3^eg%8jsr@||>6qaEKc-m6xDvNU+o7~x8$ z9M<=c0fgtmldCuYp??Sxw_414 zPEjb^#(qc5oV$_qtJ;fAcHiJu&MGb2{#Hhj(ey8}iv431;v`69WXf9^?@<@tV}pPDQZP9kE-vE593=>7aDIgD-%?!zLR@SgljS)4 z2e{45y?RZnZf@QJJJo+|2R(6I#Fw-#pbeLVxn{^yWU0Fu^AqKhDwjlt-c?Js^&4#A zvto`eZOL%dgF2-Q#6Yw-dz&}$9Tu^{wK#Gon21M4JM8{aJx#6&^=l0$3M$w!_hyLc zDrAY_>;TFX3Z#wW_dr*E_WH6g$7Hub5mCty@J2cjA&D}zcnuylVEgA z#{E8muq7mT9HiVRRCgn?WPA{Ifn@OX^b#v$E@)ynP?9pA`~igUVU`F#h^5;?Wdd@w z22dlbzNzdc?ZBkcCvAU^e9*&1f{_f(M^8i$p(Do+(~^+wypGJh9>iV{o-;fNE)3KD zo5R<;Ref>Xe;0|Ou~Wk)( zn1R1k+sm=(o(=W*p{hjAuzSY$NK1n&b7~Ogdz+AwvMc#SoHq9pB~4!`PAew1p8*oK zHR-LqEB1j&zKss+G_U7)xFjUnAK+KWEhmEqSLq^C_2xU64*Gc3&8?JF6tsyn&vg{4+=k<>7r@@&fZ6UkqW5oJwZ=1S1f{qhxLVn2$At(o``4T8nTs$>!bYkI z!erB>SQYcO0>gA6H#^BtCiW~LqiHqde2R#0)^FiYJyidE^W?G>)j+~!aVj$Ag#T5I z+)T0bkJfB8LP0<(J^-T&StY#*<$eg*>r>kOA6F#nI84sABVRO5mVInq zN&HJf-cA)gah$oL75>a7Li%D@m1p-G5U=Nkcudaxg)FxNtlU|eT_ReYuIq9XR?vDy znqONIuj86ro1OhYCSC_Zio!UyRVVBkF4B1@R%NcSa?i#21-4y)?thB!8XLed*ixa* z#1^g%vpLJ`CKXSib+4C;_WfUYeCQr}f`Sl+l7kM}-$43B+;zNo1X-Euo}ATgDRXaq z(pDhY-C3hk+dR;}29A%%T9MbXv#o&n8&aqg4)m2f$eWzSu#6(p2>KdJ69xC*g3Ey# zc5C&nqqXd1HX8v}gA}@%3uH9!H;AI#VoPt601t1RQ~GR#B7<0#f4nN7#TAF;wmYOX zzN=kHZ&Np+O>F z-68lP5c=r+xza^*f}4Xrb~+XXII+oN3l=u)x4BB`VYQ)xkRO6k0p!^YxHNiueW4Pd zl-TJzE!}lV^fc3Exp<1T2VUoYl^dbhLoIU=VU~S>+qaeo#LY5=kL);rRsr_p zDs2}~skUx6D?%Jb?0;JX;S(~xGPaxb6hBRNqwK)WSI27`MW}?%H9#zh3?uwwbfWyb z1G1o8qt!h+HVOU9MLhgQgqEp^tFAmKY>GY`_O3Dac^oG&{&5wl^3d&JJ>CknujdRU&)d`--Lk#q5s=U{{xhNb<(}> zM0;3|!%~x9bJO_CEKhY6|9aNDg8xtDJzg!(h+h4~=f1{T8=E+w!hY>~|KW`)nYS}I z*DHohyZQACXTDCM2+Ilj-`aqrxBuT-&n`X_u_5RUFI=Cn)WYDuSrAn^w>XX~M=3b>;^I*VNZZ+eH>%$Y zDY{rzJL_*JSfYm<(nYozrX|CN*znu)AD#b4n)_#A#46Gck3_~n11Ju%>-{1n-HrZB z$p5+8byE5JZ_pb)IK+QItB?No05wO0he7@Kfc?wWtPc_eo)a5~{_ncm|9wCFe@pzY zzdQV&7Aut=ps(*DmBmrdjk7y9e&xMg{4n)o_MiuD=STCI_E-MZxPe~SN&*oPf!){)Sz5RzDt z84YBH>-Hk3!l<9IH#VHpgrU^j^kY)3uf^O;%CXO`O_|9J3D*F~O*Z&Yhwre~eC4Y< zhThF_^utd@+t9H-LaX4>ZUbZC6rA0Y;!E$TEr0FO7VoRRK=Q{ zAOE5)sL^z0X-L@zk`<4d|5d4r0nGEKo{5zxA?eCgLm$7;Sb;7s=t_bzkCo>9$14pb5D@qe)vOcUq4OMlEO3hb zyJDtVbk1UN9Uy4~E$nw%p*)*lgNk?Wu!V?yfY@p`q6-iE9koUoTa_TZc!lOlr&1bj&IS02otB`E z5X$iOyzH(R@cjE++0_hgkJ!AFu8vylXlTvavs%Tng$6kY88-cw$o38WYg3)#KjXZ< zjlGe&5-5$zDQLE=+Ie@{W>Kik^(O^w^j%|Y88W`apl>Lp#dOxLtKdsfw|BX|?jr_? z3GexqxHLi;ulol^!HQphP_q3J_TDV-@B>72BpJ9QfS~w49v^NQH%sVb z%vp0CFe5S-GVV{2CBk}ufW?*P5@phzo5;cpD6n2zUksd123DbYex$|KSrjschjEo! zYC#pp#N0(6-P{1y%i#k5v7e2O8PsA8ng8ikJI5bj#glul{JgGNzAE!{J>#%c3aq)% zWzePMx^HGiwWTZ7%i?XELaaBr_*hWl=gSM~A4*8^wAN8LJp5-ps;C?~B&8s9IkFZ~ z;Nq`bfLR!ju7{}3Y&$)7%7`Kn$*u^R<^so=mx;oy`>1HHrz&Dkrkp9yRGw>FM$`>45o;Zu3uJ;MhKCN7l+?p zhSn?A3>;V>J>sVNSLMjZRz=VuKvwVG^C!A+al~%=i#RPx*4J|Dd5B z=aqpIeCQ!;htp7#>-nzRZLiTJV3)aU?_Z3v{JJq!3}A-CpN|0>Ap#kPU~!z z&4q};XV*H$5(;}(>&cNAuwZBKqD6Aw&G0iFcZ=*{M9!dY>Va6Hbr-hn^0vd=rR ztH;XS7CBv*o-kxGz@!nrKja`FJMGgyLPGRfmJ6HJcq~`J? zb_?xd`^1O~HtVr%Rx-C|<}-b)$LXqBQ~@GZ+iH}TKc>1D8jaUr%!pmAJBSomd9pP1 zbwuupo3F|DYrh;ja`L29hF8bzG-O4_TWq$NMh#{c0_Y=LNAcrTi-hG#1_@Ma#o zw;tL}VwlX8@-nzUTB@|f=kZb_Ifh~?#PcIuZYeATgwp(4P?4Ka$E)PF#H=^37^Dga zs4TL9e=b5USwhS4#}TVjuoTsH5hxmtcBn)`mg%Tl@grkK>@0Nrc*6KH)I2zgW39)Y zcbtM9>6}gf$B#KQx4p13lgaDj1K+!GcSf32`5H|W3EwZ1su6fdi0xA&hG8sedK@hd zZnKwtoVi0W+h3jCtS^|!zd=)Ow;iflh7>BUTupRDj0_5x-8YS3{CiG~gaqT7!NbJt zrtWB^dF84Q5$!)?Tz9kT$c;H#UDksAIpBT8u$W)(i_&Q8MX!vFqf#MEi-cFnt$ZF? z+XxwqOBt|#xG%he1BYs2&|-Ft%%AhaFuy{sS`#bCuxd>(o_Igg*wN*c(MYIHZ4Yqb z#P2M!QX_g=;yAr~UMI`A2g8&y7Z=zz&pTFr_@|#6w5IH@bG5apS=hMl&rE`{3qxIX z4j?3M7snonw!atKI0m6(W|IO#od?GN+_Z6*7;|5KCGaAMt8i4ml&r4q1cwghGmawDLL0p=wu~)QD3^ zz3U>lWRgd{&a0FpTgUq`VAP0ttvNOPW|(PX(J&jkwqiG0vZfoL@ zUv`B#w^F;v21HDfxz)^Y?~^?6<+ID{mxC1$6Dl@!p7rh_9MY!%HLdWgsmzVlpuUz| z55!V?PN65&dmpQ~-8A3Z9qw$1 z3bFa&w5w5~w|7`@CNwuEEn+ZwQ*G(p?4b@KQ2q_@#vQSzJHEZNHB!vf0 zV^OE*O$H_bwI(JzVaSsT{@j*Ib4;=(t{Q*ODnU1X5tJ8!eb88&W)XRiCi;?f>wpZg z_;>;s_a%$&%z3DACq6%gf<33JMkT#cmH=%p2tLi_!Loj9TJormgwll2S zxQgRIOZXTl-WqHqjXA;;0>+?P3B1iHJaG79bRd4A>$nH5_haA;NNs6(+VZOxvl7?w!IR~*H4?e7 z&)H1^oSJO5azH57(`Dc?y5d9@b z|M0p+cu}@HMdCW5FT@T^R^CMmnEZOMInEq9*-y?5n~p_9AW&gZbZJ0(&aaRPIr|7I zeH0)Z-GF-kxpv9}9l6MgSWHO@Kgvk-Xb2SVgJpZ|t3t_|R7j2h#N{v4&C}t9H z)X;AUt5N&W&Y&b@i$CwuAZpRX&&4UVy7MXg zA#pl$e23^b-7I?F##&?9b~AMp{sT--{fSq&F*J0bH4tq}L|uI#L@NID(NKbdjzR73 zpW0RymUaCBbIuoU)P~#z!8afjoaFHolXJqx7lIhv?WlQn zUjy-f#x$~j2sShF_(ai>OY!dcJ{a*W#P>_%p|GgHi6sAKEKo_hx|3rV?|%QFRLsG+ zHU%bPV}MzyZn;2YWWVXha+#uMfB~BR4b*msS!@lswL@{#-BniONkXVwkf&cgZ*;#& zBuqz*Nf0GAaEPUiIy$E%6em>fY8?!At|aob6JMGgngh?alt4Vgu7R%vzU=X^ORDxo z?jX0H9anVp&`8-^sD<_z)8R|Tj=>v1ufG~(rvXhEdPLq`%sJO$VOX~5Hh2kTs~a7C z6j7|~(}Ik8K4Nt^h$&ohUP*e0rK{SJTk2|z%=_dPMfp}mLbZ*3#1-~|6Myf;!IhJVlP z616-8!VecFqTd%XltrZ!D?iGEe@5@K9KZPdtEltzkv_wbQJ~rMFi1*&sJD4aWgIh( zDAe0z_|s7-Az(CXMc|O3+Nd>KFkvK%1~ zr+&KZ$_{6C_wFT0|M+#aPcLBc_@wDlbTSEMRHy>)C%xce1b=Nw9?}??ZX!aCOncQT zRo%k#CAOK6nav#+u=) zqa26jqL^_BmpNEaNL=dFgLh`%#FgV#77VZR*P!x-5Z8;2CrB%XOK_34Z^?5+HR{|+?M9_gq!w6ZTiM{9tm&HRUm5N6g}twkYp zoc;RQ{!)eJr!g)xp6Iaqv~$j|!g#DXJ2^Hf8aT0?O9nlH697pHMVqV{?D|`+ZJDQ_ z*%~(|xIV*mN{mW$KJiIRCBSP5`P#n8Ukd5QUPHieLg1scx5&a!qmh|Z96_1C4rAp= zE_{dVaA^{=ef+D>VzeNUQE7DFnAuc%v~=C>O(c>pvZHhXTA`K}6TKj#Q+wSujkVYoVrTqk+H z!?dcQ$ZE~QPAg+ak>@%p1f0ODSf`?msY*36WU&j%%40tl&5u6BGuX$Z@xh!nyzbf^*qA4FHPsmc1E z;;zBse*mylwfmut2@q}<{(0iJJj>%UOyIK>4X{66xsSqi zZCcgw!CkwDei}h0JYb!;NRSkp+0A?_YR4y2Tz~H2d-O5}Tr*|D6g@|7A?`3!?&@7I z-CjVbVv)dWfLr0$`mGRMu64^xZY~pI!EVWr?&8*|0t40dL3Ojm(zzlAW z`=Tf8X-H;-W{GU_m}JIi=spZ3><(`Zrf$E6?1P>wK1#M~hp1bg5c!A_Hze8B$s@rg zhtbe^(KvNe{}+z7p1$u$Ba<_52F6KTV!-C^+UIkg)(az zn68uQKzF4&1^U5nj{$>Faa=EVioq1^CnQyZ&gP(_wUXV6V-bT!2jsB6v)xg;J4slh z4%q7~z66k&pHdo?@VC7~ns}?t)QI)^nh~p~1!h*lE z(~^M$@?Kpjg~}Pv_CQ3PWtT^KcD-QJ&%j|*1qt5^xa311sGsS(iKUF|O)`#ncqH^| zFLM;>ZUn*rTqt#hx0^XYE?lfOvXqX2>HO-XUQ`|#ha@$osW7yN&6wyerjN#7^GmmWzHeh6nd9f`}b=7{7V`f#8v`~c}B{Zm4$b@_r6Xhkp);1l?PfUO;r?HIQaAU_k zzMiBFqb33>hV@N!GUY)^;!5V%mQv{~StcN4@(W1lKMo);bayZmHVyHNYS?lTD#<45 zUeFRVOb)m1YqFpLlPe(v3^&6sP9&X{h`hr7VWSG%&OB5dlu6>3;sgzXqf5Oh*xXODUqDxh5Cqv)iS+l>6eI-CbICk2#n ze2Mkd5u>n)A!~b8k)?+>Rn=c*?iUPRW9?RwMXtF&?$TLed#qP^2pvn9L{tJ;LVk1| zW-n83&0LfuTuDmbL8C}eZGM+<&pRQWV#Dn6i$skk1LuFjUl)tESg+8RT)qRgeJ?C? zmK<_fiW7SW6pqC;m3Netf~QTQ|t?eC(GAo!8KEJSEe zGE38MeOBk4kBrhB*B!8l%7S}}qIO~SW}1ngz{|qwxBj;6f|~4?nLsK;b5NiZF}Kn{ z(ga~zN8PXma2|7szES*o0J4|hw(yKx92h;F9s)Mrn8;OV^RFrP6IrcVm~;A>!{HPEOl>D612A!5 zx+?dp(EA^8R`q7 z!%Xcf{RPIwBTro*tLmS!PSN2q^OhC9aN!g51A^?3=cZU2!M+?PJdnj-$O>v1JWWfL z=}23O^IvZS`Bu!gl%AlTCLKF8IW^fXX%l3yF};J8<4}~g^whY2hY~9nznvWMCBZgQ zbwhBm3I605=YiTQV${mq(Ok_K`IU(+G(!PFHtbtR4`2m52n})UvbK!v6_0CUaABjp zUUayg^fC0GU-dK7HBz;Ps5wp;v7*%4XQthE6ez}BibO3cdQR3-kULBDc)s7iKc9mI zDm#ge8h%HM=0n~hVM|%=0J3EkbB2-aNsh08M&Wf_S+3`gOOU9M(P&lY>?Nf0hM$n}L(Yre4%=?0+6po^U%w|48W@Nzl5;^T zE|)+?_+mp=Sk$$gpiujrzGF^rAiP57zeR29C39OQG*O#4N|TwWOotW#B_IgDsra+U z!qTEl`Y;Z{y^F$r&_YR{G#01{V6VHm__R<)OwL+3wFvzla`foh?b~W5Ls5l~Lko6M zjjqy?Lk*2qX8j!R=wUU@R}M29a}DwLBT{8SVKsijTJiBquJ@4M;n!T!t8K7%sa1h% zRr&;LE~%GtypZ!WN*k_{l%oe0^n_Fs^qerNGwI;P>+w=~VPh(Zhxb`G6GCpxrIRQ& zFUqLl;SL!Fs;ONuDA$B5p9uymX37P{5DMPL3K;WB#RBGiFhryA`$wgYLhzM3X6EdmN@%%e+nF> zsf@t)z5aA#L(FR<32K%tfu^d_e0yoPFs*890!lC{;u0xdo8n=BGS1z^QAx3yyY3km z76d6x)eAdZIoaqtb`TB$+tp^JBT|9sY)Pdprzt%9?+2|QpbW7EmxHEuM2(X=um)6y zLH_KjBKmtWi1lxD_$w znff)qAr_N-{nR}PyGt%swJqBsn1$#Uk zO}o}Ek~=mvMom7em9n9Tv?oLM_IcL&t+E*5($QE?E}3VyfABp0e2@GS-;!wY5}Y*0 z!l!{BMXpw?r&>NAcu3yO0FDn_0|d1wIxr8W1BASBrUsjh55&92%#Yxk!EbI!?`Ofv ze&@ogz>SU<3mRv~BQV&fkPN_0E12-ghhi<5J{L4hE|Y^@kNbr_HMz+<_w0t$^s^D< z!T46AM+3E<& zdxUxhXrS@tE0}-_y4(gF?y`=hjnX(f3D6j$rS8l}AiK+IF;4X|4XIyJo9(y=H7f(; zEQ*&FRP$+Nk85uhf3J6|zizX1#iawpxV5zVYCge`W|$9Y+8n`qrPYF^Ysw90p&;=z z-l1`N?^S-L-E5;kby+L;h%0z_s*S!ojc&CC(CgWYQB65CVG8>)K79a9w_VDxz3)qX zn=lPZXI3^isnE<92SzK!a2ZHZVPmGCJtU@++3s-2u0_Z>uINH^gJhvr8GeKt6OAJM z=X+AN#W7!KlyTLoE&pPooJVvzv|+yk&&|ihYB9=VD;c|!?mpX};wFA3t!mUg{yh+* z2`E0ZEJ1I;DO;!D*^ns3)<+IwGN`A9oMh2M99g0l!@il!V14EaLyFM;5)$jVh$dB3xc`x0gDyto=e(ZgkXeT(TIHeI`I!I7vZ{#Z(B_)(l$ia?y{qkT=^DUm)nrogR zWP+oHoa|kAVmTw-SE<9OHd!-Yx|f5jkYD^;PBRjY943z_cu!;lKCm3_nN9W213CM} z)&|sV2J44M9A{nS|MS*;X_QM7WP8;|^b<;m$rn!ZYUr)&t{Cm1;2^RPP^DX>2mC8l zvuP`A-&4M*dc=O-|w`e*1kdnz7NG z`A4VfEj>zH8M8&WyXPUq`4-IC`l)A^m*PpiV0COVXsuAhI5GvO6j(7X{l@n?9PHDJ z!4{6b-RBBFnRn@|8CN*R9>X`&li=o|WjP((^*GiH zik{2w)aGJB9M@buK}d{Pmeuuk67#8G`X{0-SOkZLKbX%x%=@p(1OaM7z&k>j(Qhz z-N_O6%*7fTW+dGybFnJo0`RR3&Z8%3pH$kiXa}P-ggbOMn1`JJJUG*lkVEJFQ?XF& zjvwahe;na=SUrKP^7k?RZC=u)_vAu%bbFYn<>V$1S+WQ z(C9pOzlA?F)w9qCM{V7kUZOHKG-v_}jaRA@pGUxQz`Lnn<_o~6}tt=pxm z0tb&rco1MiLyo2W;GM23c4#Ztn);-GnID~e@dk4nR6qE*Y3N{f=gkj}J41UGR>t!; z4tW#+ivt1)$*WU~tUm=YXHbSh{*_A7o8%1ATP}X}1;?}Xw(>X`AjIcKwY!nW-B3H@ zeXls#qhL(wRe6H7DCa7ItadAr<02K(Ce_aRe5Q{;?I-dzTbeQQNLsvYCMskWAJ9u^ zO3Bs>I0%*w@eLL?M9pUN7KeANYzEe#s@C^Fu=@Ry9RSN0ItA?Gc4ljb!D}yols{CC z(m!`&@`tMQ2{#F^9ZZoF+f>N(L(ex%R_ixpO*sI%o%qb#y+RdY*HFJ%K_laj#e(cB z-4F%g@BFg0q{Dfk< zkP^Ak{Y}i@%*FH=zGUvWy9!d6VUtZ-$zW7z8Ok~x9>{#u1qd|q^Id6mlnh#cG(r&2 z0`m)db4Ao5k0RdBQd1f^w)xYXhM^a(cW#qib5Rb?ufhxc!G^b0GGhPmT$^+yMFbuoq*&Z=!{FPMae4H^R@jR;+>r&`>8vGk4&-AJ1Y4g2Yq{Z?(J_gC3; zdH7(I=A5Xl6d+*Z^q*Jxho>P<>_$7x8`R197yw5+Q7v~;s=Dip zomISBwT@(AOJL7Uchjh|3%_Rl;T^x0mK=zP_v2vKbl*$H&|F26P=dMW;S>>|^Nxj) zqyCnLH97^F4mOn|y9MQ5edW~Y9oo0w=8}m_c%`K++Ue|{)+hmKymF^$8uDDb9!EG< zajeye?M^?MZ|reV@EHX4VI{vRpcAJ)7Pd6rc+;zsi{a3_az-5LWZM`Dw?)AwAhcb-4ng zAtk*v-vq>L@gQaGK78&n6!5$NGfP{Q7^TT>dkP0noBW-fHoEr%A8=*e(oB?fhAh`6 z6;zR829KPxLPMu^3K9TwLh6XPl6Gj37=0(iV9Icf^@pJF!2suFG=&OV7&vbnkRSuD z002EGah!0yOyi79;rRMtW`FjPkkFChWnYdwpT0y>jxcOPC@B=YUCWo)^`bGP^&5neRrP`V-81HbWTv2^OAIaVXJzb#|22 z;sZfFq{$|!hB*h@)`A*ttSkmSO1U=P7m0a(Ua7^Ipip2ug{R8*kX_7x!eb=%^hQdn zQohBSwXpsmkcVd_Nl8I@R^mwpZldED=Ea~ah2&rj6Lg$1G-}e{)cXQk0^eZLLGxCT zS+|_^R)ffnUoxW{r_1y_UYQ!GLsyIdFC+8CoDyF;c<~RwN0(g1(p& zxrM*Uv7ERJ!9IdHZnJogCGgA(>xq7nZ55# zMIXYx+>vE&u@(<2CJl9GUlnS;4}^kc3=+q8{(g6-6^u3oS?(D z74Y|sj+Or7cb5zz*p5wBZo=GgtIh^@d#pc+AY74nWobJ+TZ_ewH8Fe(QP0%?9T^U9 z+^}H{2~zr}q_l*VqTDMsWd)Y0+g>T@zx22wn6Ou{8?_blc@5_uMW-3j&;^OGR`umI z!t1ew<=ykAQJ9~k?K|5h`|&8$&Jjc)lcnFX`(DZc@jn}Zz!)YLsN)~@T^K0Zm8C4T zVCh{1(14Y>JFdaj2NHQYKxLa)P!;pClw@yr)%P?3$1UD6((1v9=e!$)zLLpLp# z0i&6svVn6-bDzQljkU7R7%Y?0WZ9n$0!tYRl4a$l-G1Z1eB%5aZ{@VVx}|I{#pg?$ z1fD4idlbVk8BlG&TIW^~e_Hi+G-X_HOLTuzw}b}r$JV2GhT7s8Hby7v2ci6O5R@^$ zu7ny_t0h~7!s{?f7Rs$S7H$63y*@3$r6Ryfv?j9DA<9Vg`v*md$$SbZfS;aHZVRaMte#2)7 zuM0ORi>Bzx-0T4j-7-1?7I54!=F{+1`TPC^5FRjenN)m)p zC{oC$R_`HAze#S@ZnHSBBLm(v<)Wo92tumgvuPIt5-5-uuZh&BItO>U|MCdDfm@nzm7@?nmndfgy1L=2b zCc4havU^L)4{k}8c&tDdOaad@37`~4)SRXAj4{TT&m6L>E7bCICr0cVs%iv?#5J->C$7 z)zwHwH~zA1h9YmX%Df>Ltf@GD6?T-Zbgq9QBR_7XjCYL0Wki6z{vb03YWzSfpv#B8 zm;xAnod4$7{rEG%bURl7i!-CgoVz~#oIIUp8Fm)=N>W%?8yAxCJ$`wsEtRZhFHJ_b zAbHC@3r*e6M%(B*9okivch%BXvNVVh!>4w{p)!wfTxwH`JEko9W#G`&qGuht7ifeELkgs zF<&vBn1bwRtp%YfFQ#=u3?(Gp)?};s{m5r`%_u-sck-D;ak9OAgF27PU2>!1PXAW+MrKJHetLIyYzICngEIkcLTLLH7_aM3ODw!c}ohD7_ z_aiWCiwwIzaTd&wAH6PN3Cf#G^#qEbzEejyl7l++5nMLM^Kw8ISuP>EobZwV6Tt?f z^<+pdr6u!1+Pk>YFlh(|p>8zK7KDhW8!0oQyfjeD8K%{@@AvoL8Tg*wcs8vX>N9@+ z-Kd8@3f%Gg1Vths-$Y_#kbj8JwnB%dwNlVO!$8!dj&+DnJ=bkWxzZU%Cs}x^Z7T%m z0d36G6w;t-lMGhB2nAaob>@G+9uw}K_R7vrzFH5*Bomv^@|y{(SDTO~AGUJMC&81# z_!Cu#mOlR-(9=^?jEt`L4ueiF z9DXMt5_2)(lIBxHyb`a5U3qD9sdE&GXY?$8w*K+kjBT47iX$o2v*T7XexR!*0N}3{ zampIVUt59Vlgr~2vf?$Xu|9NaDMLCImf$JLt#Q#&z{ zY(KHoSE*|JyG& zxLR74_I$o>Vwf{}pCK*|w4NOtA85COy9bB#vbK$kS zbYkanr1eEKIdR=DS!4EA%x{>5!BcNnER@#Fuzu_#w!*rWKW+?NPCKVAj(?CqkV!vY zMhpBtUtCGGv^@g;0@X!C#%2_rv!LV(3z-W+Cy2G0qor1*Z zPD$w)-Q6)7L8KdFh#(=|-Q77lM^CyR-~aGlAJ2bqopCHP(|l&`qmw=m{M0;EH-($> z$#T*eVixmHFDDuJJDe+y&RIe)F9>7ll(nP1gT!ub)9LY55*>G zy93(T88=-HEY82}L=UW0Kv-iy$B%P3n~ivE-$6$WO5*zLdudH0jUr^QwkFVea9`8l za5hI4D;B()x3DZ!dlxu2Nna>^KtSt%82modYe8W?!(YqL$0h46{Eq&T{}W>U`ww{1 zZT0{LwYh$G++gouguvg20yRwY%YhoZv4{IMKpz5<^#ag2LxDN5k2T8}c*$e>&N*18 zFPhpyAX?IVwM^I>qHPj1d*u_;ZA1A^{MIqsvHJjLt)aEsLRu+w{CfPofhs)6uNSF` zGmsz2WP99m0qTa#%^hjc2$@h3C#&&=j|h9D6vjicliA+=Xn<0vi_QU!t-3Q!e~Y(& z!Ju;C$JD)LrxT6-MTF%F1HyK?dA$Bq&cZ33*<+IGjR0#fuNXpu49AR$zUVDfSz}W5 zxGFOzaB`h=mk{$!L6j9`+yWaWnx#+IBGe6BF{5Gx_FnF=x)~o9h)P0EMjas$J^k&K zSC5_R+OC2@6^?Xqelwk#Gr#;1f$#km9m#}lZvj&g}GDZl@$vOj2NyMR~NPdo0vmXc)3iMTQu(niXxOeuHVnFem zgiZ^V8BSKUP$mN2g}An{1(oc?*^NHnl|a_(#bs0%Nk|d3*Lce5Le({^SxUm@+784P zmYin9f^ppau-Rc?GFCT|L7`xi7c`sOg!$8%5~%$rmz>Dk!0n82Rpxpj6wk&hm`S1` z1XM1lr;~`KNk%EGaO(Cr>`7ndvnaTqqw?!x>P)Z}6fBuP9oqz)uzq8Xq9<^a=-F&q?G!S{pu2YN5A>^U+MdyaY37X8^K8Z#V%WYOy_RpK0$dy$%Y!0jp}FW!`WL=F4!>8Ox_Xqv9PDf7}x3k&N3}A-vSMf3%32{E&ANq@rVP1szah z7UJ9i_G{f7+f=oXJpX5?``hyi`z*cdsEeXC1F;8@Ni~ZPz-#y#XkeJw{muVfO%VU` zaC!A8^j=Xbv#1@eupIIFH=DSW6fZBxsNW72yp9Dz%>JG_;CgpFs>81vM}dw!yPu5M zGscTtURRW50WYE=PtB=I^e#r3{582{(pG9JziGWNYo@3+UzeJ0AGt;nBE!-I7B;?L!79UZg@o0MjSJp_6V$Z*m;;oH0{uEkjv0Qu{S z9C=dM0p;mYVT`#87xOfrwl;ejogTB?PH&NE4QS1Kb0G*}kAM7^_G`ZNXNweUgFK-pEXHMZkQO!?MzFL8YC52q2uYXCz^VcqUa4Wi@@G@Kt#TFJYFn>C6{ z>dNt7CJ3>$>vmVWp>jY#;mvTPfIhvI(!{P^x|vSf{~$pzf_wC&$6Fj@hZUdSpCGi! zBP7gs<~P-#Rd7_sHY}=RzK?vmaTv`#&RLObxU2BwbIF(qq_B;kH@+e)OpI;|s5Hqs zCAi$M%)e$ELmesj?kOfA8=rRnT#zU4-g-fk>jPc0#f9%Seb2jhOwCQh5V-7}Xup;a8P4n^~?tnVxldpi2=AT9ATd4@tVxZKyS4MQG-{ zO;D`=*Ka{V*58^j0C{#;04(z)m47bP&Bki0v&SJl zx@RYPkwy~LZI|m;2zzNAAFhP&Q?#7+I?vj8o(8kTRlS^&Y5xj(l0l;Cw832V&-4DU z*j>6mNFgN+^7?`Q5?*?6`oC|#kj`6vO;R@ZWT27y{Y-71)6GIn(0O%XF}5PPXa5hd zlzCx=(=k5eY%<1Z`mkRnJzkLT0efPevZ5rv66vV~ep8){-(weR7V>nlO{^F(@RQ1Unlwj!Df0 z%PT9)AwmPrR9>2joOt#H5BmF^kezV3KPcy6<)=Ny3z6Q)xp-k@cBAh-nV+H9{^hlQ z1Ht^%?i8BAQl})gz~lJfbN!~|8ea98q&9vHHc5~dA)=!}7^^3;xHj}fbK60#twpKW zo>o%1#aqF2P$dh;`W!X5Vdz%IeNR+!xjPZmiDi(^>=A{Lt-Qvs>hJO^`O6}dl zac!p-mRo}@t>c6anHBKhWfoSiBr8T?~Hoa+UH9qhpA9j%oH;>@Iy1htM9 zpJ@8ZAEVpqJemIYeJahB7S=RyR{SxTsPUVYWS%~&aro!aaa2U%rryBM3A}|2WYpwl z0jb*()cT&I8^a;Ge(LkAXV=L5 zeDg?hA*%MR zR|0rU-hrD=aoq>fiNhbi>4qLsH$PBCaFdURsu$8d-KgV0)$u}oDd_^d2BQt3>A%*9 zh(b&noX)Bl{v~hf|NcY1oWLVOZ4ol9d2R^9n-b-20KoKFpOCFpI%d12J}a&0&!LfP zsOX;_MIP*hQ*zLcyALL_;Is04Iu@HXZwbl`BORaa&}0Q2u`J~#ORvr)>p4^r>ParV z8DSBKW=tZsE$~X3gRbR69z)k!F}CAiJwx9qE3JGeO$DV2-pcmfF&wrk95#qUGG)Li4?c& zT~WR&QSKJyaQbbgTLyi3a}uqaDbLD|OvvGOU)#(0qizf)nLI)fKCttr{X-uD*5m}~>*pRQ8`pmpxrUYs00hi95(^~ZG zz2`e1OEh62tf`NRBHmV}HC_6M%Q#3p0}MjS*i@Yd>86dzP~t&NTGJglVkoquBUsYD zpcPVZFF%;D?k%^=`Ggt_jCQg~e9JM>qLGJnXNOj=>$IAZ{!{6XqjV1zT|{ELkz1@TgrQKl)bFiuhNlUStq(!{U@Dl;5qJ4@$+V?v7}FXwQg!2hfZU5+S?L5WLDLV8=!4`5iW->4=^ L&p*z@|JVHwvD{&> literal 0 HcmV?d00001 From d426070cedeb202f297693b18508bc9ac46d9678 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Sat, 1 Aug 2026 11:10:10 +0300 Subject: [PATCH 3/5] fix(replication): read replica source from response header Typed Hono context only exposes requestId via c.get(); use the X-Database-Source header set by getPgClient instead. Co-authored-by: Cursor --- supabase/functions/_backend/public/replication.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/supabase/functions/_backend/public/replication.ts b/supabase/functions/_backend/public/replication.ts index dcf2820e5d..8b2f642f7b 100644 --- a/supabase/functions/_backend/public/replication.ts +++ b/supabase/functions/_backend/public/replication.ts @@ -524,7 +524,7 @@ app.get('/', async (c) => { try { replicaPgClient = getPgClient(c, true) - const replicaSource = String(c.get('databaseSource') ?? c.res.headers.get('X-Database-Source') ?? '') + const replicaSource = c.res.headers.get('X-Database-Source') ?? '' if (!isReplicaDatabaseSource(replicaSource)) { subscription = skippedSubscription('no_replica_connection') dataCanary = skippedDataCanary('no_replica_connection') From 312b2a1c5f7700539f0c49e241f1aad718247677 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Sat, 1 Aug 2026 17:37:35 +0300 Subject: [PATCH 4/5] fix(replication): address review feedback on canary checks Require every enabled subscription to be healthy, split worker/receipt signals, parallelize replica checks, bound cache writes, and fix DaisyUI badge classes. Co-authored-by: Cursor --- src/pages/admin/dashboard/replication.vue | 11 ++-- .../functions/_backend/public/replication.ts | 44 ++++++++++------ tests/replication-data-canary.unit.test.ts | 50 +++++++++++++++++-- 3 files changed, 82 insertions(+), 23 deletions(-) diff --git a/src/pages/admin/dashboard/replication.vue b/src/pages/admin/dashboard/replication.vue index dc16b7c9ed..1f250561ca 100644 --- a/src/pages/admin/dashboard/replication.vue +++ b/src/pages/admin/dashboard/replication.vue @@ -40,6 +40,7 @@ interface SubscriptionHealth { subname: string enabled: boolean has_apply_worker: boolean + has_recent_receipt: boolean apply_lag_seconds: number | null last_msg_receipt_time: string | null status: 'ok' | 'ko' @@ -168,6 +169,8 @@ const canarySubtitle = computed(() => { return canary.reasons.join(', ') const primary = canary.primary_count == null ? '-' : formatNumberValue(canary.primary_count) const replica = canary.replica_count == null ? '-' : formatNumberValue(canary.replica_count) + if (canary.status === 'skipped') + return canary.reasons.join(', ') || 'skipped' const cacheLabel = canary.cached ? 'cached' : 'fresh' return `${canary.table}: ${primary} / ${replica} (${cacheLabel})` }) @@ -361,7 +364,7 @@ displayStore.defaultBack = '/dashboard' {{ sub.apply_lag_seconds == null ? '-' : formatNumberValue(sub.apply_lag_seconds, { maximumFractionDigits: 1 }) }} - + {{ sub.status.toUpperCase() }} @@ -400,7 +403,7 @@ displayStore.defaultBack = '/dashboard'
Cache - {{ data.data_canary?.cached ? 'cached' : 'fresh' }} + {{ data.data_canary?.status === 'skipped' ? 'skipped' : (data.data_canary?.cached ? 'cached' : 'fresh') }}
Notes @@ -478,8 +481,8 @@ displayStore.defaultBack = '/dashboard' {{ slot.status.toUpperCase() }} diff --git a/supabase/functions/_backend/public/replication.ts b/supabase/functions/_backend/public/replication.ts index 8b2f642f7b..d87693493c 100644 --- a/supabase/functions/_backend/public/replication.ts +++ b/supabase/functions/_backend/public/replication.ts @@ -37,6 +37,7 @@ export interface SubscriptionWorkerRow { subname: string subenabled: boolean has_apply_worker: boolean + has_recent_receipt: boolean apply_lag_seconds: number | null last_msg_receipt_time: string | null } @@ -49,6 +50,7 @@ export interface SubscriptionHealthResult { subname: string enabled: boolean has_apply_worker: boolean + has_recent_receipt: boolean apply_lag_seconds: number | null last_msg_receipt_time: string | null status: SlotStatus @@ -131,17 +133,23 @@ export function evaluateSubscriptionHealth( const subscriptions = rows.map((row) => { const reasons: string[] = [] - if (!row.subenabled) + if (!row.subenabled) { reasons.push('subscription_disabled') - if (!row.has_apply_worker) - reasons.push('no_apply_worker') - if (row.apply_lag_seconds !== null && row.apply_lag_seconds > thresholdSeconds) - reasons.push('apply_lag_threshold_exceeded') + } + else { + if (!row.has_apply_worker) + reasons.push('no_apply_worker') + if (!row.has_recent_receipt) + reasons.push('no_recent_receipt') + if (row.apply_lag_seconds !== null && row.apply_lag_seconds > thresholdSeconds) + reasons.push('apply_lag_threshold_exceeded') + } return { subname: row.subname, enabled: row.subenabled, has_apply_worker: row.has_apply_worker, + has_recent_receipt: row.has_recent_receipt, apply_lag_seconds: row.apply_lag_seconds, last_msg_receipt_time: row.last_msg_receipt_time, status: (reasons.length > 0 ? 'ko' : 'ok') as SlotStatus, @@ -149,17 +157,19 @@ export function evaluateSubscriptionHealth( } }) - const hasHealthy = subscriptions.some(sub => sub.status === 'ok') - const reasons = hasHealthy + // Disabled leftovers are ignored. Every enabled subscription must be healthy. + const enabled = subscriptions.filter(sub => sub.enabled) + const allEnabledHealthy = enabled.length > 0 && enabled.every(sub => sub.status === 'ok') + const reasons = allEnabledHealthy ? [] - : [...new Set(subscriptions.flatMap(sub => sub.reasons))] + : [...new Set((enabled.length ? enabled : subscriptions).flatMap(sub => sub.reasons))] return { - status: hasHealthy ? 'ok' : 'ko', + status: allEnabledHealthy ? 'ok' : 'ko', checked_at: checkedAt, threshold_seconds: thresholdSeconds, subscriptions, - reasons: reasons.length > 0 ? reasons : (hasHealthy ? [] : ['subscription_unhealthy']), + reasons: reasons.length > 0 ? reasons : ['subscription_unhealthy'], } } @@ -299,7 +309,8 @@ async function querySubscriptionHealth( SELECT s.subname, s.subenabled, - COALESCE(bool_or(ss.pid IS NOT NULL AND ss.last_msg_receipt_time IS NOT NULL), false) AS has_apply_worker, + COALESCE(bool_or(ss.pid IS NOT NULL), false) AS has_apply_worker, + COALESCE(bool_or(ss.last_msg_receipt_time IS NOT NULL), false) AS has_recent_receipt, MAX(EXTRACT(EPOCH FROM (now() - ss.last_msg_receipt_time))) FILTER (WHERE ss.last_msg_receipt_time IS NOT NULL) AS apply_lag_seconds, MAX(ss.last_msg_receipt_time) AS last_msg_receipt_time @@ -313,6 +324,7 @@ async function querySubscriptionHealth( subname: String(row.subname), subenabled: Boolean(row.subenabled), has_apply_worker: Boolean(row.has_apply_worker), + has_recent_receipt: Boolean(row.has_recent_receipt), apply_lag_seconds: toNumber(row.apply_lag_seconds), last_msg_receipt_time: row.last_msg_receipt_time ? new Date(row.last_msg_receipt_time).toISOString() @@ -406,7 +418,7 @@ async function getCachedDataCanary( reasons: result.reasons, } setDataCanaryMemoryEntry(cacheKey, cacheEntry) - await cacheHelper.putJson(cacheRequest, cacheEntry, DATA_CANARY_TTL_SECONDS) + await cacheHelper.putJson(cacheRequest, cacheEntry, DATA_CANARY_TTL_SECONDS, { timeoutMs: DATA_CANARY_CACHE_TIMEOUT_MS }) const { expiresAt: _expiresAt, ...payload } = cacheEntry return { ...payload, cached: false } }) @@ -531,8 +543,12 @@ app.get('/', async (c) => { } else { const replicaDrizzle = getDrizzleClient(replicaPgClient) - subscription = await querySubscriptionHealth(replicaDrizzle, thresholdSeconds) - dataCanary = await getCachedDataCanary(c, drizzleClient, replicaDrizzle, replicaSource) + const [subscriptionResult, canaryResult] = await Promise.all([ + querySubscriptionHealth(replicaDrizzle, thresholdSeconds), + getCachedDataCanary(c, drizzleClient, replicaDrizzle, replicaSource), + ]) + subscription = subscriptionResult + dataCanary = canaryResult } } catch (error) { diff --git a/tests/replication-data-canary.unit.test.ts b/tests/replication-data-canary.unit.test.ts index 220a9a4420..ff1a9a65b9 100644 --- a/tests/replication-data-canary.unit.test.ts +++ b/tests/replication-data-canary.unit.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, it, vi } from 'vitest' +import { afterEach, describe, expect, it } from 'vitest' import { clearDataCanaryCacheForTests, evaluateAppVersionsCanary, @@ -38,6 +38,7 @@ describe('replication data canary evaluation', () => { subname: 'capgo_google_eu_2', subenabled: false, has_apply_worker: false, + has_recent_receipt: false, apply_lag_seconds: null, last_msg_receipt_time: null, }, @@ -45,6 +46,7 @@ describe('replication data canary evaluation', () => { subname: 'capgo_google_eu_2_sub', subenabled: true, has_apply_worker: true, + has_recent_receipt: true, apply_lag_seconds: 2, last_msg_receipt_time: '2026-08-01T07:00:00.000Z', }, @@ -54,15 +56,37 @@ describe('replication data canary evaluation', () => { expect(result.subscriptions.find(s => s.subname === 'capgo_google_eu_2_sub')?.status).toBe('ok') }) - it('marks enabled subscription without apply worker as ko', () => { - vi.useFakeTimers() - vi.setSystemTime(new Date('2026-08-01T07:00:00Z')) + it('fails when any enabled subscription is unhealthy', () => { + const result = evaluateSubscriptionHealth([ + { + subname: 'capgo_google_eu_2', + subenabled: true, + has_apply_worker: false, + has_recent_receipt: false, + apply_lag_seconds: null, + last_msg_receipt_time: null, + }, + { + subname: 'capgo_google_eu_2_sub', + subenabled: true, + has_apply_worker: true, + has_recent_receipt: true, + apply_lag_seconds: 2, + last_msg_receipt_time: '2026-08-01T07:00:00.000Z', + }, + ]) + expect(result.status).toBe('ko') + expect(result.reasons).toContain('no_apply_worker') + }) + + it('marks enabled subscription without apply worker as ko', () => { const result = evaluateSubscriptionHealth([ { subname: 'capgo_google_eu_2_sub', subenabled: true, has_apply_worker: false, + has_recent_receipt: false, apply_lag_seconds: null, last_msg_receipt_time: null, }, @@ -70,7 +94,22 @@ describe('replication data canary evaluation', () => { expect(result.status).toBe('ko') expect(result.reasons).toContain('no_apply_worker') - vi.useRealTimers() + }) + + it('marks enabled subscription with pid but no receipt as ko', () => { + const result = evaluateSubscriptionHealth([ + { + subname: 'capgo_google_eu_2_sub', + subenabled: true, + has_apply_worker: true, + has_recent_receipt: false, + apply_lag_seconds: null, + last_msg_receipt_time: null, + }, + ]) + + expect(result.status).toBe('ko') + expect(result.reasons).toContain('no_recent_receipt') }) it('marks apply lag above threshold as ko', () => { @@ -79,6 +118,7 @@ describe('replication data canary evaluation', () => { subname: 'capgo_google_eu_2_sub', subenabled: true, has_apply_worker: true, + has_recent_receipt: true, apply_lag_seconds: 400, last_msg_receipt_time: '2026-08-01T06:50:00.000Z', }, From bd25412ad13a26be5ef7640da5e98d89f477ad34 Mon Sep 17 00:00:00 2001 From: Martin Donadieu Date: Sun, 2 Aug 2026 15:21:30 +0300 Subject: [PATCH 5/5] fix(replication): clarify disabled and empty-replica canary states Mark disabled subscription rows as disabled (not KO), fail empty replica counts, and expose slot_status while keeping overall health inclusive of subscription and canary checks. Co-authored-by: Cursor --- src/pages/admin/dashboard/replication.vue | 11 ++++-- .../functions/_backend/public/replication.ts | 37 ++++++++++++------- tests/replication-data-canary.unit.test.ts | 8 ++++ 3 files changed, 39 insertions(+), 17 deletions(-) diff --git a/src/pages/admin/dashboard/replication.vue b/src/pages/admin/dashboard/replication.vue index 1f250561ca..16c00a4452 100644 --- a/src/pages/admin/dashboard/replication.vue +++ b/src/pages/admin/dashboard/replication.vue @@ -43,7 +43,7 @@ interface SubscriptionHealth { has_recent_receipt: boolean apply_lag_seconds: number | null last_msg_receipt_time: string | null - status: 'ok' | 'ko' + status: 'ok' | 'ko' | 'disabled' reasons: string[] }> reasons: string[] @@ -364,8 +364,11 @@ displayStore.defaultBack = '/dashboard' {{ sub.apply_lag_seconds == null ? '-' : formatNumberValue(sub.apply_lag_seconds, { maximumFractionDigits: 1 }) }} - - {{ sub.status.toUpperCase() }} + + {{ sub.enabled ? sub.status.toUpperCase() : 'DISABLED' }} @@ -403,7 +406,7 @@ displayStore.defaultBack = '/dashboard'
Cache - {{ data.data_canary?.status === 'skipped' ? 'skipped' : (data.data_canary?.cached ? 'cached' : 'fresh') }} + {{ !data.data_canary ? '-' : data.data_canary.status === 'skipped' ? 'skipped' : (data.data_canary.cached ? 'cached' : 'fresh') }}
Notes diff --git a/supabase/functions/_backend/public/replication.ts b/supabase/functions/_backend/public/replication.ts index d87693493c..21f861c774 100644 --- a/supabase/functions/_backend/public/replication.ts +++ b/supabase/functions/_backend/public/replication.ts @@ -53,7 +53,7 @@ export interface SubscriptionHealthResult { has_recent_receipt: boolean apply_lag_seconds: number | null last_msg_receipt_time: string | null - status: SlotStatus + status: SlotStatus | 'disabled' reasons: string[] }> reasons: string[] @@ -103,7 +103,7 @@ export function evaluateAppVersionsCanary( const diffPercent = baseline === 0 ? 0 : diff / baseline const reasons: string[] = [] - if (primaryCount > 0 && replicaCount === 0) + if (replicaCount === 0) reasons.push('replica_empty') else if (diffPercent > thresholdPercent) reasons.push('count_mismatch') @@ -132,22 +132,30 @@ export function evaluateSubscriptionHealth( } const subscriptions = rows.map((row) => { - const reasons: string[] = [] if (!row.subenabled) { - reasons.push('subscription_disabled') - } - else { - if (!row.has_apply_worker) - reasons.push('no_apply_worker') - if (!row.has_recent_receipt) - reasons.push('no_recent_receipt') - if (row.apply_lag_seconds !== null && row.apply_lag_seconds > thresholdSeconds) - reasons.push('apply_lag_threshold_exceeded') + return { + subname: row.subname, + enabled: false, + has_apply_worker: row.has_apply_worker, + has_recent_receipt: row.has_recent_receipt, + apply_lag_seconds: row.apply_lag_seconds, + last_msg_receipt_time: row.last_msg_receipt_time, + status: 'disabled' as const, + reasons: ['subscription_disabled'], + } } + const reasons: string[] = [] + if (!row.has_apply_worker) + reasons.push('no_apply_worker') + if (!row.has_recent_receipt) + reasons.push('no_recent_receipt') + if (row.apply_lag_seconds !== null && row.apply_lag_seconds > thresholdSeconds) + reasons.push('apply_lag_threshold_exceeded') + return { subname: row.subname, - enabled: row.subenabled, + enabled: true, has_apply_worker: row.has_apply_worker, has_recent_receipt: row.has_recent_receipt, apply_lag_seconds: row.apply_lag_seconds, @@ -569,10 +577,12 @@ app.get('/', async (c) => { subscription.status === 'ko', dataCanary.status === 'ko', ] + // Intentionally includes subscription + canary: /replication is the admin health probe. const overallStatus: SlotStatus = failingChecks.some(Boolean) ? 'ko' : 'ok' const response = { status: overallStatus, + slot_status: slotStatus, estimation_source: mode, threshold_seconds: thresholdSeconds, threshold_minutes: Number((thresholdSeconds / 60).toFixed(2)), @@ -609,6 +619,7 @@ app.get('/', async (c) => { max_lag_minutes: null, max_lag_slot: null, slots: [], + slot_status: 'ko', subscription: skippedSubscription('replication_lag_error'), data_canary: skippedDataCanary('replication_lag_error'), }, 500) diff --git a/tests/replication-data-canary.unit.test.ts b/tests/replication-data-canary.unit.test.ts index ff1a9a65b9..487b6705fa 100644 --- a/tests/replication-data-canary.unit.test.ts +++ b/tests/replication-data-canary.unit.test.ts @@ -25,6 +25,13 @@ describe('replication data canary evaluation', () => { }) }) + it('rejects empty replica even when primary is empty', () => { + expect(evaluateAppVersionsCanary(0, 0)).toMatchObject({ + status: 'ko', + reasons: ['replica_empty'], + }) + }) + it('rejects count drift above threshold', () => { expect(evaluateAppVersionsCanary(1000, 900)).toMatchObject({ status: 'ko', @@ -53,6 +60,7 @@ describe('replication data canary evaluation', () => { ]) expect(result.status).toBe('ok') + expect(result.subscriptions.find(s => s.subname === 'capgo_google_eu_2')?.status).toBe('disabled') expect(result.subscriptions.find(s => s.subname === 'capgo_google_eu_2_sub')?.status).toBe('ok') })