Serialize the target object to JSON as a single line
(&self, v: T, required: bool)
| 212 | |
| 213 | /// Serialize the target object to JSON as a single line |
| 214 | pub(crate) async fn send_impl<T: Serialize>(&self, v: T, required: bool) -> Result<()> { |
| 215 | let mut guard = self.inner.lock().await; |
| 216 | // Check if we have an inner value; if not, nothing to do. |
| 217 | let Some(inner) = guard.as_mut() else { |
| 218 | return Ok(()); |
| 219 | }; |
| 220 | |
| 221 | // If this is our first message, emit the Start message |
| 222 | if !inner.sent_start { |
| 223 | inner.sent_start = true; |
| 224 | let start = Event::Start { |
| 225 | version: API_VERSION.into(), |
| 226 | }; |
| 227 | Self::send_impl_inner(inner, &start).await?; |
| 228 | } |
| 229 | |
| 230 | // For messages that can be dropped, if we already sent an update within this cycle, discard this one. |
| 231 | // TODO: Also consider querying the pipe buffer and also dropping if we can't do this write. |
| 232 | let now = Instant::now(); |
| 233 | if !required { |
| 234 | const REFRESH_MS: u32 = 1000 / REFRESH_HZ as u32; |
| 235 | if let Some(elapsed) = inner.last_write.map(|w| now.duration_since(w)) { |
| 236 | if elapsed.as_millis() < REFRESH_MS.into() { |
| 237 | return Ok(()); |
| 238 | } |
| 239 | } |
| 240 | } |
| 241 | |
| 242 | Self::send_impl_inner(inner, &v).await?; |
| 243 | // Update the last write time |
| 244 | inner.last_write = Some(now); |
| 245 | Ok(()) |
| 246 | } |
| 247 | |
| 248 | /// Send an event. |
| 249 | pub(crate) async fn send(&self, event: Event<'_>) { |
no test coverage detected