|
| 1 | +require "test_helper" |
| 2 | + |
| 3 | +class DashboardRollupConcurrencyTest < ActiveSupport::TestCase |
| 4 | + self.use_transactional_tests = false |
| 5 | + include ActiveJob::TestHelper |
| 6 | + |
| 7 | + test "a correction during refresh leaves a coherent snapshot and a pending refresh" do |
| 8 | + original_cache = Rails.cache |
| 9 | + original_adapter = ActiveJob::Base.queue_adapter |
| 10 | + Rails.cache = ActiveSupport::Cache::MemoryStore.new |
| 11 | + ActiveJob::Base.queue_adapter = :test |
| 12 | + user = create(:user) |
| 13 | + create(:heartbeat, user: user, project: "api", language: "Ruby", time: 1_700_000_000.0) |
| 14 | + heartbeat = create(:heartbeat, user: user, project: "api", language: "Ruby", time: 1_700_000_060.0) |
| 15 | + DashboardRollupRefreshService.new(user: user).call |
| 16 | + clear_enqueued_jobs |
| 17 | + Rails.cache.clear |
| 18 | + DashboardRollupRefreshJob.schedule_for(user.id) |
| 19 | + main_thread = Thread.current |
| 20 | + corrected = false |
| 21 | + subscriber = lambda do |*args| |
| 22 | + payload = args.last |
| 23 | + if Thread.current == main_thread && !corrected && payload[:sql].include?('SELECT COUNT(*) FROM "heartbeats"') |
| 24 | + corrected = true |
| 25 | + Thread.new do |
| 26 | + ActiveRecord::Base.connection_pool.with_connection do |
| 27 | + Heartbeat.find(heartbeat.id).update!(language: "Python") |
| 28 | + end |
| 29 | + end.value |
| 30 | + end |
| 31 | + end |
| 32 | + |
| 33 | + ActiveSupport::Notifications.subscribed(subscriber, "sql.active_record") do |
| 34 | + DashboardRollupRefreshJob.perform_now(user.id) |
| 35 | + end |
| 36 | + |
| 37 | + assert corrected, "the concurrent correction must run during the aggregate reads" |
| 38 | + assert DashboardRollup.dirty?(user.id), "the newer invalidation must survive refresh completion" |
| 39 | + buckets = DashboardRollup.where(user: user, dimension: "language").to_h { |row| [ row.bucket, row.total_seconds ] } |
| 40 | + assert_equal({ "Ruby" => 60 }, buckets) |
| 41 | + assert_enqueued_jobs 2, only: DashboardRollupRefreshJob |
| 42 | + |
| 43 | + DashboardRollupRefreshJob.perform_now(user.id) |
| 44 | + assert_not DashboardRollup.dirty?(user.id) |
| 45 | + assert_equal 60, DashboardRollup.find_by!(user: user, dimension: "language", bucket_value: "Python").total_seconds |
| 46 | + DashboardRollup.mark_dirty(user.id) |
| 47 | + Rails.cache.clear |
| 48 | + assert DashboardRollup.dirty?(user.id), "invalidations must survive cache loss" |
| 49 | + ensure |
| 50 | + DashboardRollup.where(user: user).delete_all if user |
| 51 | + Heartbeat.with_deleted.where(user: user).delete_all if user |
| 52 | + user&.destroy! |
| 53 | + clear_enqueued_jobs |
| 54 | + Rails.cache = original_cache |
| 55 | + ActiveJob::Base.queue_adapter = original_adapter |
| 56 | + end |
| 57 | +end |
0 commit comments