///|
pub(all) struct CausalBuffer {
delivered : VectorClock
pending : Array[CausalEvent]
max_pending : Int
dropped : Int
}
///|
pub fn CausalBuffer::new() -> CausalBuffer {
CausalBuffer::with_capacity(1024)
}
///|
/// Creates a bounded buffer. When full, a newly pending event is rejected and
/// `dropped_count` records the overflow; callers can use it for backpressure.
pub fn CausalBuffer::with_capacity(max_pending : Int) -> CausalBuffer {
{
delivered: VectorClock::new(),
pending: [],
max_pending: if max_pending > 0 {
max_pending
} else {
1
},
dropped: 0,
}
}
///|
pub fn CausalBuffer::with_delivered(delivered : VectorClock) -> CausalBuffer {
{ delivered, pending: [], max_pending: 1024, dropped: 0 }
}
///|
pub fn CausalBuffer::ready(self : CausalBuffer, event : CausalEvent) -> Bool {
let origin_next = self.delivered.get(event.origin) + 1
if event.clock.get(event.origin) != origin_next {
return false
}
for entry in event.clock.entries {
if entry.node != event.origin &&
entry.counter > self.delivered.get(entry.node) {
return false
}
}
true
}
///|
pub fn CausalBuffer::deliver(
self : CausalBuffer,
event : CausalEvent,
) -> CausalBuffer {
{ ..self, delivered: self.delivered.merge(event.clock) }
}
///|
pub fn CausalBuffer::push(
self : CausalBuffer,
event : CausalEvent,
) -> CausalBuffer {
if self.ready(event) {
self.deliver(event)
} else if self.pending.length() >= self.max_pending {
{ ..self, dropped: self.dropped + 1 }
} else {
let next = self.pending.copy()
next.push(event)
{ ..self, pending: next }
}
}
///|
pub fn CausalBuffer::flush_ready(self : CausalBuffer) -> CausalBuffer {
let mut delivered = self.delivered
let mut pending = self.pending
let mut changed = true
while changed {
changed = false
let next_pending = []
for event in pending {
let buffer = {
delivered,
pending: [],
max_pending: self.max_pending,
dropped: self.dropped,
}
if buffer.ready(event) {
delivered = delivered.merge(event.clock)
changed = true
} else {
next_pending.push(event)
}
}
pending = next_pending
}
{ delivered, pending, max_pending: self.max_pending, dropped: self.dropped }
}
///|
pub fn CausalBuffer::pending_count(self : CausalBuffer) -> Int {
self.pending.length()
}
///|
pub fn CausalBuffer::dropped_count(self : CausalBuffer) -> Int {
self.dropped
}