|
@@ -17,6 +17,7 @@ WorkflowController::WorkflowController(storage::StorageClient& storage,
|
|
|
WebSocketServer& ws_server,
|
|
WebSocketServer& ws_server,
|
|
|
WorkflowScheduler& scheduler,
|
|
WorkflowScheduler& scheduler,
|
|
|
DatabaseWatcher& db_watcher,
|
|
DatabaseWatcher& db_watcher,
|
|
|
|
|
+ auth::AccessControl& access,
|
|
|
nodes::NodeStore& node_store)
|
|
nodes::NodeStore& node_store)
|
|
|
: storage_(storage)
|
|
: storage_(storage)
|
|
|
, middleware_(middleware)
|
|
, middleware_(middleware)
|
|
@@ -25,6 +26,7 @@ WorkflowController::WorkflowController(storage::StorageClient& storage,
|
|
|
, ws_server_(ws_server)
|
|
, ws_server_(ws_server)
|
|
|
, scheduler_(scheduler)
|
|
, scheduler_(scheduler)
|
|
|
, db_watcher_(db_watcher)
|
|
, db_watcher_(db_watcher)
|
|
|
|
|
+ , access_(access)
|
|
|
, node_store_(node_store) {}
|
|
, node_store_(node_store) {}
|
|
|
|
|
|
|
|
|
|
|
|
@@ -147,6 +149,36 @@ void WorkflowController::registerRoutes(httplib::Server& server) {
|
|
|
});
|
|
});
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+
|
|
|
|
|
+bool WorkflowController::loadAllowed(const httplib::Request& req, httplib::Response& res,
|
|
|
|
|
+ const auth::AuthContext& ctx, const std::string& id,
|
|
|
|
|
+ auth::Action action, nlohmann::json& out) {
|
|
|
|
|
+ (void)req;
|
|
|
|
|
+ auto stored = storage_.get("workflows", id);
|
|
|
|
|
+ if (stored.failed()) {
|
|
|
|
|
+ sendError(res, "Workflow not found", 404);
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ const std::string project = auth::AccessControl::projectOf(stored.value());
|
|
|
|
|
+
|
|
|
|
|
+ if (!access_.allowed(ctx, project, auth::Action::Read)) {
|
|
|
|
|
+ // Not "forbidden". Whether a workflow exists is itself something only
|
|
|
|
|
+ // people who can reach it should learn.
|
|
|
|
|
+ sendError(res, "Workflow not found", 404);
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (action != auth::Action::Read && !access_.allowed(ctx, project, action)) {
|
|
|
|
|
+ sendError(res, action == auth::Action::Run
|
|
|
|
|
+ ? "You can see this workflow but not run it"
|
|
|
|
|
+ : "You can see this workflow but not change it", 403);
|
|
|
|
|
+ return false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ out = stored.value();
|
|
|
|
|
+ return true;
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
void WorkflowController::listWorkflows(const httplib::Request& req, httplib::Response& res,
|
|
void WorkflowController::listWorkflows(const httplib::Request& req, httplib::Response& res,
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
storage::QueryOptions options;
|
|
storage::QueryOptions options;
|
|
@@ -202,6 +234,24 @@ void WorkflowController::listWorkflows(const httplib::Request& req, httplib::Res
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // Only the projects this caller can reach. Filtering here rather than in
|
|
|
|
|
+ // the query because the database cannot be asked "any of these ids", and
|
|
|
|
|
+ // a workflows collection is the handful of things somebody built.
|
|
|
|
|
+ const auto reachable = access_.projectsFor(ctx);
|
|
|
|
|
+ {
|
|
|
|
|
+ nlohmann::json kept = nlohmann::json::array();
|
|
|
|
|
+ for (const auto& workflow : result.value().documents) {
|
|
|
|
|
+ if (reachable.contains(auth::AccessControl::projectOf(workflow))) {
|
|
|
|
|
+ kept.push_back(workflow);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ result.value().documents = kept;
|
|
|
|
|
+ // The totals have to describe what is being returned, or the pager
|
|
|
|
|
+ // promises pages that are not there.
|
|
|
|
|
+ result.value().total_count = static_cast<int64_t>(kept.size());
|
|
|
|
|
+ result.value().has_more = false;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
nlohmann::json workflows = result.value().documents;
|
|
nlohmann::json workflows = result.value().documents;
|
|
|
int64_t total = result.value().total_count;
|
|
int64_t total = result.value().total_count;
|
|
|
bool has_more = result.value().has_more;
|
|
bool has_more = result.value().has_more;
|
|
@@ -241,13 +291,8 @@ void WorkflowController::getWorkflow(const httplib::Request& req, httplib::Respo
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
std::string id = req.matches[1];
|
|
std::string id = req.matches[1];
|
|
|
|
|
|
|
|
- auto result = storage_.get("workflows", id);
|
|
|
|
|
- if (result.failed()) {
|
|
|
|
|
- sendError(res, "Workflow not found", 404);
|
|
|
|
|
- return;
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- auto workflow = result.value();
|
|
|
|
|
|
|
+ nlohmann::json workflow;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Read, workflow)) return;
|
|
|
|
|
|
|
|
// Whether the workflow that runs differs from the workflow on screen.
|
|
// Whether the workflow that runs differs from the workflow on screen.
|
|
|
//
|
|
//
|
|
@@ -285,6 +330,29 @@ void WorkflowController::createWorkflow(const httplib::Request& req, httplib::Re
|
|
|
body["ownerId"] = ctx.user_id;
|
|
body["ownerId"] = ctx.user_id;
|
|
|
body["active"] = body.value("active", false);
|
|
body["active"] = body.value("active", false);
|
|
|
|
|
|
|
|
|
|
+ // Everything belongs to a project. Where the caller asked for, if they
|
|
|
|
|
+ // may write there; their own project otherwise - never nowhere, which
|
|
|
|
|
+ // would make it reachable only by an admin.
|
|
|
|
|
+ std::string project = body.value("projectId", "");
|
|
|
|
|
+ if (!project.empty() && !access_.allowed(ctx, project, auth::Action::Write)) {
|
|
|
|
|
+ sendError(res, "You cannot create a workflow in that project", 403);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (project.empty()) {
|
|
|
|
|
+ auto personal = access_.personalProjectFor(ctx.user_id);
|
|
|
|
|
+ if (personal.failed()) {
|
|
|
|
|
+ auto created = access_.ensurePersonalProject(ctx.user_id, ctx.username);
|
|
|
|
|
+ if (created.failed()) {
|
|
|
|
|
+ sendError(res, "You have no project to put this in", 500);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ project = created.value().value("_id", "");
|
|
|
|
|
+ } else {
|
|
|
|
|
+ project = personal.value();
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+ body["projectId"] = project;
|
|
|
|
|
+
|
|
|
// Validate required fields
|
|
// Validate required fields
|
|
|
if (!body.contains("name") || body["name"].get<std::string>().empty()) {
|
|
if (!body.contains("name") || body["name"].get<std::string>().empty()) {
|
|
|
sendError(res, "Name is required", 400);
|
|
sendError(res, "Name is required", 400);
|
|
@@ -320,8 +388,24 @@ void WorkflowController::updateWorkflow(const httplib::Request& req, httplib::Re
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
try {
|
|
try {
|
|
|
std::string id = req.matches[1];
|
|
std::string id = req.matches[1];
|
|
|
|
|
+
|
|
|
|
|
+ nlohmann::json existing;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Write, existing)) return;
|
|
|
|
|
+
|
|
|
auto body = nlohmann::json::parse(req.body);
|
|
auto body = nlohmann::json::parse(req.body);
|
|
|
|
|
|
|
|
|
|
+ // Moving it to another project is a change of who can reach it, so it
|
|
|
|
|
+ // needs the right to write in the project it is going to as well as
|
|
|
|
|
+ // the one it is leaving.
|
|
|
|
|
+ if (body.contains("projectId") && body["projectId"].is_string()) {
|
|
|
|
|
+ const std::string target = body["projectId"];
|
|
|
|
|
+ if (target != auth::AccessControl::projectOf(existing) &&
|
|
|
|
|
+ !access_.allowed(ctx, target, auth::Action::Write)) {
|
|
|
|
|
+ sendError(res, "You cannot move a workflow into a project you cannot write to", 403);
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// What the editor had loaded when the user started changing it. Two
|
|
// What the editor had loaded when the user started changing it. Two
|
|
|
// people with the same workflow open used to mean whoever pressed save
|
|
// people with the same workflow open used to mean whoever pressed save
|
|
|
// second silently threw the other's work away, with nothing anywhere to
|
|
// second silently threw the other's work away, with nothing anywhere to
|
|
@@ -466,6 +550,9 @@ void WorkflowController::deleteWorkflow(const httplib::Request& req, httplib::Re
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
std::string id = req.matches[1];
|
|
std::string id = req.matches[1];
|
|
|
|
|
|
|
|
|
|
+ nlohmann::json existing;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Write, existing)) return;
|
|
|
|
|
+
|
|
|
// Every trigger it had, not just the scheduled ones - a deleted workflow
|
|
// Every trigger it had, not just the scheduled ones - a deleted workflow
|
|
|
// that kept its database watch would go on being started by changes, with
|
|
// that kept its database watch would go on being started by changes, with
|
|
|
// nothing left to run.
|
|
// nothing left to run.
|
|
@@ -487,14 +574,8 @@ void WorkflowController::executeWorkflow(const httplib::Request& req, httplib::R
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
std::string workflow_id = req.matches[1];
|
|
std::string workflow_id = req.matches[1];
|
|
|
|
|
|
|
|
- // Get workflow
|
|
|
|
|
- auto workflow_result = storage_.get("workflows", workflow_id);
|
|
|
|
|
- if (workflow_result.failed()) {
|
|
|
|
|
- sendError(res, "Workflow not found", 404);
|
|
|
|
|
- return;
|
|
|
|
|
- }
|
|
|
|
|
-
|
|
|
|
|
- auto& workflow = workflow_result.value();
|
|
|
|
|
|
|
+ nlohmann::json workflow;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, workflow_id, auth::Action::Run, workflow)) return;
|
|
|
|
|
|
|
|
// Select runner
|
|
// Select runner
|
|
|
auto runner = load_balancer_.selectRunner();
|
|
auto runner = load_balancer_.selectRunner();
|
|
@@ -585,11 +666,9 @@ void WorkflowController::listWorkflowVersions(const httplib::Request& req, httpl
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- auto current = storage_.get("workflows", id);
|
|
|
|
|
- if (current.failed()) {
|
|
|
|
|
- sendError(res, "Workflow not found", 404);
|
|
|
|
|
- return;
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ nlohmann::json current_doc;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Read, current_doc)) return;
|
|
|
|
|
+ auto current = common::Result<nlohmann::json>(current_doc);
|
|
|
|
|
|
|
|
auto listed = storage_.listVersions("workflows", id, limit, 0);
|
|
auto listed = storage_.listVersions("workflows", id, limit, 0);
|
|
|
if (listed.failed()) {
|
|
if (listed.failed()) {
|
|
@@ -632,6 +711,9 @@ void WorkflowController::getWorkflowVersion(const httplib::Request& req, httplib
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ nlohmann::json current_doc;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Read, current_doc)) return;
|
|
|
|
|
+
|
|
|
auto doc = storage_.getVersion("workflows", id, version);
|
|
auto doc = storage_.getVersion("workflows", id, version);
|
|
|
if (doc.failed()) {
|
|
if (doc.failed()) {
|
|
|
sendError(res, "That version of the workflow is no longer stored", 404);
|
|
sendError(res, "That version of the workflow is no longer stored", 404);
|
|
@@ -644,11 +726,9 @@ void WorkflowController::publishWorkflow(const httplib::Request& req, httplib::R
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
std::string id = req.matches[1];
|
|
std::string id = req.matches[1];
|
|
|
|
|
|
|
|
- auto current = storage_.get("workflows", id);
|
|
|
|
|
- if (current.failed()) {
|
|
|
|
|
- sendError(res, "Workflow not found", 404);
|
|
|
|
|
- return;
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ nlohmann::json current_doc;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Write, current_doc)) return;
|
|
|
|
|
+ auto current = common::Result<nlohmann::json>(current_doc);
|
|
|
|
|
|
|
|
// Publishing what is already published would store a new version whose only
|
|
// Publishing what is already published would store a new version whose only
|
|
|
// difference is the act of publishing it, and leave the button offering to
|
|
// difference is the act of publishing it, and leave the button offering to
|
|
@@ -715,6 +795,9 @@ void WorkflowController::activateWorkflow(const httplib::Request& req, httplib::
|
|
|
// An active workflow runs what was published, so switching one on without
|
|
// An active workflow runs what was published, so switching one on without
|
|
|
// anything published would leave its triggers with nothing to run. Turning
|
|
// anything published would leave its triggers with nothing to run. Turning
|
|
|
// it on is a clear enough statement that what is there now should go live.
|
|
// it on is a clear enough statement that what is there now should go live.
|
|
|
|
|
+ nlohmann::json guard_doc;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Write, guard_doc)) return;
|
|
|
|
|
+
|
|
|
auto before = storage_.get("workflows", id);
|
|
auto before = storage_.get("workflows", id);
|
|
|
nlohmann::json patch = {{"active", true},
|
|
nlohmann::json patch = {{"active", true},
|
|
|
{"consecutiveFailures", 0},
|
|
{"consecutiveFailures", 0},
|
|
@@ -751,6 +834,9 @@ void WorkflowController::deactivateWorkflow(const httplib::Request& req, httplib
|
|
|
const auth::AuthContext& ctx) {
|
|
const auth::AuthContext& ctx) {
|
|
|
std::string id = req.matches[1];
|
|
std::string id = req.matches[1];
|
|
|
|
|
|
|
|
|
|
+ nlohmann::json guard_doc;
|
|
|
|
|
+ if (!loadAllowed(req, res, ctx, id, auth::Action::Write, guard_doc)) return;
|
|
|
|
|
+
|
|
|
auto result = storage_.update("workflows", id, {{"active", false}, {"updatedAt", TimeUtils::nowMs()}}, 0, true);
|
|
auto result = storage_.update("workflows", id, {{"active", false}, {"updatedAt", TimeUtils::nowMs()}}, 0, true);
|
|
|
if (result.failed()) {
|
|
if (result.failed()) {
|
|
|
sendError(res, "Workflow not found", 404);
|
|
sendError(res, "Workflow not found", 404);
|