Skip to content

Shash sriv/stream 1807 discard tasks from incorrect applications - #803

Open
ShashSriv wants to merge 5 commits into
mainfrom
ShashSriv/stream-1807-discard-tasks-from-incorrect-applications
Open

ShashSriv wants to merge 5 commits into
mainfrom
ShashSriv/stream-1807-discard-tasks-from-incorrect-applications

Conversation

@ShashSriv

Copy link
Copy Markdown

A pool can only run tasks for the applications its workers serve. Nothing checked
that, so a misrouted task never died:

  • pull mode nobody claimed it and it sat inthe DB until it paged us
  • push mode it was claimed, found no worker, and got re-fetched forever.

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 to
upkeep's existing discard/deadletter path.

Push pools derive the allowed set from worker_map and are covered on deploy. Pull
pools have no worker_map, so they stay unrestricted until ops sets applications
on them — follow-up PR in ops. Nothing changes behaviour until then.

ShashSriv and others added 3 commits October 1, 2026 15:53
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>
@ShashSriv
ShashSriv requested a review from a team as a code owner October 2, 2026 19:50
@linear-code

linear-code Bot commented Oct 2, 2026

Copy link
Copy Markdown

STREAM-1807

@cursor cursor Bot left a comment •

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Stale Bugbot comment from a previous run.

Comment thread src/config/mod.rs
Comment thread src/grpc/server.rs
Comment thread src/grpc/server.rs
…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.

@cursor cursor Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Cursor Bugbot has reviewed your changes and found 1 potential issue.

Fix All in Cursor

❌ 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.

Comment thread src/kafka/activation_batcher.rs
metrics::histogram!("consumer.unknown_application_failures", "topic" => topic)
.record((attempts - successes) as f64);

self.deadletter_batch.clear();

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

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.

@ShashSriv ShashSriv Oct 2, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

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>
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.

1 participant