
Http.pudu
Pudu55 lines2.3 KB
1/** @Log.Sinks.Http.Sink — batches of events posted to an endpoint */2module PuduLangLog.Sinks.Http34import Std.Http.Client as Client5import Std.Http as Web6import PuduLangLog.Formatting.Compact as Compact7import PuduLangLog as Log8import PuduLangLog.Sink as Sink9import PuduLangLog.Sinks.Batching as Batching1011/** @Log.Sinks.Http.Options — endpoint, headers, body format, batching, and transport */12export type Options = {13 endpoint: Str,14 headers: Array[(Str, Str)],15 contentType: Str,16 formatter: fn(&Log.Event) -> Str,17 batching: Batching.Options,18 transport: fn(&Web.Request) -> Result[Int, Str]19}202122export fn defaults(endpoint: Str) -> Options {23 Options{endpoint: endpoint, headers: [], contentType: "application/x-ndjson", formatter: Compact.compact(), batching: Batching.defaults(), transport: transport(Client.limits())}24}25262728export fn transport(limits: Client.Limits) -> fn(&Web.Request) -> Result[Int, Str] {29 fn(request: &Web.Request) -> Result[Int, Str] {30 match Client.send(request, &limits) {31 case Ok(response) => Ok(response.status.code)32 case Err(problem) => Err(Client.explain(&problem))33 }34 }35}363738export fn requestOf(options: &Options, batch: &Array[Log.Event]) -> Web.Request {39 var request = Web.withBody(&Web.request(Web.Post, options.endpoint), batch.map(|event: Log.Event| (options.formatter)(&event)).join(""))40 request = Web.withHeader(&request, "content-type", options.contentType)41 for header in options.headers { request = Web.withHeader(&request, header[0], header[1]) }42 request43}444546export fn sink(options: Options) -> Sink.Sink {47 Batching.sink(Batching.Target {48 emitBatch: fn(batch: &Array[Log.Event]) -> Result[(), Str] {49 let status = (options.transport)(&requestOf(&options, batch)) ?50 if status >= 200 && status < 300 { Ok(()) } else { Err("the endpoint answered " + show(status)) }51 },52 onEmptyBatch: fn() -> () {}53 }, options.batching)54}55