Bläddra i källkod

feat: a project's retention reaches the workflows that inherit it

The project endpoint accepted a name and a description and nothing else,
so the setting a workflow inherits could not be set in the first place.
It takes settings now, and a retention change cascades to the workflows
that inherit it - and only those. A workflow that set its own has decided
for itself; having the project overwrite that would make the override
meaningless.

Validated rather than trusted. A retention stored as a string or a
negative number reads back as "nothing declared", which silently means
keep for ever - the failure would be invisible until somebody noticed the
data was never ageing out.

Two things the testing turned up.

The retire job took 4.7 seconds for a workflow with one execution. Each
pass re-queried the executions collection, which is 9,572 documents, to
be told what the previous pass had already shown. A short page is the
last page - restamping only removes documents, so nothing can shrink into
view behind it - and stopping there halves the work. What remains is one
filtered scan of a large collection, which is the size of the data rather
than a mistake in the loop, and it runs on a worker.

And the retention service was being dereferenced six lines before it was
constructed, exactly like the crash earlier on this branch: an empty
unique_ptr that fails inside the first request to use it rather than at
the mistake. Constructed before its controllers now, with the reason
written down where the next person will reorder it.

Verified: setting a retention on a project with two workflows cascaded to
the one that inherits and left the one that decided for itself alone,
which the API then reports as inherited true/false; a negative value and
a string are both refused by name. The test suite now leaves no
collections behind at all, because deleting a workflow removes its
storage - it used to accumulate one per run. 67 passed, 0 failed.
fszontagh 1 månad sedan
förälder
incheckning
a3313198db

+ 61 - 3
src/webserver/api/project_controller.cpp

@@ -1,4 +1,5 @@
 #include "project_controller.hpp"
+#include "storage/retention.hpp"
 
 #include "common/time_utils.hpp"
 #include "common/uuid.hpp"
@@ -27,8 +28,10 @@ const std::vector<std::string>& ownedCollections() {
 ProjectController::ProjectController(storage::StorageClient& storage,
                                      auth::AuthMiddleware& middleware,
                                      auth::AccessControl& access,
-                                     auth::AuthStore& auth_store)
-    : storage_(storage), middleware_(middleware), access_(access), auth_store_(auth_store) {}
+                                     auth::AuthStore& auth_store,
+                                     retention::RetentionService& retention)
+    : storage_(storage), middleware_(middleware), access_(access), auth_store_(auth_store),
+      retention_(retention) {}
 
 void ProjectController::registerRoutes(httplib::Server& server) {
     server.Get("/api/v1/projects", [this](const httplib::Request& req, httplib::Response& res) {
@@ -227,8 +230,40 @@ void ProjectController::updateProject(const httplib::Request& req, httplib::Resp
         patch["name"] = body["name"];
     }
     if (body.contains("description")) patch["description"] = body["description"];
+
+    // How long the project's workflows keep their data, unless a workflow says
+    // otherwise for itself. Validated here rather than trusted: a retention
+    // stored as a string or a negative number would read back as "nothing
+    // declared" and silently mean "keep for ever".
+    const auto before_retention =
+        storage::declaredTtlSeconds(existing.value().value("settings", nlohmann::json::object()));
+    std::optional<int64_t> after_retention = before_retention;
+
+    if (body.contains("settings")) {
+        if (!body["settings"].is_object()) {
+            sendError(res, "settings must be an object", 400);
+            return;
+        }
+        nlohmann::json settings = existing.value().value("settings", nlohmann::json::object());
+        for (const auto& [key, value] : body["settings"].items()) {
+            settings[key] = value;
+        }
+        if (settings.contains(storage::kRetentionKey)) {
+            const auto& retention = settings[storage::kRetentionKey];
+            if (!retention.is_object() || !retention.contains(storage::kTtlSecondsKey) ||
+                !retention[storage::kTtlSecondsKey].is_number_integer() ||
+                retention[storage::kTtlSecondsKey].get<int64_t>() < 0) {
+                sendError(res, "retention.ttlSeconds must be a whole number of seconds, "
+                               "0 or more. 0 means keep for ever", 400);
+                return;
+            }
+            after_retention = retention[storage::kTtlSecondsKey].get<int64_t>();
+        }
+        patch["settings"] = settings;
+    }
+
     if (patch.empty()) {
-        sendError(res, "Nothing to change - send a name or a description", 400);
+        sendError(res, "Nothing to change - send a name, a description or settings", 400);
         return;
     }
 
@@ -238,6 +273,29 @@ void ProjectController::updateProject(const httplib::Request& req, httplib::Resp
         return;
     }
 
+    // A project's retention is inherited, so changing it has to reach the
+    // workflows that inherit it - and only those. A workflow that set its own
+    // has decided for itself, and having a project setting overwrite that would
+    // make the override meaningless.
+    if (after_retention != before_retention && after_retention) {
+        int64_t cascaded = 0;
+        storage::QueryOptions opts;
+        opts.filters.emplace_back("projectId", id);
+        opts.page_size = 500;
+        auto workflows = storage_.query("workflows", opts);
+        if (workflows.ok()) {
+            for (const auto& workflow : workflows.value().documents) {
+                const auto own = storage::declaredTtlSeconds(
+                    workflow.value("settings", nlohmann::json::object()));
+                if (own) continue;  // it decided for itself
+                retention_.apply(workflow.value("_id", ""), *after_retention);
+                cascaded++;
+            }
+        }
+        LOG_INFO("Retention: project {} set to {}s - {} workflows inherit it",
+                 id, *after_retention, cascaded);
+    }
+
     auto after = storage_.get(PROJECTS, id);
     sendJson(res, describe(after.ok() ? after.value() : existing.value(), ctx));
 }

+ 4 - 1
src/webserver/api/project_controller.hpp

@@ -4,6 +4,7 @@
 #include <nlohmann/json.hpp>
 
 #include "../auth/access.hpp"
+#include "../retention/retention_service.hpp"
 #include "../auth/auth_middleware.hpp"
 #include "../auth/auth_store.hpp"
 #include "storage/storage_client.hpp"
@@ -23,7 +24,8 @@ public:
     ProjectController(storage::StorageClient& storage,
                       auth::AuthMiddleware& middleware,
                       auth::AccessControl& access,
-                      auth::AuthStore& auth_store);
+                      auth::AuthStore& auth_store,
+                      retention::RetentionService& retention);
 
     void registerRoutes(httplib::Server& server);
 
@@ -61,6 +63,7 @@ private:
     auth::AuthMiddleware& middleware_;
     auth::AccessControl& access_;
     auth::AuthStore& auth_store_;
+    retention::RetentionService& retention_;
 };
 
 } // namespace smartbotic::webserver::api

+ 14 - 2
src/webserver/retention/retention_service.cpp

@@ -75,8 +75,13 @@ void RetentionService::worker() {
             // 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());
+                // A workflow that never stored anything has no collection to
+                // drop, which is the ordinary case rather than a fault. Only
+                // worth a word when there was something there.
+                if (documents > 0) {
+                    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);
             }
@@ -215,6 +220,13 @@ int64_t RetentionService::restamp(const std::string& collection,
 
         // Nothing new on a full pass means every document has been seen.
         if (handled_this_pass == 0) break;
+
+        // A short page is the last page: there was nothing after it to shrink
+        // into view, and restamping only ever removes documents. Without this
+        // the common case - a handful of executions - pays for a second scan of
+        // the whole collection to be told what the first one already showed,
+        // which is most of the four seconds a small job was taking.
+        if (result.value().documents.size() < static_cast<size_t>(kPageSize)) break;
     }
 
     return restamped;

+ 7 - 5
src/webserver/webserver_service.cpp

@@ -277,14 +277,16 @@ void WebServerService::setupRoutes() {
         *auth_store_, *auth_middleware_, *access_, *storage_);
     user_ctrl_->registerRoutes(server);
 
+    // Before every controller that holds a reference to it. Dereferencing an
+    // empty unique_ptr here does not fail here - it fails later, inside the
+    // first request that uses it, as a segfault nowhere near the mistake. That
+    // has happened once on this branch already.
+    retention_ = std::make_unique<retention::RetentionService>(*storage_);
+
     project_ctrl_ = std::make_unique<api::ProjectController>(
-        *storage_, *auth_middleware_, *access_, *auth_store_);
+        *storage_, *auth_middleware_, *access_, *auth_store_, *retention_);
     project_ctrl_->registerRoutes(server);
 
-    // Constructed before the controller that holds a reference to it, and
-    // destroyed after: it owns a worker thread that touches storage_.
-    retention_ = std::make_unique<retention::RetentionService>(*storage_);
-
     workflow_ctrl_ = std::make_unique<api::WorkflowController>(
         *storage_, *auth_middleware_, *runner_registry_, *load_balancer_, *ws_server_,
         *scheduler_, *db_watcher_, *access_, *node_store_, *retention_);