|
|
@@ -60,14 +60,62 @@ WebhookController::~WebhookController() {
|
|
|
}
|
|
|
|
|
|
void WebhookController::stop() {
|
|
|
+ std::vector<std::string> discarded_workflow_ids;
|
|
|
{
|
|
|
std::lock_guard<std::mutex> lock(dispatch_mutex_);
|
|
|
if (!dispatch_running_) {
|
|
|
return; // already stopped
|
|
|
}
|
|
|
dispatch_running_ = false;
|
|
|
+
|
|
|
+ // Whatever is still queued never started, so it never reached
|
|
|
+ // scheduler_.notifyExecutionStarted (that happens once a worker pops
|
|
|
+ // it, inside the task) - dropping it here leaks no slot. But every
|
|
|
+ // one of these submitters was already told "thanks, your answer was
|
|
|
+ // received", so silently discarding is exactly the kind of promise
|
|
|
+ // this project keeps getting bitten by breaking - log it instead of
|
|
|
+ // draining it. Draining it would mean shutdown waits for the entire
|
|
|
+ // backlog rather than just what is already running, which is what
|
|
|
+ // made the previous version of this able to block for hours.
|
|
|
+ while (!dispatch_queue_.empty()) {
|
|
|
+ discarded_workflow_ids.push_back(std::move(dispatch_queue_.front().workflow_id));
|
|
|
+ dispatch_queue_.pop();
|
|
|
+ }
|
|
|
+ }
|
|
|
+ if (!discarded_workflow_ids.empty()) {
|
|
|
+ LOG_WARN("Shutting down with {} queued immediate-mode form submission(s) never started; "
|
|
|
+ "discarded for workflows: {}",
|
|
|
+ discarded_workflow_ids.size(), StringUtils::join(discarded_workflow_ids, ", "));
|
|
|
}
|
|
|
dispatch_cv_.notify_all();
|
|
|
+
|
|
|
+ // 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
|
|
|
+ // 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));
|
|
|
+
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(active_mutex_);
|
|
|
+ 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();
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ // TryCancel unblocks this process's own client call quickly regardless
|
|
|
+ // of whether the runner notices - see the comment above registerActiveTask
|
|
|
+ // in the header, and the follow-up note in the task report about whether
|
|
|
+ // the runner treats a cancelled RPC as a reason to report
|
|
|
+ // execution.cancelled/.failed. Either way, nothing can still be running
|
|
|
+ // against this controller once join returns, which is the lifetime
|
|
|
+ // guarantee that matters here.
|
|
|
for (auto& worker : dispatch_workers_) {
|
|
|
if (worker.joinable()) {
|
|
|
worker.join();
|
|
|
@@ -78,34 +126,49 @@ void WebhookController::stop() {
|
|
|
|
|
|
void WebhookController::dispatchWorkerLoop() {
|
|
|
while (true) {
|
|
|
- std::function<void()> task;
|
|
|
+ QueuedDispatch item;
|
|
|
{
|
|
|
std::unique_lock<std::mutex> lock(dispatch_mutex_);
|
|
|
dispatch_cv_.wait(lock, [this] {
|
|
|
return !dispatch_running_ || !dispatch_queue_.empty();
|
|
|
});
|
|
|
- if (!dispatch_running_ && dispatch_queue_.empty()) {
|
|
|
+ // Stop pulling from the queue the moment shutdown starts, even if
|
|
|
+ // work remains - stop() itself discards and logs whatever is left
|
|
|
+ // rather than this loop draining it, which would mean shutdown
|
|
|
+ // waits for the whole backlog instead of only what already
|
|
|
+ // started.
|
|
|
+ if (!dispatch_running_) {
|
|
|
return;
|
|
|
}
|
|
|
- task = std::move(dispatch_queue_.front());
|
|
|
+ item = std::move(dispatch_queue_.front());
|
|
|
dispatch_queue_.pop();
|
|
|
}
|
|
|
- task();
|
|
|
+ item.task();
|
|
|
}
|
|
|
}
|
|
|
|
|
|
-bool WebhookController::tryEnqueueDispatch(std::function<void()> task) {
|
|
|
+bool WebhookController::tryEnqueueDispatch(const std::string& workflow_id, std::function<void()> task) {
|
|
|
{
|
|
|
std::lock_guard<std::mutex> lock(dispatch_mutex_);
|
|
|
if (!dispatch_running_ || dispatch_queue_.size() >= dispatch_queue_capacity_) {
|
|
|
return false;
|
|
|
}
|
|
|
- dispatch_queue_.push(std::move(task));
|
|
|
+ dispatch_queue_.push(QueuedDispatch{workflow_id, std::move(task)});
|
|
|
}
|
|
|
dispatch_cv_.notify_one();
|
|
|
return true;
|
|
|
}
|
|
|
|
|
|
+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});
|
|
|
+}
|
|
|
+
|
|
|
+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; });
|
|
|
+}
|
|
|
+
|
|
|
void WebhookController::registerRoutes(httplib::Server& server) {
|
|
|
// Match any HTTP method for webhooks
|
|
|
auto webhook_handler = [this](const httplib::Request& req, httplib::Response& res) {
|
|
|
@@ -321,20 +384,39 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
|
|
|
// scheduler's own single thread (joined in WorkflowScheduler::stop())
|
|
|
// and relies on the client-side deadline lapsing; this runs on a
|
|
|
// pool sized and owned for exactly this purpose.
|
|
|
- bool accepted = tryEnqueueDispatch([this, channel, workflow_id, execution_id, grpc_req]() {
|
|
|
+ bool accepted = tryEnqueueDispatch(workflow_id,
|
|
|
+ [this, channel, workflow_id, execution_id, grpc_req]() {
|
|
|
scheduler_.notifyExecutionStarted(workflow_id, execution_id);
|
|
|
|
|
|
auto bg_stub = proto::RunnerService::NewStub(channel);
|
|
|
proto::ExecuteWorkflowResponse bg_res;
|
|
|
grpc::ClientContext bg_ctx;
|
|
|
- // Nobody is waiting on this response, so the deadline only needs
|
|
|
- // to bound a runner that never answers at all.
|
|
|
+ // Nobody is waiting on this response under normal operation, so
|
|
|
+ // the deadline only needs to bound a runner that never answers at
|
|
|
+ // all. At shutdown, stop() calls TryCancel() on this context
|
|
|
+ // directly instead of waiting for the deadline - see stop() for
|
|
|
+ // why an hour-long deadline would otherwise make shutdown itself
|
|
|
+ // take up to an hour.
|
|
|
bg_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::hours(1));
|
|
|
|
|
|
+ registerActiveTask(&bg_ctx, workflow_id);
|
|
|
auto bg_status = bg_stub->ExecuteWorkflow(&bg_ctx, grpc_req, &bg_res);
|
|
|
+ unregisterActiveTask(&bg_ctx);
|
|
|
|
|
|
if (!bg_status.ok()) {
|
|
|
- if (bg_status.error_code() != grpc::StatusCode::DEADLINE_EXCEEDED) {
|
|
|
+ // Neither a lapsed deadline nor a shutdown-time TryCancel
|
|
|
+ // means the run stopped - both only end this process's own
|
|
|
+ // wait. The runner's ExecuteWorkflow handler never checks
|
|
|
+ // whether its caller cancelled (confirmed by reading
|
|
|
+ // runner_service.cpp: no ServerContext::IsCancelled() call
|
|
|
+ // anywhere), so the workflow keeps running server-side
|
|
|
+ // regardless and will still report its own
|
|
|
+ // execution.completed/failed/waiting event when it finishes -
|
|
|
+ // releasing the slot here would say idle while it is still
|
|
|
+ // working. Making the runner honor cancellation is a
|
|
|
+ // legitimate follow-up, left for its own change.
|
|
|
+ if (bg_status.error_code() != grpc::StatusCode::DEADLINE_EXCEEDED &&
|
|
|
+ bg_status.error_code() != grpc::StatusCode::CANCELLED) {
|
|
|
scheduler_.notifyExecutionFinished(workflow_id, execution_id);
|
|
|
}
|
|
|
LOG_ERROR("Immediate-mode webhook execution failed for workflow {}: {}",
|