|
|
@@ -0,0 +1,132 @@
|
|
|
+#include "database_watcher.hpp"
|
|
|
+
|
|
|
+#include <algorithm>
|
|
|
+
|
|
|
+#include "logging/logger.hpp"
|
|
|
+
|
|
|
+namespace smartbotic::webserver {
|
|
|
+
|
|
|
+DatabaseWatcher::DatabaseWatcher(storage::StorageClient& storage) : storage_(storage) {}
|
|
|
+
|
|
|
+DatabaseWatcher::~DatabaseWatcher() {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ // Dropping the handles closes the streams.
|
|
|
+ watches_.clear();
|
|
|
+}
|
|
|
+
|
|
|
+void DatabaseWatcher::setExecuteCallback(ExecuteCallback callback) {
|
|
|
+ execute_callback_ = std::move(callback);
|
|
|
+}
|
|
|
+
|
|
|
+void DatabaseWatcher::watch(const std::string& workflow_id,
|
|
|
+ const std::string& workflow_name,
|
|
|
+ const std::string& trigger_node_id,
|
|
|
+ const std::string& collection,
|
|
|
+ const std::vector<std::string>& event_types) {
|
|
|
+ if (workflow_id.empty() || collection.empty()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ Watch entry;
|
|
|
+ entry.workflow_id = workflow_id;
|
|
|
+ entry.workflow_name = workflow_name;
|
|
|
+ entry.trigger_node_id = trigger_node_id;
|
|
|
+ entry.collection = collection;
|
|
|
+ entry.event_types = event_types;
|
|
|
+
|
|
|
+ // Subscribing happens outside the lock: the callback can fire on the
|
|
|
+ // client's thread the moment the stream opens, and it takes the same lock.
|
|
|
+ auto subscription = storage_.subscribe(
|
|
|
+ {collection},
|
|
|
+ [this, workflow_id](const std::string& changed_collection,
|
|
|
+ const std::string& id,
|
|
|
+ const std::string& event_type,
|
|
|
+ const nlohmann::json& document) {
|
|
|
+ onEvent(workflow_id, changed_collection, id, event_type, document);
|
|
|
+ });
|
|
|
+
|
|
|
+ if (!subscription) {
|
|
|
+ LOG_WARN("Database watch for workflow {} on {} could not be opened",
|
|
|
+ workflow_id, collection);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ entry.subscription = std::move(subscription);
|
|
|
+
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ // Assigning replaces whatever was there, and the old handle goes with
|
|
|
+ // it - which is how a workflow that changed collections stops listening
|
|
|
+ // to the one it used to watch.
|
|
|
+ watches_[workflow_id] = std::move(entry);
|
|
|
+ }
|
|
|
+
|
|
|
+ LOG_INFO("Watching {} for workflow {} ({})", collection, workflow_name, workflow_id);
|
|
|
+}
|
|
|
+
|
|
|
+void DatabaseWatcher::unwatch(const std::string& workflow_id) {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ if (watches_.erase(workflow_id) > 0) {
|
|
|
+ LOG_INFO("Stopped watching for workflow {}", workflow_id);
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+size_t DatabaseWatcher::watchCount() const {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ return watches_.size();
|
|
|
+}
|
|
|
+
|
|
|
+nlohmann::json DatabaseWatcher::describe() const {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ nlohmann::json out = nlohmann::json::array();
|
|
|
+ for (const auto& [id, w] : watches_) {
|
|
|
+ out.push_back({
|
|
|
+ {"workflowId", id},
|
|
|
+ {"workflowName", w.workflow_name},
|
|
|
+ {"collection", w.collection},
|
|
|
+ {"triggerNodeId", w.trigger_node_id},
|
|
|
+ {"eventTypes", w.event_types},
|
|
|
+ });
|
|
|
+ }
|
|
|
+ return out;
|
|
|
+}
|
|
|
+
|
|
|
+void DatabaseWatcher::onEvent(const std::string& workflow_id,
|
|
|
+ const std::string& collection,
|
|
|
+ const std::string& id,
|
|
|
+ const std::string& event_type,
|
|
|
+ const nlohmann::json& document) {
|
|
|
+ std::string trigger_node_id;
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ auto it = watches_.find(workflow_id);
|
|
|
+ // The workflow may have been switched off between the event being sent
|
|
|
+ // and this running.
|
|
|
+ if (it == watches_.end()) return;
|
|
|
+
|
|
|
+ const auto& wanted = it->second.event_types;
|
|
|
+ if (!wanted.empty() &&
|
|
|
+ std::find(wanted.begin(), wanted.end(), event_type) == wanted.end()) {
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ trigger_node_id = it->second.trigger_node_id;
|
|
|
+ }
|
|
|
+
|
|
|
+ if (!execute_callback_) return;
|
|
|
+
|
|
|
+ nlohmann::json event = {
|
|
|
+ {"collection", collection},
|
|
|
+ {"documentId", id},
|
|
|
+ {"eventType", event_type},
|
|
|
+ {"document", document},
|
|
|
+ };
|
|
|
+
|
|
|
+ LOG_DEBUG("Database event {} on {}/{} starting workflow {}",
|
|
|
+ event_type, collection, id, workflow_id);
|
|
|
+
|
|
|
+ // Straight into the callback, which hands the run to a runner over gRPC.
|
|
|
+ // This is the database client's own thread, so nothing here may block for
|
|
|
+ // long or events queue up behind it.
|
|
|
+ execute_callback_(workflow_id, trigger_node_id, "database-change", event);
|
|
|
+}
|
|
|
+
|
|
|
+} // namespace smartbotic::webserver
|