Event-time watermark, lateness, and window assignment primitives for MoonBit.
let state = WatermarkState::new(default_policy())
.observe(StreamEvent::new("sensor-a", 120_000, 1, 10))
let assignment = assign_event(state, minute_window(), StreamEvent::new("sensor-a", 122_000, 2, 11))
println(assignment.status.name())moon fmt --check
moon info && git diff --exit-code -- '*.mbti'
moon check --target all
moon build --target all
moon test --target all
moon run cmd/mainmoon add Oyc996/moonwatermarkkitpub(all) struct EventAssignment {
event : StreamEvent
status : EventStatus
watermark_ms : Int
windows : Array[TimeWindow]
} derive(Eq, Debug)pub(all) struct PartitionWatermark {
partition : String
state : WatermarkState
idle : Bool
} derive(Eq, Debug)fn StreamEvent::new(key : String, event_time_ms : Int, arrival_time_ms : Int, value : Int) -> StreamEventfn TimeWindow::is_finalized(self : TimeWindow, watermark_ms : Int, allowed_lateness_ms : Int) -> Boolpub(all) struct WatermarkCoordinator {
policy : WatermarkPolicy
partitions : Array[PartitionWatermark]
} derive(Eq, Debug)fn WatermarkCoordinator::observe(self : WatermarkCoordinator, partition : String, event : StreamEvent) -> WatermarkCoordinatorfn WatermarkCoordinator::refresh_idle(self : WatermarkCoordinator, now_ms : Int) -> WatermarkCoordinatorfn WatermarkPolicy::new(max_out_of_order_ms : Int, allowed_lateness_ms : Int, idle_timeout_ms : Int) -> WatermarkPolicypub(all) struct WatermarkState {
policy : WatermarkPolicy
max_event_time_ms : Int
current_watermark_ms : Int
last_arrival_ms : Int
initialized : Bool
} derive(Eq, Debug)pub(all) struct WindowAggregate {
window : TimeWindow
key : String
count : Int
sum : Int
late : Int
dropped : Int
} derive(Eq, Debug)Event-time watermark, lateness, and window assignment primitives for MoonBit.