Skip to content

Commit 6f0388b

Browse files
committed
fix(webapp): engine writes settle before any failure is reported
Both bounds' engine writes now settle before a sync step fails, so no write is still in flight when the compensating re-sync runs; a late sibling can never land after the compensation and leave one bound stale.
1 parent c3158e8 commit 6f0388b

1 file changed

Lines changed: 18 additions & 5 deletions

File tree

‎apps/webapp/app/v3/services/concurrencyLimitsSystem.server.ts‎

Lines changed: 18 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -300,7 +300,7 @@ function syncResetToEngine(
300300
? updateQueueTotalConcurrencyLimits(environment, row.name, totalTarget)
301301
: removeQueueTotalConcurrencyLimits(environment, row.name);
302302

303-
return fromPromise(Promise.all([perKeySync, totalSync]), (error) => ({
303+
return fromPromise(settleBothEngineWrites(perKeySync, totalSync), (error) => ({
304304
type: "sync_limit_to_engine_failed" as const,
305305
cause: error,
306306
})).map(() => row);
@@ -326,6 +326,19 @@ function resetLimitOverrides(db: PrismaClientOrTransaction, row: TaskQueue) {
326326
return guardedLimitUpdate(db, row, data);
327327
}
328328

329+
/**
330+
* Both engine writes settle before a failure is reported, so no write is still in
331+
* flight when a caller's compensation runs — a late sibling can never land after
332+
* the compensating re-sync and leave one bound stale.
333+
*/
334+
async function settleBothEngineWrites(a: Promise<unknown>, b: Promise<unknown>): Promise<void> {
335+
const results = await Promise.allSettled([a, b]);
336+
const failed = results.find((result) => result.status === "rejected");
337+
if (failed && failed.status === "rejected") {
338+
throw failed.reason;
339+
}
340+
}
341+
329342
/**
330343
* Optimistic update: the where clause carries the row's updatedAt plus both override
331344
* markers as read, so ANY concurrent write — another override or reset, or a deploy
@@ -379,7 +392,7 @@ function compensateEngineFromFreshRow(
379392
if (!fresh || fresh.updatedAt.getTime() === lastSyncedAt) {
380393
return;
381394
}
382-
await Promise.all([
395+
await settleBothEngineWrites(
383396
typeof fresh.concurrencyLimit === "number"
384397
? updateQueueConcurrencyLimits(environment, fresh.name, fresh.concurrencyLimit)
385398
: removeQueueConcurrencyLimits(environment, fresh.name),
@@ -389,8 +402,8 @@ function compensateEngineFromFreshRow(
389402
fresh.name,
390403
fresh.totalConcurrencyLimit
391404
)
392-
: removeQueueTotalConcurrencyLimits(environment, fresh.name),
393-
]);
405+
: removeQueueTotalConcurrencyLimits(environment, fresh.name)
406+
);
394407
lastSyncedAt = fresh.updatedAt.getTime();
395408
}
396409
})(),
@@ -414,7 +427,7 @@ function syncLimitToEngine(environment: AuthenticatedEnvironment, row: TaskQueue
414427
? updateQueueTotalConcurrencyLimits(environment, row.name, row.totalConcurrencyLimit)
415428
: removeQueueTotalConcurrencyLimits(environment, row.name);
416429

417-
return fromPromise(Promise.all([perKeySync, totalSync]), (error) => ({
430+
return fromPromise(settleBothEngineWrites(perKeySync, totalSync), (error) => ({
418431
type: "sync_limit_to_engine_failed" as const,
419432
cause: error,
420433
})).map(() => row);

0 commit comments

Comments
 (0)