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