Make follower sync asynchronous
This commit is contained in:
parent
133ebae847
commit
bbb1ad97ef
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
@ -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}"
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -0,0 +1,5 @@
|
|||
# frozen_string_literal: true
|
||||
|
||||
Rails.application.config.after_initialize do
|
||||
TournamentSyncWorker.start unless Rails.env.test?
|
||||
end
|
||||
|
|
@ -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
|
||||
|
|
@ -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?
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Reference in New Issue