@@ -818,30 +818,43 @@ async function createWorkerQueue(
818818 /**
819819 * The optimistic markers only guard the Postgres write; an override or reset can
820820 * still land between that write and the engine sync above, which would leave the
821- * engine holding this deploy's stale values. Re-read the markers and re-sync once
822- * from the fresh row when they moved: every actor writes Postgres before its own
823- * engine sync, so whoever syncs last is syncing the freshest row.
821+ * engine holding this deploy's stale values. Re-read the markers and re-sync from
822+ * the fresh row until they stop moving (bounded): every actor writes Postgres
823+ * before its own engine sync, so re-syncing whatever is freshest converges. A
824+ * marker moving after the final read is healed by that actor's own engine sync
825+ * or the next deploy.
824826 */
825- const freshQueue = await prisma . taskQueue . findFirst ( {
826- where : { id : taskQueue . id } ,
827- select : {
828- name : true ,
829- paused : true ,
830- concurrencyLimit : true ,
831- totalConcurrencyLimit : true ,
832- concurrencyLimitOverriddenAt : true ,
833- totalConcurrencyLimitOverriddenAt : true ,
834- } ,
835- } ) ;
827+ let syncedMarkers = {
828+ concurrency : taskQueue . concurrencyLimitOverriddenAt ?. getTime ( ) ,
829+ total : taskQueue . totalConcurrencyLimitOverriddenAt ?. getTime ( ) ,
830+ } ;
831+
832+ for ( let i = 0 ; i < 3 ; i ++ ) {
833+ const freshQueue = await prisma . taskQueue . findFirst ( {
834+ where : { id : taskQueue . id } ,
835+ select : {
836+ name : true ,
837+ paused : true ,
838+ concurrencyLimit : true ,
839+ totalConcurrencyLimit : true ,
840+ concurrencyLimitOverriddenAt : true ,
841+ totalConcurrencyLimitOverriddenAt : true ,
842+ } ,
843+ } ) ;
844+
845+ if (
846+ ! freshQueue ||
847+ ( freshQueue . concurrencyLimitOverriddenAt ?. getTime ( ) === syncedMarkers . concurrency &&
848+ freshQueue . totalConcurrencyLimitOverriddenAt ?. getTime ( ) === syncedMarkers . total )
849+ ) {
850+ break ;
851+ }
836852
837- if (
838- freshQueue &&
839- ( freshQueue . concurrencyLimitOverriddenAt ?. getTime ( ) !==
840- taskQueue . concurrencyLimitOverriddenAt ?. getTime ( ) ||
841- freshQueue . totalConcurrencyLimitOverriddenAt ?. getTime ( ) !==
842- taskQueue . totalConcurrencyLimitOverriddenAt ?. getTime ( ) )
843- ) {
844853 await syncQueueLimitsToEngine ( freshQueue ) ;
854+ syncedMarkers = {
855+ concurrency : freshQueue . concurrencyLimitOverriddenAt ?. getTime ( ) ,
856+ total : freshQueue . totalConcurrencyLimitOverriddenAt ?. getTime ( ) ,
857+ } ;
845858 }
846859
847860 return taskQueue ;
0 commit comments