11import type { TaskQueue , User } from "@trigger.dev/database" ;
2- import { errAsync , fromPromise , okAsync } from "neverthrow" ;
2+ import { errAsync , fromPromise , okAsync , type ResultAsync } from "neverthrow" ;
33import { Prisma , type PrismaClientOrTransaction } from "~/db.server" ;
44import type { AuthenticatedEnvironment } from "~/services/apiAuth.server" ;
55import {
@@ -134,12 +134,18 @@ export class ConcurrencyLimitsSystem {
134134 } ,
135135 reset : ( environment : AuthenticatedEnvironment , name : string ) => {
136136 return findLimitByName ( this . db , environment , name )
137- . andThen ( ( row ) => syncResetToEngine ( environment , row ) )
137+ . andThen ( ( row ) =>
138+ syncResetToEngine ( environment , row ) . orElse ( ( error ) =>
139+ compensateEngineFromFreshRow ( this . db , environment , row . id )
140+ . orElse ( ( ) => okAsync ( undefined ) )
141+ . andThen ( ( ) => errAsync ( error ) )
142+ )
143+ )
138144 . andThen ( ( row ) =>
139145 resetLimitOverrides ( this . db , row ) . orElse ( ( error ) =>
140- compensateEngineFromFreshRow ( this . db , environment , row . id ) . andThen ( ( ) =>
141- errAsync ( error )
142- )
146+ compensateEngineFromFreshRow ( this . db , environment , row . id )
147+ . orElse ( ( ) => okAsync ( undefined ) )
148+ . andThen ( ( ) => errAsync ( error ) )
143149 )
144150 )
145151 . andThen ( ( row ) =>
@@ -266,7 +272,13 @@ function applyLimitOverride(
266272 * override markers clear, so an engine failure leaves the markers set and a retry
267273 * converges instead of being rejected while the overridden limit stays enforced.
268274 */
269- function syncResetToEngine ( environment : AuthenticatedEnvironment , row : TaskQueue ) {
275+ function syncResetToEngine (
276+ environment : AuthenticatedEnvironment ,
277+ row : TaskQueue
278+ ) : ResultAsync <
279+ TaskQueue ,
280+ { type : "limit_not_overridden" } | { type : "sync_limit_to_engine_failed" ; cause : unknown }
281+ > {
270282 if ( row . concurrencyLimitOverriddenAt === null && row . totalConcurrencyLimitOverriddenAt === null ) {
271283 return errAsync ( { type : "limit_not_overridden" as const } ) ;
272284 }
@@ -315,10 +327,12 @@ function resetLimitOverrides(db: PrismaClientOrTransaction, row: TaskQueue) {
315327}
316328
317329/**
318- * Optimistic update: the where clause carries the row's updatedAt as read, so ANY
319- * concurrent write — another override or reset, or a deploy refreshing the declared
320- * values — makes this update miss (P2025) and the caller gets a conflict instead of
321- * persisting values computed from a stale row.
330+ * Optimistic update: the where clause carries the row's updatedAt plus both override
331+ * markers as read, so ANY concurrent write — another override or reset, or a deploy
332+ * refreshing the declared values — makes this update miss (P2025) and the caller
333+ * gets a conflict instead of persisting values computed from a stale row. The
334+ * markers narrow the same-millisecond updatedAt window to writes that also leave
335+ * both markers untouched.
322336 */
323337function guardedLimitUpdate (
324338 db : PrismaClientOrTransaction ,
@@ -330,6 +344,8 @@ function guardedLimitUpdate(
330344 where : {
331345 id : row . id ,
332346 updatedAt : row . updatedAt ,
347+ concurrencyLimitOverriddenAt : row . concurrencyLimitOverriddenAt ,
348+ totalConcurrencyLimitOverriddenAt : row . totalConcurrencyLimitOverriddenAt ,
333349 } ,
334350 data,
335351 } ) ,
0 commit comments