|
|
@@ -2,6 +2,8 @@
|
|
|
#include "common/uuid.hpp"
|
|
|
#include "common/time_utils.hpp"
|
|
|
#include "common/config_defaults.hpp"
|
|
|
+#include "storage/retention.hpp"
|
|
|
+#include "storage/workflow_collection.hpp"
|
|
|
#include "logging/logger.hpp"
|
|
|
#include <algorithm>
|
|
|
#include <functional>
|
|
|
@@ -56,10 +58,68 @@ std::string executionStatusToString(ExecutionStatus status) {
|
|
|
}
|
|
|
}
|
|
|
|
|
|
+std::optional<storage::Retention> WorkflowEngine::retentionFor(const Workflow& workflow) {
|
|
|
+ // A workflow that says nothing about retention inherits its project's. With
|
|
|
+ // no project there is nothing to inherit from - which is the case for
|
|
|
+ // anything created before projects existed.
|
|
|
+ if (const auto own = storage::declaredTtlSeconds(workflow.settings)) {
|
|
|
+ return storage::Retention{*own, false};
|
|
|
+ }
|
|
|
+ if (workflow.project_id.empty()) {
|
|
|
+ return std::nullopt;
|
|
|
+ }
|
|
|
+
|
|
|
+ nlohmann::json project_settings = nlohmann::json::object();
|
|
|
+ const int64_t now = TimeUtils::nowMs();
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(project_settings_mutex_);
|
|
|
+ auto it = project_settings_cache_.find(workflow.project_id);
|
|
|
+ if (it != project_settings_cache_.end() &&
|
|
|
+ now - it->second.fetched_at_ms < kProjectSettingsCacheMs) {
|
|
|
+ return storage::declaredRetention(workflow.settings, it->second.settings);
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
+ auto project = storage_.get("projects", workflow.project_id);
|
|
|
+ if (project.ok()) {
|
|
|
+ project_settings = project.value().value("settings", nlohmann::json::object());
|
|
|
+ } else {
|
|
|
+ // A project that cannot be read must not quietly become "keep for
|
|
|
+ // ever": say so, and cache the empty answer only briefly so a
|
|
|
+ // transient failure does not pin the wrong retention.
|
|
|
+ LOG_WARN("Retention: cannot read project {} for workflow {} ({}); this run keeps its "
|
|
|
+ "data on the default lifetime",
|
|
|
+ workflow.project_id, workflow.id, project.error().message());
|
|
|
+ }
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(project_settings_mutex_);
|
|
|
+ project_settings_cache_[workflow.project_id] = {project_settings, now};
|
|
|
+ }
|
|
|
+ return storage::declaredRetention(workflow.settings, project_settings);
|
|
|
+}
|
|
|
+
|
|
|
+void WorkflowEngine::ensureWorkflowCollection(const std::string& collection) {
|
|
|
+ {
|
|
|
+ std::lock_guard<std::mutex> lock(created_collections_mutex_);
|
|
|
+ if (created_workflow_collections_.contains(collection)) return;
|
|
|
+ }
|
|
|
+ // Not conditional on a listCollections check: creating one that already
|
|
|
+ // exists is harmless, and asking first costs a round trip on every run.
|
|
|
+ auto created = storage_.createCollection(collection);
|
|
|
+ if (created.failed()) {
|
|
|
+ // Not fatal - the insert that follows reports the real failure if the
|
|
|
+ // collection genuinely is not there.
|
|
|
+ LOG_DEBUG("Storage: could not create {} ({})", collection, created.error().message());
|
|
|
+ }
|
|
|
+ std::lock_guard<std::mutex> lock(created_collections_mutex_);
|
|
|
+ created_workflow_collections_.insert(collection);
|
|
|
+}
|
|
|
+
|
|
|
Workflow Workflow::fromJson(const nlohmann::json& j) {
|
|
|
Workflow wf;
|
|
|
wf.id = j.value("_id", j.value("id", ""));
|
|
|
wf.name = j.value("name", "");
|
|
|
+ wf.project_id = j.value("projectId", j.value("project_id", ""));
|
|
|
wf.active = j.value("active", false);
|
|
|
wf.settings = j.value("settings", nlohmann::json::object());
|
|
|
|
|
|
@@ -342,6 +402,7 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
nlohmann::json snapshot;
|
|
|
snapshot["id"] = workflow.id;
|
|
|
snapshot["name"] = workflow.name;
|
|
|
+ snapshot["projectId"] = workflow.project_id;
|
|
|
snapshot["settings"] = workflow.settings;
|
|
|
snapshot["nodes"] = nlohmann::json::array();
|
|
|
for (const auto& node : workflow.nodes) {
|
|
|
@@ -365,6 +426,14 @@ Result<ExecutionResult> WorkflowEngine::execute(const Workflow& workflow,
|
|
|
}
|
|
|
result.workflow_snapshot = snapshot;
|
|
|
|
|
|
+ // Resolved once, at the start, and carried on the result: the record is
|
|
|
+ // written several times over a run's life - before the walk, on pause, at
|
|
|
+ // the end - and every one of those writes has to agree about how long the
|
|
|
+ // record is kept, or the last one silently changes it.
|
|
|
+ if (const auto retention = retentionFor(workflow)) {
|
|
|
+ result.retention_ttl_ms = retention->ttlMs();
|
|
|
+ }
|
|
|
+
|
|
|
// Write the record before the walk starts, not only when it ends.
|
|
|
//
|
|
|
// Until this, an execution existed in the database only once it had
|
|
|
@@ -1301,6 +1370,13 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
return result;
|
|
|
}
|
|
|
|
|
|
+ // Nothing declared leaves documents on the lifetime they have always had -
|
|
|
+ // for ever - so an installation that has not chosen a retention keeps
|
|
|
+ // behaving exactly as it did.
|
|
|
+ const auto declared_retention = retentionFor(workflow);
|
|
|
+ const int64_t retention_ttl_ms = declared_retention ? declared_retention->ttlMs() : 0;
|
|
|
+ const std::string own_collection = storage::workflowCollectionName(workflow.id);
|
|
|
+
|
|
|
// Helper to get per-workflow storage permission for a collection
|
|
|
// Returns: "none", "read-only", or "read-write"
|
|
|
auto getWorkflowAccess = [&workflow, this](const std::string& collection) -> std::string {
|
|
|
@@ -1309,6 +1385,27 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
return "none";
|
|
|
}
|
|
|
|
|
|
+ // A private collection is settled before the settings are consulted at
|
|
|
+ // all. Another workflow's is refused whatever the settings say - which
|
|
|
+ // matters most for defaultAccess, since a workflow set to read-write by
|
|
|
+ // default would otherwise reach every other workflow's private data,
|
|
|
+ // across projects, just by naming it.
|
|
|
+ if (storage::isWorkflowCollection(collection)) {
|
|
|
+ if (!storage::isWorkflowCollectionOf(collection, workflow.id)) {
|
|
|
+ return "none";
|
|
|
+ }
|
|
|
+ // Its owner holds it read-write without asking. An entry in the
|
|
|
+ // settings still wins, so the grant can be given up deliberately.
|
|
|
+ if (workflow.settings.contains("storagePermissions")) {
|
|
|
+ const auto& perms = workflow.settings["storagePermissions"];
|
|
|
+ if (perms.contains("collections") && perms["collections"].is_object() &&
|
|
|
+ perms["collections"].contains(collection)) {
|
|
|
+ return perms["collections"][collection].get<std::string>();
|
|
|
+ }
|
|
|
+ }
|
|
|
+ return "read-write";
|
|
|
+ }
|
|
|
+
|
|
|
// Check workflow settings for storage permissions
|
|
|
if (workflow.settings.contains("storagePermissions")) {
|
|
|
const auto& storage_perms = workflow.settings["storagePermissions"];
|
|
|
@@ -1359,51 +1456,94 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
}
|
|
|
};
|
|
|
|
|
|
+ // "@self" means this workflow's own collection, whoever is running. Every
|
|
|
+ // storage callback resolves it the same way and before the permission
|
|
|
+ // check, so the name cannot be used to reach anything else.
|
|
|
+ auto resolveCollection = [own_collection](const std::string& collection,
|
|
|
+ std::string& resolved) -> common::Result<void> {
|
|
|
+ if (!storage::isSelfCollection(collection)) {
|
|
|
+ resolved = collection;
|
|
|
+ return {};
|
|
|
+ }
|
|
|
+ if (own_collection.empty()) {
|
|
|
+ return common::Error(common::ErrorCode::InvalidArgument,
|
|
|
+ "\"@self\" means this workflow's own storage, and this run has no workflow id");
|
|
|
+ }
|
|
|
+ resolved = own_collection;
|
|
|
+ return {};
|
|
|
+ };
|
|
|
+
|
|
|
// Storage API callbacks with per-workflow permission checks
|
|
|
- ctx.storage_get_doc = [this, canReadCollection](const std::string& collection, const std::string& id)
|
|
|
+ ctx.storage_get_doc = [this, canReadCollection, resolveCollection](
|
|
|
+ const std::string& collection, const std::string& id)
|
|
|
-> common::Result<nlohmann::json> {
|
|
|
- if (!canReadCollection(collection)) {
|
|
|
+ std::string target;
|
|
|
+ if (auto r = resolveCollection(collection, target); r.failed()) return r.error();
|
|
|
+ if (!canReadCollection(target)) {
|
|
|
return common::Error(common::ErrorCode::PermissionDenied,
|
|
|
- "No read access to collection: " + collection);
|
|
|
+ "No read access to collection: " + target);
|
|
|
}
|
|
|
- return storage_.get(collection, id);
|
|
|
+ return storage_.get(target, id);
|
|
|
};
|
|
|
|
|
|
- ctx.storage_insert = [this, canWriteCollection](const std::string& collection, const nlohmann::json& data,
|
|
|
+ ctx.storage_insert = [this, canWriteCollection, resolveCollection, retention_ttl_ms, own_collection](
|
|
|
+ const std::string& collection, const nlohmann::json& data,
|
|
|
const std::string& id, int64_t ttl_ms)
|
|
|
-> common::Result<std::string> {
|
|
|
- if (!canWriteCollection(collection)) {
|
|
|
+ std::string target;
|
|
|
+ if (auto r = resolveCollection(collection, target); r.failed()) return r.error();
|
|
|
+ if (!canWriteCollection(target)) {
|
|
|
return common::Error(common::ErrorCode::PermissionDenied,
|
|
|
- "No write access to collection: " + collection);
|
|
|
- }
|
|
|
- return storage_.insert(collection, data, id, ttl_ms);
|
|
|
+ "No write access to collection: " + target);
|
|
|
+ }
|
|
|
+ // A workflow's own collection is made on first write rather than when
|
|
|
+ // the workflow is saved: a workflow that never stores anything should
|
|
|
+ // not leave an empty collection behind, and one written to by a runner
|
|
|
+ // that has never seen it before must still work.
|
|
|
+ if (storage::collectionBaseName(target) == own_collection) {
|
|
|
+ ensureWorkflowCollection(target);
|
|
|
+ }
|
|
|
+ // A TTL the node asked for wins - it knows what it wrote. Otherwise the
|
|
|
+ // workflow's retention applies, so data written by a node that never
|
|
|
+ // heard of retention still ages out.
|
|
|
+ const int64_t effective_ttl = ttl_ms > 0 ? ttl_ms : retention_ttl_ms;
|
|
|
+ return storage_.insert(target, data, id, effective_ttl);
|
|
|
};
|
|
|
|
|
|
- ctx.storage_update = [this, canWriteCollection](const std::string& collection, const std::string& id,
|
|
|
+ ctx.storage_update = [this, canWriteCollection, resolveCollection](
|
|
|
+ const std::string& collection, const std::string& id,
|
|
|
const nlohmann::json& data, int64_t expected_version, bool partial)
|
|
|
-> common::Result<int64_t> {
|
|
|
- if (!canWriteCollection(collection)) {
|
|
|
+ std::string target;
|
|
|
+ if (auto r = resolveCollection(collection, target); r.failed()) return r.error();
|
|
|
+ if (!canWriteCollection(target)) {
|
|
|
return common::Error(common::ErrorCode::PermissionDenied,
|
|
|
- "No write access to collection: " + collection);
|
|
|
+ "No write access to collection: " + target);
|
|
|
}
|
|
|
- return storage_.update(collection, id, data, expected_version, partial);
|
|
|
+ return storage_.update(target, id, data, expected_version, partial);
|
|
|
};
|
|
|
|
|
|
- ctx.storage_delete = [this, canWriteCollection](const std::string& collection, const std::string& id,
|
|
|
+ ctx.storage_delete = [this, canWriteCollection, resolveCollection](
|
|
|
+ const std::string& collection, const std::string& id,
|
|
|
int64_t expected_version)
|
|
|
-> common::Result<void> {
|
|
|
- if (!canWriteCollection(collection)) {
|
|
|
+ std::string target;
|
|
|
+ if (auto r = resolveCollection(collection, target); r.failed()) return r.error();
|
|
|
+ if (!canWriteCollection(target)) {
|
|
|
return common::Error(common::ErrorCode::PermissionDenied,
|
|
|
- "No write access to collection: " + collection);
|
|
|
+ "No write access to collection: " + target);
|
|
|
}
|
|
|
- return storage_.remove(collection, id, expected_version);
|
|
|
+ return storage_.remove(target, id, expected_version);
|
|
|
};
|
|
|
|
|
|
- ctx.storage_query = [this, canReadCollection](const std::string& collection, const engine::StorageQueryOptions& options)
|
|
|
+ ctx.storage_query = [this, canReadCollection, resolveCollection](
|
|
|
+ const std::string& collection, const engine::StorageQueryOptions& options)
|
|
|
-> common::Result<engine::StorageQueryResult> {
|
|
|
- if (!canReadCollection(collection)) {
|
|
|
+ std::string target;
|
|
|
+ if (auto r = resolveCollection(collection, target); r.failed()) return r.error();
|
|
|
+ if (!canReadCollection(target)) {
|
|
|
return common::Error(common::ErrorCode::PermissionDenied,
|
|
|
- "No read access to collection: " + collection);
|
|
|
+ "No read access to collection: " + target);
|
|
|
}
|
|
|
|
|
|
// Convert engine query options to storage query options
|
|
|
@@ -1414,7 +1554,7 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
storage_opts.page = options.page;
|
|
|
storage_opts.page_size = options.page_size;
|
|
|
|
|
|
- auto result = storage_.query(collection, storage_opts);
|
|
|
+ auto result = storage_.query(target, storage_opts);
|
|
|
if (result.failed()) {
|
|
|
return common::Error(result.error().code(), result.error().message());
|
|
|
}
|
|
|
@@ -1426,16 +1566,30 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
return engine_result;
|
|
|
};
|
|
|
|
|
|
- ctx.storage_list_collections = [this, canReadCollection](bool include_system)
|
|
|
+ ctx.storage_list_collections = [this, canReadCollection, own_collection](bool include_system)
|
|
|
-> common::Result<std::vector<std::string>> {
|
|
|
auto all_collections = storage_.listCollections();
|
|
|
std::vector<std::string> accessible;
|
|
|
|
|
|
+ // A workflow's own storage is offered before it exists - it is created
|
|
|
+ // on first write, and a collection that only appears once you have
|
|
|
+ // already written to it cannot be picked from a list.
|
|
|
+ if (!own_collection.empty()) {
|
|
|
+ accessible.push_back(storage::kSelfCollection);
|
|
|
+ }
|
|
|
+
|
|
|
for (const auto& coll : all_collections) {
|
|
|
// Skip system collections unless explicitly requested
|
|
|
if (!include_system && collection_permissions_->isSystemCollection(coll)) {
|
|
|
continue;
|
|
|
}
|
|
|
+ // Private workflow storage is never offered under its raw name:
|
|
|
+ // another workflow's is not reachable anyway, and this workflow's
|
|
|
+ // own is already in the list as "@self". Offering both would put
|
|
|
+ // the same collection in the list twice under two names.
|
|
|
+ if (storage::isWorkflowCollection(coll)) {
|
|
|
+ continue;
|
|
|
+ }
|
|
|
// Only include collections the workflow has at least read access to
|
|
|
if (canReadCollection(coll)) {
|
|
|
accessible.push_back(coll);
|
|
|
@@ -1448,7 +1602,7 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
// Files sit outside the collection namespace, so they are gated on a reserved
|
|
|
// "files" permission entry. Deny-by-default, consistent with collections: a
|
|
|
// workflow must declare settings.storagePermissions.collections.files.
|
|
|
- ctx.storage_upload_file = [this, canWriteCollection](const std::vector<uint8_t>& data,
|
|
|
+ ctx.storage_upload_file = [this, canWriteCollection, retention_ttl_ms](const std::vector<uint8_t>& data,
|
|
|
const engine::ScriptFileMeta& meta)
|
|
|
-> common::Result<engine::ScriptFileInfo> {
|
|
|
if (!canWriteCollection("files")) {
|
|
|
@@ -1464,7 +1618,10 @@ NodeExecutionResult WorkflowEngine::executeNode(const WorkflowNode& node,
|
|
|
upstream.is_public = meta.is_public;
|
|
|
upstream.metadata = meta.metadata;
|
|
|
|
|
|
- auto result = storage_.uploadFile(data, upstream);
|
|
|
+ // The file expires with the document that points at it. A stored file
|
|
|
+ // whose record has aged out is unreachable weight, and a record whose
|
|
|
+ // file has gone is a broken link - they have to share one lifetime.
|
|
|
+ auto result = storage_.uploadFile(data, upstream, retention_ttl_ms);
|
|
|
if (result.failed()) {
|
|
|
return common::Error(result.error().code(), result.error().message());
|
|
|
}
|
|
|
@@ -2165,8 +2322,19 @@ void WorkflowEngine::storeExecution(const ExecutionResult& result) {
|
|
|
// deadline", not "no need to extend" - it is given the 30-day ceiling the
|
|
|
// approval node caps at, so a permanent approval does not fall back to
|
|
|
// the ordinary seven-day log TTL and evaporate with no trace.
|
|
|
+ //
|
|
|
+ // A workflow that sets its own retention overrides all of that: its history
|
|
|
+ // and the data it wrote age out together, which is the point of the setting.
|
|
|
+ // An explicit "keep for ever" has to survive the waiting floor below, so it
|
|
|
+ // is checked rather than compared - a 0 fed into a `wanted > ttl_ms` test
|
|
|
+ // reads as the shortest possible lifetime instead of the longest.
|
|
|
int64_t ttl_ms = 7 * 24 * 60 * 60 * 1000;
|
|
|
- if (result.status == ExecutionStatus::Waiting) {
|
|
|
+ if (result.retention_ttl_ms >= 0) {
|
|
|
+ ttl_ms = result.retention_ttl_ms;
|
|
|
+ }
|
|
|
+ const bool keep_forever = (ttl_ms == 0);
|
|
|
+
|
|
|
+ if (result.status == ExecutionStatus::Waiting && !keep_forever) {
|
|
|
const int64_t grace = 24 * 60 * 60 * 1000;
|
|
|
const int64_t never_ttl = 30LL * 24 * 60 * 60 * 1000;
|
|
|
int64_t wanted = never_ttl;
|