diff --git a/lib/leopard/nats_jetstream_consumer.rb b/lib/leopard/nats_jetstream_consumer.rb index 9695e38..e94b854 100644 --- a/lib/leopard/nats_jetstream_consumer.rb +++ b/lib/leopard/nats_jetstream_consumer.rb @@ -69,33 +69,27 @@ 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. + # 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. # # @return [Object] The JetStream pull subscription. def build_subscription(endpoint) - ensure_consumer(endpoint) @jetstream.pull_subscribe( endpoint.subject, - endpoint.durable, - stream: endpoint.stream, + endpoint.durable || endpoint.name, + config: consumer_config(endpoint), ) 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) - rescue NATS::JetStream::Error::NotFound - @jetstream.add_consumer(endpoint.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/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 }, diff --git a/test/lib/nats_jetstream_consumer_test.rb b/test/lib/nats_jetstream_consumer_test.rb index f31d693..a991edf 100644 --- a/test/lib/nats_jetstream_consumer_test.rb +++ b/test/lib/nats_jetstream_consumer_test.rb @@ -41,8 +41,25 @@ def test_consumer_config_keeps_safe_string_key_overrides assert_equal 100, string_key_config['max_ack_pending'] 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 symbol_key_config @symbol_key_config ||= @consumer.send(:consumer_config, endpoint_with_consumer(symbol_key_overrides)) end