Current section

Files

Jump to
fluentd_forwarder lib fluentd_forwarder handler.ex.orig
Raw

lib/fluentd_forwarder/handler.ex.orig

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