///|
pub(all) struct ScheduledEvent {
due_ms : Int
subscription_id : String
reason : String
event : Envelope
} derive(Eq, @debug.Debug)
///|
pub(all) struct Outbox {
scheduled : Array[ScheduledEvent]
} derive(Eq, @debug.Debug)
///|
pub(all) struct OutboxDrain {
due : Array[ScheduledEvent]
remaining : Outbox
} derive(Eq, @debug.Debug)
///|
pub fn outbox() -> Outbox {
{ scheduled: [] }
}
///|
pub fn scheduled_event(
due_ms : Int,
subscription_id : StringView,
reason : StringView,
event : Envelope,
) -> ScheduledEvent {
{
due_ms,
subscription_id: subscription_id.to_owned(),
reason: reason.to_owned(),
event,
}
}
///|
pub fn Outbox::len(self : Outbox) -> Int {
self.scheduled.length()
}
///|
pub fn Outbox::is_empty(self : Outbox) -> Bool {
self.scheduled.length() == 0
}
///|
pub fn Outbox::push(self : Outbox, item : ScheduledEvent) -> Outbox {
let scheduled = self.scheduled.copy()
scheduled.push(item)
scheduled.sort_by(compare_scheduled_event)
{ scheduled, }
}
///|
pub fn Outbox::from_report(
self : Outbox,
event : Envelope,
report : PublishReport,
now_ms : Int,
) -> Outbox {
let mut box = self
for delivery in report.deliveries {
match delivery.status {
RetryScheduled(reason) =>
box = box.push(
scheduled_event(
now_ms + delivery.delay_ms,
delivery.subscription_id,
reason,
event.next_attempt().add_trace("retry:\{delivery.subscription_id}"),
),
)
Delivered(_) | Dropped(_) | DeadLettered(_) => ()
}
}
box
}
///|
pub fn Outbox::due(self : Outbox, now_ms : Int) -> Array[ScheduledEvent] {
let items : Array[ScheduledEvent] = []
for item in self.scheduled {
if item.due_ms <= now_ms {
items.push(item)
}
}
items
}
///|
pub fn Outbox::drain_due(self : Outbox, now_ms : Int) -> OutboxDrain {
let due : Array[ScheduledEvent] = []
let future : Array[ScheduledEvent] = []
for item in self.scheduled {
if item.due_ms <= now_ms {
due.push(item)
} else {
future.push(item)
}
}
{ due, remaining: { scheduled: future } }
}
///|
pub fn Outbox::drop_subscription(
self : Outbox,
subscription_id : StringView,
) -> Outbox {
let wanted = subscription_id.to_owned()
let scheduled : Array[ScheduledEvent] = []
for item in self.scheduled {
if item.subscription_id != wanted {
scheduled.push(item)
}
}
{ scheduled, }
}
///|
pub fn Outbox::reschedule_all(self : Outbox, delta_ms : Int) -> Outbox {
let scheduled : Array[ScheduledEvent] = []
for item in self.scheduled {
scheduled.push({ ..item, due_ms: item.due_ms + delta_ms })
}
scheduled.sort_by(compare_scheduled_event)
{ scheduled, }
}
///|
pub fn Outbox::manifest(self : Outbox) -> String {
self.scheduled.map(item => item.to_manifest_line()).join("\n")
}
///|
pub fn ScheduledEvent::to_manifest_line(self : ScheduledEvent) -> String {
"due=\{self.due_ms};sub=\{self.subscription_id};reason=\{escape_wire_text(self.reason)};event=\{self.event.id}"
}
///|
pub fn OutboxDrain::summary(self : OutboxDrain) -> String {
"due=\{self.due.length()} remaining=\{self.remaining.len()}"
}
///|
fn compare_scheduled_event(
left : ScheduledEvent,
right : ScheduledEvent,
) -> Int {
if left.due_ms != right.due_ms {
left.due_ms - right.due_ms
} else if left.subscription_id != right.subscription_id {
left.subscription_id.compare(right.subscription_id)
} else {
left.event.id.compare(right.event.id)
}
}