From bbb1ad97ef9a68b71df238dc46fa80a62079cd14 Mon Sep 17 00:00:00 2001 From: Malaber Date: Mon, 13 Apr 2026 20:54:28 +0200 Subject: [PATCH] Make follower sync asynchronous --- app/controllers/match_scores_controller.rb | 4 +- app/controllers/matches_controller.rb | 4 +- app/controllers/stages_controller.rb | 4 +- app/controllers/teams_controller.rb | 4 +- app/controllers/tournaments_controller.rb | 4 +- app/models/tournament_sync_queue_entry.rb | 63 +++++++++++++++++++ app/services/tournament_sync_enqueue.rb | 31 +++++++++ app/services/tournament_sync_processor.rb | 33 ++++++++++ app/services/tournament_sync_pusher.rb | 10 ++- app/services/tournament_sync_worker.rb | 36 +++++++++++ config/initializers/tournament_sync_worker.rb | 5 ++ ...00_create_tournament_sync_queue_entries.rb | 22 +++++++ spec/e2e/http/tournament_follow_sync_spec.rb | 20 +++++- .../tournament_sync_processor_spec.rb | 25 ++++++++ 14 files changed, 246 insertions(+), 19 deletions(-) create mode 100644 app/models/tournament_sync_queue_entry.rb create mode 100644 app/services/tournament_sync_enqueue.rb create mode 100644 app/services/tournament_sync_processor.rb create mode 100644 app/services/tournament_sync_worker.rb create mode 100644 config/initializers/tournament_sync_worker.rb create mode 100644 db/migrate/20260413210000_create_tournament_sync_queue_entries.rb create mode 100644 spec/services/tournament_sync_processor_spec.rb diff --git a/app/controllers/match_scores_controller.rb b/app/controllers/match_scores_controller.rb index 1468cf9..d239aef 100644 --- a/app/controllers/match_scores_controller.rb +++ b/app/controllers/match_scores_controller.rb @@ -35,8 +35,6 @@ class MatchScoresController < ApplicationController end def push_sync_if_needed!(tournament) - TournamentSyncPusher.push!(tournament) - rescue TournamentSyncPusher::SyncFailed => e - logger.warn("Tournament sync push failed for #{tournament.id}: #{e.message}") + TournamentSyncEnqueue.call(tournament) end end diff --git a/app/controllers/matches_controller.rb b/app/controllers/matches_controller.rb index 1a7bd67..370662d 100644 --- a/app/controllers/matches_controller.rb +++ b/app/controllers/matches_controller.rb @@ -162,8 +162,6 @@ class MatchesController < ApplicationController end def push_sync_if_needed!(tournament) - TournamentSyncPusher.push!(tournament) - rescue TournamentSyncPusher::SyncFailed => e - logger.warn("Tournament sync push failed for #{tournament.id}: #{e.message}") + TournamentSyncEnqueue.call(tournament) end end diff --git a/app/controllers/stages_controller.rb b/app/controllers/stages_controller.rb index c02dffc..8760686 100644 --- a/app/controllers/stages_controller.rb +++ b/app/controllers/stages_controller.rb @@ -79,8 +79,6 @@ class StagesController < ApplicationController end def push_sync_if_needed!(tournament) - TournamentSyncPusher.push!(tournament) - rescue TournamentSyncPusher::SyncFailed => e - logger.warn("Tournament sync push failed for #{tournament.id}: #{e.message}") + TournamentSyncEnqueue.call(tournament) end end diff --git a/app/controllers/teams_controller.rb b/app/controllers/teams_controller.rb index 42e9b25..7b5d7b5 100644 --- a/app/controllers/teams_controller.rb +++ b/app/controllers/teams_controller.rb @@ -32,8 +32,6 @@ class TeamsController < ApplicationController end def push_sync_if_needed!(tournament) - TournamentSyncPusher.push!(tournament) - rescue TournamentSyncPusher::SyncFailed => e - logger.warn("Tournament sync push failed for #{tournament.id}: #{e.message}") + TournamentSyncEnqueue.call(tournament) end end diff --git a/app/controllers/tournaments_controller.rb b/app/controllers/tournaments_controller.rb index 36a4612..d2d7da3 100644 --- a/app/controllers/tournaments_controller.rb +++ b/app/controllers/tournaments_controller.rb @@ -253,9 +253,7 @@ class TournamentsController < ApplicationController end def push_sync_if_needed!(tournament) - TournamentSyncPusher.push!(tournament) - rescue TournamentSyncPusher::SyncFailed => e - logger.warn("Tournament sync push failed for #{tournament.id}: #{e.message}") + TournamentSyncEnqueue.call(tournament) end end diff --git a/app/models/tournament_sync_queue_entry.rb b/app/models/tournament_sync_queue_entry.rb new file mode 100644 index 0000000..297cd4b --- /dev/null +++ b/app/models/tournament_sync_queue_entry.rb @@ -0,0 +1,63 @@ +# frozen_string_literal: true + +class TournamentSyncQueueEntry < ApplicationRecord + LOCK_TTL = 2.minutes + MAX_BACKOFF = 5.minutes + + belongs_to :tournament + + validates :snapshot_json, presence: true + validates :status, presence: true + + scope :due, -> { where(status: 'pending').where('next_attempt_at <= ?', Time.current) } + scope :unlocked, lambda { + where(locked_at: nil).or(where('locked_at < ?', Time.current - LOCK_TTL)) + } + + def snapshot + JSON.parse(snapshot_json, symbolize_names: true) + end + + def schedule_retry!(error_message) + update!( + status: 'pending', + last_error: error_message, + attempts: attempts + 1, + last_attempt_at: Time.current, + next_attempt_at: Time.current + retry_delay, + locked_at: nil, + lock_token: nil + ) + end + + def mark_synced! + update!( + status: 'synced', + last_error: nil, + last_attempt_at: Time.current, + next_attempt_at: Time.current, + attempts: 0, + locked_at: nil, + lock_token: nil + ) + end + + def acquire_lock! + token = SecureRandom.hex(8) + updated = self.class + .where(id: id) + .where(locked_at: nil) + .or(self.class.where(id: id).where('locked_at < ?', Time.current - LOCK_TTL)) + .update_all(locked_at: Time.current, lock_token: token) + return nil if updated.zero? + + reload + token + end + + private + + def retry_delay + [2**attempts, MAX_BACKOFF].min.seconds + end +end diff --git a/app/services/tournament_sync_enqueue.rb b/app/services/tournament_sync_enqueue.rb new file mode 100644 index 0000000..f4bd3a8 --- /dev/null +++ b/app/services/tournament_sync_enqueue.rb @@ -0,0 +1,31 @@ +# frozen_string_literal: true + +class TournamentSyncEnqueue + def self.call(tournament) + new(tournament).call + end + + def initialize(tournament) + @tournament = tournament + end + + def call + return false unless tournament.sync_push_enabled? + + payload = TournamentSnapshotBuilder.build(tournament) + entry = TournamentSyncQueueEntry.find_or_initialize_by(tournament: tournament) + entry.snapshot_json = JSON.dump(payload) + entry.status = 'pending' + entry.next_attempt_at = Time.current + entry.last_error = nil + entry.locked_at = nil + entry.lock_token = nil + entry.save! + TournamentSyncWorker.start + true + end + + private + + attr_reader :tournament +end diff --git a/app/services/tournament_sync_processor.rb b/app/services/tournament_sync_processor.rb new file mode 100644 index 0000000..515b07b --- /dev/null +++ b/app/services/tournament_sync_processor.rb @@ -0,0 +1,33 @@ +# frozen_string_literal: true + +class TournamentSyncProcessor + def self.process_due! + new.process_due! + end + + def process_due! + loop do + entry = next_entry + break if entry.nil? + + process_entry(entry) + end + end + + private + + def next_entry + TournamentSyncQueueEntry.due.unlocked.order(:next_attempt_at, :id).first + end + + def process_entry(entry) + lock_token = entry.acquire_lock! + return if lock_token.nil? + + entry.reload + TournamentSyncPusher.push_snapshot!(entry.tournament, entry.snapshot) + entry.mark_synced! + rescue TournamentSyncPusher::SyncFailed => e + entry.reload.schedule_retry!(e.message) + end +end diff --git a/app/services/tournament_sync_pusher.rb b/app/services/tournament_sync_pusher.rb index bb110bb..aad9fea 100644 --- a/app/services/tournament_sync_pusher.rb +++ b/app/services/tournament_sync_pusher.rb @@ -11,6 +11,10 @@ class TournamentSyncPusher new(tournament).push! end + def self.push_snapshot!(tournament, snapshot) + new(tournament).push_snapshot!(snapshot) + end + def initialize(tournament) @tournament = tournament end @@ -18,7 +22,11 @@ class TournamentSyncPusher def push! return false unless tournament.sync_push_enabled? - response = perform_request(snapshot: TournamentSnapshotBuilder.build(tournament)) + push_snapshot!(TournamentSnapshotBuilder.build(tournament)) + end + + def push_snapshot!(snapshot) + response = perform_request(snapshot: snapshot) unless response.is_a?(Net::HTTPSuccess) raise SyncFailed, "sync push failed with status #{response.code}: #{response.body}" diff --git a/app/services/tournament_sync_worker.rb b/app/services/tournament_sync_worker.rb new file mode 100644 index 0000000..36afd5e --- /dev/null +++ b/app/services/tournament_sync_worker.rb @@ -0,0 +1,36 @@ +# frozen_string_literal: true + +class TournamentSyncWorker + POLL_INTERVAL = 2.seconds + + class << self + def start + return if Rails.env.test? + + mutex.synchronize do + return if running? + + @thread = Thread.new do + Thread.current.name = 'tournament-sync-worker' if Thread.current.respond_to?(:name=) + loop do + TournamentSyncProcessor.process_due! + sleep POLL_INTERVAL + rescue StandardError => e + Rails.logger.warn("Tournament sync worker error: #{e.message}") + sleep POLL_INTERVAL + end + end + end + end + + def running? + @thread&.alive? + end + + private + + def mutex + @mutex ||= Mutex.new + end + end +end diff --git a/config/initializers/tournament_sync_worker.rb b/config/initializers/tournament_sync_worker.rb new file mode 100644 index 0000000..b0407a9 --- /dev/null +++ b/config/initializers/tournament_sync_worker.rb @@ -0,0 +1,5 @@ +# frozen_string_literal: true + +Rails.application.config.after_initialize do + TournamentSyncWorker.start unless Rails.env.test? +end diff --git a/db/migrate/20260413210000_create_tournament_sync_queue_entries.rb b/db/migrate/20260413210000_create_tournament_sync_queue_entries.rb new file mode 100644 index 0000000..b06c3f9 --- /dev/null +++ b/db/migrate/20260413210000_create_tournament_sync_queue_entries.rb @@ -0,0 +1,22 @@ +# frozen_string_literal: true + +class CreateTournamentSyncQueueEntries < ActiveRecord::Migration[7.0] + def change + create_table :tournament_sync_queue_entries do |t| + t.references :tournament, null: false, foreign_key: { on_delete: :cascade }, index: { unique: true } + t.text :snapshot_json, null: false + t.datetime :next_attempt_at, null: false + t.datetime :last_attempt_at + t.datetime :locked_at + t.string :lock_token + t.integer :attempts, null: false, default: 0 + t.string :status, null: false, default: 'pending' + t.string :last_error + + t.timestamps + end + + add_index :tournament_sync_queue_entries, :next_attempt_at + add_index :tournament_sync_queue_entries, :status + end +end diff --git a/spec/e2e/http/tournament_follow_sync_spec.rb b/spec/e2e/http/tournament_follow_sync_spec.rb index f4e417b..fa44fdb 100644 --- a/spec/e2e/http/tournament_follow_sync_spec.rb +++ b/spec/e2e/http/tournament_follow_sync_spec.rb @@ -38,7 +38,7 @@ RSpec.describe 'Tournament follower sync HTTP E2E' do expect(configure_sync[:status]).to eq(200) source = fetch_tournament(client: source_anonymous_client, tournament_id: source.fetch(:id)) - follower = fetch_tournament(client: follower_anonymous_client, tournament_id: follower.fetch(:id)) + follower = wait_for_tournament_sync!(source_tournament_id: source.fetch(:id), follower_tournament_id: follower.fetch(:id)) expect(tournament_signature(follower)).to eq(tournament_signature(source)) update_cutoff = source_owner_client.patch("/tournaments/#{source.fetch(:id)}", body: { @@ -53,7 +53,7 @@ RSpec.describe 'Tournament follower sync HTTP E2E' do play_group_with_decider_lifecycle!(source_id: source.fetch(:id), groups: source_group_stage.fetch(:groups).sort_by { |group| group.fetch(:number) }) source = fetch_tournament(client: source_anonymous_client, tournament_id: source.fetch(:id)) - follower = fetch_tournament(client: follower_anonymous_client, tournament_id: follower.fetch(:id)) + follower = wait_for_tournament_sync!(source_tournament_id: source.fetch(:id), follower_tournament_id: follower.fetch(:id)) expect(tournament_signature(follower)).to eq(tournament_signature(source)) expect(follower.fetch(:stages).map { |stage| stage.fetch(:level) }).to include(-1, 0, 1) @@ -66,7 +66,7 @@ RSpec.describe 'Tournament follower sync HTTP E2E' do finish_playoff_bracket!(source_id: source.fetch(:id)) source = fetch_tournament(client: source_anonymous_client, tournament_id: source.fetch(:id)) - follower = fetch_tournament(client: follower_anonymous_client, tournament_id: follower.fetch(:id)) + follower = wait_for_tournament_sync!(source_tournament_id: source.fetch(:id), follower_tournament_id: follower.fetch(:id)) expect(tournament_signature(follower)).to eq(tournament_signature(source)) disable_follower = follower_owner_client.patch("/tournaments/#{follower.fetch(:id)}", body: { read_only_mode: false }) @@ -127,6 +127,20 @@ RSpec.describe 'Tournament follower sync HTTP E2E' do response.fetch(:json) end + def wait_for_tournament_sync!(source_tournament_id:, follower_tournament_id:, timeout: 30) + deadline = Time.now + timeout + + loop do + source = fetch_tournament(client: source_anonymous_client, tournament_id: source_tournament_id) + follower = fetch_tournament(client: follower_anonymous_client, tournament_id: follower_tournament_id) + return follower if tournament_signature(source) == tournament_signature(follower) + + raise "Follower did not catch up within #{timeout}s" if Time.now >= deadline + + sleep 0.5 + end + end + def play_group_with_decider_lifecycle!(source_id:, groups:) groups.each_with_index do |group, index| if index.zero? diff --git a/spec/services/tournament_sync_processor_spec.rb b/spec/services/tournament_sync_processor_spec.rb new file mode 100644 index 0000000..f3b54ec --- /dev/null +++ b/spec/services/tournament_sync_processor_spec.rb @@ -0,0 +1,25 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe TournamentSyncProcessor do + describe '.process_due!' do + it 'retries failed pushes asynchronously from queue state' do + tournament = create(:tournament, + sync_target_url: 'https://remote.example.com/tournaments/1/sync_state', + sync_auth_token: 'shared-secret') + TournamentSyncEnqueue.call(tournament) + entry = TournamentSyncQueueEntry.find_by!(tournament: tournament) + + allow(TournamentSyncPusher).to receive(:push_snapshot!).and_raise(TournamentSyncPusher::SyncFailed, 'network down') + + described_class.process_due! + + entry.reload + expect(entry.status).to eq('pending') + expect(entry.last_error).to eq('network down') + expect(entry.attempts).to eq(1) + expect(entry.next_attempt_at).to be > Time.current + end + end +end