
PlaceholderResilienceTest.pudu
Pudu148 lines7.9 KB
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 Placeholder272829fn statusOf(outcome: HttpClient.Outcome[Response.Response]) -> Result[Int, HttpClient.Failure] {30 Result.map(outcome, |response: Response.Response| response.status.code)31}323334fn isRejected(outcome: &Result[Int, HttpClient.Failure]) -> Bool {35 match outcome {36 case Err(HttpClient.Rejected(_)) => true37 case _ => false38 }39}404142fn 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