katie/http-syslog-transformer

parses and forwards Nginx syslog data to VictoriaLogs

Katie KlossUse envoker for config parsing4b60201

main
3.5 KiB154 linesraw
1import envoker/config
2import gleam/float
3import gleam/http
4import gleam/http/request
5import gleam/http/response
6import gleam/httpc
7import gleam/int
8import gleam/io
9import gleam/json
10import gleam/option
11import gleam/order
12import gleam/result
13import gleam/string
14import gleam/time/timestamp
15import neon/net
16import neon/udp
17import parser
18
19type HandleError {
20  Garbage(parser.ParseError)
21  SendFailed(httpc.HttpError)
22  StorageFailed(String)
23}
24
25type Config {
26  Config(url: String, port: Int)
27}
28
29fn vl_url(config: Config, url: String) {
30  case request.to(url <> "/blah") {
31    Error(_) -> panic as "VL_URL is malformed"
32    _ -> Nil
33  }
34  Config(..config, url:)
35}
36
37fn port(config: Config, port: Int) {
38  case int.compare(port, 0), int.compare(port, 65_535) {
39    order.Lt, _ | order.Eq, _ | _, order.Gt ->
40      panic as "SYSLOG_PORT is out of bounds"
41    _, _ -> Nil
42  }
43  Config(..config, port:)
44}
45
46pub fn main() {
47  let assert Ok(config) =
48    config.load(Config(url: "", port: 0), [
49      config.required_string("VL_URL", option.None, vl_url),
50      config.required_int("SYSLOG_PORT", option.None, port),
51    ])
52
53  let parser = parser.new()
54
55  let assert Ok(addr) = net.parse_ip_address("::")
56  let assert Ok(port) = net.port(config.port)
57  let assert Ok(socket) =
58    udp.new(port)
59    |> udp.ip_address(addr)
60    |> udp.open
61
62  receive_loop(config, parser, socket)
63}
64
65fn receive_loop(
66  config: Config,
67  parser: parser.Parser,
68  socket: udp.Udp,
69) -> udp.UdpError {
70  case socket |> udp.receive(0, net.infinity) {
71    Ok(datagram) -> {
72      let r =
73        parser.parse_data(parser, datagram.payload)
74        |> result.map_error(fn(e) { Garbage(e) })
75        |> result.try(handle_nginx_line(config, _))
76
77      case r {
78        Error(e) -> {
79          io.println(
80            string.inspect(e)
81            <> " from "
82            <> net.ip_address_to_string(datagram.ip_address),
83          )
84          Nil
85        }
86        _ -> Nil
87      }
88
89      receive_loop(config, parser, socket)
90    }
91    Error(err) -> err
92  }
93}
94
95fn handle_nginx_line(
96  config: Config,
97  line: parser.NginxLine,
98) -> Result(Nil, HandleError) {
99  let body = line_to_json(line)
100
101  let assert Ok(req) =
102    request.to(
103      config.url <> "/insert/jsonline?_stream_fields=app_name,hostname",
104    )
105
106  let resp =
107    request.set_body(req, body)
108    |> request.set_method(http.Post)
109    |> httpc.send
110
111  case resp {
112    Ok(response.Response(status: 200, ..)) -> {
113      Ok(Nil)
114    }
115    Ok(r) -> {
116      Error(StorageFailed("HTTP " <> int.to_string(r.status) <> ": " <> r.body))
117    }
118    Error(e) -> Error(SendFailed(e))
119  }
120}
121
122fn line_to_json(line: parser.NginxLine) -> String {
123  let time =
124    line.time
125    |> timestamp.to_unix_seconds
126    |> float.to_string
127
128  let msg =
129    [
130      int.to_string(line.response_code),
131      line.method,
132      line.uri,
133      line.remote_addr,
134      line.user_agent,
135    ]
136    |> string.join(" ")
137
138  json.object([
139    #("app_name", json.string(line.tag)),
140    #("hostname", json.string(line.hostname)),
141    #("_time", json.string(time)),
142    #("_msg", json.string(msg)),
143    #("remote_addr", json.string(line.remote_addr)),
144    #("remote_user", json.string(line.remote_user)),
145    #("method", json.string(line.method)),
146    #("uri", json.string(line.uri)),
147    #("version", json.string(line.version)),
148    #("response_code", json.string(int.to_string(line.response_code))),
149    #("size", json.string(int.to_string(line.size))),
150    #("referer", json.string(line.referer)),
151    #("user_agent", json.string(line.user_agent)),
152  ])
153  |> json.to_string
154}