|
|
@@ -0,0 +1,223 @@
|
|
|
+#include "retention_service.hpp"
|
|
|
+
|
|
|
+#include "common/time_utils.hpp"
|
|
|
+#include "logging/logger.hpp"
|
|
|
+#include "storage/workflow_collection.hpp"
|
|
|
+
|
|
|
+namespace smartbotic::webserver::retention {
|
|
|
+
|
|
|
+namespace {
|
|
|
+
|
|
|
+constexpr const char* kExecutions = "executions";
|
|
|
+constexpr const char* kExecutionWorkflowField = "workflowId";
|
|
|
+// Nothing links a document in a workflow's own collection back to the workflow:
|
|
|
+// it does not need one, because the collection belongs to that workflow and
|
|
|
+// nothing else writes to it. An empty field name means "every document here".
|
|
|
+constexpr const char* kNoWorkflowField = "";
|
|
|
+
|
|
|
+constexpr int32_t kPageSize = 200;
|
|
|
+
|
|
|
+} // namespace
|
|
|
+
|
|
|
+RetentionService::RetentionService(storage::StorageClient& storage) : storage_(storage) {
|
|
|
+ worker_ = std::thread([this] { worker(); });
|
|
|
+}
|
|
|
+
|
|
|
+RetentionService::~RetentionService() {
|
|
|
+ stopping_ = true;
|
|
|
+ cv_.notify_all();
|
|
|
+ if (worker_.joinable()) worker_.join();
|
|
|
+}
|
|
|
+
|
|
|
+void RetentionService::apply(const std::string& workflow_id, int64_t ttl_seconds) {
|
|
|
+ if (workflow_id.empty()) return;
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ queue_.push_back(Job{workflow_id, ttl_seconds, /*drop_collection_after=*/false});
|
|
|
+ }
|
|
|
+ cv_.notify_one();
|
|
|
+}
|
|
|
+
|
|
|
+void RetentionService::retire(const std::string& workflow_id) {
|
|
|
+ if (workflow_id.empty()) return;
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(mutex_);
|
|
|
+ queue_.push_back(Job{workflow_id, kRetireTtlSeconds, /*drop_collection_after=*/true});
|
|
|
+ }
|
|
|
+ cv_.notify_one();
|
|
|
+}
|
|
|
+
|
|
|
+void RetentionService::worker() {
|
|
|
+ while (!stopping_) {
|
|
|
+ Job job;
|
|
|
+ {
|
|
|
+ std::unique_lock<std::mutex> lock(mutex_);
|
|
|
+ cv_.wait(lock, [this] { return stopping_ || !queue_.empty(); });
|
|
|
+ if (stopping_) return;
|
|
|
+ job = queue_.front();
|
|
|
+ queue_.pop_front();
|
|
|
+ }
|
|
|
+
|
|
|
+ const std::string own = storage::workflowCollectionName(job.workflow_id);
|
|
|
+
|
|
|
+ const int64_t executions =
|
|
|
+ restamp(kExecutions, kExecutionWorkflowField, job.workflow_id, job.ttl_seconds);
|
|
|
+ const int64_t documents = restamp(own, kNoWorkflowField, job.workflow_id, job.ttl_seconds);
|
|
|
+
|
|
|
+ LOG_INFO("Retention: workflow {} restamped to {}s - {} executions, {} documents",
|
|
|
+ job.workflow_id, job.ttl_seconds, executions, documents);
|
|
|
+
|
|
|
+ if (job.drop_collection_after) {
|
|
|
+ // The documents are not gone yet - they expire on their own shortly.
|
|
|
+ // Dropping the collection now would take them with it, which is the
|
|
|
+ // same outcome sooner; what it must not do is drop a collection
|
|
|
+ // whose documents are still being restamped, so it happens here,
|
|
|
+ // after the pass, on the same worker.
|
|
|
+ auto dropped = storage_.dropCollection(own);
|
|
|
+ if (dropped.failed()) {
|
|
|
+ LOG_WARN("Retention: could not drop {} ({}); its documents still expire on their own",
|
|
|
+ own, dropped.error().message());
|
|
|
+ } else {
|
|
|
+ LOG_INFO("Retention: dropped {} for deleted workflow {}", own, job.workflow_id);
|
|
|
+ }
|
|
|
+ }
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+RetentionService::Impact RetentionService::measure(const std::string& collection,
|
|
|
+ const std::string& workflow_field,
|
|
|
+ const std::string& workflow_id,
|
|
|
+ int64_t ttl_seconds) {
|
|
|
+ Impact impact;
|
|
|
+ const int64_t now = common::TimeUtils::nowMs();
|
|
|
+ const int64_t ttl_ms = ttl_seconds * 1000;
|
|
|
+
|
|
|
+ int32_t page = 1;
|
|
|
+ while (true) {
|
|
|
+ storage::QueryOptions opts;
|
|
|
+ opts.page = page;
|
|
|
+ opts.page_size = kPageSize;
|
|
|
+ if (!workflow_field.empty()) {
|
|
|
+ opts.filters.emplace_back(workflow_field, workflow_id);
|
|
|
+ }
|
|
|
+ // Only the creation stamp is needed to answer this. Asking for the whole
|
|
|
+ // document would pull every execution's node output through the wire to
|
|
|
+ // count them.
|
|
|
+ opts.fields = {"_created_at"};
|
|
|
+
|
|
|
+ auto result = storage_.query(collection, opts);
|
|
|
+ if (result.failed()) {
|
|
|
+ // A collection that is not there yet holds nothing to lose. Anything
|
|
|
+ // else is worth saying out loud rather than reporting as zero.
|
|
|
+ if (result.error().code() != common::ErrorCode::CollectionNotFound) {
|
|
|
+ LOG_WARN("Retention: cannot measure {} ({})", collection, result.error().message());
|
|
|
+ }
|
|
|
+ return impact;
|
|
|
+ }
|
|
|
+
|
|
|
+ for (const auto& doc : result.value().documents) {
|
|
|
+ impact.total++;
|
|
|
+ const int64_t created = common::TimeUtils::documentStamp(doc, "_created_at");
|
|
|
+ if (created <= 0) continue;
|
|
|
+ const int64_t age = now - created;
|
|
|
+ if (age > impact.oldest_age_ms) impact.oldest_age_ms = age;
|
|
|
+ // Keeping for ever deletes nothing, whatever the ages are.
|
|
|
+ if (ttl_seconds > 0 && age >= ttl_ms) impact.expiring_now++;
|
|
|
+ }
|
|
|
+
|
|
|
+ if (!result.value().has_more) break;
|
|
|
+ page++;
|
|
|
+ }
|
|
|
+ return impact;
|
|
|
+}
|
|
|
+
|
|
|
+std::optional<storage::Retention> RetentionService::effectiveFor(const nlohmann::json& workflow) {
|
|
|
+ const auto settings = workflow.value("settings", nlohmann::json::object());
|
|
|
+ if (const auto own = storage::declaredTtlSeconds(settings)) {
|
|
|
+ return storage::Retention{*own, false};
|
|
|
+ }
|
|
|
+
|
|
|
+ const std::string project_id = workflow.value("projectId", std::string{});
|
|
|
+ if (project_id.empty()) return std::nullopt;
|
|
|
+
|
|
|
+ auto project = storage_.get("projects", project_id);
|
|
|
+ if (project.failed()) return std::nullopt;
|
|
|
+ return storage::declaredRetention(
|
|
|
+ settings, project.value().value("settings", nlohmann::json::object()));
|
|
|
+}
|
|
|
+
|
|
|
+RetentionService::Preview RetentionService::preview(const std::string& workflow_id,
|
|
|
+ int64_t ttl_seconds) {
|
|
|
+ Preview preview;
|
|
|
+ preview.executions =
|
|
|
+ measure(kExecutions, kExecutionWorkflowField, workflow_id, ttl_seconds);
|
|
|
+ preview.documents = measure(storage::workflowCollectionName(workflow_id), kNoWorkflowField,
|
|
|
+ workflow_id, ttl_seconds);
|
|
|
+ return preview;
|
|
|
+}
|
|
|
+
|
|
|
+int64_t RetentionService::restamp(const std::string& collection,
|
|
|
+ const std::string& workflow_field,
|
|
|
+ const std::string& workflow_id,
|
|
|
+ int64_t ttl_seconds) {
|
|
|
+ int64_t restamped = 0;
|
|
|
+
|
|
|
+ // Always page 1. Restamping with a shorter retention deletes documents as it
|
|
|
+ // goes, so the result set shrinks underneath a cursor and advancing the page
|
|
|
+ // number would step over the documents that moved up into the space. Reading
|
|
|
+ // the first page repeatedly, and stopping when a pass changes nothing,
|
|
|
+ // cannot skip a document.
|
|
|
+ //
|
|
|
+ // Ids already handled are remembered for the same reason in reverse: with a
|
|
|
+ // longer retention nothing is deleted, so page 1 returns the same documents
|
|
|
+ // for ever and the loop would not end.
|
|
|
+ std::unordered_set<std::string> done;
|
|
|
+
|
|
|
+ while (!stopping_) {
|
|
|
+ storage::QueryOptions opts;
|
|
|
+ opts.page = 1;
|
|
|
+ opts.page_size = kPageSize;
|
|
|
+ if (!workflow_field.empty()) {
|
|
|
+ opts.filters.emplace_back(workflow_field, workflow_id);
|
|
|
+ }
|
|
|
+
|
|
|
+ auto result = storage_.query(collection, opts);
|
|
|
+ if (result.failed()) {
|
|
|
+ if (result.error().code() != common::ErrorCode::CollectionNotFound) {
|
|
|
+ LOG_WARN("Retention: cannot read {} ({})", collection, result.error().message());
|
|
|
+ }
|
|
|
+ return restamped;
|
|
|
+ }
|
|
|
+ if (result.value().documents.empty()) break;
|
|
|
+
|
|
|
+ int64_t handled_this_pass = 0;
|
|
|
+ for (const auto& doc : result.value().documents) {
|
|
|
+ if (stopping_) return restamped;
|
|
|
+ const std::string id = doc.value("_id", "");
|
|
|
+ if (id.empty() || done.contains(id)) continue;
|
|
|
+
|
|
|
+ // upsert rather than update: only insert and upsert carry a TTL, and
|
|
|
+ // upsert swaps the expiry entry in one locked server-side step. It
|
|
|
+ // preserves _created_at, so a restamped record still says when it
|
|
|
+ // was made - verified against the live database, because a retention
|
|
|
+ // change that quietly re-dated every execution would be worse than
|
|
|
+ // the growth it was meant to fix.
|
|
|
+ auto written = storage_.upsert(collection, doc, id, ttl_seconds * 1000);
|
|
|
+ if (written.failed()) {
|
|
|
+ LOG_WARN("Retention: could not restamp {}/{} ({})", collection, id,
|
|
|
+ written.error().message());
|
|
|
+ } else {
|
|
|
+ restamped++;
|
|
|
+ }
|
|
|
+ done.insert(id);
|
|
|
+ handled_this_pass++;
|
|
|
+ }
|
|
|
+
|
|
|
+ // Nothing new on a full pass means every document has been seen.
|
|
|
+ if (handled_this_pass == 0) break;
|
|
|
+ }
|
|
|
+
|
|
|
+ return restamped;
|
|
|
+}
|
|
|
+
|
|
|
+} // namespace smartbotic::webserver::retention
|