diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..b9584d5 --- /dev/null +++ b/.github/workflows/ci.yml @@ -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 diff --git a/.travis.yml b/.travis.yml deleted file mode 100644 index 6fb8618..0000000 --- a/.travis.yml +++ /dev/null @@ -1,4 +0,0 @@ -language: ruby -rvm: - - 2.2.3 -before_install: gem install bundler -v 1.10.6 diff --git a/lib/leveret.rb b/lib/leveret.rb index d5f302c..fd779a6 100644 --- a/lib/leveret.rb +++ b/lib/leveret.rb @@ -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' diff --git a/lib/leveret/configuration.rb b/lib/leveret/configuration.rb index 9a20802..cbfc339 100644 --- a/lib/leveret/configuration.rb +++ b/lib/leveret/configuration.rb @@ -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 @@ -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 diff --git a/lib/leveret/worker.rb b/lib/leveret/worker.rb index 651c054..704520b 100644 --- a/lib/leveret/worker.rb +++ b/lib/leveret/worker.rb @@ -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] All of the queues this worker is going to subscribe to # @!attribute consumers @@ -103,6 +108,8 @@ 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 @@ -110,6 +117,48 @@ def fork_and_run(incoming_message) 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 diff --git a/spec/configuration_spec.rb b/spec/configuration_spec.rb index 6141dad..e144878 100644 --- a/spec/configuration_spec.rb +++ b/spec/configuration_spec.rb @@ -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') @@ -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 diff --git a/spec/spec_helper.rb b/spec/spec_helper.rb index 99626ba..e057d0d 100644 --- a/spec/spec_helper.rb +++ b/spec/spec_helper.rb @@ -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 } @@ -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' diff --git a/spec/worker_spec.rb b/spec/worker_spec.rb index 4439f63..3529fb0 100644 --- a/spec/worker_spec.rb +++ b/spec/worker_spec.rb @@ -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