diff --git a/app/services/tournament_sync_enqueue.rb b/app/services/tournament_sync_enqueue.rb index 056ab18..f2694b7 100644 --- a/app/services/tournament_sync_enqueue.rb +++ b/app/services/tournament_sync_enqueue.rb @@ -16,12 +16,15 @@ class TournamentSyncEnqueue entry = TournamentSyncQueueEntry.find_or_initialize_by(tournament: tournament) entry.snapshot_json = JSON.dump(payload) entry.status = 'pending' + entry.attempts = 0 + entry.last_attempt_at = nil entry.next_attempt_at = Time.current entry.last_error = nil entry.locked_at = nil entry.lock_token = nil entry.save! TournamentSyncWorker.start + TournamentSyncWorker.wake true end diff --git a/app/services/tournament_sync_worker.rb b/app/services/tournament_sync_worker.rb index 36afd5e..87b733d 100644 --- a/app/services/tournament_sync_worker.rb +++ b/app/services/tournament_sync_worker.rb @@ -14,15 +14,21 @@ class TournamentSyncWorker Thread.current.name = 'tournament-sync-worker' if Thread.current.respond_to?(:name=) loop do TournamentSyncProcessor.process_due! - sleep POLL_INTERVAL + wait_for_wake_or_timeout rescue StandardError => e Rails.logger.warn("Tournament sync worker error: #{e.message}") - sleep POLL_INTERVAL + wait_for_wake_or_timeout end end end end + def wake + mutex.synchronize do + condition.broadcast + end + end + def running? @thread&.alive? end @@ -32,5 +38,15 @@ class TournamentSyncWorker def mutex @mutex ||= Mutex.new end + + def condition + @condition ||= ConditionVariable.new + end + + def wait_for_wake_or_timeout + mutex.synchronize do + condition.wait(mutex, POLL_INTERVAL) + end + end end end diff --git a/spec/services/tournament_sync_enqueue_spec.rb b/spec/services/tournament_sync_enqueue_spec.rb new file mode 100644 index 0000000..77cc3ad --- /dev/null +++ b/spec/services/tournament_sync_enqueue_spec.rb @@ -0,0 +1,38 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe TournamentSyncEnqueue do + describe '.call' do + it 'resets retry state for newly enqueued snapshots' do + tournament = create(:tournament, + sync_target_url: 'https://remote.example.com/tournaments/1/sync_state', + sync_auth_token: 'shared-secret') + entry = TournamentSyncQueueEntry.create!( + tournament: tournament, + snapshot_json: '{"old":true}', + status: 'pending', + attempts: 4, + last_error: 'network down', + last_attempt_at: 5.minutes.ago, + next_attempt_at: 3.minutes.from_now, + locked_at: 1.minute.ago, + lock_token: 'stale-lock' + ) + + allow(TournamentSyncWorker).to receive(:start) + allow(TournamentSyncWorker).to receive(:wake) + + described_class.call(tournament) + + entry.reload + expect(entry.status).to eq('pending') + expect(entry.attempts).to eq(0) + expect(entry.last_error).to be_nil + expect(entry.last_attempt_at).to be_nil + expect(entry.next_attempt_at).to be <= Time.current + expect(entry.locked_at).to be_nil + expect(entry.lock_token).to be_nil + end + end +end