@@ -34,66 +34,83 @@ export function calculateNextRunDate(current: Date, frequency: string): Date {
3434 return next ;
3535}
3636
37+ let isRunning = false ;
38+
3739export async function runDueScheduledTransactions ( ) : Promise < void > {
38- const now = new Date ( ) ;
40+ if ( isRunning ) return ;
41+ isRunning = true ;
42+
43+ try {
44+ const now = new Date ( ) ;
3945
40- const due = await db . query . scheduledTransactions . findMany ( {
41- where : and (
42- eq ( scheduledTransactions . status , 'active' ) ,
43- lte ( scheduledTransactions . nextRunDate , now )
44- ) ,
45- } ) ;
46+ const due = await db . query . scheduledTransactions . findMany ( {
47+ where : and (
48+ eq ( scheduledTransactions . status , 'active' ) ,
49+ lte ( scheduledTransactions . nextRunDate , now )
50+ ) ,
51+ } ) ;
4652
47- if ( due . length === 0 ) return ;
53+ if ( due . length === 0 ) return ;
4854
49- logger . info ( 'Scheduled runner: processing due transactions' , { count : due . length } ) ;
55+ logger . info ( 'Scheduled runner: processing due transactions' , { count : due . length } ) ;
5056
51- let succeeded = 0 ;
52- let failed = 0 ;
57+ let succeeded = 0 ;
58+ let failed = 0 ;
5359
54- for ( const scheduled of due ) {
55- try {
56- // Catch-up loop: create one transaction per missed occurrence, capped at 12
57- // to avoid flooding months of backlog in one shot.
60+ for ( const scheduled of due ) {
5861 let baseDate = new Date ( scheduled . nextRunDate ) ;
5962 let created = 0 ;
6063 const MAX_CATCH_UP = 12 ;
6164
6265 while ( baseDate <= now && created < MAX_CATCH_UP ) {
63- // Use the actual scheduled date, not today — so "gaji tgl 25" is recorded on the 25th
64- await createTransaction ( {
65- userId : scheduled . userId ,
66- type : scheduled . type as 'income' | 'expense' | 'transfer' ,
67- amount : String ( scheduled . amount ) ,
68- categoryId : scheduled . categoryId ,
69- accountId : scheduled . accountId ?? undefined ,
70- toAccountId : scheduled . toAccountId ?? undefined ,
71- description : scheduled . description ?? undefined ,
72- date : baseDate . toISOString ( ) . split ( 'T' ) [ 0 ] ,
73- } ) ;
66+ try {
67+ // Use the actual scheduled date so "gaji tgl 25" records on the 25th
68+ await createTransaction ( {
69+ userId : scheduled . userId ,
70+ type : scheduled . type as 'income' | 'expense' | 'transfer' ,
71+ amount : String ( scheduled . amount ) ,
72+ categoryId : scheduled . categoryId ,
73+ accountId : scheduled . accountId ?? undefined ,
74+ toAccountId : scheduled . toAccountId ?? undefined ,
75+ description : scheduled . description ?? undefined ,
76+ date : baseDate . toISOString ( ) . split ( 'T' ) [ 0 ] ,
77+ } ) ;
7478
75- baseDate = calculateNextRunDate ( baseDate , scheduled . frequency ) ;
76- created ++ ;
77- }
79+ const nextDate = calculateNextRunDate ( baseDate , scheduled . frequency ) ;
80+
81+ // Persist progress after each occurrence — prevents re-duplicating
82+ // already-created transactions if a later occurrence throws
83+ await db
84+ . update ( scheduledTransactions )
85+ . set ( { nextRunDate : nextDate , updatedAt : new Date ( ) } )
86+ . where ( eq ( scheduledTransactions . id , scheduled . id ) ) ;
7887
79- // Advance nextRunDate to first future occurrence
80- await db
81- . update ( scheduledTransactions )
82- . set ( { nextRunDate : baseDate , updatedAt : new Date ( ) } )
83- . where ( eq ( scheduledTransactions . id , scheduled . id ) ) ;
88+ baseDate = nextDate ;
89+ created ++ ;
90+ succeeded ++ ;
91+ } catch ( err ) {
92+ failed ++ ;
93+ logger . error ( 'Scheduled runner: failed to process transaction' , {
94+ scheduledId : scheduled . id ,
95+ userId : scheduled . userId ,
96+ frequency : scheduled . frequency ,
97+ occurrenceDate : baseDate . toISOString ( ) ,
98+ error : err instanceof Error ? err . message : String ( err ) ,
99+ } ) ;
100+ break ; // stop catch-up for this item; nextRunDate already points to the failed date
101+ }
102+ }
84103
85- succeeded += created ;
86- } catch ( err ) {
87- failed ++ ;
88- logger . error ( 'Scheduled runner: failed to process transaction' , {
89- scheduledId : scheduled . id ,
90- userId : scheduled . userId ,
91- frequency : scheduled . frequency ,
92- nextRunDate : scheduled . nextRunDate ,
93- error : err instanceof Error ? err . message : String ( err ) ,
94- } ) ;
104+ if ( created === MAX_CATCH_UP && baseDate <= now ) {
105+ logger . warn ( 'Scheduled runner: hit catch-up cap, will resume next run' , {
106+ scheduledId : scheduled . id ,
107+ remaining : baseDate . toISOString ( ) ,
108+ } ) ;
109+ }
95110 }
96- }
97111
98- logger . info ( 'Scheduled runner: done' , { succeeded, failed } ) ;
112+ logger . info ( 'Scheduled runner: done' , { succeeded, failed } ) ;
113+ } finally {
114+ isRunning = false ;
115+ }
99116}
0 commit comments