
HttpTest.pudu
Pudu85 lines3.7 KB
1/** @Test.Sinks.Http.Suite — batches posted to an endpoint */2module PuduLangLog.Sinks.HttpTest34import Std.Concurrent as Concurrent5import Std.Http.Client as Client6import Std.Http.Server.Route as Route7import Std.Http.Server as Server8import Std.Http as Web9import Std.Io as Io10import Std.Sync as Sync11import Std.Test as Test12import PuduLangLog.Event as Event13import PuduLangLog as Log14import PuduLangLog.Sinks.Http as Http15import PuduLangLog.Value as Value161718fn event(message: Str) -> Log.Event { Event.create(Log.Timestamp{millis: 0, offset: 0}, Log.Information, message, []) }192021fn main() -> Int {22 let noRequests: Array[Web.Request] = []23 let sent = Sync.cell(noRequests)24 let status = Sync.cell(202)25 let fake = fn(request: &Web.Request) -> Result[Int, Str] {26 if let Ok(held) = Sync.get(&sent) { let _stored = Sync.set(&sent, held.push(*request)) }27 match Sync.get(&status) {28 case Ok(code) => Ok(code)29 case Err(_) => Err("no status")30 }31 }32 let options = Http.Options{..Http.defaults("https://logs.example/ingest"), headers: [("X-Api-Key", "k1")], transport: fake}33 let request = Http.requestOf(&options, &[event("a"), Event.withProperty(&event("b"), "N", Value.int(1))])34 let sink = Http.sink(options)35 let _queued = (sink.emit)(&event("posted"))36 (sink.flush)()37 let _refusing = Sync.set(&status, 500)38 let _failed = (sink.emit)(&event("refused"))39 (sink.flush)()40 let requests = match Sync.get(&sent) {41 case Ok(held) => held42 case Err(_) => []43 }44 (sink.close)()4546 let noBodies: Array[Str] = []47 let received = Sync.cell(noBodies)48 let router = Route.routing(&[Route.post("/ingest", fn(incoming: Route.Request) -> Web.Response {49 if let Ok(held) = Sync.get(&received) { let _stored = Sync.set(&received, held.push(Route.body(&incoming))) }50 Web.Response{status: Web.status(201), headers: [], body: "", binaryBody: None}51 })])52 let server = Server.server(&router)53 let delivered = match Server.start("127.0.0.1", 0) {54 case Ok(running) => {55 let serving = Concurrent.start(fn() -> () { let _served = Server.run(&server, &running, 1) })56 let live = Http.sink(Http.Options{..Http.defaults("http://127.0.0.1:" + show(Server.portOf(&running)) + "/ingest"), transport: Http.transport(Client.limits().reaching(&["127.0.0.1"]))})57 let _sentLive = (live.emit)(&event("over the wire"))58 (live.flush)()59 (live.close)()60 if let Ok(worker) = serving { let _joined = Concurrent.join(&worker) }61 let _stopped = Server.stop(&running)62 match Sync.get(&received) {63 case Ok(held) => held64 case Err(_) => []65 }66 }67 case Err(_) => ["the test server could not start"]68 }6970 let checks = Test.suite("Sinks.Http", &[71 Test.equals("a batch is one post", &request.method, &Web.Post),72 Test.equals("to the endpoint", &request.target, &"https://logs.example/ingest"),73 Test.equals("carrying one compact line per event", &request.body, &"\{\"@t\":\"1970-01-01T00:00:00.0000000Z\",\"@mt\":\"a\"\}\n\{\"@t\":\"1970-01-01T00:00:00.0000000Z\",\"@mt\":\"b\",\"N\":1\}\n"),74 Test.equals("with the content type and headers", &(Web.header(&request.headers, "content-type"), Web.header(&request.headers, "X-Api-Key")), &(Some("application/x-ndjson"), Some("k1"))),75 Test.equals("an accepted batch is sent once", &requests.length(), &2),76 Test.equals("a refused batch is kept for a retry", &requests[1].body.contains("refused"), &true),77 Test.equals("a live endpoint receives the batch", &delivered, &["\{\"@t\":\"1970-01-01T00:00:00.0000000Z\",\"@mt\":\"over the wire\"\}\n"])78 ])79 let ran = Test.run(&checks)80 for failure in Test.failuresOf(&ran) {81 let _reported = Io.writeErrorLine(failure)82 }83 Test.report(&ran)84}85