Skip to content
Merged
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
41 changes: 25 additions & 16 deletions lib/outboxer/publisher.rb
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,10 @@ def terminating?
@status == Status::TERMINATING
end

def publishing?
@status == Status::PUBLISHING
end

# Parses command line arguments to configure the publisher.
# @param args [Array<String>] The arguments passed via the command line.
# @return [Hash] The parsed options.
Expand Down Expand Up @@ -549,28 +553,33 @@ def create_publisher_thread(id:, index:,
# Thread.current.report_on_exception = true

while !terminating?
begin
published_message = Message.publish(logger: logger) do |message|
block.call({ id: id }, message)
end
rescue StandardError => error
logger.error(
"#{error.class}: #{error.message}\n" \
"#{error.backtrace.join("\n")}")

Publisher.sleep(
poll_interval,
tick_interval: tick_interval,
process: process,
kernel: kernel)
else
if published_message.nil?
if publishing?
begin
published_message = Message.publish(logger: logger) do |message|
block.call({ id: id }, message)
end
rescue StandardError => error
logger.error(
"#{error.class}: #{error.message}\n" \
"#{error.backtrace.join("\n")}")

Publisher.sleep(
poll_interval,
tick_interval: tick_interval,
process: process,
kernel: kernel)
else
if published_message.nil?
Publisher.sleep(
poll_interval,
tick_interval: tick_interval,
process: process,
kernel: kernel)
end
end
else
Publisher.sleep(tick_interval, tick_interval: tick_interval,
process: process, kernel: kernel)
end
end
rescue ::Exception => error
Expand Down