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
74 changes: 74 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
name: CI

on:
push:
branches: [master]
pull_request:

jobs:
rspec:
name: rspec (ruby ${{ matrix.ruby }})${{ matrix.advisory && ' [advisory]' || '' }}
runs-on: ${{ matrix.os }}
# 3.4 reports but does not gate, so a failure there cannot block a fix for the 2.7
# that actually runs in production. Promote it to required once consumers have moved.
continue-on-error: ${{ matrix.advisory || false }}

strategy:
fail-fast: false
matrix:
include:
# The Ruby DeployHQ runs today. This one gates.
- ruby: '2.7'
os: ubuntu-22.04 # setup-ruby ships no 2.7 build for ubuntu-24.04
# 2.4.22 is the last bundler line that supports Ruby 2.7; anything newer
# requires >= 3.2 and fails to install. This is the same pin the consuming
# app documents. Do NOT drop to 1.17.3 -- it has a default-gem activation
# bug that has caused production deploy outages there.
bundler: '2.4.22'
# Where the app is heading; advisory until that upgrade lands.
- ruby: '3.4'
os: ubuntu-24.04
bundler: 'latest'
advisory: true

services:
rabbitmq:
# The suite is not hermetic -- spec_helper purges real queues through Bunny, so a
# broker has to be up before any example runs.
image: rabbitmq:3-management-alpine
ports:
- 5672:5672
env:
# NOT guest. RabbitMQ restricts guest to loopback, and a service container is
# reached over the docker bridge, so guest authentication is refused.
RABBITMQ_DEFAULT_USER: leveret
RABBITMQ_DEFAULT_PASS: leveret
options: >-
--health-cmd "rabbitmq-diagnostics -q ping"
--health-interval 5s
--health-timeout 5s
--health-retries 20

steps:
- uses: actions/checkout@v4

- uses: ruby/setup-ruby@v1
with:
ruby-version: ${{ matrix.ruby }}
# Let setup-ruby install bundler, so the version is pinned per Ruby (see matrix)
# rather than resolved at run time. A bare `gem install bundler` picks the newest
# release, which on 2.7 is one that refuses to install.
bundler: ${{ matrix.bundler }}
# No bundler-cache: this gem deliberately ships no Gemfile.lock, so there is no
# stable key to cache against. The dependency set is small.
bundler-cache: false

- name: Install dependencies
run: |
bundle --version
bundle install --jobs 4 --retry 3

- name: rspec
env:
LEVERET_AMQP_URL: amqp://leveret:leveret@localhost:5672
run: bundle exec rspec --format documentation
4 changes: 0 additions & 4 deletions .travis.yml

This file was deleted.

7 changes: 7 additions & 0 deletions lib/leveret.rb
Original file line number Diff line number Diff line change
@@ -1,6 +1,13 @@
require 'bunny'
require 'json'
require 'logger'
# Queue, Worker and DelayQueue all `extend Forwardable`, but nothing required it. Loading the gem
# standalone (its own spec suite) therefore failed with an uninitialized-constant error; under
# Rails it only ever worked because ActiveSupport happens to require forwardable first.
require 'forwardable'
# Worker#run_before_child_exit_hook bounds the host application's hook. Required explicitly
# rather than relying on bunny pulling it in transitively.
require 'timeout'

require 'leveret/configuration'
require 'leveret/delay_queue'
Expand Down
11 changes: 10 additions & 1 deletion lib/leveret/configuration.rb
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,19 @@ module Leveret
# @return [Proc] A proc which will be executed in a child after forking to process a message. Default: +proc {}+
# @!attribute error_handler
# @return [Proc] A proc which will be called if a job raises an exception. Default: +proc {|ex| ex }+
# @!attribute before_child_exit
# @return [Proc] A proc executed in the child immediately before it exits, after the message has been
# acknowledged. The child leaves via +Kernel#exit!+, which skips +at_exit+ and every buffer with it, so
# any sink that batches in memory (an HTTP log shipper, a metrics client) silently loses whatever the
# job wrote last. Use this hook to flush those. Bounded by
# {Worker::CHILD_EXIT_HOOK_TIMEOUT} and rescued, so a slow sink cannot stop a child exiting.
# Default: +proc {}+
# @!attribute concurrent_fork_count
# @return [Integer] The number of jobs that can be processes simultanously. Default: +1+
class Configuration
attr_accessor :amqp, :exchange_name, :queue_name_prefix, :log_file, :log_level, :default_queue_name, :after_fork,
:error_handler, :concurrent_fork_count, :delay_exchange_name, :delay_queue_name, :delay_time
:error_handler, :concurrent_fork_count, :delay_exchange_name, :delay_queue_name, :delay_time,
:before_child_exit

# Create a new instance of Configuration with a set of sane defaults.
def initialize
Expand All @@ -54,6 +62,7 @@ def assign_defaults
self.delay_queue_name = 'leveret_delay_queue'
self.delay_time = 10_000
self.after_fork = proc {}
self.before_child_exit = proc {}
self.error_handler = proc { |ex| ex }
self.concurrent_fork_count = 1
end
Expand Down
49 changes: 49 additions & 0 deletions lib/leveret/worker.rb
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,11 @@ module Leveret
class Worker
extend Forwardable

# Wall-clock bound on Configuration#before_child_exit. Generous enough for a log shipper to
# complete one synchronous HTTP delivery, short enough that an unreachable sink delays a
# child's exit by seconds rather than pinning a fork slot.
CHILD_EXIT_HOOK_TIMEOUT = 5

# @!attribute queues
# @return [Array<Queue>] All of the queues this worker is going to subscribe to
# @!attribute consumers
Expand Down Expand Up @@ -103,13 +108,57 @@ def fork_and_run(incoming_message)
result_handler.handle(result)

log.info "[#{incoming_message.delivery_tag}] Exiting child process #{pid}"
run_before_child_exit_hook
flush_own_log
exit!(0)
end

# Master doesn't need to know how it all went down, the worker will report it's own status back to the queue
Process.detach(pid)
end

# Give the host application a chance to flush anything it has buffered before the child
# leaves via #exit!.
#
# #exit! is deliberate -- it skips at_exit handlers that were registered in the PARENT and
# inherited across the fork, which must not run once per job. But it skips ALL of them, and
# an in-memory buffer flushed by an at_exit handler goes with them. An HTTP log shipper is
# the common case: a batching sink relies on `at_exit { close }` to deliver its tail, so the
# LAST lines a job writes -- the ones saying whether it succeeded -- are the ones most
# reliably lost. Long jobs hide this, because a periodic flush ships everything except the
# final batch; short jobs can lose their entire output.
#
# Runs AFTER the acknowledgement, so a hook that hangs cannot cause redelivery. Bounded and
# rescued for the same reason: an unreachable sink must never stop a child exiting, or forks
# accumulate until the host runs out of processes. A failed flush costs log lines; a wedged
# child costs the worker.
def run_before_child_exit_hook
Timeout.timeout(CHILD_EXIT_HOOK_TIMEOUT) { configuration.before_child_exit.call }
rescue Exception => e # rubocop:disable Lint/RescueException
# Timeout::Error is not a StandardError on older rubies, and this runs microseconds before
# exit! -- there is nothing left to protect by letting anything propagate.
log.warn "before_child_exit hook failed: #{e.class}: #{e.message}"
end

# Flush OUR OWN log before exit! discards it. This gem's log_file defaults to STDOUT, and a
# redirected STDOUT is block-buffered, so `Job returned ...` and `Exiting child process ...`
# -- both written in the child, microseconds before exit! -- are usually never written out.
#
# This is not theoretical. On one production host over two days the log held 4,769
# "Forked to child" lines (written by the PARENT, which exits normally and flushes) against
# 14 "Job returned" lines (written by the CHILD). 0.3%. The result of virtually every job
# this gem has ever run was discarded by its own exit path, which makes the log useless for
# the one question it is most often asked: did that job succeed?
#
# Must run LAST, so it also flushes whatever before_child_exit logged.
def flush_own_log
device = log.instance_variable_get(:@logdev)
io = device.respond_to?(:dev) ? device.dev : nil
io.flush if io.respond_to?(:flush)
rescue Exception # rubocop:disable Lint/RescueException
# Nothing can be reported here -- reporting is what just failed -- and exit! is next.
end

# Constantize the class name in the payload and execute the job with parameters
#
# @param [Parameters] payload The job name and parameters the job requires
Expand Down
3 changes: 2 additions & 1 deletion spec/configuration_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
it 'has a default set of configuration params' do
config = Leveret::Configuration.new

expect(config.amqp).to eq("amqp://guest:guest@localhost:5672")
expect(config.amqp).to eq('amqp://guest:guest@localhost:5672')
expect(config.exchange_name).to eq('leveret_exch')
expect(config.queue_name_prefix).to eq('leveret_queue')
expect(config.delay_exchange_name).to eq('leveret_delay_exch')
Expand All @@ -14,6 +14,7 @@
expect(config.log_level).to eq(Logger::DEBUG)
expect(config.default_queue_name).to eq('standard')
expect(config.after_fork).to be_a(Proc)
expect(config.before_child_exit).to be_a(Proc)
expect(config.error_handler).to be_a(Proc)
expect(config.concurrent_fork_count).to eq(1)
end
Expand Down
6 changes: 6 additions & 0 deletions spec/spec_helper.rb
Original file line number Diff line number Diff line change
@@ -1,5 +1,8 @@
$LOAD_PATH.unshift File.expand_path('../../lib', __FILE__)
require 'leveret'
# Three spec files build doubles with OpenStruct. ostruct is no longer loaded implicitly, so
# without this every example in those files errors before it runs -- 14 of them.
require 'ostruct'

Dir[File.join(File.dirname(__FILE__), 'support/**/*.rb')].each { |f| require f }

Expand All @@ -8,6 +11,9 @@

c.before(:all) do
Leveret.configure do |conf|
# Overridable so CI can point at a broker that does not accept the loopback-only
# `guest` account. Defaults to the same local broker developers already use.
conf.amqp = ENV.fetch('LEVERET_AMQP_URL', 'amqp://guest:guest@localhost:5672')
conf.log_level = Logger::ERROR
conf.queue_name_prefix = 'leveret_test_queue'
conf.default_queue_name = 'test'
Expand Down
94 changes: 93 additions & 1 deletion spec/worker_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,102 @@
end

it 'can use custom queue names' do
queue_names = %w(test other)
queue_names = %w[test other]

worker = Leveret::Worker.new(*queue_names)
expect(worker.queues.map(&:name)).to eq(queue_names)
end
end

# The child leaves via exit!, which skips at_exit and every buffer flushed by one. This hook is
# the only opportunity a batching sink gets to deliver what the job just wrote.
describe '#run_before_child_exit_hook' do
subject(:worker) { Leveret::Worker.new }

def run_hook
worker.send(:run_before_child_exit_hook)
end

around do |example|
original = Leveret.configuration.before_child_exit
example.run
Leveret.configuration.before_child_exit = original
end

it 'calls the configured hook' do
called = false
Leveret.configuration.before_child_exit = proc { called = true }

run_hook

expect(called).to be(true)
end

it 'is a no-op with the default hook' do
expect { run_hook }.not_to raise_error
end

it 'swallows an exception raised by the hook' do
Leveret.configuration.before_child_exit = proc { raise 'sink unavailable' }

expect { run_hook }.not_to raise_error
end

# Timeout::Error does not descend from StandardError on the rubies this gem supports, so a
# bare `rescue StandardError` would let it escape and stop the child exiting.
it 'swallows a non-StandardError raised by the hook' do
Leveret.configuration.before_child_exit = proc { raise Timeout::Error, 'too slow' }

expect { run_hook }.not_to raise_error
end

it 'still flushes our own log when the hook itself blows up' do
Leveret.configuration.before_child_exit = proc { raise 'sink unavailable' }

expect do
run_hook
worker.send(:flush_own_log)
end.not_to raise_error
end

it 'gives up on a hanging hook instead of blocking the exit forever' do
stub_const("#{described_class}::CHILD_EXIT_HOOK_TIMEOUT", 0.1)
Leveret.configuration.before_child_exit = proc { sleep 5 }

started = Time.now
expect { run_hook }.not_to raise_error

expect(Time.now - started).to be < 2
end
end

# The child writes "Job returned ..." and "Exiting child process ..." microseconds before
# exit!, and a redirected STDOUT is block-buffered, so without this both are usually lost.
describe '#flush_own_log' do
# Built BEFORE the logger is stubbed: Worker.new logs while connecting to its queues, so a
# stub installed first would be exercised by construction rather than by the method itself.
let!(:worker) { Leveret::Worker.new }

it "flushes the logger's underlying IO" do
io = StringIO.new
allow(Leveret).to receive(:log).and_return(Logger.new(io))
expect(io).to receive(:flush).at_least(:once)

worker.send(:flush_own_log)
end

it 'does not raise when the logger exposes no flushable device' do
allow(Leveret).to receive(:log).and_return(double('logger'))

expect { worker.send(:flush_own_log) }.not_to raise_error
end

it 'does not raise when flushing itself fails' do
io = StringIO.new
allow(io).to receive(:flush).and_raise(IOError, 'stream closed')
allow(Leveret).to receive(:log).and_return(Logger.new(io))

expect { worker.send(:flush_own_log) }.not_to raise_error
end
end
end
Loading