nexus_acto_rs/actor/event_stream/
event_stream_process.rs1use std::any::Any;
2
3use async_trait::async_trait;
4
5use crate::actor::actor::ExtendedPid;
6use crate::actor::actor_system::ActorSystem;
7use crate::actor::message::unwrap_envelope;
8use crate::actor::message::MessageHandle;
9use crate::actor::process::Process;
10
11#[derive(Debug, Clone)]
12pub struct EventStreamProcess {
13 system: ActorSystem,
14}
15
16impl EventStreamProcess {
17 pub fn new(system: ActorSystem) -> Self {
18 EventStreamProcess { system }
19 }
20}
21
22#[async_trait]
23impl Process for EventStreamProcess {
24 async fn send_user_message(&self, _: Option<&ExtendedPid>, message_handle: MessageHandle) {
25 let (_, msg, _) = unwrap_envelope(message_handle);
26 self.system.get_event_stream().await.publish(msg).await;
27 }
28
29 async fn send_system_message(&self, _: &ExtendedPid, _: MessageHandle) {}
30
31 async fn stop(&self, _: &ExtendedPid) {}
32
33 fn set_dead(&self) {}
34
35 fn as_any(&self) -> &dyn Any {
36 self
37 }
38}