Fix AccountsStatusesCleanupScheduler not spreading deletes across accounts correctly (#24607)

This commit is contained in:
Claire 2023-04-23 22:25:40 +02:00 committed by Tarrien
parent 1f969fe448
commit e92cd220f7
2 changed files with 49 additions and 103 deletions

View File

@ -42,7 +42,7 @@ class Scheduler::AccountsStatusesCleanupScheduler
num_processed_accounts = 0 num_processed_accounts = 0
scope = AccountStatusesCleanupPolicy.where(enabled: true) scope = AccountStatusesCleanupPolicy.where(enabled: true)
scope.where(Account.arel_table[:id].gt(first_policy_id)) if first_policy_id.present? scope = scope.where(id: first_policy_id...) if first_policy_id.present?
scope.find_each(order: :asc) do |policy| scope.find_each(order: :asc) do |policy|
num_deleted = AccountStatusesCleanupService.new.call(policy, [budget, PER_ACCOUNT_BUDGET].min) num_deleted = AccountStatusesCleanupService.new.call(policy, [budget, PER_ACCOUNT_BUDGET].min)
num_processed_accounts += 1 unless num_deleted.zero? num_processed_accounts += 1 unless num_deleted.zero?
@ -77,14 +77,14 @@ class Scheduler::AccountsStatusesCleanupScheduler
end end
def last_processed_id def last_processed_id
redis.get('account_statuses_cleanup_scheduler:last_account_id') redis.get('account_statuses_cleanup_scheduler:last_policy_id')
end end
def save_last_processed_id(id) def save_last_processed_id(id)
if id.nil? if id.nil?
redis.del('account_statuses_cleanup_scheduler:last_account_id') redis.del('account_statuses_cleanup_scheduler:last_policy_id')
else else
redis.set('account_statuses_cleanup_scheduler:last_account_id', id, ex: 1.hour.seconds) redis.set('account_statuses_cleanup_scheduler:last_policy_id', id, ex: 1.hour.seconds)
end end
end end
end end

View File

@ -1,16 +1,17 @@
# frozen_string_literal: true
require 'rails_helper' require 'rails_helper'
describe Scheduler::AccountsStatusesCleanupScheduler do describe Scheduler::AccountsStatusesCleanupScheduler do
subject { described_class.new } subject { described_class.new }
let!(:account_alice) { Fabricate(:account, domain: nil, username: 'alice') } let!(:account1) { Fabricate(:account, domain: nil) }
let!(:account_bob) { Fabricate(:account, domain: nil, username: 'bob') } let!(:account2) { Fabricate(:account, domain: nil) }
let!(:account_chris) { Fabricate(:account, domain: nil, username: 'chris') } let!(:account3) { Fabricate(:account, domain: nil) }
let!(:account_dave) { Fabricate(:account, domain: nil, username: 'dave') } let!(:account4) { Fabricate(:account, domain: nil) }
let!(:account_erin) { Fabricate(:account, domain: nil, username: 'erin') } let!(:remote) { Fabricate(:account) }
let!(:remote) { Fabricate(:account) }
let!(:policy1) { Fabricate(:account_statuses_cleanup_policy, account: account1) }
let!(:policy2) { Fabricate(:account_statuses_cleanup_policy, account: account3) }
let!(:policy3) { Fabricate(:account_statuses_cleanup_policy, account: account4, enabled: false) }
let(:queue_size) { 0 } let(:queue_size) { 0 }
let(:queue_latency) { 0 } let(:queue_latency) { 0 }
@ -18,18 +19,38 @@ describe Scheduler::AccountsStatusesCleanupScheduler do
[ [
{ {
'concurrency' => 2, 'concurrency' => 2,
'queues' => %w(push default), 'queues' => ['push', 'default'],
}, },
] ]
end end
before do before do
queue_stub = instance_double(Sidekiq::Queue, size: queue_size, latency: queue_latency) queue_stub = double
allow(queue_stub).to receive(:size).and_return(queue_size)
allow(queue_stub).to receive(:latency).and_return(queue_latency)
allow(Sidekiq::Queue).to receive(:new).and_return(queue_stub) allow(Sidekiq::Queue).to receive(:new).and_return(queue_stub)
allow(Sidekiq::ProcessSet).to receive(:new).and_return(process_set_stub) allow(Sidekiq::ProcessSet).to receive(:new).and_return(process_set_stub)
sidekiq_stats_stub = instance_double(Sidekiq::Stats) sidekiq_stats_stub = double
allow(Sidekiq::Stats).to receive(:new).and_return(sidekiq_stats_stub) allow(Sidekiq::Stats).to receive(:new).and_return(sidekiq_stats_stub)
# Create a bunch of old statuses
10.times do
Fabricate(:status, account: account1, created_at: 3.years.ago)
Fabricate(:status, account: account2, created_at: 3.years.ago)
Fabricate(:status, account: account3, created_at: 3.years.ago)
Fabricate(:status, account: account4, created_at: 3.years.ago)
Fabricate(:status, account: remote, created_at: 3.years.ago)
end
# Create a bunch of newer statuses
5.times do
Fabricate(:status, account: account1, created_at: 3.minutes.ago)
Fabricate(:status, account: account2, created_at: 3.minutes.ago)
Fabricate(:status, account: account3, created_at: 3.minutes.ago)
Fabricate(:status, account: account4, created_at: 3.minutes.ago)
Fabricate(:status, account: remote, created_at: 3.minutes.ago)
end
end end
describe '#under_load?' do describe '#under_load?' do
@ -49,19 +70,19 @@ describe Scheduler::AccountsStatusesCleanupScheduler do
end end
end end
describe '#compute_budget' do describe '#get_budget' do
context 'with a single thread' do context 'on a single thread' do
let(:process_set_stub) { [{ 'concurrency' => 1, 'queues' => %w(push default) }] } let(:process_set_stub) { [ { 'concurrency' => 1, 'queues' => ['push', 'default'] } ] }
it 'returns a low value' do it 'returns a low value' do
expect(subject.compute_budget).to be < 10 expect(subject.compute_budget).to be < 10
end end
end end
context 'with a lot of threads' do context 'on a lot of threads' do
let(:process_set_stub) do let(:process_set_stub) do
[ [
{ 'concurrency' => 2, 'queues' => %w(push default) }, { 'concurrency' => 2, 'queues' => ['push', 'default'] },
{ 'concurrency' => 2, 'queues' => ['push'] }, { 'concurrency' => 2, 'queues' => ['push'] },
{ 'concurrency' => 2, 'queues' => ['push'] }, { 'concurrency' => 2, 'queues' => ['push'] },
{ 'concurrency' => 2, 'queues' => ['push'] }, { 'concurrency' => 2, 'queues' => ['push'] },
@ -75,96 +96,21 @@ describe Scheduler::AccountsStatusesCleanupScheduler do
end end
describe '#perform' do describe '#perform' do
around do |example|
Timeout.timeout(30) do
example.run
end
end
before do
# Policies for the accounts
Fabricate(:account_statuses_cleanup_policy, account: account_alice)
Fabricate(:account_statuses_cleanup_policy, account: account_chris)
Fabricate(:account_statuses_cleanup_policy, account: account_dave, enabled: false)
Fabricate(:account_statuses_cleanup_policy, account: account_erin)
# Create a bunch of old statuses
4.times do
Fabricate(:status, account: account_alice, created_at: 3.years.ago)
Fabricate(:status, account: account_bob, created_at: 3.years.ago)
Fabricate(:status, account: account_chris, created_at: 3.years.ago)
Fabricate(:status, account: account_dave, created_at: 3.years.ago)
Fabricate(:status, account: account_erin, created_at: 3.years.ago)
Fabricate(:status, account: remote, created_at: 3.years.ago)
end
# Create a bunch of newer statuses
Fabricate(:status, account: account_alice, created_at: 3.minutes.ago)
Fabricate(:status, account: account_bob, created_at: 3.minutes.ago)
Fabricate(:status, account: account_chris, created_at: 3.minutes.ago)
Fabricate(:status, account: account_dave, created_at: 3.minutes.ago)
Fabricate(:status, account: remote, created_at: 3.minutes.ago)
end
context 'when the budget is lower than the number of toots to delete' do context 'when the budget is lower than the number of toots to delete' do
it 'deletes the appropriate statuses' do it 'deletes as many statuses as the given budget' do
expect(Status.count).to be > (subject.compute_budget) # Data check expect { subject.perform }.to change { Status.count }.by(-subject.compute_budget)
expect { subject.perform }
.to change(Status, :count).by(-subject.compute_budget) # Cleanable statuses
.and (not_change { account_bob.statuses.count }) # No cleanup policy for account
.and(not_change { account_dave.statuses.count }) # Disabled cleanup policy
end end
it 'eventually deletes every deletable toot given enough runs' do it 'does not delete from accounts with no cleanup policy' do
stub_const 'Scheduler::AccountsStatusesCleanupScheduler::MAX_BUDGET', 4 expect { subject.perform }.to_not change { account2.statuses.count }
expect { 3.times { subject.perform } }.to change(Status, :count).by(-cleanable_statuses_count)
end end
it 'correctly round-trips between users across several runs' do it 'does not delete from accounts with disabled cleanup policies' do
stub_const 'Scheduler::AccountsStatusesCleanupScheduler::MAX_BUDGET', 3 expect { subject.perform }.to_not change { account4.statuses.count }
stub_const 'Scheduler::AccountsStatusesCleanupScheduler::PER_ACCOUNT_BUDGET', 2
expect { 3.times { subject.perform } }
.to change(Status, :count).by(-3 * 3)
.and change { account_alice.statuses.count }
.and change { account_chris.statuses.count }
.and(change { account_erin.statuses.count })
end end
context 'when given a big budget' do it 'eventually deletes every deletable toot' do
let(:process_set_stub) { [{ 'concurrency' => 400, 'queues' => %w(push default) }] } expect { subject.perform; subject.perform; subject.perform; subject.perform }.to change { Status.count }.by(-20)
before do
stub_const 'Scheduler::AccountsStatusesCleanupScheduler::MAX_BUDGET', 400
end
it 'correctly handles looping in a single run' do
expect(subject.compute_budget).to eq(400)
expect { subject.perform }.to change(Status, :count).by(-cleanable_statuses_count)
end
end
context 'when there is no work to be done' do
let(:process_set_stub) { [{ 'concurrency' => 400, 'queues' => %w(push default) }] }
before do
stub_const 'Scheduler::AccountsStatusesCleanupScheduler::MAX_BUDGET', 400
subject.perform
end
it 'does not get stuck' do
expect(subject.compute_budget).to eq(400)
expect { subject.perform }.to_not change(Status, :count)
end
end
def cleanable_statuses_count
Status
.where(account_id: [account_alice, account_chris, account_erin]) # Accounts with enabled policies
.where('created_at < ?', 2.weeks.ago) # Policy defaults is 2.weeks
.count
end end
end end
end end