Skip to content

isolate worker completions on dedicated gRPC connections - #166

Open
Tomer Rosenthal (torosent) wants to merge 7 commits into
mainfrom
torosent-go-completion-transport-draft
Open

Tomer Rosenthal (torosent) wants to merge 7 commits into
mainfrom
torosent-go-completion-transport-draft

Conversation

@torosent

@torosent Tomer Rosenthal (torosent) commented Oct 6, 2026 •

Copy link
Copy Markdown
Member

Why

A DTS server allows 100 concurrent HTTP/2 streams per connection, and a worker's long-lived GetWorkItems stream permanently occupies one. With a single connection per worker, activity/orchestration/entity completions and abandons queue behind that limit before they reach the server, so a busy worker stalls on acknowledgements even when its local concurrency limits have room.

What

  • A worker now keeps one intake connection (Hello, GetWorkItems, everything else) and routes an exact allowlist of Complete*/Abandon* methods round-robin over a bounded group of dedicated completion connections. It is still one logical worker: one WorkerID, one executor, one intake stream, and unchanged per-kind concurrency limits.
  • Completion isolation is a required part of the worker architecture, not a mode. The budget is configurable with WithWorkerCompletionConnections(n) (1–8, default DefaultWorkerCompletionConnections = 3; the default is a starting point, not a tuned optimum).
  • Owned workers build the extra connections through the existing DTS connection factory, so auth/token refresh, TLS, authority, metadata, call options and message limits are identical to intake.
  • Lifecycle is generation-safe: accepted work is acknowledged on the connection group that accepted it, and a group is retired and closed only after its in-flight work drains. Shutdown is bounded, partial dial failures roll back, and every owned connection is closed exactly once. Caller-owned connections are never closed.
  • Graceful shutdown and run-context cancellation now keep the intake stream open until accepted work and its acknowledgements drain; new dispatch stops immediately.

Validation

  • go build ./..., go vet, gofmt, full tests and race checks pass, including minimum Go 1.25. CI passes on Go 1.25 coverage, Go 1.27 race, samples and CodeQL.
  • Real HTTP/2/TLS fixtures with a 100-stream limit verify completion isolation while preserving one worker/intake. Deterministic regressions cover method/options routing, receive-to-dispatch lease retention, blocked stream-open cancellation, reconnect retirement, saturated shutdown, ownership/rollback and terminal work-item slot release.
  • Real DTS on AKS: graceful shutdown and run-context cancellation acknowledged accepted work before client intake finished; forced shutdown and separate recovery produced correct persisted output.
  • Post-review comparison, 615abb6 versus 1d83a1f: two 600-second windows at 1,075 starts/s, with the same Go 1.25 build, fixture, observers, resources and completion budget 3. The fixed worker acknowledged and client-validated all 645,000 planned workflows; SQL independently confirmed every output with zero unfinished instances. It measured 210,208 persisted actions/s during input, and 216,057 including drain. The pre-review control measured 196,006 during input and reached 94.7% of planned starts.
  • In the earlier same-process routing experiment, mean completion pre-header wait dropped from approximately 1,421 ms to 20 ms.

Risk / reviewer notes

  • Out of scope: retry-budget redesign, receive fairness, client long polls.
  • One ordered post-review comparison is not a repeated or isolated causal performance proof: placement, order and persistence variation remain. Both windows reached the disk-bandwidth ceiling. Three completion channels are not a universal optimum.

A DTS server allows 100 concurrent streams per HTTP/2 connection, and the
long-lived GetWorkItems stream holds one. With one connection per worker,
completions and abandons for a busy worker queue behind that limit.

Workers now keep one intake connection (Hello, GetWorkItems) and route an
exact allowlist of Complete*/Abandon* methods round-robin over a bounded
group of owned completion connections (default 3, 1-8 via
WithWorkerCompletionConnections). Work is completed on the connection
generation that accepted it, and a retired generation closes only after its
in-flight work drains.

Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
@torosent
Tomer Rosenthal (torosent) marked this pull request as ready for review October 6, 2026 18:48
Copilot AI balanced review requested due to automatic review settings October 6, 2026 18:48

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

Copilot review overview

🟡 Changes recommended

Public documentation retains provisional internal-review language, and one factory validation error misidentifies intake failures.

Review effort: Balanced
Findings: 2 Low severity

Open (2)
What changed in this PR

Introduces dedicated gRPC completion connections to prevent worker acknowledgements from competing with intake streams.

Changes:

  • Adds configurable completion-connection groups with round-robin routing and generation-safe retirement.
  • Preserves intake leases during graceful shutdown and cancellation.
  • Adds extensive lifecycle, routing, ownership, and integration coverage.
File Description
README.md Documents worker connection roles and budgets.
CHANGELOG.md Records completion isolation and shutdown changes.
client/​grpc_worker.go Adds configuration and connection-group lifecycle.
client/​grpc_worker_compat.go Updates borrowed-listener requirements.
client/​grpc_worker_connection.go Implements completion routing and ownership.
client/​grpc_worker_connection_test.go Tests routing, limits, readiness, and cleanup.
client/​grpc_worker_drain_test.go Tests lease-preserving shutdown behavior.
client/​grpc_worker_processor.go Synchronizes dispatch and graceful draining.
client/​grpc_worker_transport_test.go Updates borrowed-transport coverage.
durabletaskscheduler/​README.md Documents completion isolation architecture.
durabletaskscheduler/​connection_test.go Tests compatibility-listener ownership.
durabletaskscheduler/​resource_id_test.go Updates authentication reconnect coverage.
durabletaskscheduler/​worker_completion_test.go Verifies owned connection option parity.
exporthistory/​versioning_test.go Supplies a dedicated test completion transport.
samples/​largepayloads/​largepayloads_test.go Verifies intake remains leased through drain.
samples/​worker/​README.md Documents the completion-budget sample.
samples/​worker/​worker.go Demonstrates a one-connection completion budget.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread client/grpc_worker_connection.go Outdated
Comment thread durabletaskscheduler/README.md Outdated
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 6, 2026 19:07

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

Copilot review overview

🟡 Changes recommended

Stream establishment can become uncancellable, leaving worker startup or shutdown permanently blocked.

Review effort: Balanced
Findings: 1 High severity

Open (1)
Resolved since last review (2)

Comment thread client/grpc_worker.go
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 6, 2026 19:23

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

Copilot review overview

🟡 Changes recommended

Cancellation-path abandonment can race intake stream closure because it is not tracked as pending work.

Review effort: Balanced
Findings: 1 High severity

Open (1)
Resolved since last review (1)

Comment thread client/grpc_worker_processor.go Outdated
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 6, 2026 19:36

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

Copilot review overview

🟡 Changes recommended

A receive-to-dispatch race can cancel intake before a newly received item is registered and abandoned.

Review effort: Balanced
Findings: 1 High severity

Open (1)
Resolved since last review (1)

Comment thread client/grpc_worker_processor.go Outdated
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 6, 2026 19:51

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

Copilot review overview

🟡 Changes recommended

A receive-registration race can cancel intake and release a work-item lease before its abandonment is tracked.

Review effort: Balanced
Findings: 2 High severity

Open (2)

Comment thread client/grpc_worker_processor.go Outdated
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Copilot AI balanced review requested due to automatic review settings October 6, 2026 20:16

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

Copilot review overview

🔵 Needs a closer look

The concurrency-sensitive connection ownership, reconnection, and lease-drain lifecycle warrants final human review despite comprehensive tests.

Review effort: Balanced
Findings: None

Resolved since last review (2)

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