///|
/// Scenario package for M2-SCENARIO-EVENT.
///|
pub fn package_id() -> String {
"lockwire/scenario"
}
///|
pub(all) enum EventState {
Pending
Scheduled(@core.VTime)
Processed
Cancelled(CancelReason)
Failed(String)
} derive(Eq, Debug)
///|
pub(all) enum CancelReason {
UserCancelled(String)
} derive(Eq, Debug)
///|
pub(all) enum StepResult {
StepProcessed(event_id~ : Int, time~ : @core.VTime)
Idle
} derive(Eq, Debug)
///|
pub(all) struct AnyResult[T] {
winner_id : Int
winner_index : Int
value : T
pending_ids : Array[Int]
} derive(Eq, Debug)
///|
struct EventCell[T] {
state : Ref[EventState]
value : Ref[T?]
callbacks : Array[(Event[T]) -> Unit]
}
///|
pub struct Event[T] {
priv id : Int
priv cause_id : Int
priv cell : Ref[EventCell[T]]
}
///|
pub type Timeout = Event[Unit]
///|
struct ScheduledAction {
id : Int
cause_id : Int
time : @core.VTime
run : () -> Unit
}
///|
pub struct SimEnv {
priv seed : Int
priv now_ref : Ref[@core.VTime]
priv next_event_id : Ref[Int]
priv next_cause_id : Ref[Int]
priv scheduled : Array[ScheduledAction]
priv processed_ids : Array[Int]
}
///|
pub struct ProcessHandle[T] {
priv name : String
priv event : Event[T]
}
///|
pub fn SimEnv::new(seed~ : Int) -> SimEnv {
{
seed,
now_ref: Ref(@core.VTime::from_ns(0L)),
next_event_id: Ref(1),
next_cause_id: Ref(1),
scheduled: [],
processed_ids: [],
}
}
///|
pub fn SimEnv::now(self : SimEnv) -> @core.VTime {
self.now_ref.val
}
///|
fn[T] SimEnv::alloc_event(self : SimEnv) -> Event[T] {
let id = self.next_event_id.val
self.next_event_id.val = id + 1
let cause_id = self.next_cause_id.val
self.next_cause_id.val = cause_id + 1
{
id,
cause_id,
cell: Ref({ state: Ref(Pending), value: Ref(None), callbacks: [] }),
}
}
///|
fn SimEnv::schedule_action(
self : SimEnv,
id~ : Int,
cause_id~ : Int,
time~ : @core.VTime,
run~ : () -> Unit,
) -> Unit {
self.scheduled.push({ id, cause_id, time, run })
}
///|
pub fn[T] event(env : SimEnv) -> Event[T] {
env.alloc_event()
}
///|
pub fn timeout(env : SimEnv, d : @core.Duration) -> Event[Unit] {
let ev = env.alloc_event()
let at = env.now().add(d)
ev.cell.val.state.val = Scheduled(at)
env.schedule_action(id=ev.id, cause_id=ev.cause_id, time=at, run=fn() {
if ev.is_pending() {
succeed(ev, ())
}
})
ev
}
///|
pub fn[T] Event::id(self : Event[T]) -> Int {
self.id
}
///|
pub fn[T] Event::cause_id(self : Event[T]) -> Int {
self.cause_id
}
///|
pub fn[T] Event::state(self : Event[T]) -> EventState {
self.cell.val.state.val
}
///|
pub fn[T] Event::value(self : Event[T]) -> T? {
self.cell.val.value.val
}
///|
pub fn[T] Event::is_pending(self : Event[T]) -> Bool {
match self.state() {
Pending | Scheduled(_) => true
Processed | Cancelled(_) | Failed(_) => false
}
}
///|
fn[T] Event::on_processed(
self : Event[T],
callback : (Event[T]) -> Unit,
) -> Unit {
if self.state() is Processed {
callback(self)
} else {
self.cell.val.callbacks.push(callback)
}
}
///|
fn[T] Event::notify_processed(self : Event[T]) -> Unit {
let callbacks = self.cell.val.callbacks.copy()
for callback in callbacks {
callback(self)
}
}
///|
pub fn[T] succeed(ev : Event[T], value : T) -> Unit {
if ev.is_pending() {
ev.cell.val.value.val = Some(value)
ev.cell.val.state.val = Processed
ev.notify_processed()
}
}
///|
pub fn[T] fail(ev : Event[T], message : String) -> Unit {
if ev.is_pending() {
ev.cell.val.state.val = Failed(message)
}
}
///|
pub fn[T] cancel(ev : Event[T], reason : CancelReason) -> Unit {
if ev.is_pending() {
ev.cell.val.state.val = Cancelled(reason)
}
}
///|
fn ScheduledAction::precedes(
self : ScheduledAction,
other : ScheduledAction,
) -> Bool {
if self.time < other.time {
true
} else if self.time == other.time {
self.id < other.id
} else {
false
}
}
///|
fn SimEnv::next_action_index(self : SimEnv) -> Int? {
if self.scheduled.is_empty() {
None
} else {
let mut best = 0
for i in 1.. @core.VTime? {
match self.next_action_index() {
Some(index) => Some(self.scheduled[index].time)
None => None
}
}
///|
pub fn SimEnv::step(self : SimEnv) -> StepResult {
match self.next_action_index() {
None => Idle
Some(index) => {
let action = self.scheduled.remove(index)
self.now_ref.val = action.time
(action.run)()
self.processed_ids.push(action.id)
StepProcessed(event_id=action.id, time=action.time)
}
}
}
///|
pub fn SimEnv::run(self : SimEnv) -> Unit {
while !self.scheduled.is_empty() {
ignore(self.step())
}
}
///|
pub fn SimEnv::run_until_time(self : SimEnv, t : @core.VTime) -> Unit {
while self.peek() is Some(next) && next <= t {
ignore(self.step())
}
self.now_ref.val = t
}
///|
pub fn[T] SimEnv::run_until_event(self : SimEnv, ev : Event[T]) -> T {
while ev.is_pending() && !self.scheduled.is_empty() {
ignore(self.step())
}
ev.value().unwrap()
}
///|
pub fn SimEnv::run_until_pred(self : SimEnv, pred : () -> Bool) -> Unit {
while !pred() && !self.scheduled.is_empty() {
ignore(self.step())
}
}
///|
pub fn SimEnv::checkpoint(self : SimEnv) -> @core.SimDigest {
let mut digest = @core.SimDigest::empty(seed=self.seed)
.mix(self.now().ns().to_int())
.mix(self.scheduled.length())
.mix(self.processed_ids.length())
for action in self.scheduled {
digest = digest
.mix(action.id)
.mix(action.cause_id)
.mix(action.time.ns().to_int())
}
for id in self.processed_ids {
digest = digest.mix(id)
}
digest
}
///|
pub fn[T] any_of(env : SimEnv, events : Array[Event[T]]) -> Event[AnyResult[T]] {
let out = env.alloc_event()
for index, ev in events {
ev.on_processed(fn(done) {
if out.is_pending() {
let pending_ids : Array[Int] = []
for other in events {
if other.id() != done.id() && other.is_pending() {
pending_ids.push(other.id())
}
}
succeed(out, {
winner_id: done.id(),
winner_index: index,
value: done.value().unwrap(),
pending_ids,
})
}
})
}
out
}
///|
pub fn[T] all_of(env : SimEnv, events : Array[Event[T]]) -> Event[Array[T]] {
let out = env.alloc_event()
let values : Array[T?] = Array::make(events.length(), None)
for index, ev in events {
ev.on_processed(fn(done) {
values[index] = done.value()
if out.is_pending() && values.all(fn(value) { value is Some(_) }) {
let result : Array[T] = []
for value in values {
match value {
Some(v) => result.push(v)
None => ()
}
}
succeed(out, result)
}
})
}
out
}
///|
pub fn[T] spawn(
env : SimEnv,
name~ : String,
f : (SimEnv) -> T,
) -> ProcessHandle[T] {
let ev = env.alloc_event()
let at = env.now()
ev.cell.val.state.val = Scheduled(at)
env.schedule_action(id=ev.id, cause_id=ev.cause_id, time=at, run=fn() {
if ev.is_pending() {
succeed(ev, f(env))
}
})
{ name, event: ev }
}
///|
pub fn[T] ProcessHandle::event(self : ProcessHandle[T]) -> Event[T] {
self.event
}
///|
pub fn[T] ProcessHandle::name(self : ProcessHandle[T]) -> String {
self.name
}