import envoker/config import gleam/float import gleam/http import gleam/http/request import gleam/http/response import gleam/httpc import gleam/int import gleam/io import gleam/json import gleam/option import gleam/order import gleam/result import gleam/string import gleam/time/timestamp import neon/net import neon/udp import parser type HandleError { Garbage(parser.ParseError) SendFailed(httpc.HttpError) StorageFailed(String) } type Config { Config(url: String, port: Int) } fn vl_url(config: Config, url: String) { case request.to(url <> "/blah") { Error(_) -> panic as "VL_URL is malformed" _ -> Nil } Config(..config, url:) } fn port(config: Config, port: Int) { case int.compare(port, 0), int.compare(port, 65_535) { order.Lt, _ | order.Eq, _ | _, order.Gt -> panic as "SYSLOG_PORT is out of bounds" _, _ -> Nil } Config(..config, port:) } pub fn main() { let assert Ok(config) = config.load(Config(url: "", port: 0), [ config.required_string("VL_URL", option.None, vl_url), config.required_int("SYSLOG_PORT", option.None, port), ]) let parser = parser.new() let assert Ok(addr) = net.parse_ip_address("::") let assert Ok(port) = net.port(config.port) let assert Ok(socket) = udp.new(port) |> udp.ip_address(addr) |> udp.open receive_loop(config, parser, socket) } fn receive_loop( config: Config, parser: parser.Parser, socket: udp.Udp, ) -> udp.UdpError { case socket |> udp.receive(0, net.infinity) { Ok(datagram) -> { let r = parser.parse_data(parser, datagram.payload) |> result.map_error(fn(e) { Garbage(e) }) |> result.try(handle_nginx_line(config, _)) case r { Error(e) -> { io.println( string.inspect(e) <> " from " <> net.ip_address_to_string(datagram.ip_address), ) Nil } _ -> Nil } receive_loop(config, parser, socket) } Error(err) -> err } } fn handle_nginx_line( config: Config, line: parser.NginxLine, ) -> Result(Nil, HandleError) { let body = line_to_json(line) let assert Ok(req) = request.to( config.url <> "/insert/jsonline?_stream_fields=app_name,hostname", ) let resp = request.set_body(req, body) |> request.set_method(http.Post) |> httpc.send case resp { Ok(response.Response(status: 200, ..)) -> { Ok(Nil) } Ok(r) -> { Error(StorageFailed("HTTP " <> int.to_string(r.status) <> ": " <> r.body)) } Error(e) -> Error(SendFailed(e)) } } fn line_to_json(line: parser.NginxLine) -> String { let time = line.time |> timestamp.to_unix_seconds |> float.to_string let msg = [ int.to_string(line.response_code), line.method, line.uri, line.remote_addr, line.user_agent, ] |> string.join(" ") json.object([ #("app_name", json.string(line.tag)), #("hostname", json.string(line.hostname)), #("_time", json.string(time)), #("_msg", json.string(msg)), #("remote_addr", json.string(line.remote_addr)), #("remote_user", json.string(line.remote_user)), #("method", json.string(line.method)), #("uri", json.string(line.uri)), #("version", json.string(line.version)), #("response_code", json.string(int.to_string(line.response_code))), #("size", json.string(int.to_string(line.size))), #("referer", json.string(line.referer)), #("user_agent", json.string(line.user_agent)), ]) |> json.to_string }