MCPcopy Create free account
hub / github.com/It4innovations/hyperqueue / send_event

Method send_event

crates/hyperqueue/src/server/event/streamer.rs:364–393  ·  view source on GitHub ↗
(
        &self,
        payload: EventPayload,
        now: Option<DateTime<Utc>>,
        forward_mode: ForwardMode,
    )

Source from the content-addressed store, hash-verified

362 }
363
364 fn send_event(
365 &self,
366 payload: EventPayload,
367 now: Option<DateTime<Utc>>,
368 forward_mode: ForwardMode,
369 ) {
370 let mut inner = self.inner.get_mut();
371 if inner.storage_sender.is_none() && inner.client_listeners.is_empty() {
372 return;
373 }
374 let event = Event {
375 time: now.unwrap_or_else(Utc::now),
376 payload,
377 };
378
379 inner.client_listeners.retain(|listener| {
380 if listener.filter.check(&event.payload) {
381 listener.sender.send(event.clone()).is_ok()
382 } else {
383 true
384 }
385 });
386
387 if let Some(ref streamer) = inner.storage_sender
388 && matches!(forward_mode, ForwardMode::StreamAndPersist)
389 && streamer.send(EventStreamMessage::Event(event)).is_err()
390 {
391 log::error!("Event streaming queue has been closed.");
392 }
393 }
394
395 pub fn on_server_stop(&self) {
396 self.send_event(

Callers 15

on_worker_addedMethod · 0.80
on_worker_lostMethod · 0.80
on_overview_receivedMethod · 0.80
on_job_openedMethod · 0.80
on_job_closedMethod · 0.80
on_job_submittedMethod · 0.80
on_job_completedMethod · 0.80
on_job_cancelMethod · 0.80
on_task_startedMethod · 0.80
on_task_finishedMethod · 0.80
on_task_canceledMethod · 0.80
on_task_abortedMethod · 0.80

Calls 6

EventClass · 0.70
get_mutMethod · 0.45
is_emptyMethod · 0.45
checkMethod · 0.45
sendMethod · 0.45
cloneMethod · 0.45

Tested by

no test coverage detected