From 1cdeaf699c15f6fffdd5de70cdf243325c5a0431 Mon Sep 17 00:00:00 2001 From: tiyabelay Date: Wed, 7 Oct 2026 11:28:00 -0700 Subject: [PATCH 1/5] fix: enable jetstream endpoint to resolve stream from subject and add error logs --- lib/leopard/nats_jetstream_consumer.rb | 30 ++++++++++++++--- test/lib/nats_jetstream_consumer_test.rb | 43 +++++++++++++++++++++++- 2 files changed, 67 insertions(+), 6 deletions(-) diff --git a/lib/leopard/nats_jetstream_consumer.rb b/lib/leopard/nats_jetstream_consumer.rb index 9695e38..ed77b76 100644 --- a/lib/leopard/nats_jetstream_consumer.rb +++ b/lib/leopard/nats_jetstream_consumer.rb @@ -69,6 +69,9 @@ def start_endpoint(endpoint) subscription = build_subscription(endpoint) subscriptions << subscription threads << @thread_factory.new { consume_endpoint(subscription, endpoint) } + rescue StandardError => e + @logger.error "JetStream endpoint #{endpoint.name} failed to start: ", e + raise end # Ensures the durable consumer exists and creates a pull subscription for it. @@ -77,23 +80,40 @@ def start_endpoint(endpoint) # # @return [Object] The JetStream pull subscription. def build_subscription(endpoint) - ensure_consumer(endpoint) + stream = stream_for(endpoint) + ensure_consumer(endpoint, stream) @jetstream.pull_subscribe( endpoint.subject, endpoint.durable, - stream: endpoint.stream, + stream:, ) end + # Finds the stream for an endpoint, looking it up from the subject when none is configured. + # + # @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to find a stream for. + # + # @raise [ConfigurationError] When no stream is configured and none captures the endpoint's subject. + # + # @return [String] The JetStream stream name. + def stream_for(endpoint) + endpoint.stream || @jetstream.find_stream_name_by_subject(endpoint.subject) + rescue NATS::JetStream::Error::NotFound + raise ConfigurationError, + "JetStream endpoint #{endpoint.name} has no stream and none captures subject + #{endpoint.subject}, " \ + 'create the stream or set `stream:`' + end + # Verifies that the durable consumer exists, creating it when missing. # # @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to ensure. # # @return [Object] Consumer metadata from `consumer_info` or `add_consumer`. - def ensure_consumer(endpoint) - @jetstream.consumer_info(endpoint.stream, endpoint.durable) + def ensure_consumer(endpoint, stream) + @jetstream.consumer_info(stream, endpoint.durable) rescue NATS::JetStream::Error::NotFound - @jetstream.add_consumer(endpoint.stream, consumer_config(endpoint)) + @jetstream.add_consumer(stream, consumer_config(endpoint)) end # Builds the JetStream consumer configuration for an endpoint. diff --git a/test/lib/nats_jetstream_consumer_test.rb b/test/lib/nats_jetstream_consumer_test.rb index f31d693..1b249e1 100644 --- a/test/lib/nats_jetstream_consumer_test.rb +++ b/test/lib/nats_jetstream_consumer_test.rb @@ -1,9 +1,9 @@ # frozen_string_literal: true require_relative '../helper' +require 'nats/client' require Rubyists::Leopard.libroot / 'leopard/nats_jetstream_consumer' require Rubyists::Leopard.libroot / 'leopard/nats_jetstream_endpoint' - class NatsJetstreamConsumerTest < Minitest::Test def setup @consumer = Rubyists::Leopard::NatsJetstreamConsumer.new( @@ -41,8 +41,49 @@ def test_consumer_config_keeps_safe_string_key_overrides assert_equal 100, string_key_config['max_ack_pending'] end + def test_stream_for_with_the_configured_stream + assert_equal 'EVENTS', consumer_with(jetstream: Object.new).send(:stream_for, endpoint_with_consumer(nil)) + end + + def test_stream_for_with_the_subject_when_no_stream_is_configured + jetstream = Object.new + jetstream.define_singleton_method(:find_stream_name_by_subject) { |_subject| 'EVENTS' } + + assert_equal 'EVENTS', consumer_with(jetstream:).send(:stream_for, endpoint_without_stream) + end + + def test_stream_for_raises_a_configuration_error_when_no_stream_has_the_subject + jetstream = Object.new + jetstream.define_singleton_method(:find_stream_name_by_subject) { |_subject| raise NATS::JetStream::Error::NotFound } + + assert_raises(Rubyists::Leopard::ConfigurationError) do + consumer_with(jetstream:).send(:stream_for, endpoint_without_stream) + end + end + + def test_start_endpoint_logs_and_reraises_when_the_endpoint_fails_to_start + logged = [] + logger = Object.new + logger.define_singleton_method(:error) { |*args| logged << args } + + assert_raises(NoMethodError) do + consumer_with(jetstream: Object.new, logger:).send(:start_endpoint, + endpoint_with_consumer(nil)) + end + assert_match(/failed to start/, logged.first.first) + end + private + def consumer_with(jetstream:, logger: Object.new) + Rubyists::Leopard::NatsJetstreamConsumer.new(jetstream:, endpoints: [], + logger:, process_message: ->(*_) {}) + end + + def endpoint_without_stream + Rubyists::Leopard::NatsJetstreamEndpoint.new(**base_endpoint_attributes, stream: nil) + end + def symbol_key_config @symbol_key_config ||= @consumer.send(:consumer_config, endpoint_with_consumer(symbol_key_overrides)) end From 909230ca55b85b42fda3703aa06cb00e4726b089 Mon Sep 17 00:00:00 2001 From: tiyabelay Date: Wed, 7 Oct 2026 13:52:57 -0700 Subject: [PATCH 2/5] fix: Remove ensure_consumer and fallback on pull_subscrbe to resolve stream --- lib/leopard/nats_jetstream_consumer.rb | 31 +----------------------- test/lib/nats_jetstream_consumer_test.rb | 25 ------------------- 2 files changed, 1 insertion(+), 55 deletions(-) diff --git a/lib/leopard/nats_jetstream_consumer.rb b/lib/leopard/nats_jetstream_consumer.rb index ed77b76..14cc45f 100644 --- a/lib/leopard/nats_jetstream_consumer.rb +++ b/lib/leopard/nats_jetstream_consumer.rb @@ -80,42 +80,13 @@ def start_endpoint(endpoint) # # @return [Object] The JetStream pull subscription. def build_subscription(endpoint) - stream = stream_for(endpoint) - ensure_consumer(endpoint, stream) @jetstream.pull_subscribe( endpoint.subject, endpoint.durable, - stream:, + stream: endpoint.stream, ) end - # Finds the stream for an endpoint, looking it up from the subject when none is configured. - # - # @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to find a stream for. - # - # @raise [ConfigurationError] When no stream is configured and none captures the endpoint's subject. - # - # @return [String] The JetStream stream name. - def stream_for(endpoint) - endpoint.stream || @jetstream.find_stream_name_by_subject(endpoint.subject) - rescue NATS::JetStream::Error::NotFound - raise ConfigurationError, - "JetStream endpoint #{endpoint.name} has no stream and none captures subject - #{endpoint.subject}, " \ - 'create the stream or set `stream:`' - end - - # Verifies that the durable consumer exists, creating it when missing. - # - # @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to ensure. - # - # @return [Object] Consumer metadata from `consumer_info` or `add_consumer`. - def ensure_consumer(endpoint, stream) - @jetstream.consumer_info(stream, endpoint.durable) - rescue NATS::JetStream::Error::NotFound - @jetstream.add_consumer(stream, consumer_config(endpoint)) - end - # Builds the JetStream consumer configuration for an endpoint. # # @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to translate. diff --git a/test/lib/nats_jetstream_consumer_test.rb b/test/lib/nats_jetstream_consumer_test.rb index 1b249e1..1d5d14d 100644 --- a/test/lib/nats_jetstream_consumer_test.rb +++ b/test/lib/nats_jetstream_consumer_test.rb @@ -1,7 +1,6 @@ # frozen_string_literal: true require_relative '../helper' -require 'nats/client' require Rubyists::Leopard.libroot / 'leopard/nats_jetstream_consumer' require Rubyists::Leopard.libroot / 'leopard/nats_jetstream_endpoint' class NatsJetstreamConsumerTest < Minitest::Test @@ -41,26 +40,6 @@ def test_consumer_config_keeps_safe_string_key_overrides assert_equal 100, string_key_config['max_ack_pending'] end - def test_stream_for_with_the_configured_stream - assert_equal 'EVENTS', consumer_with(jetstream: Object.new).send(:stream_for, endpoint_with_consumer(nil)) - end - - def test_stream_for_with_the_subject_when_no_stream_is_configured - jetstream = Object.new - jetstream.define_singleton_method(:find_stream_name_by_subject) { |_subject| 'EVENTS' } - - assert_equal 'EVENTS', consumer_with(jetstream:).send(:stream_for, endpoint_without_stream) - end - - def test_stream_for_raises_a_configuration_error_when_no_stream_has_the_subject - jetstream = Object.new - jetstream.define_singleton_method(:find_stream_name_by_subject) { |_subject| raise NATS::JetStream::Error::NotFound } - - assert_raises(Rubyists::Leopard::ConfigurationError) do - consumer_with(jetstream:).send(:stream_for, endpoint_without_stream) - end - end - def test_start_endpoint_logs_and_reraises_when_the_endpoint_fails_to_start logged = [] logger = Object.new @@ -80,10 +59,6 @@ def consumer_with(jetstream:, logger: Object.new) logger:, process_message: ->(*_) {}) end - def endpoint_without_stream - Rubyists::Leopard::NatsJetstreamEndpoint.new(**base_endpoint_attributes, stream: nil) - end - def symbol_key_config @symbol_key_config ||= @consumer.send(:consumer_config, endpoint_with_consumer(symbol_key_overrides)) end From aedb08d7760beea26cd237c646a0860b04c25a68 Mon Sep 17 00:00:00 2001 From: tiyabelay Date: Wed, 7 Oct 2026 13:54:01 -0700 Subject: [PATCH 3/5] fix: Add empty line --- test/lib/nats_jetstream_consumer_test.rb | 1 + 1 file changed, 1 insertion(+) diff --git a/test/lib/nats_jetstream_consumer_test.rb b/test/lib/nats_jetstream_consumer_test.rb index 1d5d14d..a991edf 100644 --- a/test/lib/nats_jetstream_consumer_test.rb +++ b/test/lib/nats_jetstream_consumer_test.rb @@ -3,6 +3,7 @@ require_relative '../helper' require Rubyists::Leopard.libroot / 'leopard/nats_jetstream_consumer' require Rubyists::Leopard.libroot / 'leopard/nats_jetstream_endpoint' + class NatsJetstreamConsumerTest < Minitest::Test def setup @consumer = Rubyists::Leopard::NatsJetstreamConsumer.new( From b39fe75c9142dbdb2358e01d5ba6d96ec777c938 Mon Sep 17 00:00:00 2001 From: tiyabelay Date: Wed, 7 Oct 2026 15:21:21 -0700 Subject: [PATCH 4/5] fix: Update integration test to reflect reliance on pull_subscribe --- lib/leopard/nats_jetstream_consumer.rb | 7 +++++-- test/integration/nats_jetstream_integration_test.rb | 5 +++-- 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/lib/leopard/nats_jetstream_consumer.rb b/lib/leopard/nats_jetstream_consumer.rb index 14cc45f..9fd015a 100644 --- a/lib/leopard/nats_jetstream_consumer.rb +++ b/lib/leopard/nats_jetstream_consumer.rb @@ -74,7 +74,10 @@ def start_endpoint(endpoint) raise end - # Ensures the durable consumer exists and creates a pull subscription for it. + # Creates a pull subscription for the endpoint's durable consumer. + # + # No stream is passed, so nats-pure resolves the stream from the subject and adds the durable consumer + # (using {#consumer_config}) when it does not already exist. # # @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to subscribe to. # @@ -83,7 +86,7 @@ def build_subscription(endpoint) @jetstream.pull_subscribe( endpoint.subject, endpoint.durable, - stream: endpoint.stream, + config: consumer_config(endpoint), ) end diff --git a/test/integration/nats_jetstream_integration_test.rb b/test/integration/nats_jetstream_integration_test.rb index 4e20be5..a3f97df 100644 --- a/test/integration/nats_jetstream_integration_test.rb +++ b/test/integration/nats_jetstream_integration_test.rb @@ -43,7 +43,7 @@ def build_names token = SecureRandom.hex(4) { stream: "EVENTS_#{token}", - subject: "events.#{token}", + subject: "leopard_it.events.#{token}", durable: "events_consumer_#{token}", service: "JetstreamService#{token}", } @@ -101,6 +101,8 @@ def build_worker(names, middleware: nil, &handler) worker end + # Only the stream is created here. Leopard calls pull_subscribe without a stream, so nats-pure resolves the + # stream from the subject and adds the durable consumer when it does not exist. def create_stream(names) @jetstream.add_stream(name: names[:stream], subjects: [names[:subject]]) @streams << names[:stream] @@ -121,7 +123,6 @@ def build_service_class(names, middleware: nil, &handler) def endpoint_options(names) { - stream: names[:stream], subject: names[:subject], durable: names[:durable], consumer: { ack_wait: 1, max_deliver: 5 }, From 06e88c7a520c0f6dfec0f88110a039bea9b4c56f Mon Sep 17 00:00:00 2001 From: tiyabelay Date: Wed, 7 Oct 2026 15:24:37 -0700 Subject: [PATCH 5/5] fix: Pass endpoint name when durable is not available --- lib/leopard/nats_jetstream_consumer.rb | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/leopard/nats_jetstream_consumer.rb b/lib/leopard/nats_jetstream_consumer.rb index 9fd015a..e94b854 100644 --- a/lib/leopard/nats_jetstream_consumer.rb +++ b/lib/leopard/nats_jetstream_consumer.rb @@ -85,7 +85,7 @@ def start_endpoint(endpoint) def build_subscription(endpoint) @jetstream.pull_subscribe( endpoint.subject, - endpoint.durable, + endpoint.durable || endpoint.name, config: consumer_config(endpoint), ) end