|
|
@@ -33,7 +33,8 @@ 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});
|
|
|
+ queue_.push_back(Job{workflow_id, ttl_seconds, /*drop_collection_after=*/false,
|
|
|
+ /*from_creation=*/true});
|
|
|
}
|
|
|
cv_.notify_one();
|
|
|
}
|
|
|
@@ -42,7 +43,8 @@ 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});
|
|
|
+ queue_.push_back(Job{workflow_id, kRetireTtlSeconds, /*drop_collection_after=*/true,
|
|
|
+ /*from_creation=*/false});
|
|
|
}
|
|
|
cv_.notify_one();
|
|
|
}
|
|
|
@@ -61,8 +63,9 @@ void RetentionService::worker() {
|
|
|
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);
|
|
|
+ restamp(kExecutions, kExecutionWorkflowField, job.workflow_id, job.ttl_seconds,
|
|
|
+ job.from_creation);
|
|
|
+ const int64_t documents = restamp(own, kNoWorkflowField, job.workflow_id, job.ttl_seconds, job.from_creation);
|
|
|
|
|
|
LOG_INFO("Retention: workflow {} restamped to {}s - {} executions, {} documents",
|
|
|
job.workflow_id, job.ttl_seconds, executions, documents);
|
|
|
@@ -164,69 +167,96 @@ RetentionService::Preview RetentionService::preview(const std::string& workflow_
|
|
|
int64_t RetentionService::restamp(const std::string& collection,
|
|
|
const std::string& workflow_field,
|
|
|
const std::string& workflow_id,
|
|
|
- int64_t ttl_seconds) {
|
|
|
+ int64_t ttl_seconds,
|
|
|
+ bool from_creation) {
|
|
|
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.
|
|
|
+ // Every id already handled, so a document is not restamped twice and a
|
|
|
+ // sweep can tell whether it found anything new.
|
|
|
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);
|
|
|
- }
|
|
|
+ // Full sweeps until one finds nothing new.
|
|
|
+ //
|
|
|
+ // A TTL a document has already outlived means deleting it, and deleting as
|
|
|
+ // the sweep runs shifts everything behind it forward - so paging straight
|
|
|
+ // through would step over the documents that moved into the gap. A second
|
|
|
+ // sweep finds them, at their new positions, and the id set makes revisiting
|
|
|
+ // the rest harmless. Only worth repeating when something was actually
|
|
|
+ // deleted: with nothing removed, nothing moved.
|
|
|
+ bool sweep_again = true;
|
|
|
+ while (sweep_again && !stopping_) {
|
|
|
+ sweep_again = false;
|
|
|
+ int64_t deleted_this_sweep = 0;
|
|
|
+
|
|
|
+ for (int32_t page = 1; !stopping_; ++page) {
|
|
|
+ storage::QueryOptions opts;
|
|
|
+ opts.page = page;
|
|
|
+ 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());
|
|
|
+ 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;
|
|
|
}
|
|
|
- return restamped;
|
|
|
- }
|
|
|
- if (result.value().documents.empty()) break;
|
|
|
+ if (result.value().documents.empty()) break;
|
|
|
+
|
|
|
+ const int64_t now = common::TimeUtils::nowMs();
|
|
|
+ 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;
|
|
|
+ done.insert(id);
|
|
|
+
|
|
|
+ // A retention is "kept for N after it was written", not "kept
|
|
|
+ // for another N starting now". The database measures a TTL from
|
|
|
+ // the moment it is set, so the remaining life is worked out here
|
|
|
+ // - otherwise applying a seven-day retention to a two-year-old
|
|
|
+ // document would grant it another week rather than removing it,
|
|
|
+ // and the count offered beforehand would have been a fiction.
|
|
|
+ int64_t ttl_ms = ttl_seconds * 1000;
|
|
|
+ if (from_creation && ttl_seconds > 0) {
|
|
|
+ const int64_t created = common::TimeUtils::documentStamp(doc, "_created_at");
|
|
|
+ if (created > 0) {
|
|
|
+ const int64_t remaining = created + ttl_ms - now;
|
|
|
+ if (remaining <= 0) {
|
|
|
+ // Its life is already over. Waiting for an expiry
|
|
|
+ // that is in the past would mean waiting for ever.
|
|
|
+ if (storage_.remove(collection, id).ok()) {
|
|
|
+ restamped++;
|
|
|
+ deleted_this_sweep++;
|
|
|
+ }
|
|
|
+ continue;
|
|
|
+ }
|
|
|
+ ttl_ms = remaining;
|
|
|
+ }
|
|
|
+ }
|
|
|
|
|
|
- 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++;
|
|
|
+ // 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 everything
|
|
|
+ // would be worse than the growth it was meant to fix.
|
|
|
+ auto written = storage_.upsert(collection, doc, id, ttl_ms);
|
|
|
+ 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;
|
|
|
+ // A short page is the last page.
|
|
|
+ if (result.value().documents.size() < static_cast<size_t>(kPageSize)) 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;
|
|
|
+ if (deleted_this_sweep > 0) sweep_again = true;
|
|
|
}
|
|
|
|
|
|
return restamped;
|