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
4 changes: 4 additions & 0 deletions lib/logtail/config.rb
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,10 @@ def call(severity, timestamp, progname, msg)

attr_writer :http_body_limit

def initialize
@debug_logger = nil
end

# Whether a particular {Logtail::LogEntry} should be sent to Better Stack
def send_to_better_stack?(log_entry)
!@better_stack_filters&.any? { |blocker| blocker.call(log_entry) }
Expand Down
18 changes: 16 additions & 2 deletions lib/logtail/log_devices/http.rb
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ class HTTP
DEFAULT_INGESTING_SCHEME = "https".freeze
CONTENT_TYPE = "application/msgpack".freeze
USER_AGENT = "Logtail Ruby/#{Logtail::VERSION} (HTTP)".freeze
INITIAL_RECONNECT_WAIT = 1 # second
MAX_RECONNECT_WAIT = 30 # seconds

# Instantiates a new HTTP log device that can be passed to {Logtail::Logger#initialize}.
#
Expand Down Expand Up @@ -83,6 +85,7 @@ def initialize(source_token, options = {})
@request_queue = options[:request_queue] || FlushableDroppingSizedQueue.new(25)
@successive_error_count = 0
@requests_in_flight = 0
@reconnect_wait = INITIAL_RECONNECT_WAIT
end

# Write a new log line message to the buffer, and flush asynchronously if the
Expand Down Expand Up @@ -298,15 +301,19 @@ def build_http
http
end

# Creates a loop that processes the `@request_queue` on an interval.
# Creates a loop that processes the `@request_queue` on an interval. After a failed
# connection it waits before reconnecting, twice as long after every consecutive
# failure up to {MAX_RECONNECT_WAIT}, so an unreachable host is not retried in a busy
# loop. A delivered request starts the wait over (see {#deliver_requests}).
def request_outlet
loop do
http = build_http
connection_healthy = false

begin
Logtail::Config.instance.debug { "Starting HTTP connection" }

http.start do |conn|
connection_healthy = http.start do |conn|
deliver_requests(conn)
end
rescue => e
Expand All @@ -315,6 +322,12 @@ def request_outlet
Logtail::Config.instance.debug { "Finishing HTTP connection" }
http.finish if http.started?
end

next if connection_healthy

Logtail::Config.instance.debug { "Reconnecting in #{@reconnect_wait} seconds" }
sleep(@reconnect_wait)
@reconnect_wait = [@reconnect_wait * 2, MAX_RECONNECT_WAIT].min
end
end

Expand Down Expand Up @@ -360,6 +373,7 @@ def deliver_requests(conn)
num_reqs += 1

@last_resp = resp
@reconnect_wait = INITIAL_RECONNECT_WAIT

Logtail::Config.instance.debug do
if resp.code == "202"
Expand Down
16 changes: 16 additions & 0 deletions spec/logtail/config_spec.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
require "spec_helper"

describe Logtail::Config do
describe "#debug_logger" do
# Ruby before 3.0 warns under `ruby -w` when an unset instance variable is read, and the
# HTTP device reads the debug logger on every connection attempt.
it "can be read before it is set without an uninitialized instance variable warning" do
config = described_class.send(:new)
verbose, $VERBOSE = $VERBOSE, true

expect { config.debug_logger }.not_to output.to_stderr
ensure
$VERBOSE = verbose
end
end
end
81 changes: 81 additions & 0 deletions spec/logtail/log_devices/http_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,87 @@

http.close
end

context "when connecting or delivering fails" do
# Raised from the stubs below to leave the outlet's endless loop. It is not a
# StandardError, so the outlet's own `rescue => e` lets it through.
let(:stop_outlet) { Class.new(Exception) }
let(:http_device) { described_class.new("MYKEY", flush_continuously: false, requests_per_conn: 1) }
let(:request_queue) { http_device.instance_variable_get(:@request_queue) }
let(:waits) { [] }

before do
allow(http_device).to receive(:sleep) { |seconds| waits << seconds }
end

it "waits before reconnecting, twice as long after every refused connection, up to 30 seconds" do
connection_attempts = 0
allow_any_instance_of(Net::HTTP).to receive(:start) do
connection_attempts += 1
raise stop_outlet if connection_attempts > 7
raise Errno::ECONNREFUSED
end

expect { http_device.send(:request_outlet) }.to raise_error(stop_outlet)
expect(waits).to eq([1, 2, 4, 8, 16, 30, 30])
end

it "waits the same way when a request fails on an open connection, and still drops it after 3 attempts" do
request_queue.enq(Logtail::LogDevices::HTTP::RequestAttempt.new(Net::HTTP::Post.new("/")))
request_attempts = 0
allow_any_instance_of(Net::HTTP).to receive(:request) do
request_attempts += 1
raise Errno::ECONNRESET
end
allow(http_device).to receive(:sleep) do |seconds|
waits << seconds
raise stop_outlet if waits.size == 3
end

expect { http_device.send(:request_outlet) }.to raise_error(stop_outlet)
expect(waits).to eq([1, 2, 4])
expect(request_attempts).to eq(3)
expect(request_queue.size).to eq(0)
end

it "starts over at 1 second once a request is delivered" do
request_queue.enq(Logtail::LogDevices::HTTP::RequestAttempt.new(Net::HTTP::Post.new("/")))
connection = double("connection", request: double("response", code: "202"))
connections = [:refused, :refused, :delivers, :refused, :refused]
allow_any_instance_of(Net::HTTP).to receive(:start) do |_http, &block|
case connections.shift
when :refused then raise Errno::ECONNREFUSED
when :delivers then block.call(connection)
else raise stop_outlet
end
end

expect { http_device.send(:request_outlet) }.to raise_error(stop_outlet)
expect(waits).to eq([1, 2, 1, 2])
end
end

it "lets close stop the outlet while it waits to reconnect" do
connection_attempts = 0
allow_any_instance_of(Net::HTTP).to receive(:start) do
connection_attempts += 1
raise Errno::ECONNREFUSED
end
http_device = described_class.new("MYKEY")
http_device.send(:ensure_flush_threads_are_started)
outlet = http_device.instance_variable_get(:@request_outlet_thread)
# Up to 5 seconds for the thread's first attempt, which is slow on a cold TruffleRuby.
500.times do
break if connection_attempts > 0 && outlet.status == "sleep"
sleep 0.01
end

closing = Process.clock_gettime(Process::CLOCK_MONOTONIC)
http_device.close
expect(Process.clock_gettime(Process::CLOCK_MONOTONIC) - closing).to be < 0.5
expect(outlet).not_to be_alive
expect(connection_attempts).to eq(1)
end
end

describe "#deliver_requests" do
Expand Down
Loading