moonwatermarkkit

Event-time watermark, lateness, and window assignment primitives for MoonBit.

watermark
event-time
window
lateness
streaming
moon add Oyc996/moonwatermarkkit@0.2.0
Download zip
Author
Version
0.2.0
License
Apache-2.0
Last updated
last month
Downloads
2
README

#MoonWatermarkKit

MoonWatermarkKit is a MoonBit foundation library for event-time watermarking, late-event classification, window assignment, and stream audit summaries.

It is not a message broker, storage engine, or distributed stream runtime. It provides the small decision core that stream processors, log analyzers, metrics pipelines, and IoT event systems can embed.

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())

#Verification

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/main

#Installation

moon add Oyc996/moonwatermarkkit

The package has no third-party runtime dependencies. WatermarkState is an immutable value, so a caller can keep one state per source or partition and persist it using the caller's own checkpoint mechanism.

#
EventAssignment

pub(all) struct EventAssignment {
event : StreamEvent
status : EventStatus
watermark_ms : Int
windows : Array[TimeWindow]
} derive(Eq,
Debug
)

#
EventAssignment::accepted

fn EventAssignment::accepted(self : EventAssignment) -> Bool

#
EventAssignment::to_json

fn EventAssignment::to_json(self : EventAssignment) -> String

#
EventStatus

pub(all) enum EventStatus {
OnTime
Late
TooLate
} derive(Eq,
Debug
)

#
EventStatus::name

fn EventStatus::name(self : EventStatus) -> String

#
PartitionWatermark

pub(all) struct PartitionWatermark {
partition : String
state : WatermarkState
idle : Bool
} derive(Eq,
Debug
)

#
StreamAudit

pub(all) struct StreamAudit {
total : Int
on_time : Int
late : Int
too_late : Int
accepted : Int
} derive(Eq,
Debug
)

#
StreamAudit::healthy

fn StreamAudit::healthy(self : StreamAudit) -> Bool

#
StreamAudit::to_json

fn StreamAudit::to_json(self : StreamAudit) -> String

#
StreamEvent

pub(all) struct StreamEvent {
key : String
event_time_ms : Int
arrival_time_ms : Int
value : Int
} derive(Eq,
Debug
)

#
StreamEvent::new

fn StreamEvent::new(key : String, event_time_ms : Int, arrival_time_ms : Int, value : Int) -> StreamEvent

#
TimeWindow

pub(all) struct TimeWindow {
start_ms : Int
end_ms : Int
} derive(Eq,
Debug
)

#
TimeWindow::contains

fn TimeWindow::contains(self : TimeWindow, event_time_ms : Int) -> Bool

#
TimeWindow::is_finalized

fn TimeWindow::is_finalized(self : TimeWindow, watermark_ms : Int, allowed_lateness_ms : Int) -> Bool

#
TimeWindow::new

fn TimeWindow::new(start_ms : Int, end_ms : Int) -> TimeWindow

#
TimeWindow::to_json

fn TimeWindow::to_json(self : TimeWindow) -> String

#
WatermarkCoordinator

pub(all) struct WatermarkCoordinator {
policy : WatermarkPolicy
partitions : Array[PartitionWatermark]
} derive(Eq,
Debug
)

#
WatermarkCoordinator::active_partition_count

fn WatermarkCoordinator::active_partition_count(self : WatermarkCoordinator) -> Int

#
WatermarkCoordinator::global_watermark

fn WatermarkCoordinator::global_watermark(self : WatermarkCoordinator) -> Int?

#
WatermarkCoordinator::new

#
WatermarkCoordinator::observe

fn WatermarkCoordinator::observe(self : WatermarkCoordinator, partition : String, event : StreamEvent) -> WatermarkCoordinator

#
WatermarkCoordinator::refresh_idle

fn WatermarkCoordinator::refresh_idle(self : WatermarkCoordinator, now_ms : Int) -> WatermarkCoordinator

#
WatermarkPolicy

pub(all) struct WatermarkPolicy {
max_out_of_order_ms : Int
allowed_lateness_ms : Int
idle_timeout_ms : Int
} derive(Eq,
Debug
)

#
WatermarkPolicy::new

fn WatermarkPolicy::new(max_out_of_order_ms : Int, allowed_lateness_ms : Int, idle_timeout_ms : Int) -> WatermarkPolicy

#
WatermarkState

pub(all) struct WatermarkState {
policy : WatermarkPolicy
max_event_time_ms : Int
current_watermark_ms : Int
last_arrival_ms : Int
initialized : Bool
} derive(Eq,
Debug
)

#
WatermarkState::classify

fn WatermarkState::classify(self : WatermarkState, event : StreamEvent) -> EventStatus

#
WatermarkState::deadline_for

fn WatermarkState::deadline_for(self : WatermarkState, event_time_ms : Int) -> Int

#
WatermarkState::mark_idle

fn WatermarkState::mark_idle(self : WatermarkState, now_ms : Int) -> Bool

#
WatermarkState::new

#
WatermarkState::observe

#
WindowAggregate

pub(all) struct WindowAggregate {
window : TimeWindow
key : String
count : Int
sum : Int
late : Int
dropped : Int
} derive(Eq,
Debug
)

#
WindowAggregate::add

#
WindowAggregate::avg_x100

fn WindowAggregate::avg_x100(self : WindowAggregate) -> Int

#
WindowAggregate::new

fn WindowAggregate::new(window : TimeWindow, key : String) -> WindowAggregate

#
WindowAggregate::to_json

fn WindowAggregate::to_json(self : WindowAggregate) -> String

#
WindowSpec

pub(all) struct WindowSpec {
size_ms : Int
slide_ms : Int
} derive(Eq,
Debug
)

#
WindowSpec::max_windows_per_event

fn WindowSpec::max_windows_per_event(self : WindowSpec) -> Int

#
WindowSpec::sliding

fn WindowSpec::sliding(size_ms : Int, slide_ms : Int) -> WindowSpec

#
WindowSpec::tumbling

fn WindowSpec::tumbling(size_ms : Int) -> WindowSpec

#
assign_event

fn assign_event(state : WatermarkState, spec : WindowSpec, event : StreamEvent) -> EventAssignment

#
assign_windows

fn assign_windows(spec : WindowSpec, event_time_ms : Int) -> Array[TimeWindow]

#
audit_assignments

fn audit_assignments(assignments : Array[EventAssignment]) -> StreamAudit

#
default_policy

fn default_policy() -> WatermarkPolicy

#
five_minute_sliding_window

fn five_minute_sliding_window() -> WindowSpec

#
max_int

fn max_int(a : Int, b : Int) -> Int

#
minute_window

fn minute_window() -> WindowSpec

Powered by MoonBit

Site sourceReport issuePackagesBuild queueSkillsStatistics

© 2026 mooncakes.io