From 380b60f1559eda1ed7ba03aea94ef6cac0aed2e8 Mon Sep 17 00:00:00 2001 From: Nicholas Dudfield Date: Fri, 18 Sep 2026 11:16:16 +0700 Subject: [PATCH] feat(jobqueue): claim jobs before the no-threads assert DispatchHook is empty in production. claimedQueued/claimedDropped skip the worker pool so a stepping harness can run with zero JobQueue threads. --- src/xrpld/core/JobQueue.h | 22 ++++++++++++++++++++++ src/xrpld/core/detail/JobQueue.cpp | 20 ++++++++++++++++++++ 2 files changed, 42 insertions(+) diff --git a/src/xrpld/core/JobQueue.h b/src/xrpld/core/JobQueue.h index 06f90dd1c4..c2370d27c0 100644 --- a/src/xrpld/core/JobQueue.h +++ b/src/xrpld/core/JobQueue.h @@ -29,6 +29,7 @@ #include #include // workaround for boost 1.72 bug #include // workaround for boost 1.72 bug +#include namespace ripple { @@ -140,6 +141,14 @@ public: using JobFunction = std::function; + // pass: queue on workers. claimedQueued/claimedDropped: the hook owns the + // job (enqueue on a harness scheduler, or drop). Consulted before the + // no-threads assert so a claimed job needs no workers. Empty hook = prod. + enum class JobDisposition { pass, claimedQueued, claimedDropped }; + + using DispatchHook = std::function< + JobDisposition(JobType, std::string const&, JobFunction const&)>; + JobQueue( int threadCount, beast::insight::Collector::ptr const& collector, @@ -173,6 +182,13 @@ public: return false; } + /** Install a dispatch hook. Call before job flow and do not mutate later. */ + void + setDispatchHook(DispatchHook hook) + { + dispatchHook_ = std::move(hook); + } + /** Creates a coroutine and adds a job to the queue which will run it. @param t The type of job. @@ -223,6 +239,10 @@ public: void rendezvous(); + /** True when rendezvous() would not block. */ + [[nodiscard]] bool + isIdle() const; + void stop(); @@ -247,6 +267,8 @@ private: std::uint64_t m_lastJob; std::set m_jobSet; JobCounter jobCounter_; + // Empty in production. Read without a lock in addRefCountedJob. + DispatchHook dispatchHook_; std::atomic_bool stopping_{false}; std::atomic_bool stopped_{false}; JobDataMap m_jobData; diff --git a/src/xrpld/core/detail/JobQueue.cpp b/src/xrpld/core/detail/JobQueue.cpp index 5eb1f24d43..c8b094079d 100644 --- a/src/xrpld/core/detail/JobQueue.cpp +++ b/src/xrpld/core/detail/JobQueue.cpp @@ -97,6 +97,19 @@ JobQueue::addRefCountedJob( << __func__ << " : Adding job : " << name << " : " << type; JobTypeData& data(iter->second); + if (dispatchHook_) + { + switch (dispatchHook_(type, name, func)) + { + case JobDisposition::claimedQueued: + return true; + case JobDisposition::claimedDropped: + return false; + case JobDisposition::pass: + break; + } + } + // FIXME: Workaround incorrect client shutdown ordering // do not add jobs to a queue with no threads XRPL_ASSERT( @@ -274,6 +287,13 @@ JobQueue::rendezvous() cv_.wait(lock, [this] { return m_processCount == 0 && m_jobSet.empty(); }); } +bool +JobQueue::isIdle() const +{ + std::lock_guard lock(m_mutex); + return m_processCount == 0 && m_jobSet.empty(); +} + JobTypeData& JobQueue::getJobTypeData(JobType type) {