Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-httpclient

Named HTTP clients for Pudu: a client factory, delegating handlers, pooled keep-alive connections, handler lifetimes, logging, and resilience

0.1.1Apache-2.01

InstallClose

PlaceholderResilienceTest.pudu

Pudu148 lines7.9 KB

GitHub ↗
1/** @Test.Integration.PlaceholderResilience.Suite — the placeholder service under faults, latency, and load */2module Integration.PlaceholderResilienceTest34import Std.Concurrent as Concurrent5import Std.Io as Io6import Std.Result as Result7import Std.Sync as Sync8import Std.Test as Test9import PuduLangHttpClient.Client as Client10import PuduLangHttpClient.Events as Events11import PuduLangHttpClient.Factory.Builder as Builder12import PuduLangHttpClient.Factory as Factory13import PuduLangHttpClient.Handlers.Logging as Logging14import PuduLangHttpClient.Handlers.Metrics as Metrics15import PuduLangHttpClient.Handlers.Resilience as Resilience16import PuduLangHttpClient as HttpClient17import PuduLangHttpClient.Response as Response18import PuduLangHttpClient.Utils.Shared as Shared19import PuduLangLog.Configuration as Configuration20import PuduLangLog as Log21import PuduLangLog.Sinks.Memory as Memory22import PuduLangResilience.CircuitBreaker as CircuitBreaker23import PuduLangResilience.Pipeline as Pipeline24import PuduLangResilience.Retry as Retry25import Support.Placeholder.Api as Api26import Support.Placeholder.Server as Placeholder2728/// The status or failure of a request.29fn statusOf(outcome: HttpClient.Outcome[Response.Response]) -> Result[Int, HttpClient.Failure] {30  Result.map(outcome, |response: Response.Response| response.status.code)31}3233/// Whether an outcome is a resilience rejection.34fn isRejected(outcome: &Result[Int, HttpClient.Failure]) -> Bool {35  match outcome {36    case Err(HttpClient.Rejected(_)) => true37    case _ => false38  }39}4041/// Runs the suite.42fn main() -> Int {43  let service = Placeholder.start()44  let memory = Memory.create()45  let logger = Configuration.create().minimumLevel(Log.Verbose).writeTo(Memory.sink(&memory)).createLogger()46  let meter = Metrics.meter()47  let failures = Shared.shared([])48  let changes = Shared.shared(0)49  let listeners = Events.listeners()50    .onRequestFailed(|event: Events.RequestFailed| Shared.update(&failures, |held: Array[Str]| held.push(event.client + ": " + event.failure)))51    .onPipelineChanged(|_event: Events.PipelineChanged| Shared.update(&changes, |count: Int| count + 1))52  let options = Resilience.standardOptions()53  let standard = match Resilience.standardWith(&Resilience.Standard{..options, retry: Retry.Options{..options.retry, delay: 5, useJitter: false}}) {54    case Ok(built) => built55    case Err(invalid) => panic(Pipeline.explain(&invalid))56  }57  let breaker = match Pipeline.build([CircuitBreaker.strategy(CircuitBreaker.Options{..CircuitBreaker.defaults(), shouldHandle: Resilience.transient(), failureRatio: 0.5d, minimumThroughput: 2, samplingDuration: 5000, breakDuration: 5000})]) {58    case Ok(built) => built59    case Err(invalid) => panic(Pipeline.explain(&invalid))60  }61  let factory = match Factory.buildWith(&Factory.Options{listeners: listeners}, [62      Builder.defaults().withBaseAddress(service.base + "/").withDefaultHeader("authorization", "Bearer secret-token").observedBy(Logging.observer(logger)).withHandler(Metrics.handler(&meter)),63      Builder.named("resilient").withHandler(Resilience.handler(standard)),64      Builder.named("guarded").withHandler(Resilience.handler(breaker)),65      Builder.named("impatient").withTimeout(100),66      Builder.named("rotating").withHandlerLifetime(25),67      Builder.named("small").withMaxResponseContentBufferSize(1000)68    ]) {69    case Ok(built) => built70    case Err(invalid) => panic(Factory.explain(&invalid))71  }7273  Placeholder.failNext(&service, 2)74  let before = Placeholder.handled(&service)75  let recovered: HttpClient.Outcome[Api.Post] = Api.find(&Api.Api{client: Factory.createClient(&factory, "resilient")}, "posts", 1)76  let attempts = Placeholder.handled(&service) - before7778  Placeholder.failNext(&service, 50)79  let guarded = Factory.createClient(&factory, "guarded")80  let firstFault = statusOf(Client.get(&guarded, "posts/1"))81  let secondFault = statusOf(Client.get(&guarded, "posts/1"))82  let reachedBeforeOpen = Placeholder.handled(&service)83  let rejected = statusOf(Client.get(&guarded, "posts/1"))84  let reachedAfterOpen = Placeholder.handled(&service)85  Placeholder.failNext(&service, 0)8687  let impatient = Factory.createClient(&factory, "impatient")88  let slow = statusOf(Client.get(&impatient, "posts/1?_delay=400"))89  let quick = statusOf(Client.get(&impatient, "posts/1"))9091  let pending = Factory.createClient(&factory, "resilient")92  let cancelled = Sync.cell(Ok(0))93  let worker = match Concurrent.start(fn() -> () { let _set = Sync.set(&cancelled, statusOf(Client.get(&pending, "posts?_delay=2000"))) }) {94    case Ok(started) => started95    case Err(problem) => panic(show(problem))96  }97  let _settle = Concurrent.sleep(50)98  Client.cancelPending(&pending)99  let _joined = Concurrent.join(&worker)100101  let plain = Factory.createClient(&factory, "rotating")102  let buffered = Result.map(Client.getBytes(&plain, "photos"), |bytes: Bytes| bytes.length())103  let streamedBytes = Sync.counter(0)104  let streamed = Client.stream(&plain, "photos", fn(piece: Bytes) -> Bool {105      let _added = Sync.increment(&streamedBytes, piece.length())106      true107    })108  let tooLarge = Client.getBytes(&Factory.createClient(&factory, "small"), "photos")109110  var rotations: Array[Result[Int, HttpClient.Failure]] = []111  var round = 0112  while round < 6 {113    rotations = rotations.push(statusOf(Client.get(&Factory.createClient(&factory, "rotating"), "todos/1")))114    let _paced = Concurrent.sleep(15)115    round = round + 1116  }117118  let snapshot = Metrics.snapshot(&meter)119  let messages = Memory.messages(&memory)120  Factory.dispose(&factory)121  Placeholder.stop(&service)122123  let checks = Test.suite("Integration.PlaceholderResilience", &[124      Test.equals("two injected faults are retried until the record arrives", &(Result.map(recovered, |post: Api.Post| post.id), attempts), &(Ok(1), 3)),125      Test.equals("a breaker passes faults until it has seen enough of them", &(firstFault, secondFault), &(Ok(503), Ok(503))),126      Test.that("an open breaker rejects without reaching the service", isRejected(&rejected) && reachedAfterOpen == reachedBeforeOpen),127      Test.equals("a request slower than the client's timeout times out with that timeout", &slow, &Err(HttpClient.TimedOut(100))),128      Test.equals("a quick request through the same client still answers", &quick, &Ok(200)),129      Test.equals("cancelling pending requests stops a slow one in flight", &Result.unwrapOr(Sync.get(&cancelled), Ok(0)), &Err(HttpClient.Cancelled("the client cancelled its pending requests"))),130      Test.that("a streamed body is as long as the buffered one", match buffered {131          case Ok(length) => length > 500000 && Result.unwrapOr(Sync.count(&streamedBytes), 0) == length && Result.isOk(&streamed)132          case Err(_) => false133        }),134      Test.equals("a body over the client's buffer limit is refused", &Result.map(tooLarge, |bytes: Bytes| bytes.length()), &Err(HttpClient.ContentTooLarge(1000))),135      Test.equals("requests keep answering while pipelines rotate underneath them", &rotations, &[Ok(200), Ok(200), Ok(200), Ok(200), Ok(200), Ok(200)]),136      Test.that("rotating pipelines were created, expired, and closed", Shared.current(&changes) >= 8),137      Test.that("listeners heard the timeout and the cancellation", Shared.current(&failures).filter(|line: Str| line.contains("impatient: The request timed out")).length() == 1 && Shared.current(&failures).filter(|line: Str| line.contains("cancelled")).length() == 1),138      Test.that("metrics counted successes, server errors, and failures", snapshot.successful >= 10 && snapshot.serverErrors >= 2 && snapshot.failed >= 3 && snapshot.active == 0),139      Test.that("logs list redacted credentials and never the token", messages.filter(|line: Str| line.contains("authorization: *")).length() > 0 && messages.filter(|line: Str| line.contains("secret-token")).isEmpty()),140      Test.that("logs name each client's layers", messages.filter(|line: Str| line.startsWith("Start processing HTTP request")).length() > 0)141    ])142  let ran = Test.run(&checks)143  for failure in Test.failuresOf(&ran) {144    let _reported = Io.writeErrorLine(failure)145  }146  Test.report(&ran)147}148