Reset sync retries on new snapshots
This commit is contained in:
parent
38bbdf4b04
commit
4d9255354b
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Reference in New Issue