
ComposeTest.pudu
Pudu245 lines11.9 KB
1/** @Test.Compose.Suite — the documented order of every pipeline layer */2module PuduLangMediator.ComposeTest34import Std.Concurrent.Cancel as Cancel5import Std.Io as Io6import Std.Sync as Sync7import Std.Test as Test8import PuduLangMediator.Catalog as Catalog9import PuduLangMediator.Context as Context10import PuduLangMediator.Mediator as Mediator11import PuduLangMediator.Message as Message12import PuduLangMediator as Messaging13import PuduLangMediator.Notification as Notification14import PuduLangMediator.Open as Open15import PuduLangMediator.Registration as Registration16import PuduLangMediator.Request as Request17import PuduLangMediator.Stream as Stream1819/** @Test.Compose.Tracer — an open component of every layer that records what it saw */20type Tracer = { trace: Sync.Cell[Array[Str]], label: Str }2122impl Open.Behavior for Tracer {23 24 fn handleRequest[Q, R, E](self: &Self, request: &Q, info: &Message.Info, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] {25 record(&self.trace, self.label + ">")26 let outcome = next(context)27 record(&self.trace, "<" + self.label)28 outcome29 }30}3132impl Open.PreProcessor for Tracer {33 34 fn beforeRequest[Q, E](self: &Self, request: &Q, info: &Message.Info, context: Context.Context) -> Messaging.Outcome[(), E] {35 record(&self.trace, "pre " + self.label)36 Ok(())37 }38}3940impl Open.PostProcessor for Tracer {41 42 fn afterRequest[Q, R, E](self: &Self, request: &Q, response: &R, info: &Message.Info, context: Context.Context) -> Messaging.Outcome[(), E] {43 record(&self.trace, "post " + self.label + " " + show(*response))44 Ok(())45 }46}4748impl Open.ExceptionAction for Tracer {49 50 fn onFailure[Q, E](self: &Self, request: &Q, failure: &Messaging.Failure[E], info: &Message.Info, context: Context.Context) -> () {51 record(&self.trace, "action " + self.label + " " + show(*failure))52 }53}5455impl Open.NotificationHandler for Tracer {56 57 fn handleNotification[N, E](self: &Self, notification: &N, info: &Message.Info, context: Context.Context) -> Messaging.Outcome[(), E] {58 record(&self.trace, "hear " + self.label + " " + show(*notification))59 Ok(())60 }61}6263impl Open.StreamBehavior for Tracer {64 65 fn handleStream[Q, T, E](self: &Self, request: &Q, info: &Message.Info, context: Context.Context, sink: Stream.Sink[T], next: Stream.Next[T, E]) -> Messaging.Outcome[(), E] {66 record(&self.trace, self.label + ">")67 let trace = self.trace68 let label = self.label69 let outcome = next(context, fn(item: T) -> Bool {70 record(&trace, label + " " + show(item))71 sink(item)72 })73 record(&self.trace, "<" + self.label)74 outcome75 }76}777879fn record(trace: &Sync.Cell[Array[Str]], line: Str) -> () {80 let held = match Sync.get(trace) {81 case Ok(found) => found82 case Err(_) => []83 }84 let _stored = Sync.set(trace, held.push(line))85}868788fn taken(trace: &Sync.Cell[Array[Str]]) -> Array[Str] {89 match Sync.swap(trace, []) {90 case Ok(found) => found91 case Err(_) => []92 }93}949596fn built(options: Mediator.Options, registrations: Array[Registration.Registration]) -> Mediator.Mediator {97 match Mediator.buildWith(&options, registrations) {98 case Ok(mediator) => mediator99 case Err(invalid) => panic(Mediator.explain(&invalid))100 }101}102103104fn main() -> Int {105 let trace: Sync.Cell[Array[Str]] = Sync.cell([])106 let double: Request.Kind[Int, Int, Str] = Request.kind("double")107 let open = fn(label: Str) -> Tracer { Tracer{trace: trace, label: label} }108 let closed = fn(label: Str) -> Request.Behavior[Int, Int, Str] {109 fn(request: Int, context: Context.Context, next: Request.Next[Int, Str]) -> Messaging.Outcome[Int, Str] {110 record(&trace, label + ">")111 let outcome = next(context)112 record(&trace, "<" + label)113 outcome114 }115 }116 let layered = [117 Registration.openBehavior(open("o1")),118 Registration.behavior(&double, closed("c1")),119 Registration.openPreProcessor(open("o")),120 Registration.preProcessor(&double, fn(request: Int, _context: Context.Context) -> Messaging.Outcome[(), Str] {121 record(&trace, "pre c " + show(request))122 if request < 0 { Messaging.raise("negative") } else { Ok(()) }123 }),124 Registration.postProcessor(&double, fn(_request: Int, response: Int, _context: Context.Context) -> Messaging.Outcome[(), Str] {125 record(&trace, "post c " + show(response))126 if response > 100 { Messaging.raise("too large") } else { Ok(()) }127 }),128 Registration.openPostProcessor(open("o")),129 Registration.openBehavior(open("o2")),130 Registration.behavior(&double, closed("c2")),131 Registration.openExceptionAction(open("o")),132 Registration.exceptionAction(&double, fn(_request: Int, failure: Messaging.Failure[Str], _context: Context.Context) -> () { record(&trace, "action c " + show(failure)) }),133 Registration.exceptionHandler(&double, fn(request: Int, _failure: Messaging.Failure[Str], _context: Context.Context) -> Option[Int] {134 record(&trace, "recover first")135 if request == 7 { Some(0) } else { None }136 }),137 Registration.exceptionHandler(&double, fn(request: Int, _failure: Messaging.Failure[Str], _context: Context.Context) -> Option[Int] {138 record(&trace, "recover second")139 if request == 9 { Some(99) } else { None }140 }),141 Registration.handler(&double, fn(request: Int, _context: Context.Context) -> Messaging.Outcome[Int, Str] {142 record(&trace, "handle " + show(request))143 if request == 7 || request == 9 || request == 13 { Messaging.raise("unlucky") } else { Ok(request * 2) }144 })145 ]146 let unhandledOnly = built(Mediator.defaults(), layered)147 let everyFailure = built(Mediator.Options{..Mediator.defaults(), actionScope: Catalog.ForAll}, layered)148149 let fine = Mediator.send(&unhandledOnly, &double, 4)150 let fineTrace = taken(&trace)151 let refusedEarly = Mediator.send(&unhandledOnly, &double, -1)152 let earlyTrace = taken(&trace)153 let tooLarge = Mediator.send(&unhandledOnly, &double, 60)154 let largeTrace = taken(&trace)155 let recoveredFirst = Mediator.send(&unhandledOnly, &double, 7)156 let firstTrace = taken(&trace)157 let recoveredSecond = Mediator.send(&unhandledOnly, &double, 9)158 let secondTrace = taken(&trace)159 let unrecovered = Mediator.send(&unhandledOnly, &double, 13)160 let unrecoveredTrace = taken(&trace)161 let recoveredSeen = Mediator.send(&everyFailure, &double, 7)162 let seenTrace = taken(&trace)163164 let gate: Request.Kind[Int, Int, Str] = Request.kind("gate")165 let token: Context.Key[Str] = Context.key("token")166 let gated = built(Mediator.defaults(), [167 Registration.behavior(&gate, fn(request: Int, context: Context.Context, next: Request.Next[Int, Str]) -> Messaging.Outcome[Int, Str] {168 if request == 0 { Messaging.refuse("closed") } else { next(Context.withToken(&context, Cancel.expiring(0))) }169 }),170 Registration.handler(&gate, fn(request: Int, context: Context.Context) -> Messaging.Outcome[Int, Str] {171 Context.check(&context) ?172 Ok(request + Context.getOr(&context, &token, "").length())173 })174 ])175 let shortCircuited = Mediator.send(&gated, &gate, 0)176 let passedContext = Mediator.send(&gated, &gate, 1)177178 let placed: Notification.Kind[Int, Str] = Notification.kind("placed")179 let published = built(Mediator.defaults(), [180 Registration.notificationHandler(&placed, "first", fn(notification: Int, _context: Context.Context) -> Messaging.Outcome[(), Str] {181 record(&trace, "first " + show(notification))182 Ok(())183 }),184 Registration.openNotificationHandler("open", open("o")),185 Registration.notificationHandler(&placed, "second", fn(notification: Int, _context: Context.Context) -> Messaging.Outcome[(), Str] {186 record(&trace, "second " + show(notification))187 Ok(())188 })189 ])190 let publication = Mediator.publish(&published, &placed, 3)191 let publicationTrace = taken(&trace)192 let unregistered: Notification.Kind[Str, Str] = Notification.kind("unregistered")193 let openOnly = Mediator.publish(&published, &unregistered, "x")194 let openOnlyTrace = taken(&trace)195196 let letters: Stream.Kind[Int, Str, Str] = Stream.kind("letters")197 let streamed = built(Mediator.defaults(), [198 Registration.openStreamBehavior(open("o")),199 Registration.streamBehavior(&letters, fn(_request: Int, context: Context.Context, sink: Stream.Sink[Str], next: Stream.Next[Str, Str]) -> Messaging.Outcome[(), Str] {200 next(context, fn(item: Str) -> Bool { if item == "b" { true } else { sink(item + "!") } })201 }),202 Registration.streamHandler(&letters, fn(request: Int, _context: Context.Context, sink: Stream.Sink[Str]) -> Messaging.Outcome[(), Str] {203 Stream.each(&["a", "b", "c", "d"].slice(0, request), sink)204 })205 ])206 let letterItems = Mediator.collect(&streamed, &letters, 3)207 let letterTrace = taken(&trace)208209 let checks = Test.suite("Compose", &[210 Test.equals("a request passes every layer in order", &fine, &Ok(8)),211 Test.equals("pre-processors before, behaviors in registration order, post-processors after", &fineTrace, &[212 "pre o", "pre c 4", "o1>", "c1>", "o2>", "c2>", "handle 4", "<c2", "<o2", "<c1", "<o1", "post c 8", "post o 8"213 ]),214 Test.equals("a failing pre-processor stops the request", &refusedEarly, &Err(Messaging.Raised("negative"))),215 Test.equals("the rest of the pre-processors, the behaviors, and the handler never run", &earlyTrace, &[216 "pre o", "pre c -1", "recover first", "recover second", "action o Raised(\"negative\")", "action c Raised(\"negative\")"217 ]),218 Test.equals("a failing post-processor replaces the response", &tooLarge, &Err(Messaging.Raised("too large"))),219 Test.equals("later post-processors do not run after one fails", &largeTrace.contains("post o 120"), &false),220 Test.equals("the first exception handler that answers wins", &recoveredFirst, &Ok(0)),221 Test.equals("a recovered failure is never seen by exception actions", &firstTrace.filter(|line: Str| line.startsWith("action") || line.startsWith("recover")), &["recover first"]),222 Test.equals("exception handlers are tried in registration order", &(recoveredSecond, secondTrace.filter(|line: Str| line.startsWith("recover"))), &(223 Ok(99), ["recover first", "recover second"]224 )),225 Test.equals("an unrecovered failure reaches every exception action, open first", &(unrecovered, unrecoveredTrace.filter(|line: Str| line.startsWith("action"))), &(226 Err(Messaging.Raised("unlucky")), ["action o Raised(\"unlucky\")", "action c Raised(\"unlucky\")"]227 )),228 Test.equals("with every failure in scope, actions see a failure that is then recovered", &(recoveredSeen, seenTrace.filter(|line: Str| line.startsWith("action"))), &(229 Ok(0), ["action o Raised(\"unlucky\")", "action c Raised(\"unlucky\")"]230 )),231 Test.equals("a behavior may answer without running the handler", &shortCircuited, &Err(Messaging.Refused("closed"))),232 Test.equals("a behavior may hand the rest of the pipeline another context", &passedContext, &Err(Messaging.Cancelled("the deadline passed"))),233 Test.equals("a publication reaches every handler", &publication, &Ok(())),234 Test.equals("open and own notification handlers run in registration order", &publicationTrace, &["first 3", "hear o 3", "second 3"]),235 Test.equals("a kind without handlers still reaches the open ones", &(openOnly, openOnlyTrace), &(Ok(()), ["hear o \"x\""])),236 Test.equals("stream behaviors wrap the sink in registration order", &letterItems, &Ok(["a!", "c!"])),237 Test.equals("an open stream behavior sees what the inner ones pass on", &letterTrace, &["o>", "o \"a!\"", "o \"c!\"", "<o"])238 ])239 let ran = Test.run(&checks)240 for failure in Test.failuresOf(&ran) {241 let _reported = Io.writeErrorLine(failure)242 }243 Test.report(&ran)244}245