Skip to content

fix(redis-worker): stop fair queue leaking concurrency slots - #4540

Open
matt-aitken wants to merge 10 commits into
mainfrom
fix/fair-queue-concurrency-slot-leak
Open

fix(redis-worker): stop fair queue leaking concurrency slots#4540
matt-aitken wants to merge 10 commits into
mainfrom
fix/fair-queue-concurrency-slot-leak

Conversation

@matt-aitken

Copy link
Copy Markdown
Member

Summary

A tenant could reach a state where none of its queued work ever ran again. The fair queue's per-tenant concurrency limiter was holding slots for messages that had already finished, and nothing reclaimed them. Once the leaked slots reached the tenant's limit, every queue belonging to that tenant stalled permanently, with items sitting in Redis untouched.

Root cause

completeMessage and releaseMessage both read the message's in-flight record to build a queue descriptor, then skipped the concurrency release entirely when that read came back empty:

if (this.concurrencyManager && storedMessage) {
  await this.concurrencyManager.release(descriptor, messageId);
}

The release does not need that record. The descriptor built on the line above already falls back to extractTenantId(queueId), and tenantId is the only field release() reads.

Ordering made it worse: the release ran after visibilityManager.complete() had deleted the in-flight record, so the reclaim loop, the one thing that can recover a slot, had already lost its handle on the message. Any failure between those two awaits stranded the slot with nothing able to free it.

The fix is to always release the slot, and to release it before the in-flight record is removed, so an interruption leaves a reclaimable message rather than an orphaned slot.

The regression test deletes the in-flight record before completing a message on a tenant with a limit of 1. Without the fix the next message never gets a slot and the test times out, which is exactly the observed failure.

A slot was only released when the message's in-flight record could still be
read, and the release ran after that record had already been deleted. Any
failure in between left the slot held with nothing able to reclaim it. Once a
tenant leaked its whole limit, every queue it owned stalled for good.
@changeset-bot

changeset-bot Bot commented Aug 8, 2026

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 72947fe

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 27 packages
Name Type
@trigger.dev/redis-worker Patch
@internal/run-engine Patch
@internal/schedule-engine Patch
@trigger.dev/build Patch
@trigger.dev/core Patch
@trigger.dev/python Patch
@trigger.dev/react-hooks Patch
@trigger.dev/rsc Patch
@trigger.dev/schema-to-json Patch
@trigger.dev/sdk Patch
@trigger.dev/database Patch
@trigger.dev/otlp-importer Patch
@trigger.dev/rbac Patch
@trigger.dev/sso Patch
trigger.dev Patch
@internal/dashboard-agent Patch
@internal/cache Patch
@internal/clickhouse Patch
@internal/llm-model-catalog Patch
@internal/metrics-pipeline Patch
@internal/redis Patch
@internal/replication Patch
@internal/run-store Patch
@internal/testcontainers Patch
@internal/tracing Patch
@internal/tsql Patch
@internal/sdk-compat-tests Patch

Not sure what this means? Click here to learn what changesets are.

Click here if you're a maintainer who wants to add another changeset to this PR

devin-ai-integration[bot]

This comment was marked as resolved.

@coderabbitai

coderabbitai Bot commented Aug 8, 2026

Copy link
Copy Markdown
Contributor

Review Change Stack

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review

Walkthrough

Fair queue consumers now release concurrency slots before completing, requeueing, retrying, or moving messages to the DLQ. Fallback queue metadata supports cleanup when in-flight records are missing or invalid. Timed-out reclamation releases slots per message and removes dangling in-flight entries. Regression tests cover completion, failure, reclaim, and timeout cleanup paths. Batch queues now use configurable visibility timeouts and heartbeat long-running callbacks. A patch changeset documents the fair queue fix.

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Description check ⚠️ Warning The description explains the root cause and fix, but it omits the required checklist, testing, changelog, screenshots, and issue reference sections. Add the template sections, complete the checklist, document testing steps, provide a changelog entry, and include the issue reference and screenshots section.
✅ Passed checks (4 passed)
Check name Status Explanation
Title check ✅ Passed The title clearly and concisely describes the main fix: preventing fair queue concurrency slot leaks in redis-worker.
Docstring Coverage ✅ Passed No functions found in the changed files to evaluate docstring coverage. Skipping docstring coverage check.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch fix/fair-queue-concurrency-slot-leak

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

coderabbitai[bot]

This comment was marked as resolved.

…lease path

Concurrency groups can key on queue metadata, not just the tenant. The fallback
descriptor dropped metadata and skipped the descriptor cache, so those groups
released against the wrong Redis set and leaked exactly as before. Prefer the
cached descriptor in both branches, and move the retry path's release ahead of
the re-queue so a redelivery cannot have its fresh reservation deleted by the
previous holder.
devin-ai-integration[bot]

This comment was marked as resolved.

@pkg-pr-new

pkg-pr-new Bot commented Aug 8, 2026

Copy link
Copy Markdown

Open in StackBlitz

@trigger.dev/build

npm i https://pkg.pr.new/@trigger.dev/build@72947fe

trigger.dev

npm i https://pkg.pr.new/trigger.dev@72947fe

@trigger.dev/core

npm i https://pkg.pr.new/@trigger.dev/core@72947fe

@trigger.dev/python

npm i https://pkg.pr.new/@trigger.dev/python@72947fe

@trigger.dev/react-hooks

npm i https://pkg.pr.new/@trigger.dev/react-hooks@72947fe

@trigger.dev/redis-worker

npm i https://pkg.pr.new/@trigger.dev/redis-worker@72947fe

@trigger.dev/rsc

npm i https://pkg.pr.new/@trigger.dev/rsc@72947fe

@trigger.dev/schema-to-json

npm i https://pkg.pr.new/@trigger.dev/schema-to-json@72947fe

@trigger.dev/sdk

npm i https://pkg.pr.new/@trigger.dev/sdk@72947fe

commit: 72947fe

coderabbitai[bot]

This comment was marked as resolved.

… queue

Four paths could strand a slot with nothing left to reclaim it.

Reclaim released the slot only after requeuing, and only for messages that
survived the requeue, so a message whose requeue threw kept its slot forever.
It now releases before the message becomes claimable, which also stops a
redelivery having its fresh reservation deleted by the previous holder.

The dead-letter branch completed the message before releasing, so anything
throwing in between stranded the slot. Release now happens first.

failMessage returned early when the stored message was missing or unparseable
without releasing, completing or requeuing, leaving the slot held.

The requeue script returned without removing the in-flight entry when the
payload was gone, so the member was rescanned on every reclaim tick forever.
Enough of them fill the scan window and starve a whole shard's reclaim.

reclaimTimedOut takes an optional pre-requeue hook; no existing caller changes.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick comments (2)
packages/redis-worker/src/fair-queue/index.ts (2)

1302-1311: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Consider extracting the repeated descriptor construction.

Four sites now build the same fallback QueueDescriptor: completeMessage (Lines 1249-1253), releaseMessage (Lines 1302-1306), #releaseOrphanedConcurrency (Lines 1346-1350), and failMessage (Lines 1389-1393). The precedence rule is identical: cache first, then stored-message tenant and metadata, then extractTenantId(queueId). A single private helper keeps that rule in one place and prevents the sites from drifting apart.

The release-before-requeue ordering itself is correct.

♻️ Suggested helper
`#resolveDescriptor`(
  queueId: string,
  storedMessage?: { tenantId: string; metadata?: Record<string, unknown> } | null
): QueueDescriptor {
  return (
    this.queueDescriptorCache.get(queueId) ?? {
      id: queueId,
      tenantId: storedMessage?.tenantId ?? this.keys.extractTenantId(queueId),
      metadata: storedMessage?.metadata ?? {},
    }
  );
}

1547-1574: 🧹 Nitpick | 🔵 Trivial

Consider a metric for the swallowed release failure.

The catch block logs and continues, which is the right choice here. One failure must not stop the shard scan. However, when the release fails, the message is still requeued and the concurrency slot stays held. That is a silent slot leak that only appears in logs.

A counter incremented in the catch block would make the condition alertable. Keep the attributes low cardinality, for example a bounded outcome label only, and do not attach messageId, queueId, or tenantId.

As per coding guidelines, OTEL metric attributes must use only enums, booleans, bounded error codes, or bounded shard IDs, and must not use IDs such as runId or projectId.

Source: Coding guidelines


ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository UI

Review profile: CHILL

Plan: Pro Plus

Run ID: b6d6bee1-d525-4333-866f-fadfa29629b4

📥 Commits

Reviewing files that changed from the base of the PR and between 35568a3 and b2ebf18.

📒 Files selected for processing (4)
  • packages/redis-worker/src/fair-queue/index.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
📜 Review details
⏰ Context from checks skipped due to timeout. (21)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (10, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (8, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (12, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (6, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (5, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (9, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (1, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (11, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (4, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (3, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (7, 12)
  • GitHub Check: webapp / 🧪 Unit Tests: Webapp (2, 12)
  • GitHub Check: packages / 🧪 Unit Tests: Packages (2, 3)
  • GitHub Check: packages / 🧪 Unit Tests: Packages (3, 3)
  • GitHub Check: packages / 🧪 Unit Tests: Packages (1, 3)
  • GitHub Check: typecheck / typecheck
  • GitHub Check: internal / 🧪 Unit Tests: Internal
  • GitHub Check: runops-guard / runops-guard
  • GitHub Check: e2e-webapp / 🧪 E2E Tests: Webapp
  • GitHub Check: code-quality / code-quality
  • GitHub Check: Analyze (javascript-typescript)
🧰 Additional context used
📓 Path-based instructions (8)
**/*.{ts,tsx}

📄 CodeRabbit inference engine (.github/copilot-instructions.md)

**/*.{ts,tsx}: Use types over interfaces for TypeScript
Avoid using enums; prefer string unions or const objects instead

**/*.{ts,tsx}: Prefer static imports over dynamic import(); use dynamic imports only for unresolvable circular dependencies, genuine performance code splitting, or conditional runtime loading.
Import Trigger.dev tasks from @trigger.dev/sdk; never use @trigger.dev/sdk/v3 or deprecated client.defineJob.
Add agentcrumbs while writing code using approved namespaces; mark lines with // @Crumbs or blocks with `// `#region` `@crumbs, and strip them before merging.

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
**/*.{ts,tsx,js,jsx}

📄 CodeRabbit inference engine (.github/copilot-instructions.md)

Use function declarations instead of default exports

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
**/*.{test,spec}.{ts,tsx}

📄 CodeRabbit inference engine (.github/copilot-instructions.md)

Use vitest for all tests in the Trigger.dev repository

**/*.{test,spec}.{ts,tsx}: Use Vitest exclusively and never mock dependencies; use Testcontainers for integration dependencies.
Place test files next to the source files they test.

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
**/*.ts

📄 CodeRabbit inference engine (.cursor/rules/otel-metrics.mdc)

**/*.ts: When creating or editing OTEL metrics (counters, histograms, gauges), ensure metric attributes have low cardinality by using only enums, booleans, bounded error codes, or bounded shard IDs
Do not use high-cardinality attributes in OTEL metrics such as UUIDs/IDs (envId, userId, runId, projectId, organizationId), unbounded integers (itemCount, batchSize, retryCount), timestamps (createdAt, startTime), or free-form strings (errorMessage, taskName, queueName)
When exporting OTEL metrics via OTLP to Prometheus, be aware that the exporter automatically adds unit suffixes to metric names (e.g., 'my_duration_ms' becomes 'my_duration_ms_milliseconds', 'my_counter' becomes 'my_counter_total'). Account for these transformations when writing Grafana dashboards or Prometheus queries

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
packages/redis-worker/src/fair-queue/**/*.ts

📄 CodeRabbit inference engine (packages/redis-worker/CLAUDE.md)

Keep the fair dequeueing algorithm within src/fair-queue/.

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
packages/redis-worker/**/*.{ts,tsx}

📄 CodeRabbit inference engine (packages/redis-worker/CLAUDE.md)

Use @trigger.dev/redis-worker for all background jobs in the webapp.

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
packages/redis-worker/**/*.{test,spec}.{ts,tsx}

📄 CodeRabbit inference engine (packages/redis-worker/CLAUDE.md)

Test Redis worker behavior with ioredis and testcontainers for Redis.

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
packages/**/*.{ts,tsx}

📄 CodeRabbit inference engine (AGENTS.md)

For public packages, use build for verification.

Files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
🧠 Learnings (12)
📚 Learning: 2026-03-22T13:26:12.060Z
Learnt from: ericallam
Repo: triggerdotdev/trigger.dev PR: 3244
File: apps/webapp/app/components/code/TextEditor.tsx:81-86
Timestamp: 2026-03-22T13:26:12.060Z
Learning: In the triggerdotdev/trigger.dev codebase, do not flag `navigator.clipboard.writeText(...)` calls for `missing-await`/`unhandled-promise` issues. These clipboard writes are intentionally invoked without `await` and without `catch` handlers across the project; keep that behavior consistent when reviewing TypeScript/TSX files (e.g., usages like in `apps/webapp/app/components/code/TextEditor.tsx`).

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-03-22T19:24:14.403Z
Learnt from: matt-aitken
Repo: triggerdotdev/trigger.dev PR: 3187
File: apps/webapp/app/v3/services/alerts/deliverErrorGroupAlert.server.ts:200-204
Timestamp: 2026-03-22T19:24:14.403Z
Learning: In the triggerdotdev/trigger.dev codebase, webhook URLs are not expected to contain embedded credentials/secrets (e.g., fields like `ProjectAlertWebhookProperties` should only hold credential-free webhook endpoints). During code review, if you see logging or inclusion of raw webhook URLs in error messages, do not automatically treat it as a credential-leak/secrets-in-logs issue by default—first verify the URL does not contain embedded credentials (for example, no username/password in the URL, no obvious secret/token query params or fragments). If the URL is credential-free per this project’s conventions, allow the logging.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-05-18T08:21:27.694Z
Learnt from: d-cs
Repo: triggerdotdev/trigger.dev PR: 3632
File: apps/webapp/sentry.server.ts:4-21
Timestamp: 2026-05-18T08:21:27.694Z
Learning: When handling Prisma error P1001 ("Can't reach database server") in TypeScript, don’t assume a single error shape. Prisma can surface P1001 via two different error classes/fields: `PrismaClientKnownRequestError` exposes it as `err.code === "P1001"` (common during mid-query connection drops), while `PrismaClientInitializationError` exposes it as `err.errorCode === "P1001"` (common on client startup failure). Therefore, predicates should use `err.code === "P1001" || err.errorCode === "P1001"`. Do not flag `err.code === "P1001"` as “unreachable/never matches,” as it is expected in production.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-05-18T08:21:27.694Z
Learnt from: d-cs
Repo: triggerdotdev/trigger.dev PR: 3632
File: apps/webapp/sentry.server.ts:4-21
Timestamp: 2026-05-18T08:21:27.694Z
Learning: When handling Prisma errors for P1001 ("Can't reach database server"), do not assume it only appears under a single property name. Prisma may surface P1001 via either `PrismaClientKnownRequestError` (`err.code === "P1001"`, e.g., mid-query connection drops) or `PrismaClientInitializationError` (`err.errorCode === "P1001"`, e.g., client startup connection failure). To reliably detect the condition, check `err.code === "P1001" || err.errorCode === "P1001"`, and avoid review rules that would incorrectly flag `err.code === "P1001"` as unreachable/never-matching.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-06-13T19:53:13.759Z
Learnt from: ericallam
Repo: triggerdotdev/trigger.dev PR: 3937
File: packages/trigger-sdk/skills/realtime-and-frontend/SKILL.md:258-260
Timestamp: 2026-06-13T19:53:13.759Z
Learning: When reviewing code that uses `trigger.dev/react-hooks`’s `useRealtimeRun`, preserve the call signature where the first argument is the full realtime handle object (not `handle.id`). This is intentional to maintain type-safety and is consistent with the official docs; do not suggest changing the first argument from the handle object to `handle.id`.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-06-17T17:13:49.929Z
Learnt from: matt-aitken
Repo: triggerdotdev/trigger.dev PR: 3948
File: apps/webapp/app/routes/_app.orgs.$organizationSlug.projects.$projectParam.env.$envParam.bulk-actions.$bulkActionParam/route.tsx:48-62
Timestamp: 2026-06-17T17:13:49.929Z
Learning: In triggerdotdev/trigger.dev, within `dashboardLoader`/`dashboardAction` (or similar context resolver code) whenever you resolve an organization ID from an organization slug for RBAC/enterprise authorization scope, always read from the primary Prisma client (`prisma`), not `$replica`. Using `$replica` can hit replica-lag and cause the RBAC lookup/authorization to run without the correct org scope (bypassing intended role enforcement). Implement the slug→org lookup with `prisma.organization.findFirst(...)` (or equivalent primary-client query) and add an inline comment documenting why the primary client is required (replica lag could lead to unscoped RBAC checks).

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-06-23T13:04:21.413Z
Learnt from: carderne
Repo: triggerdotdev/trigger.dev PR: 4023
File: apps/webapp/app/services/upsertBranch.server.ts:14-18
Timestamp: 2026-06-23T13:04:21.413Z
Learning: In TypeScript, it’s valid to `import { type X }` and then use `typeof X` in a type-only position, e.g. `type Alias = z.infer<typeof X>`. The `type` modifier suppresses the runtime import, but the type checker still has the full exported type so `z.infer<typeof X>` can resolve correctly. In code reviews, don’t flag this as a TypeScript compile error as long as `typeof X` is used in a type context (e.g., with `z.infer`, `type` aliases, generics), not as a runtime value.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-05-18T14:40:02.173Z
Learnt from: ericallam
Repo: triggerdotdev/trigger.dev PR: 3658
File: packages/core/src/v3/realtimeStreams/manager.test.ts:1-147
Timestamp: 2026-05-18T14:40:02.173Z
Learning: In this repo’s trigger.dev codebase, the “never mock — use testcontainers” guideline should only be applied to integration tests that talk to real external services (e.g., Redis, Postgres, S2). For unit tests that validate in-memory logic (e.g., deduplication/cache behavior in StandardRealtimeStreamsManager and similar module-boundary call counting), it is allowed to use Vitest mocks like `vi.fn()` and to stub/mock `ApiClient` objects to count calls or simulate in-process collaborators. Do not flag `vi.fn()`-based mocks as policy violations in these unit-test scenarios; reserve the rule for true external-service integration tests.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
📚 Learning: 2026-05-18T14:40:02.173Z
Learnt from: ericallam
Repo: triggerdotdev/trigger.dev PR: 3658
File: packages/core/src/v3/realtimeStreams/manager.test.ts:1-147
Timestamp: 2026-05-18T14:40:02.173Z
Learning: In the triggerdotdev/trigger.dev repo, the policy “Never mock anything — use testcontainers instead” should only be enforced for integration tests that interact with real external services (e.g., Redis, Postgres) via actual infrastructure. For unit tests that exercise pure in-memory logic (e.g., cache semantics) it is OK to stub collaborators such as `ApiClient` using Vitest (`vi.fn()`) to assert call counts or control behavior. Do not flag `vi.fn()`-based `ApiClient` stubs in unit tests as violations of the testcontainers policy.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
📚 Learning: 2026-06-04T18:16:35.386Z
Learnt from: nicktrn
Repo: triggerdotdev/trigger.dev PR: 3836
File: apps/supervisor/src/backpressure/backpressureMonitor.ts:3-5
Timestamp: 2026-06-04T18:16:35.386Z
Learning: When reviewing TypeScript in this repo, apply the rule “prefer type aliases over interfaces” only to data/object shapes and union/intersection type modeling. If an interface is being used as a behavioral contract for collaborators to implement (e.g., method-shape interfaces that define required behavior, such as `BackpressureLogger` / `BackpressureSignalSource` in `apps/supervisor/src/backpressure/backpressureMonitor.ts`), keep it as an `interface` and do not flag it as a type-alias-vs-interface violation.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-06-09T17:58:04.699Z
Learnt from: 0ski
Repo: triggerdotdev/trigger.dev PR: 3879
File: apps/webapp/app/models/vercelIntegration.server.ts:619-630
Timestamp: 2026-06-09T17:58:04.699Z
Learning: In this codebase, outbound raw `fetch` calls should typically rely on Node/undici’s default request timeout (about ~300s) rather than adding a per-call `AbortController` + `setTimeout` wrapper inside individual functions (e.g. in files like `apps/webapp/app/models/vercelIntegration.server.ts`). During code review, do not flag the absence of a per-call timeout on a single `fetch` as an issue; if per-call timeouts are needed, they should be implemented via a codebase-wide convention (e.g., a shared fetch wrapper or documented pattern) rather than ad-hoc per-function changes.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
  • packages/redis-worker/src/fair-queue/visibility.ts
  • packages/redis-worker/src/fair-queue/index.ts
📚 Learning: 2026-06-16T09:19:47.637Z
Learnt from: d-cs
Repo: triggerdotdev/trigger.dev PR: 3960
File: apps/webapp/test/prismaInfrastructureErrorCapture.test.ts:0-0
Timestamp: 2026-06-16T09:19:47.637Z
Learning: In this repo’s Vitest setup, `vitest.config.ts` uses `globals: true`, so identifiers like `vi`, `describe`, `it`, and `expect` are available as globals in Vitest test files. During code review, do not flag missing `vi`/`describe`/`it`/`expect` imports as a runtime error or correctness issue when they’re used in `*.test.ts/tsx` or `*.spec.ts/tsx` files. Explicit imports are still preferred for consistency, but they’re not required for runtime behavior.

Applied to files:

  • packages/redis-worker/src/fair-queue/tests/visibility.test.ts
  • packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts
🔇 Additional comments (13)
packages/redis-worker/src/fair-queue/visibility.ts (4)

401-402: LGTM!


443-451: LGTM!


719-724: LGTM!


795-799: LGTM!

packages/redis-worker/src/fair-queue/index.ts (6)

28-28: LGTM!


1249-1262: LGTM!


1335-1354: LGTM!


1371-1387: LGTM!


1433-1473: LGTM!


1576-1594: LGTM!

packages/redis-worker/src/fair-queue/tests/fairQueue.test.ts (2)

1520-1590: LGTM!


1592-1659: LGTM!

packages/redis-worker/src/fair-queue/tests/visibility.test.ts (1)

916-979: LGTM!

…ite fails

The message was removed from in-flight before the dead-letter entry was written,
and the write went through a pipeline whose per-command errors were never
inspected. A failed write therefore lost the message silently: gone from
in-flight, absent from the dead-letter queue, invisible to the reclaim loop.

Write the entry first and only complete the message once it lands. On failure
the message stays in-flight, so the reclaim loop picks it up and it is retried
rather than dropped. A persistently failing write now loops visibly instead of
discarding work.
devin-ai-integration[bot]

This comment was marked as resolved.

coderabbitai[bot]

This comment was marked as resolved.

…s stranded or run twice

Releasing a slot could silently do nothing. The release pipeline never inspected
its per-command errors, and ioredis resolves a pipeline whose commands failed, so
a failed SREM reported success and the caller went on to destroy the in-flight
record. Release now throws on a failed command.

Reclaim then acted on that false success: it released per message, swallowed any
error, and requeued regardless, which is exactly how a slot ends up held with the
message gone. It now frees the whole timed-out batch in one pipeline before any of
them is requeued, and a failure aborts the requeue so the messages stay in-flight
for the next tick. That also restores the single round trip per shard.

Requeuing removed the message from in-flight before writing it back to the queue.
Redis does not roll back a script that fails partway, so a failed queue write lost
the message outright. The order is now write-then-remove.

Batch items were never heartbeated, so any item slower than the visibility timeout
was reclaimed and handed to a second consumer while the first was still running it,
executing the item twice. Items are now heartbeated for as long as their callback
runs, and the visibility timeout is configurable so this is testable.
devin-ai-integration[bot]

This comment was marked as resolved.

coderabbitai[bot]

This comment was marked as resolved.

… stop one bad key stalling a shard

Each beat extended the deadline by the same amount as the gap between beats, so
timer drift and round-trip latency left the message briefly past its deadline on
every cycle, and a reclaim scan landing in that window still handed a running item
to a second consumer. Beats now extend by a full visibility timeout while the tick
stays at a third of it.

Reclaim released the whole timed-out batch in one pipeline and aborted every
requeue if any command failed, so a single unusable concurrency key stalled reclaim
for every other tenant in that shard. A failed batch now falls back to releasing
message by message and only the messages whose slot could not be freed are held
back.
devin-ai-integration[bot]

This comment was marked as resolved.

… is lost

The heartbeat already learns when another consumer has taken the item over: the
in-flight member is gone, so the extend returns false. That signal was discarded,
leaving the original consumer free to finish and complete over the new owner's
in-flight record. It now stops and discards its result instead, so the consumer
that actually owns the item is the one that finishes it.

Also treat a discarded release pipeline as a failure rather than a success, since
a null result means none of the commands ran.
devin-ai-integration[bot]

This comment was marked as resolved.

The in-flight member is keyed only by message and queue id, so once another
consumer re-claims a reclaimed item the member exists again and an extend from
the previous consumer succeeds. The signal therefore only catches the window
where the item is back on the queue and unclaimed, which is narrower than the
comment and the changeset claimed.
The heartbeat fixes a separate pre-existing bug (items slower than the visibility
timeout are redelivered and executed twice) and shares no files with the
concurrency slot fixes, so it ships on its own.
devin-ai-integration[bot]

This comment was marked as resolved.

The test waited for the message to be queued and absent from in-flight at the
same instant, a state that only lasts about one dispatch interval before the
message is re-claimed. It now waits for the stuck in-flight entry's deadline to
move forward instead, which is monotonic once reclaim has run.

@devin-ai-integration devin-ai-integration Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Devin Review found 1 new potential issue.

Open in Devin Review

Comment on lines +5 to +7
Fair queue consumers no longer leak the concurrency slots that gate a tenant's throughput. Slots were held by messages that had already finished, were never reclaimed, and once enough of them accumulated every queue belonging to that tenant stopped being served. Slots are now freed on the paths that previously skipped them, freed before the record needed to recover them is discarded, and released before a reclaimed message goes back on the queue. A failed release is now surfaced instead of being silently treated as success.

Concurrency groups keyed on queue metadata rather than the tenant can still resolve to the wrong group when a consumer completes a message it did not enqueue, so this does not yet cover that case.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Release note text explains internal implementation instead of user impact

The release note added for this change (.changeset/fair-queue-concurrency-slot-leak.md:5-7) describes internal mechanics and a maintainer-facing caveat instead of one plain sentence about what changed for the user, which is what the repository requires because this text ships verbatim in user-visible release notes.
Impact: Users reading the release notes get implementation detail and an internal caveat rather than a clear statement of the behaviour change.

Rule in AGENTS.md on changeset wording

AGENTS.md ("Changesets and Server Changes") states: "Write the description for users, not maintainers. Both changesets and .server-changes/ notes ship verbatim in user-visible release notes. Lead with what changed for the user - one plain sentence describing behavior, not implementation, and never naming internal tools or infra."

The note instead says slots are "freed before the record needed to recover them is discarded, and released before a reclaimed message goes back on the queue. A failed release is now surfaced instead of being silently treated as success.", and the second paragraph documents a remaining internal limitation about "concurrency groups keyed on queue metadata".

Suggested change
Fair queue consumers no longer leak the concurrency slots that gate a tenant's throughput. Slots were held by messages that had already finished, were never reclaimed, and once enough of them accumulated every queue belonging to that tenant stopped being served. Slots are now freed on the paths that previously skipped them, freed before the record needed to recover them is discarded, and released before a reclaimed message goes back on the queue. A failed release is now surfaced instead of being silently treated as success.
Concurrency groups keyed on queue metadata rather than the tenant can still resolve to the wrong group when a consumer completes a message it did not enqueue, so this does not yet cover that case.
Fixes a bug where a tenant's queued work could stop running permanently after messages finished processing.
Open in Devin Review

Was this helpful? React with 👍 or 👎 to provide feedback.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants