use crate::actor::Actor;
use crate::actor::ActorHandle;
use crate::genes::DefaultGene;
use crate::genes::Gene;
use crate::message::Message;
use crate::message::MessageEnvelope;
use async_trait::async_trait;
use std::collections::HashMap;
use time::OffsetDateTime;
use tokio::sync::mpsc;
pub struct StateActor {
pub receiver: mpsc::Receiver<MessageEnvelope>,
pub output: Option<ActorHandle>,
pub state: HashMap<i32, f64>,
pub path: String,
pub gene: Box<dyn Gene + Send + Sync>,
}
#[async_trait]
impl Actor for StateActor {
async fn stop(&self) {}
async fn handle_envelope(&mut self, envelope: MessageEnvelope) {
let MessageEnvelope {
message,
respond_to,
stream_from,
..
} = envelope;
match message {
Message::InitCmd { .. } => {
log::trace!("{} init started...", self.path);
if let Some(mut stream_from) = stream_from {
log::debug!("{} init", self.path);
let mut count = 0;
while let Some(message) = stream_from.recv().await {
if !self.update_state(message.clone()) {
stream_from.close();
} else {
count += 1;
}
}
log::debug!("{} finished init from {} events", self.path, count);
if let Some(respond_to) = respond_to {
respond_to
.send(Ok(Message::EndOfStream {}))
.expect("can not reply to init");
}
}
}
Message::Update { .. } => {
log::trace!("{} handling update", self.path);
self.update_state(message.clone());
if let Some(respond_to) = respond_to {
respond_to
.send(Ok(self.get_state_rpt()))
.expect("can not reply to ask");
}
}
Message::Query { .. } => {
if let Some(respond_to) = respond_to {
respond_to
.send(Ok(self.get_state_rpt()))
.expect("can not reply to ask");
}
}
m => {
log::warn!("unexpected message: {m:?}");
}
}
if let Some(output_handle) = &self.output {
output_handle
.tell(self.get_state_rpt())
.await
.expect("cannot tell");
}
}
}
impl StateActor {
fn update_state(&mut self, message: Message) -> bool {
let mut updated = false;
match message {
Message::Update { values, .. } => {
updated = true;
self.state.extend(&values); }
Message::Query { path: _ } => {
}
Message::EndOfStream {} => {
log::trace!("{} end of update stream", self.path);
}
m => {
log::warn!("unexpected message in update stream: {:?}", m);
}
}
updated
}
fn get_state_rpt(&self) -> Message {
Message::StateReport {
path: self.path.clone(),
values: self.state.clone(),
datetime: OffsetDateTime::now_utc(),
}
}
fn new(
path: String,
receiver: mpsc::Receiver<MessageEnvelope>,
output: Option<ActorHandle>,
) -> Self {
let gene = Box::new(DefaultGene::new());
let state = HashMap::new();
StateActor {
path,
receiver,
output,
state,
gene,
}
}
}
pub fn new(path: String, bufsz: usize, output: Option<ActorHandle>) -> ActorHandle {
async fn start<'a>(mut actor: StateActor) {
while let Some(envelope) = actor.receiver.recv().await {
actor.handle_envelope(envelope).await;
}
}
let (sender, receiver) = mpsc::channel(bufsz);
let actor = StateActor::new(path, receiver, output);
let actor_handle = ActorHandle::new(sender);
tokio::spawn(start(actor));
actor_handle
}