瀏覽代碼

feat: version history, publishing, and triggers that run the published version

Editing a live workflow used to put half-finished work into production the
moment it was saved: whatever the record said last was what the next
schedule or webhook ran. Now the record is the draft, publishing records
which version should actually run, and the two are only the same when
somebody says so.

A trigger - a schedule, a webhook, an error handler - runs the published
version. A run somebody starts by hand deliberately runs the draft:
testing a change before publishing it is the whole point of publishing
being a separate act.

Version history had to be turned on first. The collection was created
without it, so the version number climbed on every write but nothing
older was kept - there was no earlier version to publish, compare or go
back to. smartbotic-database 2.4.5 can enable it on an existing
collection through configureCollection, which the storage adapter now
exposes; createCollection cannot, because it refuses a collection that is
already there.

New endpoints: the versions of a workflow, one version's content, and
publish. Activating a workflow publishes what is there if nothing has
been published yet, because otherwise switching one on would leave its
triggers with nothing to run.

A published version that has aged out of history falls back to the draft
with a warning rather than refusing to run. Stopping a live workflow
because its history rolled over would be worse than running the newer
thing - but the two can differ, so it is said out loud.

Verified end to end on one workflow whose published version and draft
report different values: the webhook ran PUBLISHED and Execute ran DRAFT,
at the same moment, from the same record. Version history returns real
earlier content, publish marks the right version, and a version that was
never stored answers 404.
fszontagh 1 月之前
父節點
當前提交
8443f5c4ab

+ 23 - 0
lib/storage/storage_client.cpp

@@ -237,6 +237,29 @@ Result<void> StorageClient::createCollection(const std::string& name,
     return {};
 }
 
+Result<void> StorageClient::configureCollection(const std::string& collection,
+                                                const CollectionConfig& cfg) {
+    smartbotic::database::Client::CollectionConfig upstream;
+    // An empty precision means "leave it alone" upstream too, so it is passed
+    // through rather than defaulted - getting this wrong would silently rewrite
+    // the precision of a collection somebody had deliberately set to ns.
+    upstream.timestampPrecision = cfg.timestamp_precision;
+    upstream.versioningEnabled = cfg.versioning_enabled;
+
+    if (!impl_->client_->configureCollection(collection, upstream)) {
+        return Error(ErrorCode::DatabaseError, "configureCollection failed: " + collection);
+    }
+    return {};
+}
+
+Result<CollectionConfig> StorageClient::getCollectionConfig(const std::string& collection) {
+    auto upstream = impl_->client_->getCollectionConfig(collection);
+    CollectionConfig cfg;
+    cfg.timestamp_precision = upstream.timestampPrecision;
+    cfg.versioning_enabled = upstream.versioningEnabled;
+    return cfg;
+}
+
 Result<void> StorageClient::dropCollection(const std::string& name) {
     bool ok = impl_->client_->dropCollection(name);
     if (!ok) return Error(ErrorCode::CollectionNotFound, "dropCollection failed: " + name);

+ 17 - 0
lib/storage/storage_client.hpp

@@ -2,6 +2,7 @@
 
 #include <map>
 #include <memory>
+#include <optional>
 #include <string>
 #include <vector>
 
@@ -99,6 +100,14 @@ struct FileInfo {
 // A read-only projection over a collection, enforced by the server. Querying the
 // view by name returns only the included fields, so the database never ships the
 // rest over the wire (upstream >= 2.4).
+// v2.4.5 - per-collection settings that can be changed after creation.
+// Unset fields are left alone, so versioning can be turned on without
+// disturbing the timestamp precision.
+struct CollectionConfig {
+    std::string timestamp_precision;             // "" leaves it unchanged
+    std::optional<bool> versioning_enabled;      // unset leaves it unchanged
+};
+
 struct ViewInfo {
     std::string name;
     std::string collection;
@@ -184,6 +193,14 @@ public:
 
     std::vector<ViewInfo> listViews();
 
+    // Version history can be turned on for a collection that already exists -
+    // createCollection cannot, because it fails outright when the collection is
+    // there. Disabling stops new versions being recorded and leaves the history
+    // already written readable, so the switch is reversible.
+    common::Result<void> configureCollection(const std::string& collection,
+                                             const CollectionConfig& cfg);
+    common::Result<CollectionConfig> getCollectionConfig(const std::string& collection);
+
     // File store. Kept separate from documents: images and other binaries do not
     // belong inline in a document, and upstream deduplicates by checksum.
     common::Result<FileUploadInfo> uploadFile(const std::vector<uint8_t>& data,

+ 9 - 0
proto/workflow.proto

@@ -124,6 +124,15 @@ message ExecuteWorkflowRequest {
     // Credential access is checked against it, and an empty workflow id is
     // treated as an administrative call with access to every credential.
     string inline_workflow = 6;
+
+    // Run the published version rather than whatever was saved last.
+    //
+    // Set by everything that starts a workflow on its own - a schedule, a
+    // webhook, an error handler - so that editing a live workflow does not put
+    // half-finished work into production the moment it is saved. A run somebody
+    // started by hand deliberately leaves this false: testing a change before
+    // publishing it is the whole point of being able to.
+    bool use_published = 7;
 }
 
 // Execute workflow response

+ 21 - 0
src/runner/runner_service.cpp

@@ -89,6 +89,27 @@ grpc::Status RunnerServiceImpl::ExecuteWorkflow(grpc::ServerContext* context,
 
     // Get workflow from database
     auto workflow_result = storage_.get("workflows", request->workflow_id());
+
+    // A trigger runs what was published. The record itself is the draft.
+    if (request->use_published() && workflow_result.ok()) {
+        const int64_t published = workflow_result.value().value("publishedVersion", int64_t{0});
+        const int64_t current = workflow_result.value().value("_version", int64_t{0});
+        if (published > 0 && published != current) {
+            auto pinned = storage_.getVersion("workflows", request->workflow_id(), published);
+            if (pinned.ok()) {
+                LOG_INFO("Workflow {} running published version {} (draft is {})",
+                         request->workflow_id(), published, current);
+                workflow_result = pinned;
+            } else {
+                // Refusing would stop a live workflow because its history has
+                // aged out, which is worse than running the draft - but it must
+                // be said out loud, because the two can differ.
+                LOG_WARN("Workflow {}: published version {} is no longer stored, "
+                         "running the current draft instead", request->workflow_id(), published);
+            }
+        }
+    }
+
     if (workflow_result.failed()) {
         LOG_ERROR("Failed to get workflow from database: {}", workflow_result.error().message());
         response->set_status(proto::EXECUTION_STATUS_FAILED);

+ 2 - 0
src/webserver/api/webhook_controller.cpp

@@ -160,6 +160,8 @@ void WebhookController::handleWebhook(const httplib::Request& req, httplib::Resp
     proto::ExecuteWorkflowRequest grpc_req;
     grpc_req.set_workflow_id(workflow_id);
     grpc_req.set_trigger_type("webhook");
+    // Fired by the outside world, so it runs the published version.
+    grpc_req.set_use_published(true);
     grpc_req.set_trigger_data(trigger_data.dump());
     grpc_req.set_wait_for_completion(true);  // Wait for webhook response
     grpc_req.set_timeout_ms(30000);  // 30 second timeout

+ 171 - 6
src/webserver/api/workflow_controller.cpp

@@ -87,6 +87,24 @@ void WorkflowController::registerRoutes(httplib::Server& server) {
         });
     });
 
+    server.Get(R"(/api/v1/workflows/([^/]+)/versions)", [this](const httplib::Request& req, httplib::Response& res) {
+        middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
+            listWorkflowVersions(req, res, ctx);
+        });
+    });
+
+    server.Get(R"(/api/v1/workflows/([^/]+)/versions/(\d+))", [this](const httplib::Request& req, httplib::Response& res) {
+        middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
+            getWorkflowVersion(req, res, ctx);
+        });
+    });
+
+    server.Post(R"(/api/v1/workflows/([^/]+)/publish)", [this](const httplib::Request& req, httplib::Response& res) {
+        middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
+            publishWorkflow(req, res, ctx);
+        });
+    });
+
     server.Post(R"(/api/v1/workflows/([^/]+)/activate)", [this](const httplib::Request& req, httplib::Response& res) {
         middleware_.requireAuth(req, res, [this](auto& req, auto& res, auto& ctx) {
             activateWorkflow(req, res, ctx);
@@ -438,6 +456,9 @@ void WorkflowController::executeWorkflow(const httplib::Request& req, httplib::R
     proto::ExecuteWorkflowRequest grpc_req;
     grpc_req.set_workflow_id(workflow_id);
     grpc_req.set_trigger_type("manual");
+    // Deliberately NOT the published version. Somebody pressing Execute is
+    // testing what is on their canvas; making them publish first to see whether
+    // a change works would defeat the point of publishing being deliberate.
     grpc_req.set_trigger_data(trigger_data.dump());
     grpc_req.set_wait_for_completion(false);
 
@@ -477,6 +498,132 @@ void WorkflowController::executeWorkflow(const httplib::Request& req, httplib::R
     sendJson(res, response, 202);
 }
 
+
+// ---------------------------------------------------------------------------
+// Versions and publishing
+//
+// The workflow record is the draft: it is whatever was written last. Publishing
+// records which version of it is the one that should actually run, so that
+// editing a live workflow does not put half-finished work into production the
+// moment it is saved. A manual run still uses the draft - that is how somebody
+// tests a change before publishing it.
+// ---------------------------------------------------------------------------
+
+void WorkflowController::listWorkflowVersions(const httplib::Request& req, httplib::Response& res,
+                                              const auth::AuthContext& ctx) {
+    (void)ctx;
+    std::string id = req.matches[1];
+
+    int limit = 50;
+    if (req.has_param("limit")) {
+        try {
+            limit = std::clamp(std::stoi(req.get_param_value("limit")), 1, 200);
+        } catch (const std::exception&) {
+            // A bad limit is not worth refusing the request over.
+        }
+    }
+
+    auto current = storage_.get("workflows", id);
+    if (current.failed()) {
+        sendError(res, "Workflow not found", 404);
+        return;
+    }
+
+    auto listed = storage_.listVersions("workflows", id, limit, 0);
+    if (listed.failed()) {
+        sendError(res, listed.error().message(), 500);
+        return;
+    }
+
+    const int64_t published = current.value().value("publishedVersion", int64_t{0});
+    const int64_t live = current.value().value("_version", int64_t{0});
+
+    nlohmann::json versions = nlohmann::json::array();
+    for (const auto& v : listed.value().versions) {
+        versions.push_back({
+            {"version", v.version},
+            {"createdAt", v.created_at},
+            {"isPublished", v.version == published},
+            {"isCurrent", v.version == live},
+        });
+    }
+
+    sendJson(res, {
+        {"versions", versions},
+        {"totalCount", listed.value().total_count},
+        {"hasMore", listed.value().has_more},
+        {"currentVersion", live},
+        {"publishedVersion", published},
+    });
+}
+
+void WorkflowController::getWorkflowVersion(const httplib::Request& req, httplib::Response& res,
+                                            const auth::AuthContext& ctx) {
+    (void)ctx;
+    std::string id = req.matches[1];
+
+    int64_t version = 0;
+    try {
+        version = std::stoll(req.matches[2]);
+    } catch (const std::exception&) {
+        sendError(res, "Version must be a number", 400);
+        return;
+    }
+
+    auto doc = storage_.getVersion("workflows", id, version);
+    if (doc.failed()) {
+        sendError(res, "That version of the workflow is no longer stored", 404);
+        return;
+    }
+    sendJson(res, doc.value());
+}
+
+void WorkflowController::publishWorkflow(const httplib::Request& req, httplib::Response& res,
+                                         const auth::AuthContext& ctx) {
+    std::string id = req.matches[1];
+
+    auto current = storage_.get("workflows", id);
+    if (current.failed()) {
+        sendError(res, "Workflow not found", 404);
+        return;
+    }
+
+    const int64_t version = current.value().value("_version", int64_t{0});
+    if (version <= 0) {
+        // Without a version there is nothing to point at, and publishing would
+        // be a promise the runner cannot keep.
+        sendError(res, "This workflow has no stored version to publish yet", 409);
+        return;
+    }
+
+    nlohmann::json patch;
+    patch["publishedVersion"] = version;
+    patch["publishedAt"] = common::TimeUtils::nowMs();
+    patch["publishedBy"] = ctx.username;
+
+    auto updated = storage_.update("workflows", id, patch, 0, true);
+    if (updated.failed()) {
+        sendError(res, updated.error().message(), 500);
+        return;
+    }
+
+    // Publishing is itself a write, so the record has moved on by one. The
+    // version that was published is the one the caller asked about, not the one
+    // this write produced.
+    ws_server_.broadcast("workflows." + id + ".published", {
+        {"workflowId", id},
+        {"publishedVersion", version},
+        {"publishedBy", ctx.username},
+    });
+
+    sendJson(res, {
+        {"id", id},
+        {"publishedVersion", version},
+        {"publishedAt", patch["publishedAt"]},
+        {"publishedBy", ctx.username},
+    });
+}
+
 void WorkflowController::activateWorkflow(const httplib::Request& req, httplib::Response& res,
                                            const auth::AuthContext& ctx) {
     std::string id = req.matches[1];
@@ -485,11 +632,26 @@ void WorkflowController::activateWorkflow(const httplib::Request& req, httplib::
     // why it was switched off. Leaving the count would switch it off again on
     // the very next failure, which is not what someone reaching for Activate
     // means.
-    auto result = storage_.update("workflows", id,
-                                  {{"active", true},
-                                   {"consecutiveFailures", 0},
-                                   {"deactivatedReason", ""},
-                                   {"updatedAt", TimeUtils::nowMs()}}, 0, true);
+    // An active workflow runs what was published, so switching one on without
+    // 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.
+    auto before = storage_.get("workflows", id);
+    nlohmann::json patch = {{"active", true},
+                            {"consecutiveFailures", 0},
+                            {"deactivatedReason", ""},
+                            {"updatedAt", TimeUtils::nowMs()}};
+    bool published_now = false;
+    if (before.ok() && before.value().value("publishedVersion", int64_t{0}) <= 0) {
+        const int64_t version = before.value().value("_version", int64_t{0});
+        if (version > 0) {
+            patch["publishedVersion"] = version;
+            patch["publishedAt"] = TimeUtils::nowMs();
+            patch["publishedBy"] = ctx.username;
+            published_now = true;
+        }
+    }
+
+    auto result = storage_.update("workflows", id, patch, 0, true);
     if (result.failed()) {
         sendError(res, "Workflow not found", 404);
         return;
@@ -499,7 +661,10 @@ void WorkflowController::activateWorkflow(const httplib::Request& req, httplib::
     updateScheduledTriggers(id, true);
 
     ws_server_.broadcast("workflows.activated", {{"id", id}});
-    sendJson(res, {{"success", true}, {"active", true}});
+    sendJson(res, {{"success", true},
+                   {"active", true},
+                   {"publishedVersion", patch.value("publishedVersion", int64_t{0})},
+                   {"publishedOnActivate", published_now}});
 }
 
 void WorkflowController::deactivateWorkflow(const httplib::Request& req, httplib::Response& res,

+ 7 - 0
src/webserver/api/workflow_controller.hpp

@@ -43,6 +43,13 @@ private:
                          const auth::AuthContext& ctx);
     void activateWorkflow(const httplib::Request& req, httplib::Response& res,
                           const auth::AuthContext& ctx);
+    void listWorkflowVersions(const httplib::Request& req, httplib::Response& res,
+                              const auth::AuthContext& ctx);
+    void getWorkflowVersion(const httplib::Request& req, httplib::Response& res,
+                            const auth::AuthContext& ctx);
+    void publishWorkflow(const httplib::Request& req, httplib::Response& res,
+                         const auth::AuthContext& ctx);
+
     void deactivateWorkflow(const httplib::Request& req, httplib::Response& res,
                             const auth::AuthContext& ctx);
 

+ 25 - 0
src/webserver/webserver_service.cpp

@@ -154,6 +154,25 @@ void WebServerService::start() {
     // Start workflow scheduler
     scheduler_->start();
 
+    // Keep a history of workflow documents. Without it the version number
+    // still climbs on every write but nothing older is retained, so there is no
+    // earlier version to publish, compare against or go back to.
+    //
+    // Turned on here rather than at creation because the collection already
+    // exists on every install that predates this, and createCollection refuses
+    // a collection that is already there.
+    {
+        storage::CollectionConfig cfg;
+        cfg.versioning_enabled = true;
+        auto configured = storage_->configureCollection("workflows", cfg);
+        if (configured.failed()) {
+            LOG_WARN("Could not enable version history on workflows: {}",
+                     configured.error().message());
+        } else {
+            LOG_INFO("Version history enabled on workflows");
+        }
+    }
+
     // Ensure the executions summary view exists before serving requests
     ensureExecutionsSummaryView();
 
@@ -546,6 +565,8 @@ void WebServerService::runErrorWorkflow(const std::string& failed_workflow_id,
     proto::ExecuteWorkflowRequest request;
     request.set_workflow_id(handler_id);
     request.set_trigger_type("error-workflow");
+    // Started by the system, so it runs what was published.
+    request.set_use_published(true);
     request.set_trigger_data(trigger_data.dump());
     request.set_wait_for_completion(false);
 
@@ -583,6 +604,10 @@ void WebServerService::executeScheduledWorkflow(const std::string& workflow_id,
     proto::ExecuteWorkflowRequest request;
     request.set_workflow_id(workflow_id);
     request.set_trigger_type(trigger_type);
+    // A schedule firing is not somebody testing an edit: it runs the published
+    // version, so a half-finished change on somebody's canvas never goes live
+    // just because it was saved.
+    request.set_use_published(true);
 
     nlohmann::json trigger_data;
     trigger_data["triggerNodeId"] = trigger_node_id;