diff --git a/lib/logtail/log_devices/http.rb b/lib/logtail/log_devices/http.rb index eeb6c7a..9347f8d 100755 --- a/lib/logtail/log_devices/http.rb +++ b/lib/logtail/log_devices/http.rb @@ -1,5 +1,6 @@ require "msgpack" require "net/https" +require "set" require "zlib" require "logtail/config" @@ -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 @@ -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) @@ -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') diff --git a/lib/logtail/logger.rb b/lib/logtail/logger.rb index c2dc67f..77ad07c 100755 --- a/lib/logtail/logger.rb +++ b/lib/logtail/logger.rb @@ -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) diff --git a/spec/logtail/log_devices/http_spec.rb b/spec/logtail/log_devices/http_spec.rb index 2f6ac4e..fd90e1f 100755 --- a/spec/logtail/log_devices/http_spec.rb +++ b/spec/logtail/log_devices/http_spec.rb @@ -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 @@ -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 @@ -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 diff --git a/spec/logtail/logger_spec.rb b/spec/logtail/logger_spec.rb index 25748cf..00e0b09 100755 --- a/spec/logtail/logger_spec.rb +++ b/spec/logtail/logger_spec.rb @@ -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) }