Skip to content
Open
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
109 changes: 53 additions & 56 deletions backend/src/services/scheduleExecutor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -100,6 +101,8 @@ export class ScheduleExecutor {
let failureCount = 0;

for (const scheduleRow of claimedSchedules) {
let executionResult: ExecutionResult;

try {
const schedule: Schedule = {
...scheduleRow,
Expand All @@ -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);
}
}

Expand All @@ -169,20 +171,6 @@ export class ScheduleExecutor {
}
}

/**
* Release the row-level claim after execution (success or failure).
*/
private async releaseClaim(scheduleId: number): Promise<void> {
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.
Expand Down Expand Up @@ -298,15 +286,15 @@ export class ScheduleExecutor {
* @param scheduleId - The schedule ID
* @param result - The execution result
*/
async recordExecution(scheduleId: number, result: ExecutionResult): Promise<void> {
const client = await pool.connect();
async recordExecution(
scheduleId: number,
result: ExecutionResult,
client: PoolClient
): Promise<void> {
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,
Expand All @@ -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,
Expand All @@ -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}`
);
}
}
}
Expand Down
59 changes: 31 additions & 28 deletions backend/src/services/scheduleService.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -361,14 +362,18 @@ export class ScheduleService {
async updateAfterExecution(
scheduleId: number,
executionResult: ExecutionResult,
client?: PoolClient,
): Promise<void> {
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",
Expand All @@ -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`);
Expand All @@ -388,55 +393,53 @@ 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),
updated_at = CURRENT_TIMESTAMP
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();
}
}
}
}
Expand Down