Repository navigation
[#33937] DocDB: Priority queueing for RPC worker dispatch (RpcPriorityQueue) - #25
Draft
craigsoules wants to merge 6 commits into
Draft
craigsoules wants to merge 6 commits into
craigsoules wants to merge 6 commits into
Conversation
Stage 0 of RPC priority queueing: introduce the RpcPriority enum (kHigh/kNormal/kLow) for dispatch ordering, distinct from ServicePriority which selects the executing thread pool. Plumb a per-service RpcPriority through RpcServer::RegisterService into ServicePool, and classify services at registration: - kHigh: consensus (master + tserver), master heartbeats, YSQL lease - kNormal (default): TabletServerService, PgClientService, master client/DDL/cluster services, and all unspecified services - kLow: TabletServerAdmin (index backfill), remote bootstrap, backups, CDC/xCluster, master xrepl, stateful services (test echo, auto-analyze, pg cron leader) The priority is stored (and logged) by ServicePoolImpl but not yet consumed; a follow-up change adds the RpcPriorityQueue that uses it to order dispatch to worker thread pools behind a feature flag. No behavior change. --- _automated · pi (claude-fable-5)_
Stage 1 of RPC priority queueing: the core RpcPriorityQueue class, not yet
wired into any dispatch path.
RpcPriorityQueue gates admission of RPC worker-pool work behind a single
permit budget and orders waiting work by RpcPriority. Tasks are submitted
together with the YBThreadPool they should execute on; the queue decides
when they are handed over, so one budget bounds concurrently executing RPC
work across all worker pools (the existing normal/high/tagged pool split is
kept only as execution targets).
Design:
- Packed atomic state word {closed, dispatched, queued}. While nothing is
waiting and a permit is free, Enqueue claims a permit with one CAS and
forwards the task directly (no lock, no band traffic) - the common
unsaturated case costs about what a direct YBThreadPool::Enqueue does.
- Once the budget is exhausted, tasks wait in per-priority FIFO bands
guarded by an adaptive spinlock (simple_spinlock: bounded spin, then
futex). Completions hand their permit directly to the highest-priority
waiting task (kHigh > kNormal > kLow) rather than releasing it, so a
waiting task cannot be overtaken by a concurrent new arrival. The
queued == 0 gate on the fast path preserves the same invariant.
- Slow-path admission is a single CAS (ClaimOrRegisterQueued) that either
claims a permit or increments queued, so it is linearizable with the
release CAS. With a separate re-check and increment, a completion's
release could land between them and leave a task registered as waiting
with dispatched below the budget and no completion left to hand it a
permit: a permanent stall.
- Strict priority, no preemption, no aging/reserves/caps in v1; the
considered follow-ups are documented at the dispatch site.
- A heap-allocated wrapper task observes completion via the pool's Done()
contract, so a permit is released on every path (run, pool refusal,
shutdown) and cannot leak.
- Completions use a single-drainer protocol: the caller that transitions
pending_completions_ from 0 loops over all registered completions; a
completion arriving while a drainer is active (including one triggered
synchronously inside Dispatch by a closing pool refusing the task) only
bumps the counter and returns. Without this, a pool shutdown with a deep
backlog recursed TaskFinished -> Dispatch -> Done -> TaskFinished once per
waiting task and overflowed the stack.
- Two-phase shutdown mirrors YBThreadPool: StartShutdown sets the closed
bit in the state word with a CAS (so no permit-claiming CAS can succeed
afterwards, whether the enqueue started before or after) and fails all
waiting tasks with the same Aborted status the pool uses; CompleteShutdown
waits for dispatched == 0 and pending_completions_ == 0 so the drainer's
final access to the queue is ordered before destruction.
- Invariant max_dispatched <= target_pool.max_workers, enforced at Enqueue
(DCHECK / rate-limited DFATAL). With it, every dispatched task has a
worker, so "dispatched" means "executing" for the queue's own work and a
saturated pool cannot hold permits that an idle pool's higher-priority
task is waiting on.
Flags (both NON_RUNTIME, unused until a later change):
- rpc_priority_queue_enabled (default false)
- rpc_priority_queue_max_dispatched (default 0 = rpc_workers_limit)
Metrics: per-band queued gauges, dispatched counters and wait-time stats,
plus an aborted counter.
Tests cover the fast path, strict ordering under saturation, no-overtake of
waiting tasks, shared budget across target pools, shutdown drain statuses
and CompleteShutdown blocking, permit release on pool refusal,
EnqueueRacesWithShutdown (8 producers vs concurrent shutdown; no task runs
after CompleteShutdown, every task gets exactly one Done),
PoolShutdownWithDeepBacklogDoesNotRecurse (20k backlog, pool shut down
first), HighPriorityForIdlePoolNotBlockedBySaturatedPool, and a
multi-producer stress test asserting the budget is never exceeded and no
permit leaks. Pass in release and debug; stress and race tests 100/100
(release) and 30/30 (debug). TSAN is not available in the macOS thirdparty
build; to be run on Linux.
---
_automated · pi (claude-fable-5)_
Stage 2 of RPC priority queueing: wire the queue into the inbound path, gated by rpc_priority_queue_enabled (default false; off path unchanged). - Messenger creates one RpcPriorityQueue in Init when the flag is on, with budget rpc_priority_queue_max_dispatched or, if 0, rpc_workers_limit. Every worker pool the messenger creates (default, tagged, high priority) is sized to rpc_workers_limit, so the budget is clamped to that value to preserve the max_dispatched <= pool.max_workers invariant. Exposed via Messenger::rpc_priority_queue() (null when disabled). - Messenger::ShutdownThreadPools closes the queue, shuts the pools down, then waits for the queue's dispatched count to drain. Reached from both RpcServer::Shutdown and Messenger::Shutdown; idempotent. - ServicePool takes a nullable RpcPriorityQueue*. After admission and timeout scheduling (unchanged), the final dispatch becomes queue->Enqueue(task, rpc_priority, thread_pool) when a queue is present, else thread_pool->Enqueue(task) as before. Both paths end with the task's Done() invoked exactly once, so failure handling is identical. The inline local-call fast path is untouched. All services, including consensus on the high-priority pool, go through the queue when it is enabled. - RpcServer::RegisterService passes the messenger's queue to ServicePool. Test infrastructure: TestServer uses its own small worker pool rather than the messenger's, so when the flag is on it creates its own queue sized to that pool and shuts it down in the same order as the messenger. Its RegisterService takes an RpcPriority. Factories for the two distinctly named generated test services are exposed for tests that need different priorities per service. New end-to-end test (rpc-test PriorityQueueDispatchesHighPriorityServiceFirst): a one-worker server whose permit is held by a running Sleep call dispatches a later-submitted kHigh AshTestService call ahead of earlier-submitted kLow CalculatorService calls. Verification: rpc-test, rpc_stub-test, mt-rpc-test, rpc_priority_queue-test pass with the flag off (default) and with it forced on via --rpc_priority_queue_enabled=true (89 rpc_stub tests exercised the queue in front of the test server pool). The one failure in both modes, RpcStubTest.TestIncoherence, is a local environment issue (no 127.0.0.11/12 loopback aliases on this machine) and is unrelated. --- _automated · pi (claude-fable-5)_
Cluster-level validation for rpc_priority_queue_enabled: KVTableTest gains a KVTablePriorityQueueTest fixture that enables the flag, so that each master and tserver runs a saturated RpcPriorityQueue with consensus/heartbeats (kHigh), user traffic (kNormal) and background services (kLow) all contending for dispatch. The existing LoadTest (4 writers + 4 readers, ClusterVerifier) and Restart (full cluster restart, consensus recovery) bodies are factored into RunLoadTest/RunRestartTest and run under both the default fixture and the priority-queue fixture. The priority-queue variants assert the feature rather than only re-running the load/restart assertions (which would keep passing if the Messenger integration were bypassed): - Every master and tserver messenger has a non-null RpcPriorityQueue with the configured budget. - Per-priority dispatch counters read from each server's metric entity are non-zero for the priorities the workload must produce: kHigh and kNormal on masters (heartbeats, client/DDL) and tservers (consensus, reads and writes), plus kLow on tservers after table creation (CreateTablet via TabletServerAdminService). The restart variant verifies the post-restart servers, proving the queue is recreated, and does not expect kLow since the tablets were created before the restart. - The load test asserts at least one task waited for a permit, so the budget was genuinely saturated. The contention assertion showed that rpc_workers_limit=8 alone never saturated: RPC handlers are mostly asynchronous, so 8 client threads never held 8 permits at once. The fixture keeps rpc_workers_limit=8 but sets rpc_priority_queue_max_dispatched=2, which produces sustained contention (typically ~800-900 waits across tservers per load-test run) with the cluster remaining fully healthy. 5/5 repeated runs pass; the default-fixture LoadTest and Restart are unchanged. --- _automated · pi (claude-fable-5)_
…ck reserve
Stage 3 of RPC priority queueing. Async outbound-call callbacks that today go
straight to the messenger's normal/high worker pools (InvokeCallbackMode
kThreadPoolNormal/kThreadPoolHigh) now go through the RpcPriorityQueue when
rpc_priority_queue_enabled, so they are subject to the same node-level
budget and priority ordering as inbound work instead of being unmanaged
"foreign" work in the same pools. kReactorThread callbacks and sync-RPC
latch callbacks are unchanged (they never used a pool).
Deadlock freedom (task classes). An audit found that the server's internal
YBClient is built on the server's own messenger, and that many RPC handlers
block on results delivered by its callbacks: every synchronous YBClient call
from a handler (SyncLeaderMasterRpc waits on a Synchronizer resolved by a
normal-pool callback), PgClientSession's table-cache Wait, FlushFuture().get(),
GetMetadata().get(), CancelTransaction's status futures, CDC service calls,
etc. With callbacks sharing one budget with handlers, enough simultaneously
blocked handlers would hold every permit and the callbacks that unblock them
could never run: a permanent stall of all RPC processing including consensus.
To make that impossible by construction, RpcPriorityQueue::Enqueue takes an
RpcTaskClass. kInbound tasks may hold at most
max(1, max_dispatched - callback_reserve) permits; kCallback tasks may use
the whole budget and may bypass inbound work waiting on the cap (in which
case everything waiting is inbound work that cannot use the free permit
anyway). Within a band, inbound and callback tasks otherwise dispatch in
arrival order. The packed state word gains a dispatched_inbound field
(closed:1 | dispatched:21 | dispatched_inbound:21 | queued:21); admission and
completion hand-off remain single-CAS / linearizable. Completions are
attributed per class for the single-drainer loop.
A finite reserve cannot break callback -> callback chains; those are
prevented by callbacks not blocking, which is made a contract rather than
an assumption: rpc_priority_queue_check_callbacks_do_not_wait (runtime,
default true) makes a callback task dispatched by the queue that waits on a
latch, condition variable or sleep crash the process in debug builds (where
ThreadRestrictions is compiled in). At the default budget the
callback -> callback exposure is identical to the pre-queue thread pool, so
a violated contract is not a regression there, only in small-budget
configurations. Audit with the check enabled: debug rpc-test, rpc_stub-test
and the KVTable mini-cluster tests (~17k callbacks through tserver queues:
consensus, YBClient batcher/meta-cache, master admin) found no blocking
callback. YSQL and CDC callback paths should be audited with the queue
enabled in debug before it is enabled by default anywhere.
Queued counter bound: with dispatched_inbound in the packed word, queued has
21 bits. Instead of a DCHECK that wraps to zero in release builds, admission
at the bound fails the task with ServiceUnavailable ("RPC priority queue is
full") and counts it in rpc_priority_queue_rejected_full.
ServicePoolImpl::Failure maps a ServiceUnavailable status to
ERROR_SERVER_TOO_BUSY (retry) rather than FATAL_SERVER_SHUTTING_DOWN (drop
connection). All state_ transitions go through one UpdateState CAS helper.
Shutdown as a quiescence barrier:
- Enqueue registers in active_enqueues_ (the role YBThreadPool's adding_
plays) so a call that read the state before closure cannot touch the
queue after CompleteShutdown returns. The registration is released after
Enqueue's last access to the queue and BEFORE any task-controlled code
runs (Done() on rejection paths, or Dispatch(), which can invoke Done()
synchronously when the pool refuses the task): that code may legitimately
shut the queue down, and CompleteShutdown would otherwise wait for a
registration only the enclosing Enqueue can release. This matches
YBThreadPool::Enqueue, which decrements adding_ before calling Done().
- StartShutdown sets drain_complete_ once the bands have been failed;
concurrent callers that lose the closed-bit race wait for it. It records
the draining thread so that a task failed by the drain may re-enter
StartShutdown from its Done() (a no-op) instead of spinning on a flag the
same thread sets only after the loop.
- A completion that observes the closed bit under the lock releases its
permit rather than dispatching a waiting task past the drain, so
"StartShutdown fails every waiting task" is true.
- CompleteShutdown waits for drain completion, zero active enqueues, zero
dispatched tasks and no completion in progress, and DCHECKs that nothing
is left waiting.
- A dispatched task holds its permit until its Done() returns (deliberately:
that lets task code use the queue from Done(), e.g. to enqueue a retry,
without racing a CompleteShutdown on another thread). The constraint we
accept instead is the one YBThreadPool already has for a task calling
Shutdown() on its own pool: task code may call StartShutdown, but a
CompleteShutdown called synchronously from Run()/Done() of a dispatched
task (including the Done() a closing pool invokes on refusal) or from a
task being drained cannot be satisfied. A thread-local records the queue
whose task is executing; such a call is LOG(DFATAL) and ignored, and the
header documents the contract. Enqueue's return contract: false means the
queue itself rejected admission (Done already invoked with Aborted or
ServiceUnavailable); true means the queue accepted the task, which the
target pool may still fail, synchronously if it is closing.
Wiring: ProxyContext::CallbackThreadPool becomes CallbackRecipient returning
a ThreadPoolTaskRecipient (YBThreadPool, or a PriorityQueueCallbackRecipient
adapter that enqueues at a fixed priority/target pool as kCallback).
OutboundCall/LocalOutboundCall hold the recipient; InvokeCallbackTask's
synchronous fallback on enqueue failure is unchanged. Messenger creates the
normal (-> default pool, RpcPriority::kNormal) and high (-> high pool,
RpcPriority::kHigh) recipients alongside the queue. New flag
rpc_priority_queue_callback_reserve (-1 = max(2, 5% of budget), 0 disables).
Tests:
- Unit: CallbacksBypassInboundWaitingOnCap, BlockingHandlersDoNotStarveCallbacks
(20 handlers that each submit a callback to the same queue and wait for
it; verified to deadlock with reserve 0 and complete with reserve 1),
ArrivalOrderPreservedAcrossClassesWithinBand, stress test mixing classes
with a reserve; ConcurrentStartShutdownWaitsForDrain,
QueuedCounterBoundRejectsAdmission; sync-point tests (debug only)
CompletionAfterCloseReleasesPermit and CompleteShutdownWaitsForActiveEnqueue
pinning the exact interleavings; BlockingCallbackIsFatal (death test),
NestedBlockingCallbackChainWithinReserveCompletes;
RejectedTaskMayShutDownQueueFromDone, TaskRejectedAtBoundMayShutDownQueueFromDone,
DrainedTaskMayReenterStartShutdown, PoolRefusedTaskMayStartShutdownFromDone,
CompleteShutdownFromDispatchedTaskIsRejected (death test in debug; asserts
no hang and clean external shutdown in release). 24/24 debug, 22/22
release.
- rpc-test / rpc_stub-test / mt-rpc-test pass with the flag off and on in
debug and release.
- KVTable mini-cluster tests run with budget 2 and callback reserve 1, i.e.
at most one inbound handler executing per server, and assert that
callbacks were dispatched through the queue on masters and tservers. The
cluster stays healthy under the load and restart tests; 5/5 repeated runs.
---
_automated · pi (claude-fable-5)_
Stage 4 of the RPC priority queue work: make the v1 gating boundary explicit in code. No behavior change. The RpcPriorityQueue class comment gains a "What is not gated (v1)" section enumerating every route by which work reaches an RPC worker pool without passing through the queue, why each is acceptable, and which are candidates for routing in a follow-up: inline local calls (ServicePoolImpl::Process fast path), direct Messenger::ThreadPool() submission (docdb wait-queue and object-lock waiter resumption, PgClientService session cleanup and table-query continuations, ThinClientService, CQL processor rescheduling, xCluster safe-time service), TabletPeer::Enqueue on the tablet's tagged pool, the PgClientService shared-memory exchange pool, reactor-thread callbacks, and non-RPC pools. Short cross-references are added at Messenger::ThreadPool()/ThreadPoolPtr(), the local-call fast path, and TabletPeer::Enqueue, and the rpc_priority_queue_enabled help text states the scope so operators know what the budget does and does not bound. Also folds one over-long line in rpc_priority_queue-test.cc flagged by build-support/lint.py. Validation performed for this stage (no code changes resulted): - Flag-off audit of every non-test file against origin/master: the off path is plumbing plus one null-check per inbound call, per outbound-call construction and per shutdown; YBThreadPool::Enqueue was already virtual via TaskRecipient, so the ThreadPoolTaskRecipient switch adds no indirection; the ServiceUnavailable branch in ServicePoolImpl::Failure is unreachable with the flag off. - Flag on, release: raft_consensus-itest 19-test subset 19/19 with the queue active on all daemons; tablet_server-test, master-test and client-test pass except one environmental failure each (unassignable bind address or missing initial sys catalog snapshot), reproduced with the flag off. - Flag on, debug, callback contract check on: pg_mini-test 22-test subset 22/22, no blocking callback detected. - Final sweep: unit 24/24 debug, 22/22 release; rpc-test, rpc_stub-test, mt-rpc-test with the flag off and on; KVTable mini-cluster tests. --- _automated · pi (claude-fable-5)_
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Adds
RpcPriorityQueue: a single gating queue perMessengerthat bounds concurrently executing RPC work across all RPC worker pools with one permit budget and dispatches waiting work in strict priority order (kHigh>kNormal>kLow). Both inbound service calls and async outbound-call callbacks go through it. Flag-gated (rpc_priority_queue_enabled, NON_RUNTIME, default false); with the flag off, the dispatch path is unchanged.Motivation and design: yugabyte#33937.
This PR is an internal review vehicle for the full series; it will be closed without merge and the six commits submitted upstream as a PR stack. Review commit by commit:
RpcPriorityenum;RegisterService(..., ServicePriority, RpcPriority); per-service classification in tserver and master. No behavior change.RpcPriorityQueue(src/yb/rpc/rpc_priority_queue.{h,cc}): packed atomic state, lock-free fast path, per-priority bands, permit hand-off on completion, single-drainer completion loop, two-phase shutdown, pool-size invariant. Flags, metrics, unit tests. Not yet wired.Messengerowns the queue;ServicePoolImpl::Processroutes admitted calls through it; end-to-endrpc-test.KVTablePriorityQueueTestmini-cluster tests (budget 2, callback reserve 1) asserting queue creation, per-priority dispatch counters, callback dispatch and genuine contention.ProxyContext::CallbackRecipient;RpcTaskClassand callback reserve for deadlock freedom; enforced no-blocking callback contract (debug); bounded queued counter ->SERVER_TOO_BUSY; shutdown as a quiescence barrier with documented re-entrancy contract.TabletPeer::Enqueue, shared-memory exchange pool).Key properties reviewers may want to probe:
Done(); released on run, pool refusal and shutdown.max_dispatched <= pool.max_workersis DCHECKed so "dispatched" means "executing".budget - reserve; callbacks may use the whole budget and must not block (debug-enforced viaThreadRestrictions). Audited in debug: rpc suites, KVTable mini-cluster,pg_mini-test22-test subset, no violations.StartShutdowncloses via the state word and fails waiting tasks;CompleteShutdownwaits for drain, active enqueues, dispatched tasks and in-flight completions.Enqueuereleases its registration before running any task code. Task code may callStartShutdown; a synchronousCompleteShutdownfrom task code is DFATAL + ignored (same constraintYBThreadPoolhas).YBThreadPool::Enqueuewas already virtual so theThreadPoolTaskRecipientswitch adds no indirection.Upgrade/Rollback safety
Process-local; no wire-format, catalog or on-disk changes.
rpc_priority_queue_enabledis NON_RUNTIME and defaults to false, so upgraded nodes behave identically to today unless explicitly restarted with the flag. Mixed-version clusters are unaffected: a node with the flag on orders only its own worker dispatch; peers see no difference. Rollback is a restart without the flag (or downgrade); nothing persists.Test plan
rpc_priority_queue-test: 24/24 debug, 22/22 release (sync-point and death tests are debug-only)rpc-test,rpc_stub-test,mt-rpc-test: pass with the flag off and on, debug and releasekv_table-test:LoadTest,Restartand theirWithPriorityQueuevariants (budget 2, reserve 1,rpc_workers_limit=8), release; 5/5 repeatsraft_consensus-itest19-test subset 19/19;tablet_server-test35/36,master-test68/69,client-test67/68 (each single failure is environmental on the dev machine and reproduces with the flag off)pg_mini-test22-test subset 22/22build-support/lint.pyclean against master;git diff --checkcleanrpc_priority_queue-test,rpc-test,kv_table-test) on Linux; not available in the macOS thirdparty buildpg_libpq-test, DDL-heavy and CDC/xCluster suites before any default-on