Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 9 additions & 15 deletions lib/leopard/nats_jetstream_consumer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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),
)
Comment thread
TiyaBelay marked this conversation as resolved.
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
Comment thread
TiyaBelay marked this conversation as resolved.

# Builds the JetStream consumer configuration for an endpoint.
#
# @param endpoint [NatsJetstreamEndpoint] The endpoint configuration to translate.
Expand Down
5 changes: 3 additions & 2 deletions test/integration/nats_jetstream_integration_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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}",
}
Expand Down Expand Up @@ -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]
Expand All @@ -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 },
Expand Down
17 changes: 17 additions & 0 deletions test/lib/nats_jetstream_consumer_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading