Skip to content

feat(consumergate): runtime stop/start of queue controllers#375

Open
sbalabanov wants to merge 1 commit into
mainfrom
consumer-gate-impl
Open

feat(consumergate): runtime stop/start of queue controllers#375
sbalabanov wants to merge 1 commit into
mainfrom
consumer-gate-impl

Conversation

@sbalabanov

@sbalabanov sbalabanov commented Jul 15, 2026

Copy link
Copy Markdown
Contributor

Implements the consumer-gate RFC (#374) and showcases a Request Cancellation e2e test using it.

Summary

  • Adds platform/extension/consumergate with the Gate, Entry, Admin, Factory, and configuration contracts, plus generated mocks.
  • Adds file-backed and no-op implementations. The file backend reads gate state directly from shared files for every delivery and uses atomic temp-file-plus-rename writes; it does not cache snapshots or verdicts.
  • Clears every queue delivery through its consumer-group/partition gate before controller processing. Blocked deliveries remain in flight through visibility extensions without consuming retry attempts or breaking partition ordering.
  • Fails open when gate reads, monitoring, or visibility extensions fail. Consumer shutdown instead leaves a blocked delivery unacked for normal redelivery.
  • Keeps parked records and payload files only while a delivery is actively blocked.
  • Uses polling as the portable convergence mechanism. Filesystem events may be added later as an optimization alongside polling because event delivery varies across operating systems, bind mounts, overlay/network filesystems, rootless Docker, and Docker Desktop.
  • Wires the shared file gate into gateway, orchestrator, and runway. Stovepipe uses the no-op implementation.
  • Runs application containers under the host test UID/GID with rootful Docker, while using 0:0 under rootless Docker so bind-mounted gate files remain readable by the host test process.

Deterministic cancellation E2E coverage

TestCancel_CaughtPreBatch_NeverLands closes the runway merge-conflict-check gate for the test queue before landing a request. The test observes the exact parked delivery, cancels the request to its terminal state, opens the gate, waits for the parked record to disappear, and uses a sentinel request to prove the stale message was consumed without reviving or landing the cancelled request.

Testing

  • make test — passed (76/76 Bazel test targets).
  • //test/integration/... — passed (8/8 integration targets).
  • //test/e2e/... — passed (2/2 E2E targets).
  • git diff --check — passed.

@CLAassistant

Copy link
Copy Markdown

CLA assistant check
Thank you for your submission! We really appreciate it. Like many open source projects, we ask that you sign our Contributor License Agreement before we can accept your contribution.
You have signed the CLA already but the status is still pending? Let us recheck it.

Comment thread platform/consumer/consumer.go Outdated
Comment thread service/submitqueue/orchestrator/server/main.go Outdated
Comment thread platform/extension/consumergate/file/store.go Outdated
Comment thread platform/consumer/gate.go Outdated
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from d04bf53 to fb4fde8 Compare July 16, 2026 00:38
Comment thread platform/extension/consumergate/consumergate.go Outdated
Comment thread platform/extension/consumergate/README.md Outdated
Comment thread platform/extension/consumergate/README.md Outdated
Comment thread platform/consumer/gate.go Outdated
Comment thread platform/consumer/gate.go Outdated
Comment thread platform/consumer/gate.go Outdated
Comment thread platform/consumer/gate.go Outdated
Comment thread platform/consumer/gate.go Outdated
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from fb4fde8 to 0b06f08 Compare July 16, 2026 20:43
@behinddwalls
behinddwalls changed the base branch from rfc-consumer-gate to main July 16, 2026 22:00
Comment thread platform/consumer/consumer.go Outdated
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch 2 times, most recently from 4e229a3 to 2db191e Compare July 16, 2026 22:55
Comment thread platform/consumer/consumer.go Outdated
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch 2 times, most recently from 617595b to 8c7078c Compare July 21, 2026 00:10
@sbalabanov
sbalabanov marked this pull request as ready for review July 21, 2026 00:10
@sbalabanov
sbalabanov requested review from a team and behinddwalls as code owners July 21, 2026 00:10
@sbalabanov
sbalabanov marked this pull request as draft July 21, 2026 00:11
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from 8c7078c to b17ed5a Compare July 21, 2026 00:13
@sbalabanov
sbalabanov marked this pull request as ready for review July 21, 2026 00:13
@sbalabanov
sbalabanov marked this pull request as draft July 21, 2026 00:13
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from b17ed5a to 40d682f Compare July 21, 2026 00:27
@sbalabanov
sbalabanov marked this pull request as ready for review July 21, 2026 00:27
@sbalabanov
sbalabanov marked this pull request as draft July 21, 2026 00:27
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from 40d682f to 357d933 Compare July 21, 2026 02:27
@sbalabanov
sbalabanov marked this pull request as ready for review July 21, 2026 16:09
@sbalabanov
sbalabanov marked this pull request as draft July 21, 2026 16:09
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from 357d933 to dd2a358 Compare July 22, 2026 01:29
@sbalabanov sbalabanov changed the title feat(consumergate): runtime stop/start of queue controllers + deterministic e2e cancel test feat(consumergate): runtime stop/start of queue controllers Jul 22, 2026
@sbalabanov
sbalabanov marked this pull request as ready for review July 22, 2026 01:30
@sbalabanov
sbalabanov marked this pull request as draft July 22, 2026 01:30
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch 2 times, most recently from 7622a3a to ea4e622 Compare July 22, 2026 03:20
Implement the consumer-gate RFC with a file-backed gate shared by gateway,
orchestrator, and runway consumers.

- clear deliveries through a consumer-group/partition gate before controller
  processing while extending visibility for blocked deliveries
- expose caller-owned delivery descriptors and let gate implementations stamp
  gate-owned parked-record fields
- keep parked payload files only while deliveries are actively blocked and
  remove them on every terminal watch path
- fail open on gate or visibility-extension failures and leave blocked
  deliveries unacked during shutdown for normal redelivery
- run service containers with the host test UID/GID under rootful Docker while
  preserving rootless Docker ownership mapping
- make the cancellation E2E scenario deterministic by parking the runway
  merge-conflict-check delivery until cancellation reaches a terminal state
- consolidate consumer white-box and behavioral unit tests in the consumer
  package

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
@sbalabanov
sbalabanov force-pushed the consumer-gate-impl branch from ea4e622 to bf2e709 Compare July 22, 2026 17:01
@sbalabanov
sbalabanov marked this pull request as ready for review July 22, 2026 17:01
@albertywu

albertywu commented Jul 22, 2026

Copy link
Copy Markdown
Contributor

Could you clarify what happens in the situation where there's 10 buffered messages awaiting controller X and the controller is paused?

Codex is telling me that the first message would wait at the closed gate and have visibility continually extended but the trailing 9 would remain buffered without visibility extension. Every 60 seconds, those 9 expire and are redelivered, consuming retry attempts and creating duplicate buffered delivery. And after 3 attempts they move to DLQ.

Assuming this is true, wouldn't this cause the undesirable impact of dropping the first 9 buffered messages while controller X is in a prolonged paused state?

@albertywu

albertywu commented Jul 22, 2026

Copy link
Copy Markdown
Contributor

It seems that with this design that uses file-based marking of gated controllers, the gating would not survive node crashes.

For example, let's say:

  1. Controller X is paused on Node A (so marker is left on Node A on the filesystem)
  2. Node A crashes and is replaced with Node B with fresh filesystem
  3. Controller X is unpaused now (pausing does not survive node crashes)

Similarly, if the fleet has 10 nodes, I don't think we fan-out the file write to all 10 nodes -- so the pause behavior would only impact the single node that the request landed on. Basically this design seems to work with a single-node only and assumes the filesystem is stateful and persistent across the lifetime of gate usage.

We could fix this by moving the gating implementation from file-based to a storage interface backed with a db. It seems to fit right in with existing Queue implementation which uses db storage so seems better to include as part of that vs a special-snowflake file-based implementation. Something like this could work as the core:

CREATE TABLE consumer_gates (
    consumer_group VARCHAR(255) NOT NULL,
    partition_key  VARCHAR(255) NOT NULL, -- empty means all partitions
    reason         TEXT NOT NULL,
    created_by     VARCHAR(255) NOT NULL,
    created_at_ms  BIGINT NOT NULL,
    PRIMARY KEY (consumer_group, partition_key)
);

@sbalabanov

Copy link
Copy Markdown
Contributor Author

Could you clarify what happens in the situation where there's 10 buffered messages awaiting controller X and the controller is paused?

Codex is telling me that the first message would wait at the closed gate and have visibility continually extended but the trailing 9 would remain buffered without visibility extension. Every 60 seconds, those 9 expire and are redelivered, consuming retry attempts and creating duplicate buffered delivery. And after 3 attempts they move to DLQ.

Assuming this is true, wouldn't this cause the undesirable impact of dropping the first 9 buffered messages while controller X is in a prolonged paused state?

It depends on a delivery order. For synchronized queues (only one message could be processed at a time), other messages would be held and won't be extended. For parallel queues, if there is enough controller instances on a production or e2e test system, multiple messages would be parked and their life extended as a result of this.

It is true that messages not accepted will not be extended. For e2e tests, this is not a problem as we control the number of messages to be validated.

@sbalabanov

Copy link
Copy Markdown
Contributor Author

It seems that with this design that uses file-based marking of gated controllers, the gating would not survive node crashes.

For example, let's say:

  1. Controller X is paused on Node A (so marker is left on Node A on the filesystem)
  2. Node A crashes and is replaced with Node B with fresh filesystem
  3. Controller X is unpaused now (pausing does not survive node crashes)

Similarly, if the fleet has 10 nodes, I don't think we fan-out the file write to all 10 nodes -- so the pause behavior would only impact the single node that the request landed on. Basically this design seems to work with a single-node only and assumes the filesystem is stateful and persistent across the lifetime of gate usage.

We could fix this by moving the gating implementation from file-based to a storage interface backed with a db. It seems to fit right in with existing Queue implementation which uses db storage so seems better to include as part of that vs a special-snowflake file-based implementation. Something like this could work as the core:

CREATE TABLE consumer_gates (
    consumer_group VARCHAR(255) NOT NULL,
    partition_key  VARCHAR(255) NOT NULL, -- empty means all partitions
    reason         TEXT NOT NULL,
    created_by     VARCHAR(255) NOT NULL,
    created_at_ms  BIGINT NOT NULL,
    PRIMARY KEY (consumer_group, partition_key)
);

File-based implementation is for e2e tests only running on the same machine. While it is theoretically possible to hook it with configuration delivery (to let it survive a node crash / restart), it is indeed easier to come up with a centralized database-backed implementation. The cost would be a single synchronous database call per each controller run - files are way faster and simpler too - the gate could be closed or open with basic shell commands.

@albertywu

albertywu commented Jul 23, 2026

Copy link
Copy Markdown
Contributor

It depends on a delivery order. For synchronized queues (only one message could be processed at a time), other messages would be held and won't be extended. For parallel queues, if there is enough controller instances on a production or e2e test system, multiple messages would be parked and their life extended as a result of this.
It is true that messages not accepted will not be extended. For e2e tests, this is not a problem as we control the number of messages to be validated.

Thanks, the parallel-partition case makes sense. Since the RFC documents production debugging use-case ("Operational pause during incident"), my concern is multiple messages in the same partition, where one lease owner processes serially. Only the head reaches the gate; the prefetched tail remains in flight without visibility extensions and can exhaust retries.

Potential fixes:

  1. Require BatchSize=1 for gateable subscriptions.
  2. Extend visibility for every buffered delivery in the gated partition.
  3. Prevent additional prefetch while a partition is gated.

For the E2E-only scope, option 1 seems simplest, at the cost of throughput

// written — without further interpretation, as what to do with a failed wait
// is the caller's policy. If ctx is cancelled while monitoring, the channel
// yields ctx.Err(). The implementation removes the parked record before
// yielding on every terminal path, so ListParked contains only deliveries

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

not sure if we need to document it, message gated is never removed from the queue for a given consumer and partition


// Admin is the write surface used by tests and tooling to operate gates and
// inspect what a stopped controller is holding.
type Admin interface {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

should we rename it to Manage? or Manager maybe

@behinddwalls

Copy link
Copy Markdown
Collaborator

It depends on a delivery order. For synchronized queues (only one message could be processed at a time), other messages would be held and won't be extended. For parallel queues, if there is enough controller instances on a production or e2e test system, multiple messages would be parked and their life extended as a result of this.
It is true that messages not accepted will not be extended. For e2e tests, this is not a problem as we control the number of messages to be validated.

Thanks, the parallel-partition case makes sense. Since the RFC documents production debugging use-case ("Operational pause during incident"), my concern is multiple messages in the same partition, where one lease owner processes serially. Only the head reaches the gate; the prefetched tail remains in flight without visibility extensions and can exhaust retries.

Potential fixes:

  1. Require BatchSize=1 for gateable subscriptions.
  2. Extend visibility for every buffered delivery in the gated partition.
  3. Prevent additional prefetch while a partition is gated.

For the E2E-only scope, option 1 seems simplest, at the cost of throughput

ideally we just extend for everything in the buffer...so the semantics remain usable for batched polling as well

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.

5 participants