diff --git a/backend/src/services/scheduleExecutor.ts b/backend/src/services/scheduleExecutor.ts index f34c242a..33eda18f 100644 --- a/backend/src/services/scheduleExecutor.ts +++ b/backend/src/services/scheduleExecutor.ts @@ -6,6 +6,7 @@ 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'; export class ScheduleExecutor { private cronJob: ScheduledTask | null = null; @@ -100,6 +101,8 @@ export class ScheduleExecutor { let failureCount = 0; for (const scheduleRow of claimedSchedules) { + let executionResult: ExecutionResult; + try { const schedule: Schedule = { ...scheduleRow, @@ -114,45 +117,44 @@ export class ScheduleExecutor { }; console.log(`[ScheduleExecutor] Executing schedule ID ${schedule.id} (Scheduled for: ${schedule.nextRunTimestamp.toISOString()})`); + executionResult = await this.executeSchedule(schedule); + } catch (error) { + executionResult = { + success: false, + error: { + message: error instanceof Error ? error.message : 'System error in executor', + details: error as any, + }, + }; + console.error( + `[ScheduleExecutor] Error executing schedule ID ${scheduleRow.id}:`, + error + ); + } - const executionResult = await this.executeSchedule(schedule); - - await this.recordExecution(schedule.id, executionResult); + try { + // Commit execution history and schedule state before releasing the claim. + // If commit fails, the claim stays in place for stale-claim recovery. + await this.recordExecution(scheduleRow.id, executionResult, client); if (executionResult.success) { successCount++; - console.log(`[ScheduleExecutor] Schedule ID ${schedule.id} executed successfully. Hash: ${executionResult.transactionHash}`); + console.log( + `[ScheduleExecutor] Schedule ID ${scheduleRow.id} executed successfully. Hash: ${executionResult.transactionHash}` + ); } else { failureCount++; console.error( - `[ScheduleExecutor] Schedule ID ${schedule.id} failed:`, + `[ScheduleExecutor] Schedule ID ${scheduleRow.id} failed:`, executionResult.error?.message ); } - } catch (error) { + } catch (recordError) { failureCount++; console.error( - `[ScheduleExecutor] Error processing schedule ID ${scheduleRow.id}:`, - error + `[ScheduleExecutor] Failed to finalize schedule ID ${scheduleRow.id}; claim retained for recovery:`, + recordError ); - - try { - await this.recordExecution(scheduleRow.id, { - success: false, - error: { - message: error instanceof Error ? error.message : 'System error in executor', - details: error as any, - }, - }); - } catch (recordError) { - console.error( - `[ScheduleExecutor] Failed to record execution error for schedule ID ${scheduleRow.id}:`, - recordError - ); - } - } finally { - // Always release the claim so the row is available for the next cycle - await this.releaseClaim(scheduleRow.id); } } @@ -169,20 +171,6 @@ export class ScheduleExecutor { } } - /** - * Release the row-level claim after execution (success or failure). - */ - private async releaseClaim(scheduleId: number): Promise { - try { - await pool.query( - 'UPDATE schedules SET locked_by = NULL, locked_at = NULL WHERE id = $1', - [scheduleId] - ); - } catch (error) { - console.error(`[ScheduleExecutor] Failed to release claim for schedule ID ${scheduleId}:`, error); - } - } - /** * Release stale claims held by crashed pods (locked_at older than 5 minutes). * Called once per cron cycle before claiming new schedules. @@ -298,15 +286,15 @@ export class ScheduleExecutor { * @param scheduleId - The schedule ID * @param result - The execution result */ - async recordExecution(scheduleId: number, result: ExecutionResult): Promise { - const client = await pool.connect(); + async recordExecution( + scheduleId: number, + result: ExecutionResult, + client: PoolClient + ): Promise { try { await client.query('BEGIN'); - // Determine execution status const status = result.success ? 'success' : 'failed'; - - // Insert into execution_history const insertQuery = ` INSERT INTO execution_history ( schedule_id, @@ -323,7 +311,7 @@ export class ScheduleExecutor { const insertValues = [ scheduleId, - new Date(), // executed_at + new Date(), status, result.transactionHash || null, result.success ? JSON.stringify({ hash: result.transactionHash }) : null, @@ -333,21 +321,30 @@ export class ScheduleExecutor { await client.query(insertQuery, insertValues); - // Update schedule state using ScheduleService - await scheduleService.updateAfterExecution(scheduleId, result); - - // Clear the lock now that execution is recorded - await client.query( - 'UPDATE schedules SET locked_by = NULL, locked_at = NULL WHERE id = $1', - [scheduleId] - ); + // Use the claim-owning connection so history and schedule-state mutation + // commit together without opening a nested transaction. + await scheduleService.updateAfterExecution(scheduleId, result, client); await client.query('COMMIT'); } catch (error) { await client.query('ROLLBACK'); throw error; - } finally { - client.release(); + } + + // The only unlock happens after execution state has committed, on the same + // physical connection, and only while this executor still owns the claim. + const releaseResult = await client.query( + `UPDATE schedules + SET locked_by = NULL, locked_at = NULL + WHERE id = $1 AND locked_by = $2 + RETURNING id`, + [scheduleId, this.podId] + ); + + if (releaseResult.rowCount !== 1) { + throw new Error( + `Schedule ${scheduleId} lock is no longer owned by executor ${this.podId}` + ); } } } diff --git a/backend/src/services/scheduleService.ts b/backend/src/services/scheduleService.ts index 5072bd34..c1362b16 100644 --- a/backend/src/services/scheduleService.ts +++ b/backend/src/services/scheduleService.ts @@ -1,5 +1,6 @@ import { DateTime } from 'luxon'; import { default as pool } from '../config/database.js'; +import type { PoolClient } from 'pg'; import type { Schedule, ScheduleFrequency, @@ -361,14 +362,18 @@ export class ScheduleService { async updateAfterExecution( scheduleId: number, executionResult: ExecutionResult, + client?: PoolClient, ): Promise { - const client = await pool.connect(); + const ownsClient = !client; + const dbClient = client ?? (await pool.connect()); + try { - await client.query('BEGIN'); + if (ownsClient) { + await dbClient.query('BEGIN'); + } - // Query the schedule to get its frequency and configuration const selectQuery = ` - SELECT + SELECT id, frequency, time_of_day as "timeOfDay", @@ -379,7 +384,7 @@ export class ScheduleService { WHERE id = $1 `; - const selectResult = await client.query(selectQuery, [scheduleId]); + const selectResult = await dbClient.query(selectQuery, [scheduleId]); if (selectResult.rows.length === 0) { throw new Error(`Schedule with ID ${scheduleId} not found`); @@ -388,35 +393,27 @@ export class ScheduleService { const schedule = selectResult.rows[0]; const executionTime = new Date(); - // Determine the new status and next_run_timestamp based on execution result let newStatus: string; let nextRunTimestamp: Date | null = null; if (!executionResult.success) { - // If execution failed, set status to 'failed' newStatus = 'failed'; + } else if (schedule.frequency === 'once') { + newStatus = 'completed'; } else { - // Execution succeeded - if (schedule.frequency === 'once') { - // For one-time schedules, set status to 'completed' - newStatus = 'completed'; - } else { - // For recurring schedules, calculate new next_run_timestamp and keep status 'active' - newStatus = 'active'; - nextRunTimestamp = this.calculateNextRun( - schedule.frequency, - schedule.timeOfDay, - new Date(schedule.startDate), - schedule.timezone, - executionTime, // Use execution time as lastRun - ); - } + newStatus = 'active'; + nextRunTimestamp = this.calculateNextRun( + schedule.frequency, + schedule.timeOfDay, + new Date(schedule.startDate), + schedule.timezone, + executionTime, + ); } - // Update the schedule in the database const updateQuery = ` UPDATE schedules - SET + SET last_run_timestamp = $1, status = $2, next_run_timestamp = COALESCE($3, next_run_timestamp), @@ -424,19 +421,25 @@ export class ScheduleService { WHERE id = $4 `; - await client.query(updateQuery, [ + await dbClient.query(updateQuery, [ executionTime, newStatus, nextRunTimestamp, scheduleId, ]); - await client.query('COMMIT'); + if (ownsClient) { + await dbClient.query('COMMIT'); + } } catch (error) { - await client.query('ROLLBACK'); + if (ownsClient) { + await dbClient.query('ROLLBACK'); + } throw error; } finally { - client.release(); + if (ownsClient) { + dbClient.release(); + } } } }