|
|
@@ -1,6 +1,7 @@
|
|
|
#include "workflow_controller.hpp"
|
|
|
#include "common/uuid.hpp"
|
|
|
#include "common/time_utils.hpp"
|
|
|
+#include "common/config_defaults.hpp"
|
|
|
#include "logging/logger.hpp"
|
|
|
#include <grpcpp/grpcpp.h>
|
|
|
|
|
|
@@ -23,6 +24,31 @@ WorkflowController::WorkflowController(storage::StorageClient& storage,
|
|
|
, scheduler_(scheduler)
|
|
|
, node_store_(node_store) {}
|
|
|
|
|
|
+void WorkflowController::materializeNodeConfigDefaults(nlohmann::json& body) {
|
|
|
+ if (!body.contains("nodes") || !body["nodes"].is_array()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ for (auto& node : body["nodes"]) {
|
|
|
+ if (!node.is_object()) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+
|
|
|
+ std::string node_type = node.value("type", "");
|
|
|
+ if (node_type.empty()) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto node_result = node_store_.get(node_type);
|
|
|
+ if (node_result.failed()) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+
|
|
|
+ auto config = node.value("config", nlohmann::json::object());
|
|
|
+ node["config"] = common::applyConfigDefaults(config, node_result.value().config_schema);
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
void WorkflowController::registerRoutes(httplib::Server& server) {
|
|
|
server.Get("/api/v1/workflows", [this](const httplib::Request& req, httplib::Response& res) {
|
|
|
middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
|
|
|
@@ -172,6 +198,8 @@ void WorkflowController::createWorkflow(const httplib::Request& req, httplib::Re
|
|
|
return;
|
|
|
}
|
|
|
|
|
|
+ materializeNodeConfigDefaults(body);
|
|
|
+
|
|
|
auto result = storage_.insert("workflows", body, id);
|
|
|
if (result.failed()) {
|
|
|
sendError(res, result.error().message(), 500);
|
|
|
@@ -222,6 +250,8 @@ void WorkflowController::updateWorkflow(const httplib::Request& req, httplib::Re
|
|
|
was_active = before.value().value("active", false);
|
|
|
}
|
|
|
|
|
|
+ materializeNodeConfigDefaults(body);
|
|
|
+
|
|
|
auto result = storage_.update("workflows", id, body, 0, true);
|
|
|
if (result.failed()) {
|
|
|
sendError(res, result.error().message(), 404);
|
|
|
@@ -409,8 +439,10 @@ void WorkflowController::updateScheduledTriggers(const std::string& workflow_id,
|
|
|
continue;
|
|
|
}
|
|
|
|
|
|
- // Get interval from node config
|
|
|
- auto config = node.value("config", nlohmann::json::object());
|
|
|
+ // Get interval from node config, filling in any schema defaults the
|
|
|
+ // stored config is missing.
|
|
|
+ auto config = common::applyConfigDefaults(
|
|
|
+ node.value("config", nlohmann::json::object()), node_def.config_schema);
|
|
|
int interval = config.value("pollInterval", 0);
|
|
|
|
|
|
if (interval > 0) {
|
|
|
@@ -437,7 +469,8 @@ int WorkflowController::getScheduledInterval(const nlohmann::json& workflow) {
|
|
|
continue;
|
|
|
}
|
|
|
|
|
|
- auto config = node.value("config", nlohmann::json::object());
|
|
|
+ auto config = common::applyConfigDefaults(
|
|
|
+ node.value("config", nlohmann::json::object()), node_result.value().config_schema);
|
|
|
return config.value("pollInterval", 0);
|
|
|
}
|
|
|
return 0;
|