|
@@ -40,13 +40,71 @@ WebhookController::WebhookController(storage::StorageClient& storage,
|
|
|
runners::LoadBalancer& load_balancer,
|
|
runners::LoadBalancer& load_balancer,
|
|
|
WebSocketServer& ws_server,
|
|
WebSocketServer& ws_server,
|
|
|
nodes::NodeStore& node_store,
|
|
nodes::NodeStore& node_store,
|
|
|
- WorkflowScheduler& scheduler)
|
|
|
|
|
|
|
+ WorkflowScheduler& scheduler,
|
|
|
|
|
+ DispatchConfig dispatch_config)
|
|
|
: storage_(storage)
|
|
: storage_(storage)
|
|
|
, registry_(registry)
|
|
, registry_(registry)
|
|
|
, load_balancer_(load_balancer)
|
|
, load_balancer_(load_balancer)
|
|
|
, ws_server_(ws_server)
|
|
, ws_server_(ws_server)
|
|
|
, node_store_(node_store)
|
|
, node_store_(node_store)
|
|
|
- , scheduler_(scheduler) {}
|
|
|
|
|
|
|
+ , scheduler_(scheduler)
|
|
|
|
|
+ , dispatch_queue_capacity_(dispatch_config.queue_capacity) {
|
|
|
|
|
+ dispatch_workers_.reserve(dispatch_config.threads);
|
|
|
|
|
+ for (std::size_t i = 0; i < dispatch_config.threads; ++i) {
|
|
|
|
|
+ dispatch_workers_.emplace_back(&WebhookController::dispatchWorkerLoop, this);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+WebhookController::~WebhookController() {
|
|
|
|
|
+ stop();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void WebhookController::stop() {
|
|
|
|
|
+ {
|
|
|
|
|
+ std::lock_guard<std::mutex> lock(dispatch_mutex_);
|
|
|
|
|
+ if (!dispatch_running_) {
|
|
|
|
|
+ return; // already stopped
|
|
|
|
|
+ }
|
|
|
|
|
+ dispatch_running_ = false;
|
|
|
|
|
+ }
|
|
|
|
|
+ dispatch_cv_.notify_all();
|
|
|
|
|
+ for (auto& worker : dispatch_workers_) {
|
|
|
|
|
+ if (worker.joinable()) {
|
|
|
|
|
+ worker.join();
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ dispatch_workers_.clear();
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+void WebhookController::dispatchWorkerLoop() {
|
|
|
|
|
+ while (true) {
|
|
|
|
|
+ std::function<void()> task;
|
|
|
|
|
+ {
|
|
|
|
|
+ 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()) {
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ task = std::move(dispatch_queue_.front());
|
|
|
|
|
+ dispatch_queue_.pop();
|
|
|
|
|
+ }
|
|
|
|
|
+ task();
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+bool WebhookController::tryEnqueueDispatch(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_cv_.notify_one();
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
|
|
|
void WebhookController::registerRoutes(httplib::Server& server) {
|
|
void WebhookController::registerRoutes(httplib::Server& server) {
|
|
|
// Match any HTTP method for webhooks
|
|
// Match any HTTP method for webhooks
|
|
@@ -210,15 +268,14 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
|
|
|
auto channel = grpc::CreateChannel(runner->address, grpc::InsecureChannelCredentials());
|
|
auto channel = grpc::CreateChannel(runner->address, grpc::InsecureChannelCredentials());
|
|
|
auto stub = proto::RunnerService::NewStub(channel);
|
|
auto stub = proto::RunnerService::NewStub(channel);
|
|
|
|
|
|
|
|
- // A schedule's overlap policy counts what is running, and this call blocks
|
|
|
|
|
- // until the run has finished - so by the time it returns with an execution
|
|
|
|
|
- // id there is nothing left to claim a slot for, and a scheduled tick in the
|
|
|
|
|
- // meantime saw the workflow as idle. Naming the id here rather than letting
|
|
|
|
|
- // the runner pick one is what makes the slot claimable up front; the runner
|
|
|
|
|
- // takes this as the run's id.
|
|
|
|
|
|
|
+ // A schedule's overlap policy counts what is running, and a synchronous
|
|
|
|
|
+ // dispatch blocks until the run has finished - so by the time it returns
|
|
|
|
|
+ // with an execution id there is nothing left to claim a slot for, and a
|
|
|
|
|
+ // scheduled tick in the meantime saw the workflow as idle. Naming the id
|
|
|
|
|
+ // here rather than letting the runner pick one is what makes the slot
|
|
|
|
|
+ // claimable up front; the runner takes this as the run's id.
|
|
|
const std::string execution_id = UUID::generatePrefixed("exec");
|
|
const std::string execution_id = UUID::generatePrefixed("exec");
|
|
|
trigger_data["_assignedExecutionId"] = execution_id;
|
|
trigger_data["_assignedExecutionId"] = execution_id;
|
|
|
- scheduler_.notifyExecutionStarted(workflow_id, execution_id);
|
|
|
|
|
|
|
|
|
|
proto::ExecuteWorkflowRequest grpc_req;
|
|
proto::ExecuteWorkflowRequest grpc_req;
|
|
|
grpc_req.set_workflow_id(workflow_id);
|
|
grpc_req.set_workflow_id(workflow_id);
|
|
@@ -226,30 +283,47 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
|
|
|
// Fired by the outside world, so it runs the published version.
|
|
// Fired by the outside world, so it runs the published version.
|
|
|
grpc_req.set_use_published(true);
|
|
grpc_req.set_use_published(true);
|
|
|
grpc_req.set_trigger_data(trigger_data.dump());
|
|
grpc_req.set_trigger_data(trigger_data.dump());
|
|
|
- // A form set to answer immediately must not hold the browser for the run.
|
|
|
|
|
- // The slot claimed above is still released correctly: the runner reports
|
|
|
|
|
- // execution.completed/failed/cancelled/waiting to
|
|
|
|
|
- // /api/v1/internal/execution-event, and execution_controller.cpp:632
|
|
|
|
|
- // releases it there - the same path fire-and-forget scheduled dispatch
|
|
|
|
|
- // relies on today.
|
|
|
|
|
const bool wait_for_run = form_trigger_data.is_null() || form_response_mode != "immediate";
|
|
const bool wait_for_run = form_trigger_data.is_null() || form_response_mode != "immediate";
|
|
|
grpc_req.set_wait_for_completion(wait_for_run);
|
|
grpc_req.set_wait_for_completion(wait_for_run);
|
|
|
grpc_req.set_timeout_ms(wait_for_run ? 30000 : 0); // 30 second timeout when waiting
|
|
grpc_req.set_timeout_ms(wait_for_run ? 30000 : 0); // 30 second timeout when waiting
|
|
|
|
|
|
|
|
if (!wait_for_run) {
|
|
if (!wait_for_run) {
|
|
|
// ExecuteWorkflow is a plain synchronous RPC: the server call runs the
|
|
// ExecuteWorkflow is a plain synchronous RPC: the server call runs the
|
|
|
- // workflow to completion (or until the client's deadline lapses) before
|
|
|
|
|
- // it returns anything, regardless of wait_for_completion - that field
|
|
|
|
|
- // only controls whether the response carries a result. Holding this
|
|
|
|
|
- // handler's thread on that call would hold the browser for the run no
|
|
|
|
|
- // matter how the request is configured, which is exactly what
|
|
|
|
|
- // "immediate" promises not to do. So the call is moved to a detached
|
|
|
|
|
- // thread: the browser gets the thank-you page right away, and the run
|
|
|
|
|
- // proceeds in the background. The scheduler slot claimed above is
|
|
|
|
|
- // still released correctly, the same way a scheduled dispatch's is -
|
|
|
|
|
- // via the runner's execution.completed/failed/cancelled/waiting
|
|
|
|
|
- // report to /api/v1/internal/execution-event.
|
|
|
|
|
- std::thread([this, channel, workflow_id, execution_id, grpc_req]() {
|
|
|
|
|
|
|
+ // workflow to completion (or until the client's deadline lapses)
|
|
|
|
|
+ // before it returns anything - wait_for_completion only controls
|
|
|
|
|
+ // whether the response carries a result, it does not make the server
|
|
|
|
|
+ // return early. Holding this handler's thread on that call would hold
|
|
|
|
|
+ // the browser for the run no matter how the request is configured,
|
|
|
|
|
+ // which is exactly what "immediate" promises not to do. So the call
|
|
|
|
|
+ // runs on the controller's bounded dispatch pool instead of this
|
|
|
|
|
+ // thread or a detached one of its own:
|
|
|
|
|
+ // - a detached thread has no owner, so nothing joins it, and it
|
|
|
|
|
+ // would still be running against `this` after the controller
|
|
|
|
|
+ // that it captured is destroyed - an ordinary
|
|
|
|
|
+ // `systemctl --user restart` while a run is in flight would be a
|
|
|
|
|
+ // use-after-free;
|
|
|
|
|
+ // - the form endpoint is public and unauthenticated by design, so
|
|
|
|
|
+ // an unbounded "one thread per submission" scheme is a resource
|
|
|
|
|
+ // exhaustion vector open to anyone who can reach it.
|
|
|
|
|
+ // Bounded and joined in stop()/~WebhookController fixes both: no
|
|
|
|
|
+ // task can still be running once the controller that owns the pool
|
|
|
|
|
+ // is gone, and a burst of submissions is throttled rather than
|
|
|
|
|
+ // spawning without limit.
|
|
|
|
|
+ //
|
|
|
|
|
+ // The scheduler slot is claimed inside the queued task, not here, so
|
|
|
|
|
+ // a submission that is refused for a full queue never claims one -
|
|
|
|
|
+ // claiming it up front and then refusing would leak it. Once
|
|
|
|
|
+ // claimed, it is released the same way a scheduled dispatch's is:
|
|
|
|
|
+ // the runner reports execution.completed/failed/cancelled/waiting to
|
|
|
|
|
+ // /api/v1/internal/execution-event, and execution_controller.cpp:632
|
|
|
|
|
+ // releases it there. That event-based release is the only part this
|
|
|
|
|
+ // shares with scheduled dispatch - scheduled dispatch blocks the
|
|
|
|
|
+ // 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]() {
|
|
|
|
|
+ scheduler_.notifyExecutionStarted(workflow_id, execution_id);
|
|
|
|
|
+
|
|
|
auto bg_stub = proto::RunnerService::NewStub(channel);
|
|
auto bg_stub = proto::RunnerService::NewStub(channel);
|
|
|
proto::ExecuteWorkflowResponse bg_res;
|
|
proto::ExecuteWorkflowResponse bg_res;
|
|
|
grpc::ClientContext bg_ctx;
|
|
grpc::ClientContext bg_ctx;
|
|
@@ -275,7 +349,21 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
|
|
|
scheduler_.notifyExecutionFinished(workflow_id, execution_id);
|
|
scheduler_.notifyExecutionFinished(workflow_id, execution_id);
|
|
|
scheduler_.notifyExecutionStarted(workflow_id, bg_res.execution_id());
|
|
scheduler_.notifyExecutionStarted(workflow_id, bg_res.execution_id());
|
|
|
}
|
|
}
|
|
|
- }).detach();
|
|
|
|
|
|
|
+ });
|
|
|
|
|
+
|
|
|
|
|
+ if (!accepted) {
|
|
|
|
|
+ // The queue is full: refuse honestly rather than tell someone
|
|
|
|
|
+ // "thanks, received" for a submission that will never run, and
|
|
|
|
|
+ // rather than silently drop it. No slot was claimed for this one.
|
|
|
|
|
+ LOG_WARN("Immediate-mode dispatch queue is full, refusing submission for workflow {}",
|
|
|
|
|
+ workflow_id);
|
|
|
|
|
+ res.status = 503;
|
|
|
|
|
+ res.set_content(form_renderer::renderMessage(
|
|
|
|
|
+ form_config_title,
|
|
|
|
|
+ "This form is busy right now. Please try again in a moment."),
|
|
|
|
|
+ "text/html; charset=utf-8");
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
res.status = 200;
|
|
res.status = 200;
|
|
|
res.set_content(form_renderer::renderMessage(form_config_title, form_response_message),
|
|
res.set_content(form_renderer::renderMessage(form_config_title, form_response_message),
|
|
@@ -283,6 +371,8 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ scheduler_.notifyExecutionStarted(workflow_id, execution_id);
|
|
|
|
|
+
|
|
|
proto::ExecuteWorkflowResponse grpc_res;
|
|
proto::ExecuteWorkflowResponse grpc_res;
|
|
|
grpc::ClientContext grpc_ctx;
|
|
grpc::ClientContext grpc_ctx;
|
|
|
grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(35));
|
|
grpc_ctx.set_deadline(std::chrono::system_clock::now() + std::chrono::seconds(35));
|