Skip to content

Commit 8a09f00

Browse files
committed
fix(webapp,run-engine): recoverable resets, conflict-guarded limit mutations, activity-refreshed counter TTL
Limit resets sync the engine to the declared base before clearing the markers, so an engine failure leaves the override intact and a retry converges. Both limit mutations carry the read markers in their where clause and surface a conflict instead of clobbering a concurrent change. The gate queued counter's TTL now refreshes on every delta, so an active gate's count never resets while drift from delta-less paths still clears once the gate goes quiet; the flag comment documents that total bounds enforce solely through the total-concurrency flag.
1 parent 35a5d45 commit 8a09f00

2 files changed

Lines changed: 72 additions & 19 deletions

File tree

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

Lines changed: 62 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import type { TaskQueue, User } from "@trigger.dev/database";
22
import { errAsync, fromPromise, okAsync } from "neverthrow";
3-
import type { PrismaClientOrTransaction } from "~/db.server";
3+
import { Prisma, type PrismaClientOrTransaction } from "~/db.server";
44
import type { AuthenticatedEnvironment } from "~/services/apiAuth.server";
55
import {
66
removeQueueConcurrencyLimits,
@@ -128,8 +128,8 @@ export class ConcurrencyLimitsSystem {
128128
},
129129
reset: (environment: AuthenticatedEnvironment, name: string) => {
130130
return findLimitByName(this.db, environment, name)
131+
.andThen((row) => syncResetToEngine(environment, row))
131132
.andThen((row) => resetLimitOverrides(this.db, row))
132-
.andThen((row) => syncLimitToEngine(environment, row))
133133
.andThen((row) =>
134134
fromPromise(toLimitItems(environment, [row]), (error) => ({
135135
type: "other" as const,
@@ -246,17 +246,43 @@ function applyLimitOverride(
246246
data.totalConcurrencyLimitOverriddenBy = overriddenBy?.id ?? null;
247247
}
248248

249-
return fromPromise(db.taskQueue.update({ where: { id: row.id }, data }), (error) => ({
250-
type: "limit_update_failed" as const,
251-
cause: error,
252-
}));
249+
return guardedLimitUpdate(db, row, data);
253250
}
254251

255-
function resetLimitOverrides(db: PrismaClientOrTransaction, row: TaskQueue) {
252+
/**
253+
* Enforce first, then persist: the engine syncs to the declared base BEFORE the
254+
* override markers clear, so an engine failure leaves the markers set and a retry
255+
* converges instead of being rejected while the overridden limit stays enforced.
256+
*/
257+
function syncResetToEngine(environment: AuthenticatedEnvironment, row: TaskQueue) {
256258
if (row.concurrencyLimitOverriddenAt === null && row.totalConcurrencyLimitOverriddenAt === null) {
257259
return errAsync({ type: "limit_not_overridden" as const });
258260
}
259261

262+
const perKeyTarget = row.concurrencyLimitOverriddenAt
263+
? row.concurrencyLimitBase
264+
: row.concurrencyLimit;
265+
const totalTarget = row.totalConcurrencyLimitOverriddenAt
266+
? row.totalConcurrencyLimitBase
267+
: row.totalConcurrencyLimit;
268+
269+
const perKeySync =
270+
typeof perKeyTarget === "number"
271+
? updateQueueConcurrencyLimits(environment, row.name, perKeyTarget)
272+
: removeQueueConcurrencyLimits(environment, row.name);
273+
274+
const totalSync =
275+
typeof totalTarget === "number"
276+
? updateQueueTotalConcurrencyLimits(environment, row.name, totalTarget)
277+
: removeQueueTotalConcurrencyLimits(environment, row.name);
278+
279+
return fromPromise(Promise.all([perKeySync, totalSync]), (error) => ({
280+
type: "sync_limit_to_engine_failed" as const,
281+
cause: error,
282+
})).map(() => row);
283+
}
284+
285+
function resetLimitOverrides(db: PrismaClientOrTransaction, row: TaskQueue) {
260286
const data: Record<string, unknown> = {};
261287

262288
if (row.concurrencyLimitOverriddenAt !== null) {
@@ -273,10 +299,35 @@ function resetLimitOverrides(db: PrismaClientOrTransaction, row: TaskQueue) {
273299
data.totalConcurrencyLimitOverriddenBy = null;
274300
}
275301

276-
return fromPromise(db.taskQueue.update({ where: { id: row.id }, data }), (error) => ({
277-
type: "limit_update_failed" as const,
278-
cause: error,
279-
}));
302+
return guardedLimitUpdate(db, row, data);
303+
}
304+
305+
/**
306+
* Optimistic update: the where clause carries the override markers as read, so a
307+
* concurrent override or reset makes this update miss (P2025) and the caller gets
308+
* a conflict instead of silently clobbering the newer state.
309+
*/
310+
function guardedLimitUpdate(
311+
db: PrismaClientOrTransaction,
312+
row: TaskQueue,
313+
data: Record<string, unknown>
314+
) {
315+
return fromPromise(
316+
db.taskQueue.update({
317+
where: {
318+
id: row.id,
319+
concurrencyLimitOverriddenAt: row.concurrencyLimitOverriddenAt,
320+
totalConcurrencyLimitOverriddenAt: row.totalConcurrencyLimitOverriddenAt,
321+
},
322+
data,
323+
}),
324+
(error) => {
325+
if (error instanceof Prisma.PrismaClientKnownRequestError && error.code === "P2025") {
326+
return { type: "conflict" as const };
327+
}
328+
return { type: "limit_update_failed" as const, cause: error };
329+
}
330+
);
280331
}
281332

282333
/**

‎internal-packages/run-engine/src/run-queue/index.ts‎

Lines changed: 10 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -166,10 +166,11 @@ end
166166
-- Callers gate the delta on the actual queue-zset transition (ZADD added == 1 /
167167
-- ZREM removed == 1) so re-enqueues and already-removed members never double count;
168168
-- gates sharing a base (duplicate entries, key variants) count once per run. The
169-
-- 24h absolute TTL (set at creation, never extended) re-anchors drift from paths
170-
-- without the delta (rolling deploys, stale-entry cleanup): the counter resets,
171-
-- floored decrements absorb the pre-reset backlog as it drains, and counts converge
172-
-- to exact for every run enqueued after the reset. Payload-driven and
169+
-- 24h TTL refreshes on every delta, so an ACTIVE gate's count never resets while
170+
-- drift from delta-less paths (a mixed-version rollout, a stale-entry cleanup)
171+
-- clears once the gate has been quiet for a day. The residual gap is a gate idle
172+
-- for 24h with runs still queued (e.g. paused with no new enqueues): its counter
173+
-- expires and under-counts until the backlog fully drains. Payload-driven and
173174
-- flag-independent, like release, so counts stay exact across flag flips.
174175
local function __gateQueuedDelta(gatesKeyPrefix, msg, delta)
175176
if type(msg) ~= 'table' or not msg.gates then return end
@@ -180,12 +181,11 @@ local function __gateQueuedDelta(gatesKeyPrefix, msg, delta)
180181
seenBases[base] = true
181182
local counterKey = base .. ':gateQueuedCounter'
182183
if delta > 0 then
183-
if redis.call('EXISTS', counterKey) == 0 then
184-
redis.call('SET', counterKey, '0', 'EX', '86400')
185-
end
186184
redis.call('INCRBY', counterKey, delta)
185+
redis.call('EXPIRE', counterKey, '86400')
187186
elseif tonumber(redis.call('GET', counterKey) or '0') > 0 then
188187
redis.call('DECRBY', counterKey, -delta)
188+
redis.call('EXPIRE', counterKey, '86400')
189189
end
190190
end
191191
end
@@ -347,7 +347,9 @@ export type RunQueueOptions = {
347347
/**
348348
* When true, queues maintain a per-base-queue groupConcurrency SET (total in-flight
349349
* across all key variants AND keyless runs) and enforce the queue's total concurrency
350-
* limit at admit time. Default false: admit paths are byte-identical to before, and
350+
* limit at admit time. V2 concurrency semantics (a limit's `total` bound, including
351+
* total-only declarations) are enforced solely through this flag: with it off, a
352+
* total-only limit caps nothing. Default false: admit paths are byte-identical to before, and
351353
* only the release-side SREM mirror runs (a no-op on an absent set), so the flag can
352354
* be flipped on a fleet that has fully rolled onto this build without draining queues.
353355
*

0 commit comments

Comments
 (0)