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
140 changes: 139 additions & 1 deletion lib/logtail/log_devices/http.rb
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
require "msgpack"
require "net/https"
require "set"
require "zlib"

require "logtail/config"
Expand All @@ -22,6 +23,8 @@ class HTTP
DEFAULT_INGESTING_SCHEME = "https".freeze
CONTENT_TYPE = "application/msgpack".freeze
USER_AGENT = "Logtail Ruby/#{Logtail::VERSION} (HTTP)".freeze
ENCODABLE_INTEGERS = (-2**63...2**64).freeze # the integers msgpack can encode
MAX_UNTRACKED_DEPTH = 100 # nested hashes and arrays, see #encodable_value
INITIAL_RECONNECT_WAIT = 1 # second
MAX_RECONNECT_WAIT = 30 # seconds

Expand Down Expand Up @@ -93,6 +96,8 @@ def initialize(source_token, options = {})
# size is constricted by the Logtail API. The actual application limit is a multiple
# of this. Hence the `@request_queue`.
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)

@msg_queue.enq(msg)
Expand Down Expand Up @@ -208,11 +213,144 @@ def build_request(msgs)
req['Content-Type'] = CONTENT_TYPE
req['Content-Encoding'] = 'deflate'
req['User-Agent'] = USER_AGENT
uncompressed = msgs.map { |msg| force_utf8_encoding(msg.to_hash) }.to_msgpack
# Entries are encoded one at a time, so one that can't be encoded doesn't lose the batch.
packer = MessagePack::DefaultFactory.packer
uncompressed = packer.write_array_header(msgs.size).to_s
packer.reset
msgs.each { |msg| uncompressed << encode_log_entry(msg, packer) }
req.body = Zlib::Deflate.deflate(uncompressed, Zlib::BEST_SPEED)
req
end

# Encodes a single log entry with msgpack, with the packer if given, which it leaves empty.
# An entry that still can't be encoded is replaced by one that says why, with the same
# level and time.
def encode_log_entry(msg, packer = MessagePack::DefaultFactory.packer)
packer.write(encodable_value(msg.to_hash)).to_s
rescue StandardError, SystemStackError => e
Logtail::Config.instance.debug { "Could not encode log entry: #{e.inspect}" }
error = force_utf8_encoding("#{e.class}: #{e.message}")
message = "Logtail could not encode this log line (#{error}): #{force_utf8_encoding(msg.message)}"
{
level: msg.level,
dt: msg.time.iso8601(LogEntry::DT_PRECISION),
message: message.byteslice(0, LogEntry::MESSAGE_MAX_BYTES).scrub(""),
}.to_msgpack
ensure
packer.reset
end

# Converts what msgpack can't encode, recursively, mostly into strings, and passes strings
# that aren't valid UTF-8 to {#force_utf8_encoding}. Returns the value itself when nothing
# needs to change, as for most log lines, and otherwise copies only the hashes and arrays
# that change. A hash or array that contains itself is cut off with "[circular]".
def encodable_value(value)
# The first pass doesn't keep track of the hashes and arrays it is in, and gives up when
# they nest too deep, as in a cycle. The second pass keeps track of them to find cycles.
catch(:too_deep) { return replacement_for(value, nil, 0) || value }
replacement_for(value, {}.compare_by_identity, 0) || value
end

# Returns what to send instead of the value, or nil to send the value as it is.
def replacement_for(value, parents, depth)
case value
when Hash
hash_replacement(value, parents, depth)
when String
force_utf8_encoding(value) unless value.valid_encoding? && (value.encoding == Encoding::UTF_8 || value.encoding == Encoding::US_ASCII)
when Integer
value.to_s unless value.bit_length < 64 || ENCODABLE_INTEGERS.cover?(value)
when nil, true, false, Symbol, Float
nil
when Array, Set, Struct
if parents
return "[circular]" if parents.key?(value)

parents[value] = true
elsif depth == MAX_UNTRACKED_DEPTH
throw :too_deep
end
replacement =
if value.is_a?(Array)
array_replacement(value, parents, depth + 1)
elsif value.is_a?(Set)
array_replacement(items = value.to_a, parents, depth + 1) || items
else
hash_replacement(members = value.to_h, parents, depth + 1) || members
end
parents.delete(value) if parents
replacement
else
force_utf8_encoding(converted_value(value))
end
end

# Returns a copy of the hash with the replacements for its keys and values, or nil if none
# needs one. The most common keys and values are checked right here, which is faster.
def hash_replacement(hash, parents, depth)
if parents
return "[circular]" if parents.key?(hash)

parents[hash] = true
elsif depth == MAX_UNTRACKED_DEPTH
throw :too_deep
end
copy = nil
key_changes = false
hash.each_pair do |key, item|
new_key = replacement_for(key, parents, depth + 1) unless key.is_a?(Symbol)
new_item =
if item.is_a?(String)
force_utf8_encoding(item) unless item.valid_encoding? && (item.encoding == Encoding::UTF_8 || item.encoding == Encoding::US_ASCII)
elsif item.is_a?(Hash)
hash_replacement(item, parents, depth + 1)
elsif !(item.nil? || item.is_a?(Integer) && item.bit_length < 64 || item.is_a?(Symbol) || item.is_a?(Float))
replacement_for(item, parents, depth + 1)
end
if new_key
key_changes = true
break
elsif new_item
(copy ||= Hash[hash])[key] = new_item
end
end
# A key that changes is rare, the copy is then built from scratch to keep the order of the keys
if key_changes
copy = {}
hash.each_pair { |key, item| copy[replacement_for(key, parents, depth + 1) || key] = replacement_for(item, parents, depth + 1) || item }
end
parents.delete(hash) if parents
copy
end

# Returns a copy of the array with the replacements for its items, or nil if none needs one.
def array_replacement(array, parents, depth)
copy = nil
array.each_with_index do |item, index|
new_item = replacement_for(item, parents, depth)
(copy ||= Array.new(array))[index] = new_item if new_item
end
copy
end

# Converts a value msgpack can't encode that isn't a hash, array, set or struct.
def converted_value(value)
case value
when Time, DateTime # Rails makes ActiveSupport::TimeWithZone match Time too
value.to_time.getutc.iso8601(LogEntry::DT_PRECISION)
when Date
value.iso8601
when Exception
{ class: value.class.name, message: value.message }
when Numeric
# BigDecimal#to_s would use an exponent, "0.1999e2"
defined?(::BigDecimal) && value.is_a?(::BigDecimal) ? value.to_s("F") : value.to_s
else
# The public id of a Rack::Session::SessionId is the cookie of a server-side session
value.respond_to?(:private_id) ? value.private_id : value.to_s
end
end

def force_utf8_encoding(data)
if data.respond_to?(:force_encoding)
data.dup.force_encoding('UTF-8')
Expand Down
6 changes: 6 additions & 0 deletions lib/logtail/logger.rb
Original file line number Diff line number Diff line change
Expand Up @@ -233,6 +233,12 @@ def add(severity, message = nil, progname = nil, &block)
super
end

# Logs the message as an info line. ::Logger#<< writes it to the log device as it is, which
# the HTTP device can't deliver. Rack::CommonLogger, for example, logs requests this way.
def <<(msg)
info(msg.to_s.chomp)
end

# Backwards compatibility with older ActiveSupport::Logger versions
Logger::Severity.constants.each do |severity|
class_eval(<<-EOT, __FILE__, __LINE__ + 1)
Expand Down
138 changes: 137 additions & 1 deletion spec/logtail/log_devices/http_spec.rb
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
require "spec_helper"
require "bigdecimal"
require "date"
require "set"

# Note: these tests access instance variables and private methods as a means of
# not muddying the public API. This object should expose a simple buffer like
Expand All @@ -22,7 +25,7 @@

it "should buffer the messages" do
http.write("test log message")
expect(msg_queue.flush).to eq(["test log message"])
expect(msg_queue.flush.map(&:message)).to eq(["test log message"])
http.close
end

Expand Down Expand Up @@ -108,6 +111,139 @@
end
end

# Testing a private method because it helps break down our tests
describe "#build_request" do
let(:http) { described_class.new("MYKEY", flush_continuously: false) }
let(:logger) { Logtail::Logger.new(http) }

# Flushes the buffer and decodes the request body, the way the API reads it.
def delivered_entries
http.send(:flush_async)
request = http.instance_variable_get(:@request_queue).deq.request
MessagePack.unpack(Zlib::Inflate.inflate(request.body))
end

a_proc = proc {}
an_object = Object.new
a_cyclic_hash = { name: "parent" }
a_cyclic_hash[:self] = a_cyclic_hash
a_cyclic_array = ["parent"]
a_cyclic_array << a_cyclic_array
# Like a Rack::Session::SessionId, whose public id is the cookie of a server-side session
a_session_id = Object.new
def a_session_id.private_id
"2::hashed-session-id"
end
def a_session_id.to_s
"session-cookie"
end

{
"a Time" => [Time.utc(2026, 10, 1, 12, 0, 0, 123456), "2026-10-01T12:00:00.123456Z"],
"a Time with a UTC offset" => [Time.new(2026, 10, 1, 14, 0, 0, "+02:00"), "2026-10-01T12:00:00.000000Z"],
"a DateTime" => [DateTime.new(2026, 10, 1, 14, 0, 0, "+02:00"), "2026-10-01T12:00:00.000000Z"],
"a Date" => [Date.new(2026, 10, 1), "2026-10-01"],
"a BigDecimal" => [BigDecimal("19.99"), "19.99"],
"a Rational" => [Rational(1, 3), "1/3"],
"an Integer above the 64-bit range" => [2**64, "18446744073709551616"],
"an Integer below the 64-bit range" => [-2**63 - 1, "-9223372036854775809"],
"a Set" => [Set[1, 2], [1, 2]],
"a Struct" => [Struct.new(:id, :name).new(1, "Ann"), { "id" => 1, "name" => "Ann" }],
"an Exception" => [ArgumentError.new("boom"), { "class" => "ArgumentError", "message" => "boom" }],
"a Range" => [1..2, "1..2"],
"a Class" => [String, "String"],
"a Proc" => [a_proc, a_proc.to_s],
"an arbitrary object" => [an_object, an_object.to_s],
"an object with a private id" => [a_session_id, "2::hashed-session-id"],
"a Hash that contains itself" => [a_cyclic_hash, { "name" => "parent", "self" => "[circular]" }],
"an Array that contains itself" => [a_cyclic_array, ["parent", "[circular]"]],
}.each do |description, (value, expected)|
it "delivers the whole batch when a log line holds #{description}" do
logger.info("line before")
logger.info("line with the value", value: value)
logger.info("line after")

entries = delivered_entries
expect(entries.map { |entry| entry["message"] }).to eq(["line before", "line with the value", "line after"])
expect(entries[1]["value"]).to eq(expected)
end
end

it "replaces a log line it still can't encode with one that says why, and delivers the rest" do
unencodable = Object.new
def unencodable.to_s
raise "to_s failed"
end

logger.info("line before")
logger.warn("line with the value", value: unencodable)
logger.info("line after")

entries = delivered_entries
expect(entries.map { |entry| entry["message"] }).to eq([
"line before",
"Logtail could not encode this log line (RuntimeError: to_s failed): line with the value",
"line after",
])
expect(entries[1].keys).to contain_exactly("level", "dt", "message")
expect(entries[1]["level"]).to eq("warn")
expect(entries[1]["dt"]).to match(/\A\d{4}-\d\d-\d\dT\d\d:\d\d:\d\d\.\d{6}Z\z/)
end

it "sends a log line that msgpack can encode as it is, without copying it" do
hash = { message: "caf\u00E9", count: 1, ratio: 0.5, flag: true, none: nil, level: :info, nested: { list: [1, "two"] } }

expect(http.send(:encodable_value, hash)).to be(hash)
end

it "passes strings that aren't valid UTF-8 to force_utf8_encoding, also in arrays and keys" do
in_array = "in an array \xFF".b
key = "key \xFF".b
allow(http).to receive(:force_utf8_encoding).and_call_original

logger.info("line", items: [in_array], counts: { key => 1 })
delivered_entries

expect(http).to have_received(:force_utf8_encoding).with(in_array).at_least(:once)
expect(http).to have_received(:force_utf8_encoding).with(key).at_least(:once)
end

it "keeps the order of the keys of a hash when it converts some of them" do
logger.info("line", value: { "a" => 1, Time.utc(2026, 10, 1) => 2, "c" => Date.new(2026, 10, 1) })

expect(delivered_entries[0]["value"].to_a).to eq([["a", 1], ["2026-10-01T00:00:00.000000Z", 2], ["c", "2026-10-01"]])
end

it "leaves the logged values as they are" do
value = { at: Time.utc(2026, 10, 1), nested: { on: Date.new(2026, 10, 1), list: [Set[1], "\xFF".b] } }
original = Marshal.load(Marshal.dump(value))

logger.info("line", value: value)
delivered_entries

expect(value).to eq(original)
end

it "delivers hashes and arrays nested 110 levels deep" do
value = "leaf"
110.times { |level| value = level.even? ? { "level #{level}" => value } : [value] }

logger.info("line", value: value)

expect(delivered_entries[0]["value"]).to eq(value)
end

it "delivers strings written to the device as info lines" do
http.write("written to the device\n")
::Logger.new(http).warn("logged by a plain Ruby logger")

entries = delivered_entries
expect(entries.map { |entry| entry["level"] }).to eq(["info", "info"])
expect(entries[0]["message"]).to eq("written to the device")
expect(entries[1]["message"]).to end_with("WARN -- : logged by a plain Ruby logger")
end
end

# Testing a private method because it helps break down our tests
describe "#intervaled_flush" do
it "should start a intervaled flush thread and flush on an interval" do
Expand Down
15 changes: 15 additions & 0 deletions spec/logtail/logger_spec.rb
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,21 @@
end
end

describe "#<<" do
it "logs the line at the info level" do
http_device = Logtail::LogDevices::HTTP.new("MYKEY", flush_continuously: false)
logger = Logtail::Logger.new(http_device)

logger << "appended line\n"

http_device.send(:flush_async)
request = http_device.instance_variable_get(:@request_queue).deq.request
entries = MessagePack.unpack(Zlib::Inflate.inflate(request.body))
expect(entries.size).to eq(1)
expect(entries[0]).to include("level" => "info", "message" => "appended line")
end
end

describe "#error" do
let(:io) { StringIO.new }
let(:logger) { Logtail::Logger.new(io) }
Expand Down
Loading