Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
4 changes: 2 additions & 2 deletions MIGRATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.
Expand Down
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion app/app.vue
Original file line number Diff line number Diff line change
Expand Up @@ -207,7 +207,7 @@ const currentEnvItem = getEnvironmentItem(nimiqNetwork) ?? { network: nimiqNetwo
<hr f-my-sm border-red-600>

<p f-mt-md text="f-sm red-1100/80">
<strong>Note:</strong> 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.
<strong>Note:</strong> 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.
</p>
</div>

Expand Down
4 changes: 2 additions & 2 deletions nuxt.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand Down
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
4 changes: 2 additions & 2 deletions server/tasks/cron/sync.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 * * * *',
}))
})

Expand Down
2 changes: 1 addition & 1 deletion server/tasks/cron/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
77 changes: 67 additions & 10 deletions server/tasks/sync/epochs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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 () => {
Expand All @@ -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()
}
})
})
19 changes: 15 additions & 4 deletions server/tasks/sync/epochs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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 {
Expand Down Expand Up @@ -59,6 +61,7 @@ export default defineTask({
},
async run() {
const config = useSafeRuntimeConfig()
const startedAt = Date.now()

try {
const rpcUrl = getRpcUrl()
Expand Down Expand Up @@ -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 }
Expand All @@ -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<EpochSyncOutcome, { status: 'failed' }> => outcome.status === 'failed')
const recentFailure = failures.find(outcome => outcome.epochNumber >= recentRange.fromEpoch)
const historicalFailures = failures.filter(outcome => outcome.epochNumber < recentRange.fromEpoch)
Expand All @@ -135,6 +145,7 @@ export default defineTask({
...(failureSummaries.length > 0 ? { error: failureSummaries.join('; ') } : {}),
totalSynced: epochsSynced.length,
epochsSynced,
deferredEpochs,
outcomes,
},
}
Expand Down
24 changes: 1 addition & 23 deletions server/utils/activity-sync.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
}

Expand Down Expand Up @@ -123,6 +122,7 @@ describe('synchronizeCompletedEpoch', () => {
expect(dependencies.fetchActivity).toHaveBeenCalledWith(EPOCH, expect.objectContaining({
network: 'testnet',
electionSet,
maxRetries: 1,
}))
expect(dependencies.repairCompletedEpoch).toHaveBeenCalledWith(
EPOCH,
Expand Down Expand Up @@ -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)
Expand Down
8 changes: 2 additions & 6 deletions server/utils/activity-sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,12 @@ type ActivityResult = [success: boolean, error?: string, activity?: EpochActivit
export interface ActivitySyncDependencies {
beginActivityEpochAttempt: (epochNumber: number, startedAt: string) => Promise<boolean>
fetchElectionSet: (epochNumber: number, options: { network: string }) => Promise<ElectionSetResult>
fetchActivity: (epochNumber: number, options: { network: string, electionSet: ElectionSet, maxBatchSize?: number }) => Promise<ActivityResult>
fetchActivity: (epochNumber: number, options: { network: string, electionSet: ElectionSet, maxBatchSize?: number, maxRetries?: number }) => Promise<ActivityResult>
getStoredFinalizedElectedAddresses: (epochNumber: number) => Promise<string[]>
finalizeVerifiedEpoch: (epochNumber: number, electionSet: ElectionSet, startedAt: string) => Promise<void>
repairCompletedEpoch: (epochNumber: number, electionSet: ElectionSet, activity: EpochActivity, startedAt: string) => Promise<void>
markActivityEpochFailed: (epochNumber: number, error: Error, startedAt: string) => Promise<void>
now: () => Date
allowRepair: boolean
}

async function getStoredFinalizedElectedAddresses(epochNumber: number): Promise<string[]> {
Expand Down Expand Up @@ -87,7 +86,6 @@ const defaultDependencies: ActivitySyncDependencies = {
repairCompletedEpoch: repairCompletedEpochDefault,
markActivityEpochFailed,
now: () => new Date(),
allowRepair: true,
}

export async function synchronizeCompletedEpoch(
Expand Down Expand Up @@ -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')
Expand Down
4 changes: 2 additions & 2 deletions wrangler.json
Original file line number Diff line number Diff line change
Expand Up @@ -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" }],
Expand All @@ -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" }],
Expand Down
Loading