use crate::actor::Actor;
use crate::actor::Handle;
use crate::actor::State;
use crate::genes::DefaultGene;
use crate::genes::Gene;
use crate::message::Message;
use crate::message::Envelope;
use async_trait::async_trait;
use time::OffsetDateTime;
use tokio::sync::mpsc;
pub struct StateActor {
pub receiver: mpsc::Receiver<Envelope>,
pub output: Option<Handle>,
pub state: State<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: Envelope) {
let Envelope {
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()) {
count += 1;
} else {
stream_from.close();
}
}
log::debug!("{} finished init from {} events", self.path, count);
if let Some(respond_to) = respond_to {
if let Err(err) = respond_to.send(Ok(Message::EndOfStream {})) {
log::error!("Error sending reply to init: {:?}", err);
}
}
}
}
Message::Update { .. } => {
log::trace!("{} handling update", self.path);
self.update_state(message.clone());
if let Some(respond_to) = respond_to {
if let Err(err) = respond_to.send(Ok(self.get_state_rpt())) {
log::error!("Error sending reply to ask: {:?}", err);
}
}
}
Message::Query { .. } => {
if let Some(respond_to) = respond_to {
if let Err(err) = respond_to.send(Ok(self.get_state_rpt())) {
log::error!("Error sending reply to ask: {:?}", err);
}
}
}
m => {
log::warn!("unexpected message: {m:?}");
}
}
if let Some(output_handle) = &self.output {
if let Err(err) = output_handle.tell(self.get_state_rpt()).await {
log::error!("Error telling output actor: {:?}", err);
}
}
}
}
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<Envelope>,
output: Option<Handle>,
) -> Self {
let gene = Box::new(DefaultGene::new());
let state = State::new();
Self {
receiver,
output,
state,
path,
gene,
}
}
}
#[must_use]
pub fn new(path: String, bufsz: usize, output: Option<Handle>) -> Handle {
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 = Handle::new(sender);
tokio::spawn(start(actor));
actor_handle
}