From b585152fcf40e2150527a8d344f92528ca0d6339 Mon Sep 17 00:00:00 2001
From: Patrick O'Reilly
Date: Thu, 17 Sep 2026 20:09:42 -0700
Subject: [PATCH] indy: make a stalled replication join name the loop it is
waiting for (#564)
GracefulShutdownUnresponsiveSlave wedges on aarch64 about one CI run in
eighteen. The wedge is in JoinReplicationServices(), which pops a blocking
semaphore once per replication loop that started -- so one loop that never
returns hangs shutdown forever, silently. The only syslog between the
connection drain and "TServer::Shutdown() complete" is the drain's own
deadline warning, so from the log a wedge here is indistinguishable from a
wedge anywhere else in Shutdown(), and the process is dead by the time anyone
reads CI.
83 attempts across four harnesses failed to reproduce it locally or on demand
in CI, so waiting for a reproduction is not a plan. Make the first natural
occurrence diagnose itself instead.
Each loop now sets a bit on entry and clears it on exit, alongside the
existing started-count and Exited push -- one RAII latch doing all three, so
every exit path stays covered including the exceptional ones the loops rely
on. The join polls the semaphore's fd rather than blocking on it outright,
and every 15s of no progress logs which loops have still not returned.
The wait stays unbounded on purpose. The teardown that follows frees the very
fds and collections a still-parked loop references (#440); giving up would
trade a hang for a use-after-free. Nothing about the shutdown contract
changes -- this only removes the silence.
Verified by injecting a 35s stall into RunReplicationWork: the report fires at
15s and 30s naming [RunReplicationWork], shutdown still completes afterwards,
and the fixture still passes 2/2. Without the injection there are no warnings
at all and both fixtures pass, so a healthy shutdown stays quiet.
---
changelog.d/564-name-the-wedged-loop.md | 1 +
orly/indy/manager.cc | 102 ++++++++++++++++++++++--
orly/indy/manager.h | 17 ++++
3 files changed, 114 insertions(+), 6 deletions(-)
create mode 100644 changelog.d/564-name-the-wedged-loop.md
diff --git a/changelog.d/564-name-the-wedged-loop.md b/changelog.d/564-name-the-wedged-loop.md
new file mode 100644
index 00000000..72215d18
--- /dev/null
+++ b/changelog.d/564-name-the-wedged-loop.md
@@ -0,0 +1 @@
+- **Changed**: a graceful shutdown that stalls in `TManager::JoinReplicationServices()` now says which replication loop it is waiting for, every 15s, instead of going silent. The wait itself is deliberately still unbounded -- the manager teardown that follows frees fds and collections a still-parked loop references (#440), so giving up would trade a hang for a use-after-free -- but it is no longer undiagnosable. This is where aarch64 shutdowns wedge (#564), and the only syslog between the connection drain and `TServer::Shutdown() complete` was the drain's own deadline warning, so a wedge there looked identical to a wedge anywhere else in `Shutdown()`. Each loop now carries a running-bit alongside the existing started/exited accounting, so the report names `RunReplicationQueue`, `RunReplicationWork` or `RunReplicateTransaction` specifically (#564).
diff --git a/orly/indy/manager.cc b/orly/indy/manager.cc
index fe7d58c5..c1e525b8 100644
--- a/orly/indy/manager.cc
+++ b/orly/indy/manager.cc
@@ -16,6 +16,8 @@
See the License for the specific language governing permissions and
limitations under the License. */
+#include
+
#include
#include
@@ -183,9 +185,42 @@ L0::TManager::TPtr TManager::NewFastRepo(const TUuid &repo_id,
return New(repo_id, *ttl, parent_repo, false);
}
+namespace {
+
+ /* Counts a replication loop in on entry and out on exit -- the Started
+ count and the Exited push that JoinReplicationServices() reaps (#440),
+ plus the running-bit that lets a stalled join name the loop it is still
+ waiting for (#564). RAII so that every exit path is covered, including
+ the exceptional ones the loops rely on. */
+ class TServiceLatch final {
+ NO_COPY(TServiceLatch);
+ public:
+
+ TServiceLatch(std::atomic &started, std::atomic &running,
+ Base::TEventSemaphore &exited, unsigned bit)
+ : Running(running), Exited(exited), Bit(bit) {
+ ++started;
+ Running.fetch_or(Bit);
+ }
+
+ ~TServiceLatch() {
+ Running.fetch_and(~Bit);
+ Exited.Push();
+ }
+
+ private:
+
+ std::atomic &Running;
+ Base::TEventSemaphore &Exited;
+ const unsigned Bit;
+
+ }; // TServiceLatch
+
+} // namespace
+
void TManager::RunReplicationQueue() {
- ++ReplicationServicesStarted;
- Base::TPushOnExit exit_latch(ReplicationServicesExited);
+ TServiceLatch exit_latch(ReplicationServicesStarted, ReplicationServicesRunning,
+ ReplicationServicesExited, ReplicationQueueService);
try {
epoll_event event;
int timeout = -1;
@@ -258,8 +293,8 @@ void TManager::RunReplicationQueue() {
}
void TManager::RunReplicationWork() {
- ++ReplicationServicesStarted;
- Base::TPushOnExit exit_latch(ReplicationServicesExited);
+ TServiceLatch exit_latch(ReplicationServicesStarted, ReplicationServicesRunning,
+ ReplicationServicesExited, ReplicationWorkService);
try {
epoll_event event;
int timeout = -1;
@@ -311,8 +346,8 @@ void TManager::RunReplicationWork() {
}
void TManager::RunReplicateTransaction() {
- ++ReplicationServicesStarted;
- Base::TPushOnExit exit_latch(ReplicationServicesExited);
+ TServiceLatch exit_latch(ReplicationServicesStarted, ReplicationServicesRunning,
+ ReplicationServicesExited, ReplicateTransactionService);
epoll_event event;
int timeout = -1;
void *state_alloc = alloca(Sabot::State::GetMaxStateSize());
@@ -560,12 +595,67 @@ void TManager::StopReplicationServices() {
}
}
+std::string TManager::DescribeReplicationServices(unsigned mask) {
+ static const std::pair names[] = {
+ {ReplicationQueueService, "RunReplicationQueue"},
+ {ReplicationWorkService, "RunReplicationWork"},
+ {ReplicateTransactionService, "RunReplicateTransaction"}
+ };
+ std::string ret;
+ for (const auto &item : names) {
+ if (mask & item.first) {
+ if (!ret.empty()) {
+ ret += ", ";
+ }
+ ret += item.second;
+ }
+ }
+ return ret.empty() ? std::string("none") : ret;
+}
+
void TManager::JoinReplicationServices() {
/* Reap exactly as many exits as loops that actually entered; a loop
whose fiber never got to run can't be waited for (and never touches
us). Every entered loop pushes Exited on its way out, including the
exception paths. */
+
+ /* The wait itself stays unbounded ON PURPOSE -- the manager teardown that
+ follows frees the fds and collections a still-parked loop references
+ (#440), so giving up here would trade a hang for a use-after-free. What
+ changes is that it stops being silent.
+
+ This join is where aarch64 shutdowns wedge (#564), and until now the log
+ simply stopped: the only syslog between the connection drain and
+ "TServer::Shutdown() complete" is the drain's own deadline warning, so a
+ wedge here was indistinguishable from a wedge anywhere else in Shutdown().
+ Poll the exit semaphore's fd instead of blocking on it outright, and every
+ time the poll expires say which loops have still not returned. A failure
+ now names its own culprit on the first occurrence, with no reproduction
+ needed -- which matters for a bug seen roughly once in eighteen CI runs
+ and never once locally in 83 attempts. */
+ static const int report_every_ms = 15000;
for (size_t n = ReplicationServicesStarted; n > 0; --n) {
+ for (;;) {
+ pollfd waiting;
+ Base::Zero(waiting);
+ waiting.fd = ReplicationServicesExited.GetFd();
+ waiting.events = POLLIN;
+ const int ret = poll(&waiting, 1, report_every_ms);
+ if (ret < 0) {
+ if (errno == EINTR) {
+ continue;
+ }
+ ::Util::ThrowSystemError(errno);
+ }
+ if (ret > 0) {
+ break;
+ }
+ syslog(LOG_WARNING,
+ "TManager::JoinReplicationServices() still waiting after %dms;"
+ " replication loop(s) that have not returned: [%s] (#564)",
+ report_every_ms,
+ DescribeReplicationServices(ReplicationServicesRunning.load()).c_str());
+ }
ReplicationServicesExited.Pop();
}
}
diff --git a/orly/indy/manager.h b/orly/indy/manager.h
index b2d245da..d2be70a9 100644
--- a/orly/indy/manager.h
+++ b/orly/indy/manager.h
@@ -541,6 +541,23 @@ namespace Orly {
reaps exactly that many exits (#440). */
std::atomic ReplicationServicesStarted;
Base::TEventSemaphore ReplicationServicesExited;
+
+ /* Which replication loops are inside their bodies right now, as a bit
+ per service. Purely so that a shutdown which does not finish can say
+ WHICH loop it is still waiting for: JoinReplicationServices() is an
+ unbounded wait, and when it hangs the log simply stops, with no way
+ to tell the three loops apart afterwards (#564). Set alongside
+ Started, cleared alongside the Exited push, by TServiceLatch. */
+ enum TReplicationService : unsigned {
+ ReplicationQueueService = 1u << 0,
+ ReplicationWorkService = 1u << 1,
+ ReplicateTransactionService = 1u << 2
+ };
+ std::atomic ReplicationServicesRunning{0};
+
+ /* Renders ReplicationServicesRunning as service names, for the log
+ line above. Returns "none" when the mask is empty. */
+ static std::string DescribeReplicationServices(unsigned mask);
std::chrono::steady_clock::time_point ReplicationNextTime;
std::chrono::milliseconds ReplicationDelay;