Skip to content

[#33937] DocDB: Priority queueing for RPC worker dispatch (RpcPriorityQueue) - #25

Draft
craigsoules wants to merge 6 commits into
masterfrom
craig/rpc-priority-queueing
Draft

craigsoules wants to merge 6 commits into
masterfrom
craig/rpc-priority-queueing

Conversation

@craigsoules

Copy link
Copy Markdown

Summary

Adds RpcPriorityQueue: a single gating queue per Messenger that 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:

Commit Content
Stage 0 RpcPriority enum; RegisterService(..., ServicePriority, RpcPriority); per-service classification in tserver and master. No behavior change.
Stage 1 Core 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.
Stage 2 Messenger owns the queue; ServicePoolImpl::Process routes admitted calls through it; end-to-end rpc-test.
Stage 2 tests KVTablePriorityQueueTest mini-cluster tests (budget 2, callback reserve 1) asserting queue creation, per-priority dispatch counters, callback dispatch and genuine contention.
Stage 3 Async callbacks routed via ProxyContext::CallbackRecipient; RpcTaskClass and 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.
Stage 4 Documents what the queue does not gate in v1 (inline local calls, direct pool submission sites, TabletPeer::Enqueue, shared-memory exchange pool).

Key properties reviewers may want to probe:

  • Permit accounting: single release point in the wrapper task's Done(); released on run, pool refusal and shutdown. max_dispatched <= pool.max_workers is DCHECKed so "dispatched" means "executing".
  • Deadlock freedom: inbound handlers capped at budget - reserve; callbacks may use the whole budget and must not block (debug-enforced via ThreadRestrictions). Audited in debug: rpc suites, KVTable mini-cluster, pg_mini-test 22-test subset, no violations.
  • Shutdown: StartShutdown closes via the state word and fails waiting tasks; CompleteShutdown waits for drain, active enqueues, dispatched tasks and in-flight completions. Enqueue releases its registration before running any task code. Task code may call StartShutdown; a synchronous CompleteShutdown from task code is DFATAL + ignored (same constraint YBThreadPool has).
  • Off path: one null-check per inbound call, per outbound-call construction and per shutdown; YBThreadPool::Enqueue was already virtual so the ThreadPoolTaskRecipient switch adds no indirection.

Upgrade/Rollback safety

Process-local; no wire-format, catalog or on-disk changes. rpc_priority_queue_enabled is 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 release
  • kv_table-test: LoadTest, Restart and their WithPriorityQueue variants (budget 2, reserve 1, rpc_workers_limit=8), release; 5/5 repeats
  • Flag on, release, queue active on all daemons: raft_consensus-itest 19-test subset 19/19; tablet_server-test 35/36, master-test 68/69, client-test 67/68 (each single failure is environmental on the dev machine and reproduces with the flag off)
  • Flag on, debug, callback contract check on: pg_mini-test 22-test subset 22/22
  • Flag-off audit of every non-test file against master
  • build-support/lint.py clean against master; git diff --check clean
  • TSAN (rpc_priority_queue-test, rpc-test, kv_table-test) on Linux; not available in the macOS thirdparty build
  • Debug callback audit of pg_libpq-test, DDL-heavy and CDC/xCluster suites before any default-on

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)_
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