Skip to content

Commit

Permalink
Expose queue_len (#8)
Browse files Browse the repository at this point in the history
  • Loading branch information
hpeebles authored Feb 7, 2024
1 parent 2fc5a9f commit baad3e4
Showing 1 changed file with 12 additions and 7 deletions.
19 changes: 12 additions & 7 deletions client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -102,22 +102,27 @@ impl<R: Runtime + Send + 'static> Client<R> {
pub fn flush_batch(&self) {
let mut guard = self.inner.lock().unwrap();
guard.next_flush_scheduled = None;
let max_batch_size = guard.max_batch_size;

let events = if guard.events.len() < max_batch_size {
mem::take(&mut guard.events)
} else {
guard.events.drain(..max_batch_size).collect()
};
if !guard.events.is_empty() {
let max_batch_size = guard.max_batch_size;

let events = if guard.events.len() < max_batch_size {
mem::take(&mut guard.events)
} else {
guard.events.drain(..max_batch_size).collect()
};

if !events.is_empty() {
let clone = self.clone();
guard
.runtime
.flush(events.clone(), move || clone.requeue_events(events));
}
}

pub fn queue_len(&self) -> usize {
self.inner.lock().unwrap().events.len()
}

fn requeue_events(&self, events: Vec<IdempotentEvent>) {
let mut guard = self.inner.lock().unwrap();
guard.events.extend(events);
Expand Down

0 comments on commit baad3e4

Please sign in to comment.