Parcourir la source

fix(runner): stop persisting binary payloads in execution records

Every node that touched an image echoed its base64 into its output, and
ExecutionResult::toJson wrote all of them, so one execution reached 3.4 MB:
1.76 MB from the fetch, then the same bytes again from the vision, dedupe,
branch and safety-gate nodes. One of these was written every ten minutes,
which is what drove the database into mass eviction and made documents
unrecoverable.

An execution record exists to explain what happened, not to carry data, so a
string longer than 8 KB is now recorded as a note of its size. Nodes still
receive the real values at runtime; only the stored record is summarised.

Measured on the same workflow: 3,393,162 bytes to 66,612, with the largest
remaining entry being the RSS feed itself at 22 KB.
fszontagh il y a 1 mois
Parent
commit
eb0fc8f7f1
1 fichiers modifiés avec 42 ajouts et 3 suppressions
  1. 42 3
      src/runner/workflow_engine.cpp

+ 42 - 3
src/runner/workflow_engine.cpp

@@ -67,6 +67,45 @@ Workflow Workflow::fromJson(const nlohmann::json& j) {
     return wf;
 }
 
+namespace {
+
+// Execution records are for observability, not data transport. A node that
+// carries binary content - a fetched image, a generated one - emits base64 that
+// is megabytes wide, and persisting every such output has pushed multi-megabyte
+// documents into the database on every run. Long strings are recorded as a note
+// of what was there instead.
+constexpr size_t kMaxPersistedStringBytes = 8192;
+
+nlohmann::json truncateLargeValues(const nlohmann::json& value) {
+    if (value.is_string()) {
+        const auto& text = value.get_ref<const std::string&>();
+        if (text.size() > kMaxPersistedStringBytes) {
+            return "[omitted " + std::to_string(text.size()) + " bytes]";
+        }
+        return value;
+    }
+
+    if (value.is_array()) {
+        nlohmann::json out = nlohmann::json::array();
+        for (const auto& item : value) {
+            out.push_back(truncateLargeValues(item));
+        }
+        return out;
+    }
+
+    if (value.is_object()) {
+        nlohmann::json out = nlohmann::json::object();
+        for (auto it = value.begin(); it != value.end(); ++it) {
+            out[it.key()] = truncateLargeValues(it.value());
+        }
+        return out;
+    }
+
+    return value;
+}
+
+} // namespace
+
 nlohmann::json ExecutionResult::toJson() const {
     nlohmann::json j;
     // Note: id is managed by database as _id (execution_id is passed to storage_.insert)
@@ -74,11 +113,11 @@ nlohmann::json ExecutionResult::toJson() const {
     j["workflowName"] = workflow_name;
     j["status"] = executionStatusToString(status);
     j["triggerType"] = trigger_type;
-    j["triggerData"] = trigger_data;
+    j["triggerData"] = truncateLargeValues(trigger_data);
     j["startedAt"] = started_at;
     j["finishedAt"] = finished_at;
     j["error"] = error;
-    j["output"] = final_output;
+    j["output"] = truncateLargeValues(final_output);
 
     j["nodeExecutions"] = nlohmann::json::array();
     for (const auto& [id, result] : node_results) {
@@ -89,7 +128,7 @@ nlohmann::json ExecutionResult::toJson() const {
         nr["finishedAt"] = result.finished_at;
         // Note: input is intentionally not stored to avoid data duplication
         // Each node's input can be reconstructed from upstream node outputs + connections
-        nr["output"] = result.output;
+        nr["output"] = truncateLargeValues(result.output);
         nr["error"] = result.error;
         nr["retryCount"] = result.retry_count;
         j["nodeExecutions"].push_back(nr);