@@ -84,6 +84,10 @@ export function createKyselyEventStoreConsumer({
8484 return ;
8585 }
8686
87+ // Track if all handlers succeeded in this batch
88+ let batchSucceeded = true ;
89+ let lastSuccessfulPosition = lastProcessedPosition ;
90+
8791 // Process each event
8892 for ( const row of events ) {
8993 const event : ReadEvent < Event , ReadEventMetadataWithGlobalPosition > = {
@@ -99,42 +103,87 @@ export function createKyselyEventStoreConsumer({
99103 } ,
100104 } ;
101105
106+ let eventSucceeded = true ;
107+
102108 // Call type-specific handlers
103109 const typeHandlers = eventHandlers . get ( row . message_type ) || [ ] ;
104110 for ( const handler of typeHandlers ) {
105111 try {
106112 await handler ( event ) ;
107113 } catch ( error ) {
108114 logger . error (
109- { error, event } ,
110- `Error processing event ${ row . message_type } ` ,
115+ {
116+ error,
117+ event,
118+ globalPosition : row . global_position ,
119+ consumerName,
120+ } ,
121+ `Error processing event ${ row . message_type } at position ${ row . global_position } . Consumer will retry this event on next poll.` ,
111122 ) ;
123+ eventSucceeded = false ;
124+ batchSucceeded = false ;
125+ break ; // Stop processing this event's handlers
112126 }
113127 }
114128
115- // Call all-event handlers
116- for ( const handler of allEventHandlers ) {
117- try {
118- await handler ( event ) ;
119- } catch ( error ) {
120- logger . error (
121- { error, event } ,
122- "Error processing event in all-event handler" ,
123- ) ;
129+ // Only call all-event handlers if type-specific handlers succeeded
130+ if ( eventSucceeded ) {
131+ for ( const handler of allEventHandlers ) {
132+ try {
133+ await handler ( event ) ;
134+ } catch ( error ) {
135+ logger . error (
136+ {
137+ error,
138+ event,
139+ globalPosition : row . global_position ,
140+ consumerName,
141+ } ,
142+ `Error processing event in all-event handler at position ${ row . global_position } . Consumer will retry this event on next poll.` ,
143+ ) ;
144+ eventSucceeded = false ;
145+ batchSucceeded = false ;
146+ break ; // Stop processing this event's handlers
147+ }
124148 }
125149 }
126150
127- // Update last processed position
151+ // If this event failed, stop processing the batch
152+ if ( ! eventSucceeded ) {
153+ ( logger . warn ?? logger . error ) (
154+ {
155+ failedPosition : row . global_position ,
156+ lastSuccessfulPosition,
157+ consumerName,
158+ } ,
159+ `Stopping batch processing due to handler failure. Will retry from position ${ lastSuccessfulPosition + 1n } ` ,
160+ ) ;
161+ break ;
162+ }
163+
164+ // Update last successful position only if all handlers succeeded for this event
128165 const globalPos = row . global_position ;
129166 if ( globalPos !== null ) {
130- lastProcessedPosition = BigInt ( String ( globalPos ) ) ;
167+ lastSuccessfulPosition = BigInt ( String ( globalPos ) ) ;
131168 }
132169 }
133170
134- // Update subscription tracking
135- await updateSubscriptionPosition ( ) ;
171+ // Only update the position if entire batch succeeded
172+ if ( batchSucceeded ) {
173+ lastProcessedPosition = lastSuccessfulPosition ;
174+ await updateSubscriptionPosition ( ) ;
175+ } else {
176+ // Keep the old position - will retry failed event on next poll
177+ logger . info (
178+ {
179+ currentPosition : lastProcessedPosition ,
180+ consumerName,
181+ } ,
182+ `Batch processing incomplete. Position unchanged at ${ lastProcessedPosition } . Will retry on next poll.` ,
183+ ) ;
184+ }
136185 } catch ( error ) {
137- logger . error ( { error } , "Error processing events" ) ;
186+ logger . error ( { error, consumerName } , "Error processing events" ) ;
138187 }
139188 } ;
140189
0 commit comments