From 7d86eb5b91b638c20d492182e529df8aaad28e58 Mon Sep 17 00:00:00 2001 From: Alberto Monterroso <14013679+Albermonte@users.noreply.github.com> Date: Thu, 24 Sep 2026 10:57:03 +0200 Subject: [PATCH] fix(sync): defer epochs without false failures --- MIGRATION.md | 4 +- README.md | 6 ++- app/app.vue | 2 +- nuxt.config.ts | 4 +- package.json | 2 +- server/tasks/cron/sync.test.ts | 4 +- server/tasks/cron/sync.ts | 2 +- server/tasks/sync/epochs.test.ts | 77 ++++++++++++++++++++++++++---- server/tasks/sync/epochs.ts | 19 ++++++-- server/utils/activity-sync.test.ts | 24 +--------- server/utils/activity-sync.ts | 8 +--- wrangler.json | 4 +- 12 files changed, 100 insertions(+), 56 deletions(-) diff --git a/MIGRATION.md b/MIGRATION.md index a5e3d08..5486b12 100644 --- a/MIGRATION.md +++ b/MIGRATION.md @@ -2,7 +2,7 @@ ## Background -Cloudflare Pages does not support scheduled tasks. This project requires six-hour syncing, so it uses Cloudflare Workers. +Cloudflare Pages does not support scheduled tasks. This project requires hourly syncing, so it uses Cloudflare Workers. ## Setting Up Redirects @@ -120,7 +120,7 @@ Repeat inspection and backup with `validators-api-mainnet` and without `--env te NuxtHub selects the testnet bindings during the build. Rebuild after each `NUXT_SCORE_V2_MODE` change; deploy the matching `.output` without `--env`. -5. Let the six-hour job discover epochs, repair recent activity, store the snapshot, then calculate v1 scores. +5. Let the hourly job discover epochs, repair recent activity, store the snapshot, then calculate v1 scores. 6. Set `NUXT_SCORE_V2_MODE=shadow`, deploy, and validate v1/v2 rows, activity coverage, score versions, and `current`, `stale`, or `no_score` API states. 7. Set `NUXT_SCORE_V2_MODE=active` only after shadow results pass. 8. Repeat the same inspect, backup, migrate, off, repair, shadow, validate, and active sequence for mainnet. Build mainnet with `pnpm build` and deploy with `pnpm exec wrangler --cwd .output deploy` after each mode change. diff --git a/README.md b/README.md index d81a3af..3140108 100644 --- a/README.md +++ b/README.md @@ -149,12 +149,14 @@ We also do have an UI component to visualize the range, check the status, and de ### Fetcher -The fetcher retrieves data from the Nimiq network and stores it in D1. It runs every six hours in this order: +The fetcher retrieves data from the Nimiq network and stores it in D1. It runs every hour in this order: 1. Discover completed epochs and verify or repair activity snapshots. 2. Store the current validator snapshot. 3. Calculate scores from finalized activity markers. +Each production run considers up to 50 epochs and attempts at most four full repairs, prioritizing recent gaps. After eight minutes it stops starting epochs so snapshot and score tasks can still run. Each activity batch gets one retry in production. Unstarted epochs remain available for the next run and are reported as `deferredEpochs`, without writing failure markers. Existing failure markers clear only after successful verification or repair. Scores remain stale until recent activity coverage reaches 100%. + Completed-epoch activity uses marker-backed integrity checks. Operator-facing marker states are: - `syncing`: repair attempt owns a fresh six-hour lease. @@ -259,7 +261,7 @@ This implementation does not run any remote migration or deployment. - `production`: [Validators API Mainnet](https://validators-api-main.je-cf9.workers.dev) via the mainnet build above - `testnet`: [Validators API Testnet](https://validators-api-test.je-cf9.workers.dev) via the testnet build above -Each environment has its own D1 database, KV cache, and R2 blob. Sync runs every six hours via Cloudflare cron triggers (see `server/tasks/sync/`). +Each environment has its own D1 database, KV cache, and R2 blob. Sync runs every hour via Cloudflare cron triggers (see `server/tasks/sync/`). ### Score v2 rollout diff --git a/app/app.vue b/app/app.vue index f6ccac1..8f8f178 100644 --- a/app/app.vue +++ b/app/app.vue @@ -207,7 +207,7 @@ const currentEnvItem = getEnvironmentItem(nimiqNetwork) ?? { network: nimiqNetwo

- Note: Data synchronization is handled automatically by scheduled tasks that run every six hours. A score lag of up to 1 epoch can be expected between sync cycles. + Note: Data synchronization is handled automatically by scheduled tasks that run every hour. A score lag of up to 1 epoch can be expected between sync cycles.

diff --git a/nuxt.config.ts b/nuxt.config.ts index bc44d2d..595ea53 100644 --- a/nuxt.config.ts +++ b/nuxt.config.ts @@ -162,8 +162,8 @@ export default defineNuxtConfig({ tasks: true, }, scheduledTasks: { - // Six-hour sync: wrapper task records run + executes sync tasks - '0 */6 * * *': ['cron:sync'], + // Hourly sync: wrapper task records run + executes sync tasks + '0 * * * *': ['cron:sync'], }, openAPI: { meta: { title: name, description, version }, diff --git a/package.json b/package.json index 600165c..5aca0c6 100644 --- a/package.json +++ b/package.json @@ -17,7 +17,7 @@ "generate": "nuxt generate", "worker:dev": "npx wrangler --cwd .output dev", "worker:dev:mainnet:prod": "tsx scripts/dev-mainnet-prod-worker.ts", - "worker:trigger-sync": "curl 'http://localhost:8787/__scheduled?cron=0+*/6+*+*+*'", + "worker:trigger-sync": "curl 'http://localhost:8787/__scheduled?cron=0+*+*+*+*'", "postinstall": "nuxt prepare", "prepublishOnly": "nr build", "typecheck": "nuxt typecheck && nr -r typecheck", diff --git a/server/tasks/cron/sync.test.ts b/server/tasks/cron/sync.test.ts index 3757945..9436d5b 100644 --- a/server/tasks/cron/sync.test.ts +++ b/server/tasks/cron/sync.test.ts @@ -69,11 +69,11 @@ afterAll(() => { }) describe('scheduled synchronization order', () => { - it('records the six-hour schedule', async () => { + it('records the hourly schedule', async () => { await task.run({ payload: {}, context: {} } as never) expect(mocks.insertValues).toHaveBeenCalledWith(expect.objectContaining({ - cron: '0 */6 * * *', + cron: '0 * * * *', })) }) diff --git a/server/tasks/cron/sync.ts b/server/tasks/cron/sync.ts index bc30226..1629154 100644 --- a/server/tasks/cron/sync.ts +++ b/server/tasks/cron/sync.ts @@ -4,7 +4,7 @@ import { runTask } from 'nitropack/runtime' import { runTasksBestEffort } from '~~/server/utils/cron-task-runner' import { eq, tables, useDrizzle } from '~~/server/utils/drizzle' -const CRON_EXPRESSION = '0 */6 * * *' +const CRON_EXPRESSION = '0 * * * *' const TASKS: string[] = ['sync:epochs', 'sync:snapshot', 'sync:scores'] interface FailedTask { diff --git a/server/tasks/sync/epochs.test.ts b/server/tasks/sync/epochs.test.ts index 8d82556..cefbfde 100644 --- a/server/tasks/sync/epochs.test.ts +++ b/server/tasks/sync/epochs.test.ts @@ -103,7 +103,7 @@ beforeEach(async () => { mocks.synchronizeCompletedEpoch .mockResolvedValueOnce({ epochNumber: 100, status: 'failed', error: 'first RPC failure', repairAttempted: true }) .mockResolvedValueOnce({ epochNumber: 99, status: 'verified' }) - .mockResolvedValueOnce({ epochNumber: 98, status: 'failed', error: 'repair budget exhausted', repairAttempted: false }) + .mockResolvedValueOnce({ epochNumber: 98, status: 'verified' }) .mockResolvedValueOnce({ epochNumber: 97, status: 'verified' }) .mockResolvedValueOnce({ epochNumber: 96, status: 'already_finalized' }) mocks.sendSyncFailureNotification.mockResolvedValue(undefined) @@ -151,7 +151,8 @@ describe('planned completed-epoch synchronization', () => { expect(isFullRepairAllowed(0, true)).toBe(true) expect(isFullRepairAllowed(5, true)).toBe(true) expect(isFullRepairAllowed(0, false)).toBe(true) - expect(isFullRepairAllowed(1, false)).toBe(false) + expect(isFullRepairAllowed(1, false)).toBe(true) + expect(isFullRepairAllowed(4, false)).toBe(false) }) it('limits production repairs while returning the failed epoch details', async () => { @@ -173,27 +174,83 @@ describe('planned completed-epoch synchronization', () => { limit: 50, syncingLeaseDurationMs: 6 * 60 * 60 * 1000, })) - expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(1, 100, 'testnet', { allowRepair: true }) - expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(2, 99, 'testnet', { allowRepair: false }) - expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(3, 98, 'testnet', { allowRepair: false }) - expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(4, 97, 'testnet', { allowRepair: false }) - expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(5, 96, 'testnet', { allowRepair: false }) + expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(1, 100, 'testnet') + expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(2, 99, 'testnet') + expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(3, 98, 'testnet') + expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(4, 97, 'testnet') + expect(mocks.synchronizeCompletedEpoch).toHaveBeenNthCalledWith(5, 96, 'testnet') expect(mocks.synchronizeCompletedEpoch).toHaveBeenCalledTimes(5) expect(result).toEqual({ result: { success: false, error: '[sync:epochs] epoch 100 failed: first RPC failure', - totalSynced: 2, - epochsSynced: [99, 97], + totalSynced: 3, + epochsSynced: [99, 98, 97], + deferredEpochs: [], outcomes: [ { epochNumber: 100, status: 'failed', error: 'first RPC failure', repairAttempted: true }, { epochNumber: 99, status: 'verified' }, - { epochNumber: 98, status: 'failed', error: 'repair budget exhausted', repairAttempted: false }, + { epochNumber: 98, status: 'verified' }, { epochNumber: 97, status: 'verified' }, { epochNumber: 96, status: 'already_finalized' }, ], }, }) }) + it('repairs a bounded batch and leaves the remaining epochs untouched for the next run', async () => { + mocks.planEpochSync.mockReturnValue([100, 99, 98, 97, 96, 95]) + mocks.synchronizeCompletedEpoch.mockReset().mockImplementation(async (epochNumber: number) => ({ + epochNumber, + status: 'repaired', + })) + + const { result } = await task.run() + + expect(result).toMatchObject({ + success: true, + totalSynced: 4, + epochsSynced: [100, 99, 98, 97], + deferredEpochs: [96, 95], + }) + expect(mocks.synchronizeCompletedEpoch.mock.calls.map(([epoch]) => epoch)).toEqual([100, 99, 98, 97]) + expect(mocks.sendSyncFailureNotification).not.toHaveBeenCalled() + }) + + it('counts failed repair attempts toward the batch limit and retains real failures', async () => { + mocks.synchronizeCompletedEpoch.mockReset().mockImplementation(async (epochNumber: number) => ({ + epochNumber, + status: 'failed', + error: 'RPC unavailable', + repairAttempted: true, + })) + + const { result } = await task.run() + + expect(result).toMatchObject({ success: false, totalSynced: 0, deferredEpochs: [96] }) + expect(mocks.synchronizeCompletedEpoch).toHaveBeenCalledTimes(4) + expect(mocks.sendSyncFailureNotification).toHaveBeenCalledWith('missing-epoch', expect.objectContaining({ + message: expect.stringContaining('RPC unavailable'), + })) + }) + + it('stops starting epochs after the runtime budget while preserving completed progress', async () => { + vi.useFakeTimers() + try { + vi.setSystemTime(new Date('2026-09-24T10:00:00Z')) + mocks.synchronizeCompletedEpoch.mockReset().mockImplementation(async (epochNumber: number) => { + vi.setSystemTime(new Date('2026-09-24T10:08:00Z')) + return { epochNumber, status: 'repaired' } + }) + + const { result } = await task.run() + + expect(result).toMatchObject({ success: true, epochsSynced: [100], deferredEpochs: [99, 98, 97, 96] }) + expect(mocks.synchronizeCompletedEpoch).toHaveBeenCalledTimes(1) + expect(mocks.sendSyncFailureNotification).not.toHaveBeenCalled() + } + finally { + vi.useRealTimers() + } + }) }) diff --git a/server/tasks/sync/epochs.ts b/server/tasks/sync/epochs.ts index ceed21b..bcc59e4 100644 --- a/server/tasks/sync/epochs.ts +++ b/server/tasks/sync/epochs.ts @@ -14,6 +14,8 @@ import { getOldestV2BackfillEpoch } from '~~/server/utils/scores' import { sendSyncFailureNotification } from '~~/server/utils/slack' const MAX_PRODUCTION_EPOCH_CANDIDATES_PER_RUN = 50 +const MAX_PRODUCTION_REPAIR_ATTEMPTS_PER_RUN = 4 +const PRODUCTION_EPOCH_START_BUDGET_MS = 8 * 60 * 1000 const DEVELOPMENT_SYNCING_LEASE_DURATION_MS = 5 * 60 * 1000 export function getEpochCandidateLimit(isDevelopment: boolean): number { @@ -27,7 +29,7 @@ export function getSyncingLeaseDurationMs(isDevelopment: boolean): number { } export function isFullRepairAllowed(fullRepairAttempts: number, isDevelopment: boolean): boolean { - return isDevelopment || fullRepairAttempts === 0 + return isDevelopment || fullRepairAttempts < MAX_PRODUCTION_REPAIR_ATTEMPTS_PER_RUN } function formatError(error: unknown): string { @@ -59,6 +61,7 @@ export default defineTask({ }, async run() { const config = useSafeRuntimeConfig() + const startedAt = Date.now() try { const rpcUrl = getRpcUrl() @@ -97,11 +100,15 @@ export default defineTask({ let fullRepairAttempts = 0 for (const epochNumber of plannedEpochs) { + // Leave unstarted epochs eligible for the next run instead of claiming and failing them. + if (!isFullRepairAllowed(fullRepairAttempts, import.meta.dev) + || (!import.meta.dev && Date.now() - startedAt >= PRODUCTION_EPOCH_START_BUDGET_MS)) { + break + } + let outcome: EpochSyncOutcome try { - outcome = await synchronizeCompletedEpoch(epochNumber, config.public.nimiqNetwork, { - allowRepair: isFullRepairAllowed(fullRepairAttempts, import.meta.dev), - }) + outcome = await synchronizeCompletedEpoch(epochNumber, config.public.nimiqNetwork) } catch (error) { outcome = { epochNumber, status: 'failed', error: formatError(error), repairAttempted: false } @@ -114,6 +121,9 @@ export default defineTask({ fullRepairAttempts++ } + const deferredEpochs = plannedEpochs.slice(outcomes.length) + consola.info(`[sync:epochs] finalized ${epochsSynced.length}, repair attempts ${fullRepairAttempts}, deferred ${deferredEpochs.length} planned epoch(s)`) + const failures = outcomes.filter((outcome): outcome is Extract => outcome.status === 'failed') const recentFailure = failures.find(outcome => outcome.epochNumber >= recentRange.fromEpoch) const historicalFailures = failures.filter(outcome => outcome.epochNumber < recentRange.fromEpoch) @@ -135,6 +145,7 @@ export default defineTask({ ...(failureSummaries.length > 0 ? { error: failureSummaries.join('; ') } : {}), totalSynced: epochsSynced.length, epochsSynced, + deferredEpochs, outcomes, }, } diff --git a/server/utils/activity-sync.test.ts b/server/utils/activity-sync.test.ts index 9aaf16e..857508d 100644 --- a/server/utils/activity-sync.test.ts +++ b/server/utils/activity-sync.test.ts @@ -83,7 +83,6 @@ function createDependencies() { repairCompletedEpoch: vi.fn().mockResolvedValue(undefined), markActivityEpochFailed: vi.fn().mockResolvedValue(undefined), now: vi.fn().mockReturnValue(new Date('2026-07-29T12:00:00.000Z')), - allowRepair: true, } } @@ -123,6 +122,7 @@ describe('synchronizeCompletedEpoch', () => { expect(dependencies.fetchActivity).toHaveBeenCalledWith(EPOCH, expect.objectContaining({ network: 'testnet', electionSet, + maxRetries: 1, })) expect(dependencies.repairCompletedEpoch).toHaveBeenCalledWith( EPOCH, @@ -153,28 +153,6 @@ describe('synchronizeCompletedEpoch', () => { expect(dependencies.repairCompletedEpoch).not.toHaveBeenCalled() }) - it('defers a mismatched epoch without fetching activity when repair is disallowed', async () => { - const dependencies = createDependencies() - dependencies.getStoredFinalizedElectedAddresses.mockResolvedValue([ADDRESS_A]) - dependencies.allowRepair = false - - const outcome = await synchronizeCompletedEpoch(EPOCH, 'testnet', dependencies) - - expect(outcome).toEqual({ - epochNumber: EPOCH, - status: 'failed', - error: expect.stringMatching(/repair|defer|budget/i), - repairAttempted: false, - }) - expect(dependencies.fetchActivity).not.toHaveBeenCalled() - expect(dependencies.repairCompletedEpoch).not.toHaveBeenCalled() - expect(dependencies.markActivityEpochFailed).toHaveBeenCalledWith( - EPOCH, - expect.objectContaining({ message: expect.stringMatching(/repair|defer|budget/i) }), - '2026-07-29T12:00:00.000Z', - ) - }) - it('returns already_finalized without fetching when attempt claim is refused', async () => { const dependencies = createDependencies() dependencies.beginActivityEpochAttempt.mockResolvedValue(false) diff --git a/server/utils/activity-sync.ts b/server/utils/activity-sync.ts index d1514dd..b936731 100644 --- a/server/utils/activity-sync.ts +++ b/server/utils/activity-sync.ts @@ -25,13 +25,12 @@ type ActivityResult = [success: boolean, error?: string, activity?: EpochActivit export interface ActivitySyncDependencies { beginActivityEpochAttempt: (epochNumber: number, startedAt: string) => Promise fetchElectionSet: (epochNumber: number, options: { network: string }) => Promise - fetchActivity: (epochNumber: number, options: { network: string, electionSet: ElectionSet, maxBatchSize?: number }) => Promise + fetchActivity: (epochNumber: number, options: { network: string, electionSet: ElectionSet, maxBatchSize?: number, maxRetries?: number }) => Promise getStoredFinalizedElectedAddresses: (epochNumber: number) => Promise finalizeVerifiedEpoch: (epochNumber: number, electionSet: ElectionSet, startedAt: string) => Promise repairCompletedEpoch: (epochNumber: number, electionSet: ElectionSet, activity: EpochActivity, startedAt: string) => Promise markActivityEpochFailed: (epochNumber: number, error: Error, startedAt: string) => Promise now: () => Date - allowRepair: boolean } async function getStoredFinalizedElectedAddresses(epochNumber: number): Promise { @@ -87,7 +86,6 @@ const defaultDependencies: ActivitySyncDependencies = { repairCompletedEpoch: repairCompletedEpochDefault, markActivityEpochFailed, now: () => new Date(), - allowRepair: true, } export async function synchronizeCompletedEpoch( @@ -116,14 +114,12 @@ export async function synchronizeCompletedEpoch( return { epochNumber, status: 'verified' } } - if (!dependencies.allowRepair) - throw new Error(`Repair budget exhausted; epoch ${epochNumber} deferred`) - repairAttempted = true const [activityOk, activityError, activity] = await dependencies.fetchActivity(epochNumber, { network, electionSet, maxBatchSize: import.meta.dev ? 120 : 6, + maxRetries: import.meta.dev ? 5 : 1, }) if (!activityOk || !activity) throw new Error(activityError || 'Unable to fetch activity') diff --git a/wrangler.json b/wrangler.json index 78f7223..e761509 100644 --- a/wrangler.json +++ b/wrangler.json @@ -13,7 +13,7 @@ "invocation_logs": true } }, - "triggers": { "crons": ["0 */6 * * *"] }, + "triggers": { "crons": ["0 * * * *"] }, "vars": { "NUXT_PUBLIC_NIMIQ_NETWORK": "main-albatross", "NUXT_SCORE_V2_MODE": "off" }, "d1_databases": [{ "binding": "DB", "database_id": "cc9f1d25-676b-4cb3-8af6-887e85a08baa", "database_name": "validators-api-mainnet" }], "kv_namespaces": [{ "binding": "CACHE", "id": "4be4d10e6d3444eca6e7c02cdbcd275f" }], @@ -23,7 +23,7 @@ "name": "validators-api-test", "compatibility_flags": ["nodejs_compat"], "observability": { "enabled": true, "logs": { "enabled": true, "invocation_logs": true } }, - "triggers": { "crons": ["0 */6 * * *"] }, + "triggers": { "crons": ["0 * * * *"] }, "vars": { "NUXT_PUBLIC_NIMIQ_NETWORK": "test-albatross", "NUXT_SCORE_V2_MODE": "off" }, "d1_databases": [{ "binding": "DB", "database_id": "de14e353-5028-4e52-a383-a9cc200d960d", "database_name": "validators-api-testnet" }], "kv_namespaces": [{ "binding": "CACHE", "id": "90c1598af92a4e72a767e6601090f014" }],