|
@@ -130,6 +130,44 @@ bool ExecutionController::canSeeExecution(const auth::AuthContext& ctx,
|
|
|
auth::Action::Read);
|
|
auth::Action::Read);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+bool ExecutionController::restrictToReachableWorkflows(const auth::AuthContext& ctx,
|
|
|
|
|
+ storage::QueryOptions& options,
|
|
|
|
|
+ const std::string& wanted_project) {
|
|
|
|
|
+ const auto reachable = access_.projectsFor(ctx);
|
|
|
|
|
+
|
|
|
|
|
+ storage::QueryOptions wf_opts;
|
|
|
|
|
+ wf_opts.page_size = 1000;
|
|
|
|
|
+ // Only the two fields the loop below reads. The projection is applied
|
|
|
|
|
+ // client-side by the adapter, so this does not shrink the transfer, but it
|
|
|
|
|
+ // keeps the intent honest and starts paying the moment upstream projection
|
|
|
|
|
+ // arrives.
|
|
|
|
|
+ wf_opts.fields = {"projectId"};
|
|
|
|
|
+ if (!wanted_project.empty()) {
|
|
|
|
|
+ wf_opts.filters.push_back({"projectId", wanted_project});
|
|
|
|
|
+ }
|
|
|
|
|
+ auto workflows = storage_.query("workflows", wf_opts);
|
|
|
|
|
+
|
|
|
|
|
+ std::vector<std::string> ids;
|
|
|
|
|
+ if (workflows.ok()) {
|
|
|
|
|
+ for (const auto& wf : workflows.value().documents) {
|
|
|
|
|
+ const std::string project = auth::AccessControl::projectOf(wf);
|
|
|
|
|
+ if (!reachable.contains(project)) continue;
|
|
|
|
|
+ if (!wanted_project.empty() && project != wanted_project) continue;
|
|
|
|
|
+ ids.push_back(wf.value("_id", ""));
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (ids.empty()) {
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (ids.size() == 1) {
|
|
|
|
|
+ options.filters.push_back({"workflowId", ids.front()});
|
|
|
|
|
+ } else {
|
|
|
|
|
+ options.filters.push_back({"workflowId", nlohmann::json{{"op", "in"}, {"value", ids}}});
|
|
|
|
|
+ }
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
void ExecutionController::registerRoutes(httplib::Server& server) {
|
|
void ExecutionController::registerRoutes(httplib::Server& server) {
|
|
|
server.Get("/api/v1/executions", [this](const httplib::Request& req, httplib::Response& res) {
|
|
server.Get("/api/v1/executions", [this](const httplib::Request& req, httplib::Response& res) {
|
|
|
middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
|
|
middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
|
|
@@ -193,40 +231,13 @@ void ExecutionController::listExecutions(const httplib::Request& req, httplib::R
|
|
|
// workflows in projects the caller cannot open. The workflow endpoints have
|
|
// workflows in projects the caller cannot open. The workflow endpoints have
|
|
|
// always refused those; their execution history did not.
|
|
// always refused those; their execution history did not.
|
|
|
if (!req.has_param("workflowId")) {
|
|
if (!req.has_param("workflowId")) {
|
|
|
- const auto reachable = access_.projectsFor(ctx);
|
|
|
|
|
- std::string wanted_project;
|
|
|
|
|
- if (req.has_param("projectId")) wanted_project = req.get_param_value("projectId");
|
|
|
|
|
-
|
|
|
|
|
- storage::QueryOptions wf_opts;
|
|
|
|
|
- wf_opts.page_size = 1000;
|
|
|
|
|
- if (!wanted_project.empty()) {
|
|
|
|
|
- wf_opts.filters.push_back({"projectId", wanted_project});
|
|
|
|
|
- }
|
|
|
|
|
- auto workflows = storage_.query("workflows", wf_opts);
|
|
|
|
|
-
|
|
|
|
|
- std::vector<std::string> ids;
|
|
|
|
|
- if (workflows.ok()) {
|
|
|
|
|
- for (const auto& wf : workflows.value().documents) {
|
|
|
|
|
- const std::string project = auth::AccessControl::projectOf(wf);
|
|
|
|
|
- if (!reachable.contains(project)) continue;
|
|
|
|
|
- if (!wanted_project.empty() && project != wanted_project) continue;
|
|
|
|
|
- ids.push_back(wf.value("_id", ""));
|
|
|
|
|
- }
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- if (ids.empty()) {
|
|
|
|
|
- // Nothing reachable means nothing to show. Answering with the whole
|
|
|
|
|
- // collection instead would be the failure this exists to prevent.
|
|
|
|
|
|
|
+ const std::string wanted_project =
|
|
|
|
|
+ req.has_param("projectId") ? req.get_param_value("projectId") : std::string{};
|
|
|
|
|
+ if (!restrictToReachableWorkflows(ctx, options, wanted_project)) {
|
|
|
sendJson(res, {{"executions", nlohmann::json::array()},
|
|
sendJson(res, {{"executions", nlohmann::json::array()},
|
|
|
{"total", 0}, {"page", 1}, {"pageSize", 0}, {"hasMore", false}});
|
|
{"total", 0}, {"page", 1}, {"pageSize", 0}, {"hasMore", false}});
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
- if (ids.size() == 1) {
|
|
|
|
|
- options.filters.push_back({"workflowId", ids.front()});
|
|
|
|
|
- } else {
|
|
|
|
|
- options.filters.push_back({"workflowId",
|
|
|
|
|
- nlohmann::json{{"op", "in"}, {"value", ids}}});
|
|
|
|
|
- }
|
|
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
// status and triggerType accept a comma separated list, because the useful
|
|
// status and triggerType accept a comma separated list, because the useful
|
|
@@ -547,7 +558,36 @@ void ExecutionController::listPending(const httplib::Request& req, httplib::Resp
|
|
|
storage::QueryOptions options;
|
|
storage::QueryOptions options;
|
|
|
options.filters.push_back({"status", "waiting"});
|
|
options.filters.push_back({"status", "waiting"});
|
|
|
|
|
|
|
|
- auto result = storage_.query("executions", options);
|
|
|
|
|
|
|
+ // Reachability, which was missing entirely. This listing was requireAuth and
|
|
|
|
|
+ // nothing else, so any account on the instance - including a member with no
|
|
|
|
|
+ // project at all - could read every waiting approval on it: workflow ids,
|
|
|
|
|
+ // workflow names and paused node ids across every project. The comment below
|
|
|
|
|
+ // reasons carefully about withholding the pause token while the rows around
|
|
|
|
|
+ // it were being handed out.
|
|
|
|
|
+ const std::string wanted_project =
|
|
|
|
|
+ req.has_param("projectId") ? req.get_param_value("projectId") : std::string{};
|
|
|
|
|
+ if (!restrictToReachableWorkflows(ctx, options, wanted_project)) {
|
|
|
|
|
+ sendJson(res, {{"pending", nlohmann::json::array()},
|
|
|
|
|
+ {"total", 0}, {"page", 1}, {"pageSize", 0}, {"hasMore", false}});
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ int page = std::max(1, intParam(req, "page", 1));
|
|
|
|
|
+ // Bounded, for the reason the executions listing is: page_size reaches the
|
|
|
|
|
+ // database as a limit, and this had none at all, so one request read every
|
|
|
|
|
+ // waiting execution as a full document - trigger data, node outputs and all.
|
|
|
|
|
+ int page_size = std::clamp(intParam(req, "pageSize", 50), 1, 200);
|
|
|
|
|
+ options.page = page;
|
|
|
|
|
+ options.page_size = page_size;
|
|
|
|
|
+
|
|
|
|
|
+ // Oldest first: an approval queue is worked from the end that has been
|
|
|
|
|
+ // waiting longest. Not pauseExpiresAt, which is 0 on a pause that never
|
|
|
|
|
+ // expires and would sort those to the front of every page.
|
|
|
|
|
+ options.sorts.push_back({"startedAt", true});
|
|
|
|
|
+
|
|
|
|
|
+ // The summary view, so the database never sends nodeExecutions or
|
|
|
|
|
+ // workflowSnapshot for rows that only need six fields.
|
|
|
|
|
+ auto result = storage_.query(kExecutionsSummaryView, options);
|
|
|
if (result.failed()) {
|
|
if (result.failed()) {
|
|
|
sendError(res, "Could not list pending approvals", 500);
|
|
sendError(res, "Could not list pending approvals", 500);
|
|
|
return;
|
|
return;
|
|
@@ -557,8 +597,12 @@ void ExecutionController::listPending(const httplib::Request& req, httplib::Resp
|
|
|
// answering, and anyone who could list every pending approval would
|
|
// answering, and anyone who could list every pending approval would
|
|
|
// otherwise be able to answer all of them. The token reaches an approver
|
|
// otherwise be able to answer all of them. The token reaches an approver
|
|
|
// through the node's own output on the execution detail.
|
|
// through the node's own output on the execution detail.
|
|
|
|
|
+ //
|
|
|
|
|
+ // The view cannot leak it either - it has a fixed field list that has never
|
|
|
|
|
+ // included pauseToken - so this holds even if a field is added here later.
|
|
|
nlohmann::json pending = nlohmann::json::array();
|
|
nlohmann::json pending = nlohmann::json::array();
|
|
|
for (const auto& record : result.value().documents) {
|
|
for (const auto& record : result.value().documents) {
|
|
|
|
|
+ if (!canSeeExecution(ctx, record)) continue;
|
|
|
pending.push_back({
|
|
pending.push_back({
|
|
|
{"executionId", record.value("_id", "")},
|
|
{"executionId", record.value("_id", "")},
|
|
|
{"workflowId", record.value("workflowId", "")},
|
|
{"workflowId", record.value("workflowId", "")},
|
|
@@ -569,7 +613,18 @@ void ExecutionController::listPending(const httplib::Request& req, httplib::Resp
|
|
|
});
|
|
});
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- sendJson(res, {{"pending", pending}, {"total", pending.size()}});
|
|
|
|
|
|
|
+ // The count of everything waiting, not of this page - the same distinction
|
|
|
|
|
+ // the executions listing gets wrong when it is written as pending.size().
|
|
|
|
|
+ const auto dropped = result.value().documents.size() - pending.size();
|
|
|
|
|
+ if (dropped > 0) {
|
|
|
|
|
+ LOG_WARN("Pending approvals listing dropped {} row(s) after the query; the total is "
|
|
|
|
|
+ "the query's and now overstates what is visible", dropped);
|
|
|
|
|
+ }
|
|
|
|
|
+ sendJson(res, {{"pending", pending},
|
|
|
|
|
+ {"total", result.value().total_count},
|
|
|
|
|
+ {"page", page},
|
|
|
|
|
+ {"pageSize", page_size},
|
|
|
|
|
+ {"hasMore", result.value().has_more}});
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
void ExecutionController::receiveExecutionEvent(const httplib::Request& req, httplib::Response& res) {
|
|
void ExecutionController::receiveExecutionEvent(const httplib::Request& req, httplib::Response& res) {
|