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
52 changes: 48 additions & 4 deletions lib/logtail/log_devices/http.rb
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,9 @@ def initialize(source_token, options = {})
@flush_continuously = options[:flush_continuously] != false
@flush_interval = options[:flush_interval] || 2 # 2 seconds
@requests_per_conn = options[:requests_per_conn] || 2_500
# The process that owns the queues and threads, see {#reset_if_forked}
@pid = Process.pid
@fork_lock = Mutex.new
@msg_queue = FlushableDroppingSizedQueue.new(@batch_size)
@request_queue = options[:request_queue] || FlushableDroppingSizedQueue.new(25)
@successive_error_count = 0
Expand All @@ -107,6 +110,7 @@ def write(msg)
# Strings, e.g. from a plain ::Logger writing to this device, are sent as info lines.
msg = LogEntry.new(:info, Time.now, nil, msg.to_s.chomp, nil, nil) unless msg.is_a?(LogEntry)
return unless Logtail.config.send_to_better_stack?(msg)
reset_if_forked

@msg_queue.enq(msg)
# No thread delivers what is written after #close, e.g. by an at_exit hook that runs
Expand All @@ -133,17 +137,25 @@ def write(msg)
end

# Flush all log messages in the buffer synchronously. This method will not return
# until delivery of the messages has been successful. If you want to flush
# until delivery of the messages has been successful, or about 5 seconds have passed.
# When no outlet thread runs (`flush_continuously: false`, or a forked child that hasn't
# logged yet), the messages are delivered in the calling thread. If you want to flush
# asynchronously see {#flush_async}.
def flush
reset_if_forked
flush_async
wait_on_request_queue
if @request_outlet_thread && @request_outlet_thread.alive?
wait_on_request_queue
else
deliver_synchronously(dequeue_requests)
end
true
end

# Closes the log device, cleans up, and attempts one last delivery. Closing it again does
# nothing; lines written after it are delivered right away (see {#write}).
def close
reset_if_forked
return if @closed
@closed = true

Expand Down Expand Up @@ -226,6 +238,38 @@ def ensure_flush_threads_are_started
end
end

# The queues and threads belong to the process that created them. After a fork, the
# parent still delivers the lines it buffered, so a child that kept them would send them
# again, and the parent's threads don't run in the child. The child starts over with
# empty queues and starts its own threads once it logs, also when the parent closed the
# device before forking.
def reset_if_forked
return if @pid == Process.pid

@fork_lock.synchronize do
return if @pid == Process.pid

@msg_queue = FlushableDroppingSizedQueue.new(@batch_size)
# The request queue can be a SizedQueue passed as the :request_queue option
@request_queue.respond_to?(:flush) ? @request_queue.flush : @request_queue.clear
@flush_thread = @request_outlet_thread = nil
@requests_in_flight = 0
@reconnect_wait = INITIAL_RECONNECT_WAIT
@closed = @late_delivery_failed = false
@pid = Process.pid
end
end

# Takes the queued requests off the request queue, for {#flush} when no outlet thread
# runs. It checks the size first because a SizedQueue (see :request_queue) blocks when empty.
def dequeue_requests
requests = []
while @request_queue.size > 0 && (request_attempt = @request_queue.deq)
requests << request_attempt.request
end
requests
end

# Builds an HTTP request based on the current messages queued.
def build_request(msgs)
path = '/'
Expand Down Expand Up @@ -447,8 +491,8 @@ def deliver_synchronously(requests)
# Waits on the request queue. This is used in {#flush} to ensure
# the log data has been delivered before returning.
def wait_on_request_queue
# Wait 20 seconds
40.times do |i|
# Wait 5 seconds
10.times do |i|
if @request_queue.size == 0 && @requests_in_flight == 0
Logtail::Config.instance.debug { "Request queue is empty and no requests are in flight, finish waiting" }
return true
Expand Down
15 changes: 15 additions & 0 deletions lib/logtail/logger.rb
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,21 @@ def level=(value)
super
end

# Delivers what was logged so far before it returns, waiting about 5 seconds at most with
# the HTTP log device. Call it before a process ends with `exit!`, which skips the at_exit
# hook that delivers the rest, as Resque's forked job processes do.
#
# Rails calls `flush` on Rails.logger after every request (ActiveSupport::LogSubscriber.flush_all!).
# Waiting there would hold up every request, so that call returns right away and the lines
# are delivered in the background as usual.
def flush
return true if caller_locations(1, 10).any? { |location| location.base_label == "flush_all!" }

@logdev.dev.flush if @logdev && @logdev.dev.respond_to?(:flush)
@extra_loggers.each { |logger| logger.flush if logger.respond_to?(:flush) }
true
end

# @private
def with_context(context, &block)
Logtail::CurrentContext.with(context, &block)
Expand Down
85 changes: 85 additions & 0 deletions spec/logtail/log_devices/http_fork_spec.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
require "spec_helper"

# Forking servers and job runners (Puma, Unicorn, Resque) fork after the app logged. These run
# logtail in a real process that forks, and count what reaches a local ingesting server.
describe Logtail::LogDevices::HTTP, "after a fork" do
let(:ingest) { LocalIngestServer.new }

before do
skip "needs fork" if !Process.respond_to?(:fork) || RUBY_ENGINE == "truffleruby"
end

after { ingest.stop }

it "delivers the lines the parent logged before forking once, not once per process" do
result = run_ruby(<<-RUBY)
require "logtail"
logger = Logtail::Logger.new(Logtail::LogDevices::HTTP.new("token", flush_interval: 60, #{ingest.device_options}))
logger.info("parent line before the fork")
2.times.map { |n| fork { logger.info("child line " + n.to_s) } }.each { |pid| Process.wait(pid) }
RUBY

expect(result.status).to be_success, result.stderr
expect(ingest.messages).to contain_exactly("parent line before the fork", "child line 0", "child line 1")
end

it "delivers the lines a forked child logs" do
result = run_ruby(<<-RUBY)
require "logtail"
logger = Logtail::Logger.new(Logtail::LogDevices::HTTP.new("token", flush_interval: 60, #{ingest.device_options}))
Process.wait(fork { 3.times { |n| logger.info("child line " + n.to_s) } })
logger.info("parent line after the fork")
RUBY

expect(result.status).to be_success, result.stderr
expect(ingest.messages).to contain_exactly("child line 0", "child line 1", "child line 2", "parent line after the fork")
end

it "lets a child that hasn't logged exit right away, without waiting on the parent's lines" do
result = run_ruby(<<-RUBY)
require "logtail"
logger = Logtail::Logger.new(Logtail::LogDevices::HTTP.new("token", flush_interval: 60, #{ingest.device_options}))
logger.info("parent line before the fork")
forking = Process.clock_gettime(Process::CLOCK_MONOTONIC)
Process.wait(fork {})
puts Process.clock_gettime(Process::CLOCK_MONOTONIC) - forking
RUBY

expect(result.stdout.to_f).to be < 5
expect(ingest.messages).to eq(["parent line before the fork"])
end

it "batches a child's lines as usual when the parent closed the device before forking" do
result = run_ruby(<<-RUBY)
require "logtail"
logger = Logtail::Logger.new(Logtail::LogDevices::HTTP.new("token", flush_interval: 60, #{ingest.device_options}))
logger.info("parent line")
logger.close
Process.wait(fork { 3.times { |n| logger.info("child line " + n.to_s) } })
RUBY

expect(result.status).to be_success, result.stderr
expect(ingest.messages).to contain_exactly("parent line", "child line 0", "child line 1", "child line 2")
# Lines written after close are each delivered on their own, see Logtail::LogDevices::HTTP#write
expect(ingest.batch_sizes).to eq([1, 3])
end

it "delivers a child's lines when it calls flush before leaving with exit!, which skips at_exit hooks" do
result = run_ruby(<<-RUBY)
require "logtail"
logger = Logtail::Logger.new(Logtail::LogDevices::HTTP.new("token", flush_interval: 60, #{ingest.device_options}))
logger.info("parent line")
Process.wait(fork do
begin
logger.info("child line")
logger.flush
ensure
exit!(0) # the way Resque ends a job's process
end
end)
RUBY

expect(result.status).to be_success, result.stderr
expect(ingest.messages).to contain_exactly("parent line", "child line")
end
end
87 changes: 87 additions & 0 deletions spec/logtail/log_devices/http_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,93 @@
http.send(:flush)
http.close
end

it "delivers in the calling thread when no outlet thread runs" do
messages = []
stub = stub_request(:post, "https://in.logs.betterstack.com/").to_return do |request|
messages.concat(MessagePack.unpack(Zlib::Inflate.inflate(request.body)).map { |line| line["message"] })
{ status: 202 }
end
http = described_class.new("MYKEY", flush_continuously: false)
http.write(Logtail::LogEntry.new("INFO", time, nil, "test log message 1", nil, nil))
http.write(Logtail::LogEntry.new("INFO", time, nil, "test log message 2", nil, nil))

http.flush

expect(stub).to have_been_requested.once
expect(messages).to eq(["test log message 1", "test log message 2"])
http.close
end

it "doesn't raise when delivering in the calling thread fails, whatever the error" do
http = described_class.new("MYKEY", flush_continuously: false)
http.write(Logtail::LogEntry.new("INFO", time, nil, "test log message", nil, nil))

# WebMock refuses to connect with an error that isn't a StandardError
expect { http.flush }.not_to raise_error
http.close
end

it "doesn't warn about @last_resp when it can't connect while Ruby shuts down, with warnings on" do
# Ruby 2.7 and older warn about an instance variable that is read before it's set
result = run_ruby(<<-RUBY)
$VERBOSE = true
require "logtail"
# Like Net::HTTP before Ruby 4.0 while Ruby shuts down: it can't start the thread that
# times out connecting
Net::HTTP.prepend(Module.new { def start(*); raise ThreadError, "can't alloc thread"; end })
http = Logtail::LogDevices::HTTP.new("token", flush_continuously: false, ingesting_host: "127.0.0.1", ingesting_port: 1, ingesting_scheme: "http")
logger = Logtail::Logger.new(http)
logger.info("line")
logger.flush
RUBY

expect(result.status).to be_success, result.stderr
expect(result.stderr).not_to include("@last_resp not initialized")
end

it "waits about 5 seconds at most for the outlet thread to deliver" do
allow_any_instance_of(Net::HTTP).to receive(:request) { sleep } # Better Stack never answers
http = described_class.new("MYKEY")
http.write(Logtail::LogEntry.new("INFO", time, nil, "test log message", nil, nil))

flushing = Process.clock_gettime(Process::CLOCK_MONOTONIC)
http.flush
expect(Process.clock_gettime(Process::CLOCK_MONOTONIC) - flushing).to be_between(4, 7)

http.instance_variable_get(:@flush_thread).kill.join
http.instance_variable_get(:@request_outlet_thread).kill.join
end

it "waits up to 5 seconds for a slow host when it delivers in the calling thread" do
# Takes 3 seconds to answer each request
slow_ingest = Class.new(LocalIngestServer) do
private

def serve(socket)
def socket.write(*)
sleep 3
super
end
super
end
end.new
result = run_ruby(<<-RUBY)
require "logtail"
# A request per line, which flush delivers as no outlet thread runs
http = Logtail::LogDevices::HTTP.new("token", flush_continuously: false, batch_size: 1, #{slow_ingest.device_options})
logger = Logtail::Logger.new(http)
logger.info("first line")
logger.info("second line")
logger.flush
RUBY

# A request that times out would drop the one after it
expect(result.status).to be_success, result.stderr
expect(slow_ingest.messages).to contain_exactly("first line", "second line")
ensure
slow_ingest.stop if slow_ingest
end
end

# Testing a private method because it helps break down our tests
Expand Down
82 changes: 82 additions & 0 deletions spec/logtail/logger_flush_all_spec.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
require "spec_helper"

# Rails flushes Rails.logger after every request (Rails::Rack::Logger calls
# ActiveSupport::LogSubscriber.flush_all!), and Logger#flush returns right away for that call
# instead of waiting for delivery. These run ActiveSupport's real flush_all! in another process, so
# an ActiveSupport that calls flush differently fails here. The ingesting host holds its answers
# until the direct flush, so a flush_all! that waited for delivery would take 5 seconds.
describe Logtail::Logger, "#flush called by ActiveSupport::LogSubscriber.flush_all!" do
{
"a Logtail::Logger" => "logger",
"a Logtail::Logger in ActiveSupport::TaggedLogging" => "ActiveSupport::TaggedLogging.new(logger)",
"a Logtail::Logger in an ActiveSupport::BroadcastLogger, as Rails 7.1 and later wrap it" =>
"ActiveSupport::BroadcastLogger.new(logger)",
"a Logtail::Logger in ActiveSupport::TaggedLogging in an ActiveSupport::BroadcastLogger" =>
"ActiveSupport::BroadcastLogger.new(ActiveSupport::TaggedLogging.new(logger))"
}.each do |description, rails_logger|
it "returns right away for #{description}, and flushing it directly delivers" do
if rails_logger.include?("BroadcastLogger") && Gem.loaded_specs["activesupport"].version < Gem::Version.new("7.1")
skip "ActiveSupport::BroadcastLogger was added in ActiveSupport 7.1"
end

result = run_ruby(<<-RUBY)
require "logger" # ActiveSupport 6.1 needs it loaded first since concurrent-ruby 1.3.5
require "active_support"
require "active_support/log_subscriber"
require "active_support/tagged_logging"
require "json"
require "logtail"

# The ingesting host keeps the lines of every request, but answers only once `answering` is set
delivered = []
answering = false
server = TCPServer.new("127.0.0.1", 0)
Thread.new do
loop do
Thread.new(server.accept) do |socket|
while socket.gets
length = 0
while (header = socket.gets) != "\\r\\n"
length = header.split(":")[1].to_i if header.downcase.start_with?("content-length:")
end
delivered.concat(MessagePack.unpack(Zlib::Inflate.inflate(socket.read(length))).map { |line| line["message"] })
sleep 0.01 until answering
socket.write("HTTP/1.1 202 Accepted\\r\\nContent-Length: 0\\r\\n\\r\\n")
end
end
end
end

# Logs a line for every request, like Rails' ActionController::LogSubscriber
class RequestLogSubscriber < ActiveSupport::LogSubscriber
def request(event)
info(event.payload[:message])
end
end
RequestLogSubscriber.attach_to(:app)

logger = Logtail::Logger.new(Logtail::LogDevices::HTTP.new("token", flush_interval: 60,
ingesting_host: "127.0.0.1", ingesting_port: server.addr[1], ingesting_scheme: "http"))
ActiveSupport::LogSubscriber.logger = #{rails_logger}
ActiveSupport::Notifications.instrument("request.app", message: "Completed 200 OK")

flushing = Process.clock_gettime(Process::CLOCK_MONOTONIC)
ActiveSupport::LogSubscriber.flush_all!
flush_all_seconds = Process.clock_gettime(Process::CLOCK_MONOTONIC) - flushing
delivered_by_flush_all = delivered.dup

answering = true
ActiveSupport::LogSubscriber.logger.flush
puts JSON.generate("flush_all_seconds" => flush_all_seconds, "delivered_by_flush_all" => delivered_by_flush_all,
"delivered_by_flush" => delivered.drop(delivered_by_flush_all.size))
RUBY

expect(result.status).to be_success, result.stderr
expect(JSON.parse(result.stdout)).to match(
"flush_all_seconds" => a_value < 1,
"delivered_by_flush_all" => [],
"delivered_by_flush" => ["Completed 200 OK"]
)
end
end
end
Loading
Loading