diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 52029992..9d71b0a9 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -115,4 +115,3 @@ jobs: run: npm test --if-present # Workflow run retention settings -retention-days: 30 \ No newline at end of file diff --git a/.github/workflows/contract-release.yml b/.github/workflows/contract-release.yml index 2655fcbc..3f66b1a5 100644 --- a/.github/workflows/contract-release.yml +++ b/.github/workflows/contract-release.yml @@ -29,4 +29,3 @@ jobs: release_token: ${{ secrets.GITHUB_TOKEN }} # Workflow run retention settings -retention-days: 90 \ No newline at end of file diff --git a/.github/workflows/dapp-ipfs.yml b/.github/workflows/dapp-ipfs.yml index 29d16e55..42248532 100644 --- a/.github/workflows/dapp-ipfs.yml +++ b/.github/workflows/dapp-ipfs.yml @@ -67,4 +67,3 @@ jobs: echo "- URL: ${{ steps.storacha.outputs.url }}" >> "$GITHUB_STEP_SUMMARY" # Workflow run retention settings -retention-days: 30 \ No newline at end of file diff --git a/.github/workflows/secrets-check.yml b/.github/workflows/secrets-check.yml index b8d3b7e8..c7a34807 100644 --- a/.github/workflows/secrets-check.yml +++ b/.github/workflows/secrets-check.yml @@ -22,4 +22,3 @@ jobs: run: ./scripts/check-k8s-secrets.sh # Workflow run retention settings -retention-days: 30 \ No newline at end of file diff --git a/backend/src/services/__tests__/scheduleExecutor.test.ts b/backend/src/services/__tests__/scheduleExecutor.test.ts index 75c4229e..4e25682a 100644 --- a/backend/src/services/__tests__/scheduleExecutor.test.ts +++ b/backend/src/services/__tests__/scheduleExecutor.test.ts @@ -58,11 +58,11 @@ describe('ScheduleExecutor', () => { }); describe('initialize', () => { - it('should set up cron job to run every minute', () => { + it('should probe scheduler leadership every 15 seconds', () => { executor.initialize(); expect(mockCron.schedule).toHaveBeenCalledWith( - '* * * * *', + '*/15 * * * * *', expect.any(Function) ); }); @@ -78,6 +78,53 @@ describe('ScheduleExecutor', () => { consoleSpy.mockRestore(); }); + + it('should skip the scheduler pass when another pod holds leadership', async () => { + const leaderQuery = jest.fn().mockResolvedValueOnce({ + rows: [{ acquired: false }], + }); + const leaderRelease = jest.fn(); + (mockPool.connect as jest.Mock).mockResolvedValueOnce({ + query: leaderQuery, + release: leaderRelease, + }); + const processSpy = jest.spyOn(executor, 'processDueSchedules').mockResolvedValue(); + + executor.initialize(); + const callback = (mockCron.schedule as jest.Mock).mock.calls[0][1] as () => Promise; + await callback(); + + expect(leaderQuery).toHaveBeenCalledWith( + expect.stringContaining('pg_try_advisory_lock'), + [expect.any(Number), expect.any(Number)] + ); + expect(processSpy).not.toHaveBeenCalled(); + expect(leaderRelease).toHaveBeenCalledWith(false); + }); + + it('should hold and release leadership around one scheduler pass', async () => { + const leaderQuery = jest.fn() + .mockResolvedValueOnce({ rows: [{ acquired: true }] }) + .mockResolvedValueOnce({ rows: [{ unlocked: true }] }); + const leaderRelease = jest.fn(); + (mockPool.connect as jest.Mock).mockResolvedValueOnce({ + query: leaderQuery, + release: leaderRelease, + }); + const processSpy = jest.spyOn(executor, 'processDueSchedules').mockResolvedValue(); + + executor.initialize(); + const callback = (mockCron.schedule as jest.Mock).mock.calls[0][1] as () => Promise; + await callback(); + + expect(processSpy).toHaveBeenCalledTimes(1); + expect(leaderQuery).toHaveBeenNthCalledWith( + 2, + expect.stringContaining('pg_advisory_unlock'), + [expect.any(Number), expect.any(Number)] + ); + expect(leaderRelease).toHaveBeenCalledWith(false); + }); }); describe('stop', () => { diff --git a/backend/src/services/scheduleExecutor.ts b/backend/src/services/scheduleExecutor.ts index f34c242a..94a40ae5 100644 --- a/backend/src/services/scheduleExecutor.ts +++ b/backend/src/services/scheduleExecutor.ts @@ -6,31 +6,79 @@ import { scheduleService } from './scheduleService.js'; import type { Schedule, ExecutionResult, PaymentRecipient } from '../types/schedule.js'; import { Operation, Asset, Memo, Keypair } from '@stellar/stellar-sdk'; import os from 'node:os'; +import type { PoolClient } from 'pg'; + +const LEADER_ELECTION_INTERVAL = '*/15 * * * * *'; +// Two int32 advisory-lock keys: ASCII-ish PAYD / SCHD namespaces. +const SCHEDULER_LOCK_NAMESPACE = 0x50415944; +const SCHEDULER_LOCK_KEY = 0x53434844; export class ScheduleExecutor { private cronJob: ScheduledTask | null = null; private readonly podId: string; + private runInProgress = false; constructor() { this.podId = `${os.hostname()}-${process.pid}`; } /** - * Initialize the cron job to run every minute - * Sets up node-cron job with error handling and logging + * Probe scheduler leadership every 15 seconds. + * + * The advisory lock is session-scoped, so PostgreSQL releases it automatically + * if the leader pod dies or loses its database connection. Keeping the lock on + * a dedicated client for the full scheduler pass guarantees that at most one + * pod enters processDueSchedules at a time. */ initialize(): void { - // Cron expression: run every minute - this.cronJob = cron.schedule('* * * * *', async () => { + this.cronJob = cron.schedule(LEADER_ELECTION_INTERVAL, async () => { + if (this.runInProgress) { + return; + } + + this.runInProgress = true; + let leaderClient: PoolClient | null = null; + let hasLeadership = false; + let destroyLeaderConnection = false; + try { - console.log('[ScheduleExecutor] Running scheduled task check...'); + leaderClient = await pool.connect(); + const election = await leaderClient.query<{ acquired: boolean }>( + 'SELECT pg_try_advisory_lock($1, $2) AS acquired', + [SCHEDULER_LOCK_NAMESPACE, SCHEDULER_LOCK_KEY] + ); + + hasLeadership = election.rows[0]?.acquired === true; + if (!hasLeadership) { + return; + } + await this.processDueSchedules(); } catch (error) { - console.error('[ScheduleExecutor] Error in cron job execution:', error); + console.error('[ScheduleExecutor] Error in leader scheduler execution:', error); + } finally { + if (leaderClient) { + if (hasLeadership) { + try { + const unlock = await leaderClient.query<{ unlocked: boolean }>( + 'SELECT pg_advisory_unlock($1, $2) AS unlocked', + [SCHEDULER_LOCK_NAMESPACE, SCHEDULER_LOCK_KEY] + ); + destroyLeaderConnection = unlock.rows[0]?.unlocked !== true; + } catch (error) { + destroyLeaderConnection = true; + console.error('[ScheduleExecutor] Failed to release scheduler leadership:', error); + } + } + + leaderClient.release(destroyLeaderConnection); + } + + this.runInProgress = false; } }); - console.log('[ScheduleExecutor] Cron job initialized - running every minute'); + console.log('[ScheduleExecutor] Cron job initialized - checking leadership every 15 seconds'); } /**