servo_base/generic_channel/
buffered.rs1use std::cell::RefCell;
6use std::mem;
7use std::panic::Location;
8
9use malloc_size_of_derive::MallocSizeOf;
10use serde::Serialize;
11
12use super::{GenericSender, SendResult};
13
14#[derive(MallocSizeOf)]
23pub struct GenericBufferedSender<T, U>
24where
25 T: Serialize,
26{
27 sender: GenericSender<T>,
28 buffer: RefCell<Vec<U>>,
29 #[ignore_malloc_size_of = "dyn are difficult to measure"]
30 buffering: Box<dyn Fn(Vec<U>) -> T>,
31 max_buffer: usize,
32}
33
34impl<T: Serialize, U> GenericBufferedSender<T, U> {
35 pub fn new(
41 sender: GenericSender<T>,
42 buffering: Box<dyn Fn(Vec<U>) -> T>,
43 max_buffer: usize,
44 ) -> Self {
45 Self {
46 sender,
47 buffer: RefCell::new(Vec::new()),
48 buffering,
49 max_buffer,
50 }
51 }
52
53 pub fn is_empty(&self) -> bool {
55 self.buffer.borrow().is_empty()
56 }
57
58 pub fn len(&self) -> usize {
60 self.buffer.borrow().len()
61 }
62
63 pub fn send(&self, msg: U) -> SendResult {
67 if self.buffer.borrow().len() + 1 >= self.max_buffer {
68 self.send_immediate(msg)
69 } else {
70 self.buffer.borrow_mut().push(msg);
71 Ok(())
72 }
73 }
74
75 #[inline]
76 #[track_caller]
77 pub fn send_or_warn(&self, msg: U) {
82 if let Err(error) = self.send(msg) {
83 let location = Location::caller();
84 log::warn!("Failed to send buffered messages due to `{error}` at {location:?}");
85 }
86 }
87
88 pub fn send_immediate(&self, msg: U) -> SendResult {
91 let mut buffer = self.buffer.borrow_mut();
92 buffer.push(msg);
93 let msgs = mem::take(&mut *buffer);
94 drop(buffer);
95 let packed = (self.buffering)(msgs);
96 self.sender.send(packed)
97 }
98
99 #[inline]
100 #[track_caller]
101 pub fn send_immediate_or_warn(&self, msg: U) {
106 if let Err(error) = self.send_immediate(msg) {
107 let location = Location::caller();
108 log::warn!(
109 "Failed to send (immediate) buffered messages due to `{error}` at {location:?}"
110 );
111 }
112 }
113
114 pub fn flush(&self) -> SendResult {
117 let mut buffer = self.buffer.borrow_mut();
118 if buffer.is_empty() {
119 return Ok(());
120 }
121 let msgs = mem::take(&mut *buffer);
122 drop(buffer);
123 let packed = (self.buffering)(msgs);
124 self.sender.send(packed)
125 }
126
127 #[inline]
128 #[track_caller]
129 pub fn flush_or_warn(&self) {
132 if let Err(error) = self.flush() {
133 let location = Location::caller();
134 log::warn!("Failed to flush buffered messages due to `{error}` at {location:?}");
135 }
136 }
137
138 pub fn discard(&self) {
140 self.buffer.borrow_mut().clear();
141 }
142}
143
144impl<T: Serialize, U> Drop for GenericBufferedSender<T, U> {
145 fn drop(&mut self) {
146 if !self.buffer.borrow().is_empty() {
148 let msgs = mem::take(&mut *self.buffer.borrow_mut());
149 let packed = (self.buffering)(msgs);
150 let _ = self.sender.send(packed);
151 }
152 }
153}