Skip to content

[BUG] Replay of task with bad limit expr bypasses limits, poisons step's limit cache for 5mins #4874

Description

@sgouliarmis

Describe the issue

While investigating #4612, Claude became sure that there were two bugs in the handling of limit expressions:

  • a task with a bad limit expression (i.e. CEL may compile but at least fails to evaluate it) will fail to launch, but replaying it actually works, and no limits are applied on that second run. Seems to be because insertTasks tries to evaluate and sees the evaluation fail, but replayTasks doesn't and just lets the scheduler do its job, which in turn sees no eval result row in the DB for a limit expression on that task and happily runs it without limiting it.
  • the same task will prevent the next ones, within the same step, from having limits applied to them for the next 5 minutes. This is because GetTaskRateLimits keeps a boolean, that signals whether a limit value applies, in cache for 5 minutes, and the value is keyed by step ID. For subsequent tasks in the same step, the same boolean decision will be re-used, ignoring those tasks' limit value (information which can be found in the database and is not consulted). The "bug", if that is one, is perhaps that the key is the step ID in the first place. At this point I got a bit confused, but I think the current code does not match the intended design (a limit expression belongs to a step, but since tasks each bind their own values to the expression's variables, a limit expression value belongs to a task, and information derived from that value should be keyed by task).

Impact: Limits are silently not applied in some cases. If the limits were there to limit use of a costly third-party service, the user might end up with a nasty surprise.

Environment

  • SDK: Python v1.39.0
  • Engine: Self-hosted v0.105.16

Expected behavior

A replayed task should fail if its limit expression is invalid. It should also not affect the limiting of subsequent tasks (here, within the same step).

Code to Reproduce, Logs, or Screenshots

An entirely Claude-generated Python script that reproduces the issue on my local setup (v0.105.16, sdk 1.39.0) is attached. It generates randomly named buckets and workflows to avoid the chore of cleaning up after itself before re-runs. Env vars HATCHET_CLIENT_TOKEN and HATCHET_CLIENT_TLS_STRATEGY may have to be set, as Claude indicates in the docstring.

Additional context

Bug #2 might be of the same kind as #4746, and while it involves different flags in different files, they are related: Queuer.hasRateLimits is set based on the output of GetTaskRateLimits where this issue's stepsWithRateLimits is computed.

Bug #1 (and #2) look like they could happen downstream of #4380.

The rest of this submission, until the AI disclaimer, is a Claude-generated blurb, meant to give a headstart to the next LLM handling this:

Engine v0.105.16 (tag v0.105.16 +1 commit, 7c89ec556). All line refs are that tree.

Bug 1 — replay never evaluates rate-limit expressions.
  insertTasks (pkg/repository/task.go:2206) gates on stepConfig.ExprCount > 0 (:2500),
  evaluates each "StepExpression" via ParseAndEvalStepRun (:2523), persists results to
  v1_task_expression_eval (:2757 -> :4196). replayTasks (:2788) has no counterpart block:
  it evaluates only v1_step_concurrency expressions (:2923) and never calls
  getStepExpressions (:3213, sole caller :2217). ExprCount appears at exactly :2500 and
  :3221. A task failed at insert has zero eval rows; replay re-queues it without
  re-evaluating and without re-failing it.

Bug 2 — a step-level verdict inferred from task-level evidence, then cached.
  queueRepository.GetTaskRateLimits (pkg/repository/scheduler_queue.go:582-821) builds
  rate limits solely from v1_task_expression_eval (:625). Two identifiers, do not
  conflate: stepsWithRateLimits is a call-local map[uuid.UUID]bool (:596) that never
  escapes the function; cachedStepIdHasRateLimit is a *cache.Cache field on
  queueRepository (:137, constructed :141 with a 5-minute TTL).
  stepsWithRateLimits[stepId] is set true only on seeing a DYNAMIC_RATE_LIMIT_KEY eval
  row (:644); its only other writer (:800) iterates ListRateLimitsForSteps over
  "StepRateLimit", which is always empty under v1 — every rate limit, including fully
  static ones, is stored as "StepExpression" rows (workflow.go:1085), so that path can
  never rescue a false verdict. The local map is then copied into the cache, one entry
  per step in the batch (:815-816), and the cache alone drives the early return (read
  :614, return nil,nil at :620-622). The guard requires every step in the batch to be
  both cached and false for skipRateLimiting to stay true (:611-618). The cached datum is
  a bool "does this step have a rate limit", not a limit value — so the failure mode is
  rate limiting skipped wholesale, not a wrong limit applied.

Coupling to #4746: Queuer.hasRateLimits (pkg/scheduling/v1/queuer.go:65) is assigned in
  exactly one place, :274, from len(GetTaskRateLimits(...)) > 0 — the call at :261, with a
  second call site at :856 — and it gates RequeueRateLimitedItems at :226. It is
  write-once-true and never reset. newQueuer (:98) builds its own queueRepository via
  QueueFactory().NewQueue, hence its own 5-minute cache, so a queuer and its cache are
  born together. A freshly leased or rebalanced queuer whose early batches carry no eval
  rows therefore poisons its own cache, receives nil from GetTaskRateLimits, never learns
  hasRateLimits, and never requeues parked items — #4746's symptom by a route that issue
  does not describe (it attributes the unlearned flag to there being no active items at
  all). Code reading only: only the over-admission above is reproduced. Neither fix
  subsumes the other — dropping the queue-local gate leaves GetTaskRateLimits still
  returning no units, and fixing the cache does not help when there are no active items
  to pass in.

Data model: "StepExpression" (definition, per step; kinds KEY/UNITS/VALUE/WINDOW — the
  enum has no non-rate-limit members) -> v1_task_expression_eval (per task, frozen at
  insert) -> "RateLimit" (bucket, per tenant+key).

Measurement trap: "RateLimit".value is a refilling bucket (lastRefill + window), not a
  consumption counter. It reads full again once the window rolls, so it cannot serve as
  after-the-fact evidence. Use v1_task_expression_eval row counts, or saturate the limit
  and observe task status.

Repro (proof.py): limit 1/HOUR, spent by a baseline phase, so any later COMPLETED cannot
  have been throttled. It cancels queued tasks before replaying — a rate-limited task
  sharing the replayed task's scheduling batch supplies eval rows and masks bug 2.

Scope: demonstrated on a single-task workflow. Queue name is stepConfig.ActionId
  (task.go:2231), so batches hold one step and the early return trips readily. Where
  several steps share a queue, every step in the batch must be cached false. Mechanism
  established; field frequency not measured.

Fix direction: derive hasRateLimit from step configuration (ExprCount > 0, or presence of
  "StepExpression" rows) rather than from whether sampled tasks carried eval rows. For
  bug 1, replayTasks should fail the task when a step with ExprCount > 0 yields no eval
  rows, matching insertTasks.

🤖 AI Disclosure
  • I acknowledge that an LLM was used in the creation of this Issue, in accordance with Hatchet's AI_POLICY.md.
  • Details: The analysis of this issue was done with heavy use of Opus 5, and the Python script is entirely generated by it. This post itself, besides the handover blurb, is human-made.

proof.py

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions