defmodule FluentdForwarder.Handler do @moduledoc false @behaviour :ranch_protocol @type tag :: String.t @type time :: integer @type record :: map @callback format(tag, time, record) :: any use GenServer require Logger defstruct [:socket, :transport, :handler, pending: "", partial: nil] def start_link(ref, transport, opts) do GenServer.start_link(__MODULE__, {ref, transport, opts}) end def init({ref, transport, opts}) do {:ok, %__MODULE__{transport: transport}, {:continue, {:handshake, ref}}} end def handle_continue({:handshake, ref}, %{transport: transport} = state) do {:ok, socket} = :ranch.handshake(ref) transport.setopts(socket, active: true) {:noreply, %{state | socket: socket}} end def handle_info({:tcp, _socket, data}, %{pending: pending} = state) do state = case Msgpax.unpack_slice(pending <> data) do {:ok, msg, pending} -> handle_msg(msg, %{state | pending: pending}) {:error, _} -> %{state | pending: pending <> data} end {:noreply, state} end def handle_info({:tcp_closed, _socket}, state) do {:stop, :normal, state} end def handle_msg([tag, entries, option], state) when is_list(entries) do for [time, record] <- entries do handle_msg(tag, time, record, option, state) end end def handle_msg([tag, time, record, option], state) do handle_msg(tag, time, record, option, state) end def handle_msg( tag, time, %{ "partial_message" => "true", "partial_id" => id, "partial_ordinal" => partial_ordinal, "partial_last" => partial_last, "log" => log } = record, option, %{partial: partial} = state ) do partial = case partial do nil -> {id, String.to_integer(partial_ordinal), log, time} {id, previous_ordinal, past_log, time} -> partial_ordinal = String.to_integer(partial_ordinal) log = if partial_ordinal == previous_ordinal + 1, do: past_log <> log, else: log {id, partial_ordinal, log, time} end partial = if partial_last == "true" do {_id, _partial_ordinal, log, time} = partial record = record |> Map.reject(fn {key, _} -> match?("partial_" <> _, key) end) |> Map.put("log", log) IO.inspect({tag, time, record, option}) nil else partial end %{state | partial: partial} end def handle_msg(tag, time, record, option, state) do IO.inspect({tag, time, record, option}) state end end