///|
/// Policy used when a handler asks for retry or returns failure.
pub(all) struct RetryPolicy {
max_attempts : Int
base_delay_ms : Int
backoff_factor : Int
} derive(Eq, @debug.Debug)
///|
pub fn retry_policy(
max_attempts? : Int = 3,
base_delay_ms? : Int = 100,
backoff_factor? : Int = 2,
) -> RetryPolicy {
{ max_attempts, base_delay_ms, backoff_factor }
}
///|
pub fn RetryPolicy::next_delay(self : RetryPolicy, attempt : Int) -> Int {
guard attempt > 0 else { return self.base_delay_ms }
let mut delay = self.base_delay_ms
for _ in 0.. Result[Subscription, EventRailError] {
let owned_id = id.to_owned()
guard owned_id != "" else { return Err(InvalidSubscription(owned_id)) }
match topic_pattern(pattern) {
Err(err) => Err(err)
Ok(parsed) =>
Ok({
id: owned_id,
pattern: parsed,
priority,
group: group.to_owned(),
order,
retry,
enabled,
guards: guards.to_owned(),
})
}
}
///|
pub fn Subscription::matches(
self : Subscription,
event : Envelope,
) -> Result[Bool, EventRailError] {
if !self.enabled {
return Ok(false)
}
match self.pattern.matches_topic(event.topic) {
Err(err) => Err(err)
Ok(false) => Ok(false)
Ok(true) => Ok(self.guards_pass(event))
}
}
///|
pub fn Subscription::guards_pass(self : Subscription, event : Envelope) -> Bool {
for predicate in self.guards {
if !predicate.evaluate(event).passed {
return false
}
}
true
}
///|
pub fn Subscription::first_guard_failure(
self : Subscription,
event : Envelope,
) -> RuleResult? {
for predicate in self.guards {
let result = predicate.evaluate(event)
if !result.passed {
return Some(result)
}
}
None
}
///|
pub(all) enum HandlerResult {
HandlerAck(String)
HandlerDrop(String)
HandlerRetry(String)
HandlerFail(String)
} derive(Eq, @debug.Debug)
///|
pub(all) enum DeliveryStatus {
Delivered(String)
Dropped(String)
RetryScheduled(String)
DeadLettered(String)
} derive(Eq, @debug.Debug)
///|
pub(all) struct Delivery {
event_id : String
subscription_id : String
topic : String
status : DeliveryStatus
attempt : Int
delay_ms : Int
} derive(Eq, @debug.Debug)
///|
pub(all) struct DeadLetter {
event_id : String
subscription_id : String
topic : String
reason : String
attempt : Int
trace : Array[String]
} derive(Eq, @debug.Debug)
///|
pub(all) struct PublishReport {
event_id : String
topic : String
routed : Int
deliveries : Array[Delivery]
dead_letters_added : Int
} derive(Eq, @debug.Debug)
///|
pub(all) enum DeliveryMode {
Fanout
FirstPerGroup
} derive(Eq, @debug.Debug)
///|
pub(all) struct RouteProbe {
subscription_id : String
pattern : String
group : String
matched : Bool
selected : Bool
reason : String
priority : Int
specificity : Int
} derive(Eq, @debug.Debug)
///|
pub(all) struct RouteExplanation {
event_id : String
topic : String
mode : DeliveryMode
selected : Array[String]
probes : Array[RouteProbe]
} derive(Eq, @debug.Debug)
///|
pub(all) struct BusAudit {
subscriptions : Int
enabled : Int
disabled : Int
groups : Int
guarded : Int
issues : Array[String]
} derive(Eq, @debug.Debug)
///|
pub(all) struct Bus {
subscriptions : Array[Subscription]
deliveries : Array[Delivery]
dead_letters : Array[DeadLetter]
} derive(Eq, @debug.Debug)
///|
pub fn bus() -> Bus {
{ subscriptions: [], deliveries: [], dead_letters: [] }
}
///|
pub fn Bus::size(self : Bus) -> Int {
self.subscriptions.length()
}
///|
pub fn Bus::subscribe(
self : Bus,
id : StringView,
pattern : StringView,
priority? : Int = 0,
group? : StringView = "",
retry? : RetryPolicy = retry_policy(),
enabled? : Bool = true,
guards? : ArrayView[EventPredicate] = [],
) -> Result[Bus, EventRailError] {
let owned_id = id.to_owned()
guard !self.has_subscription(owned_id) else {
return Err(DuplicateSubscription(owned_id))
}
match
subscription(
owned_id,
pattern,
priority~,
group~,
order=self.subscriptions.length(),
retry~,
enabled~,
guards~,
) {
Err(err) => Err(err)
Ok(sub) => {
let subscriptions = self.subscriptions.copy()
subscriptions.push(sub)
Ok({ ..self, subscriptions, })
}
}
}
///|
pub fn Bus::has_subscription(self : Bus, id : StringView) -> Bool {
let wanted = id.to_owned()
for sub in self.subscriptions {
if sub.id == wanted {
return true
}
}
false
}
///|
pub fn Bus::route(
self : Bus,
event : Envelope,
) -> Result[Array[Subscription], EventRailError] {
self.route_with_mode(event, Fanout)
}
///|
pub fn Bus::route_with_mode(
self : Bus,
event : Envelope,
mode : DeliveryMode,
) -> Result[Array[Subscription], EventRailError] {
match topic_segments(event.topic) {
Err(err) => Err(err)
Ok(_) => {
let routed : Array[Subscription] = []
for sub in self.subscriptions {
match sub.matches(event) {
Err(err) => return Err(err)
Ok(true) => routed.push(sub)
Ok(false) => ()
}
}
routed.sort_by(compare_subscription_route)
Ok(apply_delivery_mode(routed, mode))
}
}
}
///|
pub fn Bus::publish(
self : Bus,
event : Envelope,
handler : (Subscription, Envelope) -> HandlerResult,
) -> Result[(Bus, PublishReport), EventRailError] {
self.publish_with_mode(event, Fanout, handler)
}
///|
pub fn Bus::publish_with_mode(
self : Bus,
event : Envelope,
mode : DeliveryMode,
handler : (Subscription, Envelope) -> HandlerResult,
) -> Result[(Bus, PublishReport), EventRailError] {
match self.route_with_mode(event, mode) {
Err(err) => Err(err)
Ok(routed) => {
let deliveries = self.deliveries.copy()
let dead_letters = self.dead_letters.copy()
let report_deliveries : Array[Delivery] = []
let mut added_dead_letters = 0
for sub in routed {
let result = handler(sub, event)
let delivery = delivery_from_result(sub, event, result)
if delivery.status is DeadLettered(reason) {
dead_letters.push({
event_id: event.id,
subscription_id: sub.id,
topic: event.topic,
reason,
attempt: event.attempt,
trace: event.trace,
})
added_dead_letters += 1
}
deliveries.push(delivery)
report_deliveries.push(delivery)
}
let next = { ..self, deliveries, dead_letters }
let report = {
event_id: event.id,
topic: event.topic,
routed: routed.length(),
deliveries: report_deliveries,
dead_letters_added: added_dead_letters,
}
Ok((next, report))
}
}
}
///|
pub fn Bus::explain(
self : Bus,
event : Envelope,
mode? : DeliveryMode = Fanout,
) -> Result[RouteExplanation, EventRailError] {
match topic_segments(event.topic) {
Err(err) => Err(err)
Ok(_) => {
let selected = match self.route_with_mode(event, mode) {
Err(err) => return Err(err)
Ok(routes) => routes.map(sub => sub.id)
}
let probes : Array[RouteProbe] = []
for sub in self.subscriptions {
let probe = explain_subscription(sub, event, selected)
probes.push(probe)
}
Ok({ event_id: event.id, topic: event.topic, mode, selected, probes })
}
}
}
///|
pub fn Bus::audit(self : Bus) -> BusAudit {
let issues : Array[String] = []
let groups : Array[String] = []
let mut enabled = 0
let mut disabled = 0
let mut guarded = 0
for sub in self.subscriptions {
if sub.enabled {
enabled += 1
} else {
disabled += 1
}
if sub.guards.length() > 0 {
guarded += 1
}
if sub.group != "" && !groups.contains(sub.group) {
groups.push(sub.group)
}
}
if self.subscriptions.length() == 0 {
issues.push("bus has no subscriptions")
}
if enabled == 0 && self.subscriptions.length() > 0 {
issues.push("bus has no enabled subscriptions")
}
{
subscriptions: self.subscriptions.length(),
enabled,
disabled,
groups: groups.length(),
guarded,
issues,
}
}
///|
fn delivery_from_result(
sub : Subscription,
event : Envelope,
result : HandlerResult,
) -> Delivery {
match result {
HandlerAck(message) =>
{
event_id: event.id,
subscription_id: sub.id,
topic: event.topic,
status: Delivered(message),
attempt: event.attempt,
delay_ms: 0,
}
HandlerDrop(reason) =>
{
event_id: event.id,
subscription_id: sub.id,
topic: event.topic,
status: Dropped(reason),
attempt: event.attempt,
delay_ms: 0,
}
HandlerRetry(reason) | HandlerFail(reason) =>
if event.attempt + 1 >= sub.retry.max_attempts {
{
event_id: event.id,
subscription_id: sub.id,
topic: event.topic,
status: DeadLettered(reason),
attempt: event.attempt,
delay_ms: 0,
}
} else {
{
event_id: event.id,
subscription_id: sub.id,
topic: event.topic,
status: RetryScheduled(reason),
attempt: event.attempt,
delay_ms: sub.retry.next_delay(event.attempt),
}
}
}
}
///|
fn compare_subscription_route(left : Subscription, right : Subscription) -> Int {
if left.priority != right.priority {
right.priority - left.priority
} else if left.pattern.specificity != right.pattern.specificity {
right.pattern.specificity - left.pattern.specificity
} else {
left.order - right.order
}
}
///|
fn apply_delivery_mode(
routed : Array[Subscription],
mode : DeliveryMode,
) -> Array[Subscription] {
match mode {
Fanout => routed
FirstPerGroup => {
let selected : Array[Subscription] = []
let seen : Array[String] = []
for sub in routed {
let key = if sub.group == "" {
"subscription:\{sub.id}"
} else {
"group:\{sub.group}"
}
if !seen.contains(key) {
seen.push(key)
selected.push(sub)
}
}
selected
}
}
}
///|
fn explain_subscription(
sub : Subscription,
event : Envelope,
selected : Array[String],
) -> RouteProbe {
let (matched, reason) = if !sub.enabled {
(false, "subscription disabled")
} else {
match sub.pattern.matches_topic(event.topic) {
Err(err) => (false, err.message())
Ok(false) => (false, "topic did not match")
Ok(true) =>
match sub.first_guard_failure(event) {
Some(result) => (false, "guard failed: \{result.reason}")
None => (true, "matched")
}
}
}
{
subscription_id: sub.id,
pattern: sub.pattern.raw,
group: sub.group,
matched,
selected: selected.contains(sub.id),
reason,
priority: sub.priority,
specificity: sub.pattern.specificity,
}
}
///|
pub fn PublishReport::status_line(self : PublishReport) -> String {
"event=\{self.event_id} topic=\{self.topic} routed=\{self.routed} deliveries=\{self.deliveries.length()} dead_letters=\{self.dead_letters_added}"
}
///|
pub fn DeliveryStatus::to_wire(self : DeliveryStatus) -> String {
match self {
Delivered(message) => "delivered:\{message}"
Dropped(reason) => "dropped:\{reason}"
RetryScheduled(reason) => "retry:\{reason}"
DeadLettered(reason) => "dead:\{reason}"
}
}
///|
pub fn Bus::dead_letter_count(self : Bus) -> Int {
self.dead_letters.length()
}
///|
pub fn Bus::delivery_count(self : Bus) -> Int {
self.deliveries.length()
}