Module: Pigeon::Core

Defined in:
lib/pigeon/core.rb

Overview

Core functionality for Pigeon

Class Method Summary collapse

Class Method Details

.configDry::Configurable::Config

Get the configuration

Returns:

  • (Dry::Configurable::Config)


21
22
23
# File 'lib/pigeon/core.rb', line 21

def self.config
  Configuration.config
end

.configure {|config| ... } ⇒ Object

Configure the gem

Examples:

Pigeon.configure do |config|
  config.client_id = "my-application"
  config.kafka_brokers = ["kafka1:9092", "kafka2:9092"]
  config.max_retries = 5
end

Yields:

  • (config)

    Configuration instance



14
15
16
17
# File 'lib/pigeon/core.rb', line 14

def self.configure
  yield(Configuration.config) if block_given?
  initialize_karafka if @karafka_initialized.nil?
end

.initialize_karafkaKarafka::Producer

Initialize the Karafka producer

Returns:

  • (Karafka::Producer)


27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
# File 'lib/pigeon/core.rb', line 27

def self.initialize_karafka # rubocop:disable Metrics/AbcSize
  return @karafka_producer if @karafka_initialized

  # Configure Karafka
  begin
    Karafka::Setup::Config.setup do |karafka_config|
      karafka_config.client_id = config.client_id

      # Set required Kafka configuration if not provided
      if !config.karafka_config[:kafka] || config.karafka_config[:kafka].empty?
        karafka_config.kafka = {
          "bootstrap.servers": config.kafka_brokers.join(",")
        }
      end

      # Apply any additional Karafka configuration
      config.karafka_config.each do |key, value|
        karafka_config.public_send("#{key}=", value) if karafka_config.respond_to?("#{key}=")
      end
    end

    @karafka_initialized = true
    @karafka_producer = Karafka.producer
  rescue StandardError => e
    config.logger.error("Failed to initialize Karafka: #{e.message}")
    # Return a mock producer for testing
    @karafka_initialized = true
    @karafka_producer = MockProducer.new
  end

  @karafka_producer
end

.karafka_producerKarafka::Producer

Get the Karafka producer instance

Returns:

  • (Karafka::Producer)


62
63
64
65
# File 'lib/pigeon/core.rb', line 62

def self.karafka_producer
  initialize_karafka unless @karafka_initialized
  @karafka_producer
end

.processing?Boolean

Check if processing is running

Returns:

  • (Boolean)

    Whether processing is running



100
101
102
# File 'lib/pigeon/core.rb', line 100

def self.processing?
  @processor&.processing? || false
end

.processor(auto_start: false) ⇒ Pigeon::Processor

Create a new processor instance

Parameters:

  • auto_start (Boolean) (defaults to: false)

    Whether to automatically start processing pending messages

Returns:



76
77
78
# File 'lib/pigeon/core.rb', line 76

def self.processor(auto_start: false)
  Processor.new(auto_start: auto_start)
end

.publisherPigeon::Publisher

Create a new publisher instance

Returns:



69
70
71
# File 'lib/pigeon/core.rb', line 69

def self.publisher
  Publisher.new
end

.start_processing(batch_size: 100, interval: 5, thread_count: 2) ⇒ Boolean

Start processing pending messages

Parameters:

  • batch_size (Integer) (defaults to: 100)

    Number of messages to process in one batch

  • interval (Integer) (defaults to: 5)

    Interval in seconds between processing batches

  • thread_count (Integer) (defaults to: 2)

    Number of threads to use for processing

Returns:

  • (Boolean)

    Whether processing was started



85
86
87
88
# File 'lib/pigeon/core.rb', line 85

def self.start_processing(batch_size: 100, interval: 5, thread_count: 2)
  @processor ||= processor
  @processor.start_processing(batch_size: batch_size, interval: interval, thread_count: thread_count)
end

.stop_processingBoolean

Stop processing pending messages

Returns:

  • (Boolean)

    Whether processing was stopped



92
93
94
95
96
# File 'lib/pigeon/core.rb', line 92

def self.stop_processing
  return false unless @processor

  @processor.stop_processing
end