Conversation
A pool can only execute tasks for the applications its workers serve. Nothing checked that, so a misrouted activation was immortal: in pull mode no worker ever claimed it and it sat Pending until the lag alert paged someone, and in push mode it was claimed, found no worker, had its claim expire, and was fetched again forever. Add an `applications` allowlist to broker config and discard anything outside it in ActivationBatcher::reduce, alongside the existing killswitch and expiry filters, publishing the payload to the pool's deadletter topic so a misconfiguration is recoverable rather than silent data loss. Push pools derive the set from `worker_map`, which already lists every application they can deliver to. Pull pools are never given a `worker_map`, so an unset `applications` admits everything and they keep running until ops opts them in; a pool whose `worker_map` names an application outside the allowlist now fails to start, since its tasks would be discarded before reaching that worker. Two supporting changes for activations that are already stored, which the consumer filter cannot reach: the push thread now sets Failure instead of returning silently, handing the task to upkeep's discard/deadletter path, and get_task returns failed_precondition for an unserved application so a worker pointed at the wrong pool fails instead of polling an empty queue. The producer swap and batch send in flush are extracted into helpers and shared with the forwarding path rather than copied. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
The existing test stops at `reduce`, so the publish in `flush` was never exercised. Produce a discarded activation through a real broker and consume it back off the pool's deadletter topic, following the pattern upkeep uses for its own deadletter test. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
…orrect-applications main replaced the batcher's raw FutureProducer with ProducerBackend, which also switches on the use_arroyo_producer runtime flag. Both of its swap blocks move into `producer_for`, which now takes the flag, and `send_all` produces through the backend. The runtime config read is hoisted so the deadletter batch can see the flag too. Also restores the ActivationStatus import in push/thread.rs, which main dropped when it reworked claim undo, and parameterizes the deadletter test over the Arroyo flag to match main's forwarding test.
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 7bf2633. Configure here.
| metrics::histogram!("consumer.unknown_application_failures", "topic" => topic) | ||
| .record((attempts - successes) as f64); | ||
|
|
||
| self.deadletter_batch.clear(); |
There was a problem hiding this comment.
Bug: Deadlettered tasks are lost if publishing to the deadletter topic fails, as Kafka offsets are committed regardless of the publish outcome.
Severity: HIGH
Suggested Fix
After calling send_all, check if the number of successes equals the number of attempts. If they do not match, do not clear the deadletter_batch and return a result that prevents the consumer from committing the Kafka offsets. This will ensure that failed messages can be reprocessed.
Prompt for AI Agent
Review the code at the location below. A potential bug has been identified by an AI
agent. Verify if this is a real issue. If it is, propose a fix; if not, explain why it's
not valid.
Location: src/kafka/activation_batcher.rs#L327
Potential issue: In the `flush()` method's deadletter handling, tasks are published to a
deadletter topic via `send_all`. However, the code does not verify if all publish
attempts were successful. It unconditionally clears the `deadletter_batch` on line 327
and returns `Ok(Some(...))`. This successful return value causes the upstream Kafka
consumer to commit the message offsets. In a scenario where publishing to the deadletter
topic partially fails (e.g., due to transient network issues), the failed messages are
permanently lost. They are not in the database, failed to reach the deadletter topic,
and will not be re-consumed from Kafka because their offsets have been committed.
There was a problem hiding this comment.
Offsets commit per inflight batch, not per message, so the only way to avoid committing is to fail the flush which sends the batch down the reducer error path. A brief DLQ hiccup would stall the whole topic, which is worse than dropping a task this pool can't run anyway.
The send waits out kafka_send_timeout_ms with librdkafka retrying underneath. This also matches what the demoted-namespace forward path already does.
We'll alert on consumer.unknown_application_failures instead.
- `worker_map` is only delivered through in push mode, and `from_args` gives pull pools a default entry they never read, so the subset check rejected a pull pool opting in to any other application. Scope the check to push. - `set_task_status` claims through the same path as `get_task`, so apply the check there too. It returns no task rather than an error, because an error makes the worker re-send the status update. - Discards and forwards bypass `batch`, so `is_full` never tripped on them and a misrouted topic only flushed on the timer. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

A pool can only run tasks for the applications its workers serve. Nothing checked
that, so a misrouted task never died:
This pr allows the consumer to drop it before it's ever stored: metric, error log, and a copy
published to the pool's deadletter topic so it's recoverable. Tasks already in the
DB get the same treatment via the push thread setting
Failure, which hands them toupkeep's existing discard/deadletter path.
Push pools derive the allowed set from
worker_mapand are covered on deploy. Pullpools have no
worker_map, so they stay unrestricted until ops setsapplicationson them — follow-up PR in
ops. Nothing changes behaviour until then.