Skip to content

Notifications

Chris Michael edited this page Oct 6, 2026 · 1 revision

Pudu

Notifications

A notification kind may have any number of handlers, each registered with a name, and open notification handlers (Open.NotificationHandler) hear every kind. A publication hands every handler, in registration order, to the mediator's publishing strategy:

Strategy Runs Answers
Publishing.Sequential{} (default) one at a time, checking cancellation first the first failure; later handlers never run
Publishing.Continuing{} one at a time, every handler every failure, combined
Publishing.Parallel{} every handler on its own thread every failure, combined, once all finished
Publishing.Bounded{workers: n} every handler, at most n at once every failure, combined, once all finished

A strategy is a trait, so you can write your own:

module Ex05Notifications

import Std.Io as Io
import PuduLangMediator.Context as Context
import PuduLangMediator.Mediator as Mediator
import PuduLangMediator.Message as Message
import PuduLangMediator as Messaging
import PuduLangMediator.Notification as Notification
import PuduLangMediator.Open as Open
import PuduLangMediator.Publishing as Publishing
import PuduLangMediator.Registration as Registration

type Joined = { name: Str }

type Audit = {}

impl Open.NotificationHandler for Audit {
  fn handleNotification[N, E](self: &Self, notification: &N, info: &Message.Info, context: Context.Context) -> Messaging.Outcome[(), E] {
    let _said = Io.writeLine("audit: " + info.name + " " + show(*notification))
    Ok(())
  }
}

type Reversed = {}

impl Publishing.Strategy for Reversed {
  fn publish[E](self: &Self, executors: &Array[Publishing.Executor[E]], context: Context.Context) -> Messaging.Outcome[(), E] {
    for executor in executors.reverse() {
      let _said = Io.writeLine("running " + executor.handler)
      (executor.run)(context) ?
    }
    Ok(())
  }
}

fn registrations(joined: &Notification.Kind[Joined, Str]) -> Array[Registration.Registration] {
  [
    Registration.notificationHandler(joined, "welcome", fn(event: Joined, _context: Context.Context) -> Messaging.Outcome[(), Str] {
        let _said = Io.writeLine("welcome, " + event.name)
        Ok(())
      }),
    Registration.openNotificationHandler("audit", Audit{}),
    Registration.notificationHandler(joined, "crm", |event: Joined, _context: Context.Context| Messaging.raise("crm rejected " + event.name))
  ]
}

export fn main() -> Int {
  let joined: Notification.Kind[Joined, Str] = Notification.kind("users.joined")
  var outcomes: Array[Messaging.Outcome[(), Str]] = []
  let strategies: Array[(Str, dynamic Publishing.Strategy)] = [("sequential", Publishing.Sequential{}), ("continuing", Publishing.Continuing{}), ("reversed", Reversed{})]
  for (label, strategy) in strategies {
    let _title = Io.writeLine("-- " + label + " --")
    let mediator = match Mediator.buildWith(&Mediator.Options{..Mediator.defaults(), publisher: strategy}, registrations(&joined)) {
      case Ok(built) => built
      case Err(invalid) => panic(Mediator.explain(&invalid))
    }
    let outcome = Mediator.publish(&mediator, &joined, Joined{name: "ada"})
    let _said = Io.writeLine(Messaging.summarize(&outcome))
    outcomes = outcomes.push(outcome)
  }
  if outcomes.filter(|outcome: Messaging.Outcome[(), Str]| outcome == Err(Messaging.Raised("crm rejected ada"))).length() == 3 { 0 } else { 1 }
}

Check it and run it:

pudu check src/Ex05Notifications.pudu
pudu run src/Ex05Notifications.pudu

Output:

-- sequential --
welcome, ada
audit: users.joined Joined{name: "ada"}
The handler failed: "crm rejected ada"
-- continuing --
welcome, ada
audit: users.joined Joined{name: "ada"}
The handler failed: "crm rejected ada"
-- reversed --
running crm
The handler failed: "crm rejected ada"

Several failures combine into one Aggregate; one failure is reported as itself. Parallel and bounded publication contain each handler, so a handler whose thread stops is reported as Crashed and the others still run. Publishing a kind with no handlers succeeds, and still reaches the open handlers.

Related

Clone this wiki locally