Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-mediator

In-process messaging for Pudu: requests, notifications, streams, pipeline behaviors, processors, and exception handling

0.1.0Apache-2.01

InstallClose

ComposeTest.pudu

Pudu245 lines11.9 KB

GitHub ↗
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  /// Records entering and leaving.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  /// Records the request.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  /// Records the response.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  /// Records the failure.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  /// Records the notification.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  /// Records entering and leaving, and tags every item.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}7778/// Appends a line to a trace.79fn 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}8687/// The trace so far, emptied.88fn taken(trace: &Sync.Cell[Array[Str]]) -> Array[Str] {89  match Sync.swap(trace, []) {90    case Ok(found) => found91    case Err(_) => []92  }93}9495/// A mediator built from registrations, or the reason it was refused.96fn 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}102103/// Runs the suite.104fn 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