Explorar o código

fix: close the shutdown race between a queued dispatch and stop()'s cancel pass

stop() used to sleep a flat 2 seconds, then cancel whatever was registered
in active_tasks_ at that instant. A worker that had already popped its task
off the queue but had not yet reached registerActiveTask (called right
before the blocking gRPC call) was invisible to that pass. If it registered
after the 2 second window - scheduler mutex contention, a loaded machine -
nothing cancelled it, and stop() joined the worker against its context's
hour-long deadline instead. systemd SIGKILLs the unit at 90 seconds, which
is exactly what an earlier fix round was meant to prevent.

Fixed deterministically instead of by lengthening the sleep: a shutting_down_
flag is set under active_mutex_ at the start of the cancellation pass, and
registerActiveTask checks it under the same lock, calling TryCancel()
immediately if it is already true. A task can no longer register invisibly
during shutdown, however late it arrives. The flat sleep is replaced with a
condition variable wait bounded by the same 2 second grace period, with its
predicate satisfied the moment every active task unregisters - so an idle
shutdown returns immediately instead of always paying the full 2 seconds
(the M9/deferred-minor this also fixes).

Measured live (kill -TERM, signal-received to service-stopped):
- idle, nothing running: ~0.3-1.1s (previously a flat 2s)
- a long run in flight (immediate-mode form, 10-20s workflow): ~0.6-0.7s,
  TryCancel logged and the client call unblocks immediately even though the
  workflow keeps running server-side
- queue full plus 4 in-flight: ~0.3s, 2 queued submissions discarded and
  logged, 4 in-flight cancelled

Also corrects stop()'s header comment (M9): safe to call more than once
sequentially, not concurrently - the join/clear sequence has no lock
protecting it from a second concurrent caller.
fszontagh hai 1 mes
pai
achega
2d14a8e2ca

+ 33 - 9
src/webserver/api/webhook_controller.cpp

@@ -94,22 +94,35 @@ void WebhookController::stop() {
 
     // A worker that already popped a task is running it outside this lock,
     // so notify_all above does not reach it - dispatchWorkerLoop only checks
-    // dispatch_running_ between tasks. Give it a moment to reach
+    // dispatch_running_ between tasks. That worker still has to reach
     // registerActiveTask (called right before the blocking gRPC call starts)
-    // before looking for it below: 2 seconds is far more than that handful of
-    // synchronous calls needs, comfortably inside systemd's 90 second default
-    // TimeoutStopSec with room to spare for TryCancel to take effect and the
-    // join after it, and short enough that an ordinary restart with nothing
-    // in flight barely notices it.
-    std::this_thread::sleep_for(std::chrono::seconds(2));
-
+    // before it becomes visible here.
     {
-        std::lock_guard<std::mutex> lock(active_mutex_);
+        std::unique_lock<std::mutex> lock(active_mutex_);
+        // Setting this under the same lock registerActiveTask takes is what
+        // makes the race deterministic rather than timing-dependent: a task
+        // that registers after this point - however late, no matter how
+        // loaded the machine is - sees shutting_down_ already true and
+        // cancels itself right there (see registerActiveTask), instead of
+        // depending on a cancellation pass here having already run by the
+        // time it arrives.
+        shutting_down_ = true;
         for (const auto& active : active_tasks_) {
             LOG_WARN("Cancelling in-flight immediate-mode form dispatch for workflow {} at shutdown",
                      active.workflow_id);
             active.context->TryCancel();
         }
+        // Wait for every active task to unregister (its worker returned from
+        // the gRPC call and is on its way back to dispatchWorkerLoop, which
+        // exits immediately once dispatch_running_ is false), rather than
+        // sleeping a flat 2 seconds regardless of whether anything was
+        // running at all. Still bounded by the same grace period as before -
+        // TryCancel is expected to unblock everything well inside it - so a
+        // pathological case does not turn shutdown unbounded, it just stops
+        // waiting and lets the join below find out how long it actually
+        // takes.
+        active_cv_.wait_for(lock, std::chrono::seconds(2),
+                            [this] { return active_tasks_.empty(); });
     }
 
     // TryCancel unblocks this process's own client call quickly regardless
@@ -165,11 +178,22 @@ bool WebhookController::tryEnqueueDispatch(const std::string& workflow_id, std::
 void WebhookController::registerActiveTask(::grpc::ClientContext* context, const std::string& workflow_id) {
     std::lock_guard<std::mutex> lock(active_mutex_);
     active_tasks_.push_back(ActiveTask{context, workflow_id});
+    if (shutting_down_) {
+        // stop()'s own cancellation pass already ran by the time this task
+        // reached here - it would otherwise sit uncancelled until the
+        // bounded wait in stop() gives up, or worse, until the hour-long
+        // deadline set on this context. Cancel it the moment it becomes
+        // visible instead.
+        LOG_WARN("Cancelling in-flight immediate-mode form dispatch for workflow {} at shutdown "
+                 "(registered after shutdown began)", workflow_id);
+        context->TryCancel();
+    }
 }
 
 void WebhookController::unregisterActiveTask(::grpc::ClientContext* context) {
     std::lock_guard<std::mutex> lock(active_mutex_);
     std::erase_if(active_tasks_, [context](const ActiveTask& t) { return t.context == context; });
+    active_cv_.notify_all();
 }
 
 void WebhookController::registerRoutes(httplib::Server& server) {

+ 21 - 3
src/webserver/api/webhook_controller.hpp

@@ -50,7 +50,10 @@ public:
     // Stops accepting new immediate-mode dispatches and joins every worker.
     // Must run before the controller itself is destroyed - a worker task
     // captures `this` and calling into a freed controller would be a
-    // use-after-free. Safe to call more than once.
+    // use-after-free. Safe to call more than once sequentially - it is not
+    // safe to call concurrently from two threads at once, since the running
+    // task/join sequence has no lock protecting it from another stop()
+    // clearing dispatch_workers_ or racing the same active_tasks_ scan.
     void stop();
 
     // A signed token rather than the password itself: the password never
@@ -88,8 +91,12 @@ private:
     // A task that has been popped off the queue but has not yet registered
     // its gRPC context (registerActiveTask has not run) has already left the
     // queue by the time stop() drains it, but is not yet visible to the
-    // cancellation pass either - the grace period in stop() exists so this
-    // task reaches registerActiveTask before stop() looks for it.
+    // cancellation pass either. registerActiveTask checks shutting_down_
+    // under the same lock stop() sets it with, so this task cancels itself
+    // the moment it registers rather than depending on stop() having already
+    // run its pass - deterministic regardless of how long the task took to
+    // get here, not bounded by the grace period stop() waits before giving
+    // up on the join.
     void registerActiveTask(::grpc::ClientContext* context, const std::string& workflow_id);
     void unregisterActiveTask(::grpc::ClientContext* context);
 
@@ -140,6 +147,17 @@ private:
     // worker is not yet done draining.
     std::mutex active_mutex_;
     std::vector<ActiveTask> active_tasks_;
+    // Notified whenever active_tasks_ changes, so stop() can wait for
+    // "no active tasks" instead of a fixed sleep - fast when shutdown is
+    // idle, and still bounded by the same grace period when it is not.
+    std::condition_variable active_cv_;
+    // Set once stop() begins, under active_mutex_. A task that pops off the
+    // dispatch queue before dispatch_running_ goes false but has not yet
+    // reached registerActiveTask (see the comment there) checks this and
+    // self-cancels the moment it registers, instead of depending on stop()'s
+    // cancellation pass having already run by the time it gets there - which
+    // a loaded machine could outrun.
+    bool shutting_down_ = false;
 };
 
 } // namespace smartbotic::webserver::api