diff --git a/app/services/nostr_manager/fetch_event.rb b/app/services/nostr_manager/fetch_event.rb index 7767883..631fd88 100644 --- a/app/services/nostr_manager/fetch_event.rb +++ b/app/services/nostr_manager/fetch_event.rb @@ -34,7 +34,9 @@ module NostrManager received_event.signal end elsif msg && msg[0] == "EOSE" - Thread.current.exit + mutex.synchronize do + received_event.signal + end end end diff --git a/app/services/nostr_manager/fetch_latest_event.rb b/app/services/nostr_manager/fetch_latest_event.rb index 8046577..aebf522 100644 --- a/app/services/nostr_manager/fetch_latest_event.rb +++ b/app/services/nostr_manager/fetch_latest_event.rb @@ -9,7 +9,6 @@ module NostrManager end def call - received_events = 0 events = [] begin @@ -17,28 +16,17 @@ module NostrManager @relays.each do |url| event = NostrManager::FetchEvent.call(filter: @filter, relay_url: url) - if event.present? - events << event if events.none? { |e| e["id"] == event["id"] } - received_events += 1 - end - - if received_events >= @max_events - Rails.logger.debug "Found #{@max_events} events, ending the search" - break + if event.present? && events.none? { |e| e["id"] == event["id"] } + events << event + break if events.size >= @max_events end end - - events.min_by { |e| e["created_at"] } end rescue Timeout::Error - if events.size == 1 - Rails.logger.debug "[nostr] Timeout: only found 1 event within #{TIMEOUT} seconds for filter: #{@filter.inspect}" - events.first - else - Rails.logger.debug "[nostr] Timeout: no events found within #{TIMEOUT} seconds for filter: #{@filter.inspect}" - nil - end + Rails.logger.debug "[nostr] Timeout after #{TIMEOUT}s, found #{events.size} events" end + + events.max_by { |e| e["created_at"] } end end end diff --git a/spec/services/nostr_manager/fetch_latest_event_spec.rb b/spec/services/nostr_manager/fetch_latest_event_spec.rb new file mode 100644 index 0000000..dd323f0 --- /dev/null +++ b/spec/services/nostr_manager/fetch_latest_event_spec.rb @@ -0,0 +1,88 @@ +require 'rails_helper' + +RSpec.describe NostrManager::FetchLatestEvent, type: :model do + let(:filter) { Nostr::Filter.new(authors: ["pubkey"], kinds: [0], limit: 1) } + let(:event_v1) { { "id" => "aaa", "created_at" => 1000, "content" => "v1" } } + let(:event_v2) { { "id" => "bbb", "created_at" => 2000, "content" => "v2" } } + + describe "with no events found" do + it "returns nil" do + allow(NostrManager::FetchEvent).to receive(:call).and_return(nil) + + result = described_class.call(relays: ["wss://a", "wss://b"], filter: filter) + expect(result).to be_nil + end + end + + describe "with event on 1 relay only" do + it "returns the event" do + allow(NostrManager::FetchEvent).to receive(:call) do |args| + case args[:relay_url] + when "wss://a" then event_v1 + else nil + end + end + + result = described_class.call( + relays: ["wss://a", "wss://b", "wss://c"], filter: filter + ) + expect(result).to eq(event_v1) + end + end + + describe "with events on 2 relays, different versions" do + it "returns the latest event (max created_at)" do + allow(NostrManager::FetchEvent).to receive(:call) do |args| + case args[:relay_url] + when "wss://a" then event_v1 + when "wss://b" then event_v2 + end + end + + result = described_class.call(relays: ["wss://a", "wss://b"], filter: filter) + expect(result).to eq(event_v2) + end + end + + describe "with duplicate events (same ID) on 2 relays" do + it "returns the single event without double-counting" do + allow(NostrManager::FetchEvent).to receive(:call).and_return(event_v1) + + result = described_class.call(relays: ["wss://a", "wss://b"], filter: filter) + expect(result).to eq(event_v1) + end + end + + describe "with timeout after finding 1 event" do + it "still returns the found event" do + allow(NostrManager::FetchEvent).to receive(:call) do |args| + case args[:relay_url] + when "wss://a" then event_v1 + when "wss://b" then raise Timeout::Error + end + end + + result = described_class.call(relays: ["wss://a", "wss://b"], filter: filter) + expect(result).to eq(event_v1) + end + end + + describe "with max_events reached" do + it "stops querying further relays" do + call_count = 0 + allow(NostrManager::FetchEvent).to receive(:call) do |args| + call_count += 1 + raise "should not query wss://c" if args[:relay_url] == "wss://c" + case args[:relay_url] + when "wss://a" then event_v1 + when "wss://b" then event_v2 + end + end + + described_class.call( + relays: ["wss://a", "wss://b", "wss://c"], filter: filter + ) + expect(call_count).to eq(2) + end + end +end