Skip to content

Commit 4f5e9b3

Browse files
committed
Keep usage event records of running apps, service instances, and tasks
The scheduled usage event cleanup job used to delete every record older than the configured cutoff age, including the opening STARTED/CREATED event of a resource that is still running. Once the cleanup deleted that event, nothing was left to reconstruct what is running right now. Database::OldRecordCleanup can now optionally keep the records of running resources. Each model declares its lifecycles via usage_lifecycles: which states open a run (STARTED/CREATED/TASK_STARTED, plus the WAS_RUNNING/TASK_WAS_RUNNING baselines), which state closes it (STOPPED/DELETED/TASK_STOPPED), and which column names the resource. An old opening event is then only deleted when: * a closing event for the same resource exists later and is also old -- the run is over; or * it is neither the first opening of the current run nor the resource's latest one (again judged only against old rows). Consumers only need the first opening (the true start time) and the latest (the current size). The ones in between, written each time a running resource is scaled or updated, tell a consumer nothing it still needs -- and deleting them is what keeps the table size bounded for long-running, frequently-changed resources. The app and service usage event repositories turn this on with keep_running_records: true. Asking for it on a model without usage_lifecycles raises an error instead of silently deleting the records of running resources. Task events get their own lifecycle (TASK_STARTED/TASK_WAS_RUNNING -> TASK_STOPPED, matched by task_guid), so the start events of long-running tasks survive cleanup too. Task baselines use their own TASK_WAS_RUNNING state because task events carry an empty app_guid: if they said WAS_RUNNING, the app lifecycle would see them all as events of one app whose guid is '' and wrongly delete them (and the backfill's repair would write bogus STOPPED events for that phantom app). Deletion runs in a deliberate order: first the opening events that are safe to delete, while the events that make them safe still exist; then everything else. The reverse order could delete a closing event first and leave its opening event looking like a still-running resource. The cleanup log line now reports the row counts BatchDelete returns instead of running extra COUNT queries, and BatchDelete fetches each batch's ids in the same query that checks whether anything is left, halving the evaluations of the (potentially expensive) filtered dataset. Also renames the positional days_ago argument to a cutoff_age_in_days keyword.
1 parent 20fb68b commit 4f5e9b3

12 files changed

Lines changed: 603 additions & 33 deletions

‎app/jobs/runtime/events_cleanup.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ def initialize(cutoff_age_in_days)
99
end
1010

1111
def perform
12-
Database::OldRecordCleanup.new(Event, cutoff_age_in_days).delete
12+
Database::OldRecordCleanup.new(Event, cutoff_age_in_days:).delete
1313
end
1414

1515
def job_name_in_configuration

‎app/models/runtime/app_usage_event.rb‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,23 @@ class AppUsageEvent < Sequel::Model
99
:buildpack_guid, :buildpack_name,
1010
:package_state, :previous_package_state, :parent_app_guid,
1111
:parent_app_name, :process_type, :task_name, :task_guid
12+
13+
def self.usage_lifecycles
14+
[
15+
{
16+
beginning_states: [ProcessModel::STARTED, Repositories::AppUsageEventRepository::WAS_RUNNING_EVENT_STATE],
17+
ending_state: ProcessModel::STOPPED,
18+
guid_column: :app_guid
19+
},
20+
{
21+
beginning_states: [Repositories::AppUsageEventRepository::TASK_STARTED_EVENT_STATE,
22+
Repositories::AppUsageEventRepository::TASK_WAS_RUNNING_EVENT_STATE],
23+
ending_state: Repositories::AppUsageEventRepository::TASK_STOPPED_EVENT_STATE,
24+
guid_column: :task_guid
25+
}
26+
].freeze
27+
end
28+
1229
AppUsageEvent.dataset_module do
1330
def supports_window_functions?
1431
false

‎app/models/services/service_usage_event.rb‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,5 +7,17 @@ class ServiceUsageEvent < Sequel::Model
77
:service_plan_guid, :service_plan_name,
88
:service_guid, :service_label,
99
:service_broker_name, :service_broker_guid
10+
11+
def self.usage_lifecycles
12+
[
13+
{
14+
beginning_states: [Repositories::ServiceUsageEventRepository::CREATED_EVENT_STATE,
15+
Repositories::ServiceUsageEventRepository::UPDATED_EVENT_STATE,
16+
Repositories::ServiceUsageEventRepository::WAS_RUNNING_EVENT_STATE],
17+
ending_state: Repositories::ServiceUsageEventRepository::DELETED_EVENT_STATE,
18+
guid_column: :service_instance_guid
19+
}
20+
].freeze
21+
end
1022
end
1123
end

‎app/repositories/app_usage_event_repository.rb‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,18 @@
44
module VCAP::CloudController
55
module Repositories
66
class AppUsageEventRepository
7+
WAS_RUNNING_EVENT_STATE = 'WAS_RUNNING'.freeze
8+
TASK_STARTED_EVENT_STATE = 'TASK_STARTED'.freeze
9+
TASK_STOPPED_EVENT_STATE = 'TASK_STOPPED'.freeze
10+
# Task baselines get their own state (rather than reusing WAS_RUNNING)
11+
# because task events share the app_usage_events table with app events but
12+
# carry an empty app_guid. If task baselines said WAS_RUNNING, the cleanup
13+
# and the backfill's repair would both treat every task baseline as
14+
# belonging to a single app whose guid is '' -- the cleanup would wrongly
15+
# prune them, and the repair would write bogus STOPPED events for that
16+
# phantom app.
17+
TASK_WAS_RUNNING_EVENT_STATE = 'TASK_WAS_RUNNING'.freeze
18+
719
def find(guid)
820
AppUsageEvent.find(guid:)
921
end
@@ -152,7 +164,7 @@ def purge_and_reseed_started_apps!
152164
end
153165

154166
def delete_events_older_than(cutoff_age_in_days)
155-
Database::OldRecordCleanup.new(AppUsageEvent, cutoff_age_in_days, keep_at_least_one_record: true).delete
167+
Database::OldRecordCleanup.new(AppUsageEvent, cutoff_age_in_days: cutoff_age_in_days, keep_at_least_one_record: true, keep_running_records: true).delete
156168
end
157169

158170
private

‎app/repositories/service_usage_event_repository.rb‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ class ServiceUsageEventRepository
77
DELETED_EVENT_STATE = 'DELETED'.freeze
88
CREATED_EVENT_STATE = 'CREATED'.freeze
99
UPDATED_EVENT_STATE = 'UPDATED'.freeze
10+
WAS_RUNNING_EVENT_STATE = 'WAS_RUNNING'.freeze
1011

1112
def find(guid)
1213
ServiceUsageEvent.find(guid:)
@@ -92,7 +93,7 @@ def purge_and_reseed_service_instances!
9293
end
9394

9495
def delete_events_older_than(cutoff_age_in_days)
95-
Database::OldRecordCleanup.new(ServiceUsageEvent, cutoff_age_in_days, keep_at_least_one_record: true).delete
96+
Database::OldRecordCleanup.new(ServiceUsageEvent, cutoff_age_in_days: cutoff_age_in_days, keep_at_least_one_record: true, keep_running_records: true).delete
9697
end
9798
end
9899
end

‎lib/database/batch_delete.rb‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -11,19 +11,21 @@ def delete
1111
total_count = 0
1212

1313
loop do
14-
set = dataset.limit(amount)
15-
break if set.empty?
14+
# Fetch the batch's ids in the same query that checks for emptiness, so the
15+
# (potentially expensive) filtered dataset is evaluated once per batch.
16+
ids = dataset.limit(amount).select_map(:id)
17+
break if ids.empty?
1618

17-
total_count += delete_batch(set)
19+
total_count += delete_batch(ids)
1820
end
1921

2022
total_count
2123
end
2224

2325
private
2426

25-
def delete_batch(set)
26-
dataset.model.where(id: set.select_map(:id)).delete
27+
def delete_batch(ids)
28+
dataset.model.where(id: ids).delete
2729
end
2830
end
2931
end

‎lib/database/old_record_cleanup.rb‎

Lines changed: 118 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -3,25 +3,36 @@
33
module Database
44
class OldRecordCleanup
55
class NoCurrentTimestampError < StandardError; end
6-
attr_reader :model, :days_ago, :keep_at_least_one_record
6+
attr_reader :model, :cutoff_age_in_days, :keep_at_least_one_record, :keep_running_records
77

8-
def initialize(model, days_ago, keep_at_least_one_record: false)
8+
def initialize(model, cutoff_age_in_days:, keep_at_least_one_record: false, keep_running_records: false)
99
@model = model
10-
@days_ago = days_ago
10+
@cutoff_age_in_days = cutoff_age_in_days
1111
@keep_at_least_one_record = keep_at_least_one_record
12+
@keep_running_records = keep_running_records
13+
return unless keep_running_records && !model.respond_to?(:usage_lifecycles)
14+
15+
# Fail here, when the caller builds the object, not later when the
16+
# cleanup job runs and nobody is watching.
17+
raise ArgumentError.new("keep_running_records requires #{model} to define .usage_lifecycles")
1218
end
1319

20+
# The two options work together: keep_running_records always keeps the
21+
# beginning events that show a resource is still running (those with no
22+
# later ending event for the same resource), and keep_at_least_one_record
23+
# additionally keeps the single newest row, so the table is never fully
24+
# emptied for clients that poll the most recent event.
1425
def delete
15-
cutoff_date = current_timestamp_from_database - days_ago.to_i.days
16-
26+
cutoff_date = current_timestamp_from_database - cutoff_age_in_days.to_i.days
1727
old_records = model.dataset.where(Sequel.lit('created_at < ?', cutoff_date))
18-
if keep_at_least_one_record
19-
last_record = model.order(:id).last
20-
old_records = old_records.where(Sequel.lit('id < ?', last_record.id)) if last_record
21-
end
22-
logger.info("Cleaning up #{old_records.count} #{model.table_name} table rows")
2328

24-
Database::BatchDelete.new(old_records, 1000).delete
29+
if keep_running_records
30+
delete_keeping_running_records(old_records)
31+
else
32+
old_records = exclude_newest_record(old_records)
33+
logger.info("Cleaning up #{old_records.count} #{model.table_name} table rows")
34+
Database::BatchDelete.new(old_records, 1000).delete
35+
end
2536
end
2637

2738
private
@@ -35,5 +46,101 @@ def current_timestamp_from_database
3546
def logger
3647
@logger ||= Steno.logger('cc.old_record_cleanup')
3748
end
49+
50+
# Deletes old records while retaining a usable billing baseline for
51+
# still-running resources.
52+
#
53+
# For each lifecycle of the model, a beginning-state row (e.g.
54+
# STARTED/CREATED/WAS_RUNNING) is prunable when:
55+
# * a later ending-state row (e.g. STOPPED/DELETED) is also old -- the run is
56+
# over; or
57+
# * it is a superseded baseline: an earlier beginning of the same run and a
58+
# later beginning both exist (and are old). Consumers only need the first
59+
# beginning of the current run (the true start time) and the latest one
60+
# (the current size); the rows in between, written each time a running
61+
# resource is scaled or updated, tell a consumer nothing it still needs.
62+
#
63+
# The deletes are ordered deliberately: prunable beginning rows are removed
64+
# FIRST, while the rows that make them prunable still exist, so each beginning
65+
# stays prunable until it is itself deleted. Only then are the ending rows (and
66+
# any other, non-lifecycle states) removed. Reversing the order could strand a
67+
# beginning row whose paired ending was deleted in an earlier batch.
68+
def delete_keeping_running_records(old_records)
69+
lifecycles = model.usage_lifecycles
70+
prunable_beginnings = lifecycles.map { |lifecycle| prunable_beginnings_dataset(old_records, lifecycle) }
71+
72+
# Everything that is not a beginning-state row of some lifecycle (ending rows
73+
# plus any other, non-lifecycle states) is unconditionally prunable.
74+
all_beginning_states = lifecycles.flat_map { |lifecycle| lifecycle.fetch(:beginning_states) }
75+
unconditional_records = exclude_newest_record(old_records.exclude(state: all_beginning_states))
76+
77+
deleted_count = prunable_beginnings.sum { |dataset| Database::BatchDelete.new(dataset, 1000).delete }
78+
deleted_count += Database::BatchDelete.new(unconditional_records, 1000).delete
79+
80+
logger.info("Cleaned up #{deleted_count} #{model.table_name} table rows")
81+
end
82+
83+
# Builds the dataset of old beginning-state rows that are prunable for one
84+
# lifecycle. All the subqueries use the (state, guid, id) lifecycle index,
85+
# and a higher id always means created later. They only look at OLD rows: a
86+
# superseded beginning is kept until the row that supersedes it is itself
87+
# old, so a consumer reading within the retention window never sees the
88+
# decision change underneath it.
89+
def prunable_beginnings_dataset(old_records, lifecycle)
90+
beginning_states = lifecycle.fetch(:beginning_states)
91+
ending_state = lifecycle.fetch(:ending_state)
92+
guid_column = lifecycle.fetch(:guid_column)
93+
94+
old_beginnings = old_records.where(state: beginning_states)
95+
old_endings = old_records.where(state: ending_state)
96+
initial_records = old_beginnings.from_self(alias: :initial_records)
97+
98+
# The run is over: an ending row for the same resource was created later.
99+
matching_ending = old_endings.from_self(alias: :final_records).
100+
where(Sequel[:final_records][guid_column] => Sequel[:initial_records][guid_column]).
101+
where { Sequel[:final_records][:id] > Sequel[:initial_records][:id] }.
102+
select(1).exists
103+
104+
# Not the run's true start: an earlier beginning of the same run exists,
105+
# i.e. one with no ending event between the two.
106+
intervening_ending = old_endings.from_self(alias: :intervening_endings).
107+
where(Sequel[:intervening_endings][guid_column] => Sequel[:earlier_beginnings][guid_column]).
108+
where { Sequel[:intervening_endings][:id] > Sequel[:earlier_beginnings][:id] }.
109+
where { Sequel[:intervening_endings][:id] < Sequel[:initial_records][:id] }.
110+
select(1).exists
111+
earlier_beginning_in_same_run = old_beginnings.from_self(alias: :earlier_beginnings).
112+
where(Sequel[:earlier_beginnings][guid_column] => Sequel[:initial_records][guid_column]).
113+
where { Sequel[:earlier_beginnings][:id] < Sequel[:initial_records][:id] }.
114+
where(Sequel.~(intervening_ending)).
115+
select(1).exists
116+
117+
# Not the latest baseline: a later beginning for the same resource exists.
118+
later_beginning = old_beginnings.from_self(alias: :later_beginnings).
119+
where(Sequel[:later_beginnings][guid_column] => Sequel[:initial_records][guid_column]).
120+
where { Sequel[:later_beginnings][:id] > Sequel[:initial_records][:id] }.
121+
select(1).exists
122+
123+
superseded_baseline = Sequel.&(earlier_beginning_in_same_run, later_beginning)
124+
# exclude_newest_record removes nothing here today: a row is only
125+
# prunable because a later row for the same resource exists, so a
126+
# prunable row can never be the newest row in the table. We keep the
127+
# call anyway so keep_at_least_one_record still holds even if the
128+
# rules above change.
129+
exclude_newest_record(initial_records.where(Sequel.|(matching_ending, superseded_baseline)))
130+
end
131+
132+
# When keep_at_least_one_record is set, never delete the single newest row so
133+
# the table always retains at least one record.
134+
def exclude_newest_record(records)
135+
return records unless keep_at_least_one_record && newest_record_id
136+
137+
records.where(Sequel.lit('id < ?', newest_record_id))
138+
end
139+
140+
def newest_record_id
141+
return @newest_record_id if defined?(@newest_record_id)
142+
143+
@newest_record_id = model.order(:id).last&.id
144+
end
38145
end
39146
end

‎spec/unit/jobs/runtime/app_usage_events_cleanup_spec.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ module Jobs::Runtime
55
RSpec.describe AppUsageEventsCleanup, job_context: :worker do
66
let(:cutoff_age_in_days) { 30 }
77
let(:logger) { double(Steno::Logger, info: nil) }
8-
let!(:event_before_threshold) { create(:app_usage_event, created_at: (cutoff_age_in_days + 1).days.ago) }
8+
let!(:event_before_threshold) { create(:app_usage_event, created_at: (cutoff_age_in_days + 1).days.ago, state: 'STOPPED') }
99
let!(:event_after_threshold) { create(:app_usage_event, created_at: (cutoff_age_in_days - 1).days.ago) }
1010

1111
subject(:job) do

‎spec/unit/jobs/services/service_usage_events_cleanup_spec.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ module Jobs::Services
55
RSpec.describe ServiceUsageEventsCleanup, job_context: :worker do
66
let(:cutoff_age_in_days) { 30 }
77
let(:logger) { double(Steno::Logger, info: nil) }
8-
let!(:event_before_threshold) { create(:service_usage_event, created_at: (cutoff_age_in_days + 1).days.ago) }
8+
let!(:event_before_threshold) { create(:service_usage_event, created_at: (cutoff_age_in_days + 1).days.ago, state: 'DELETED') }
99
let!(:event_after_threshold) { create(:service_usage_event, created_at: (cutoff_age_in_days - 1).days.ago) }
1010

1111
subject(:job) do

0 commit comments

Comments
 (0)