use crate::actor::respond_or_log_error;
use crate::actor::Actor;
use crate::actor::Handle;
use crate::actor::State;
use crate::gene::Gene;
use crate::message::Envelope;
use crate::message::Message;
use crate::message::NvError;
use async_trait::async_trait;
use time::OffsetDateTime;
use tokio::sync::mpsc;
use tracing::debug;
use tracing::error;
use tracing::trace;
use tracing::warn;
pub struct StateActor {
pub receiver: mpsc::Receiver<Envelope<f64>>,
pub output: Option<Handle>,
pub state: State<f64>,
pub path: String,
pub gene: Box<dyn Gene<f64> + Send + Sync>,
}
#[async_trait]
impl Actor for StateActor {
async fn handle_envelope(&mut self, envelope: Envelope<f64>) {
let Envelope {
message,
respond_to,
stream_from,
..
} = envelope;
match message {
Message::InitCmd { .. } => {
trace!("{} init started...", self.path);
if let Some(mut stream_from) = stream_from {
debug!("{} init", self.path);
let mut count = 0;
while let Some(message) = stream_from.recv().await {
match &message {
Message::EndOfStream {} => {
break;
}
_ => {
if self.update_state(message.clone()) {
count += 1;
} else {
trace!("{} init closing stream.", self.path);
break;
}
}
}
}
trace!("{} init closing stream.", self.path);
stream_from.close();
debug!("{} finished init from {} events", self.path, count);
respond_or_log_error(respond_to, Ok(Message::EndOfStream {}));
}
}
Message::Update { .. } => {
trace!("{} handling update", self.path);
if self.update_state(message.clone()) {
respond_or_log_error(respond_to, Ok(self.get_state_rpt()));
} else {
error!("Error applying operators in ask");
respond_or_log_error(
respond_to,
Err(NvError {
reason: String::from("cannot apply operators"),
}),
);
}
}
Message::Query { .. } => {
respond_or_log_error(respond_to, Ok(self.get_state_rpt()));
}
m => {
warn!("unexpected message: {m}");
}
}
if let Some(output_handle) = &self.output {
if let Err(err) = output_handle.tell(self.get_state_rpt()).await {
error!("Error telling output actor: {err:?}");
}
}
}
async fn stop(&self) {}
async fn start(&mut self) {}
}
impl StateActor {
fn update_state(&mut self, message: Message<f64>) -> bool {
match self.gene.apply_operators(self.state.clone(), message) {
Ok(new_state) => {
self.state = new_state;
true
}
Err(e) => {
error!("Error applying operators in ask: {e:?}");
false
}
}
}
fn get_state_rpt(&self) -> Message<f64> {
Message::StateReport {
path: self.path.clone(),
values: self.state.clone(),
datetime: OffsetDateTime::now_utc(), }
}
fn new(
path: String,
receiver: mpsc::Receiver<Envelope<f64>>,
output: Option<Handle>,
gene: Box<dyn Gene<f64> + Send + Sync>,
) -> Self {
let state = State::new();
Self {
receiver,
output,
state,
path,
gene,
}
}
}
#[must_use]
pub fn new(
path: String,
bufsz: usize,
gene: Box<dyn Gene<f64> + Send + Sync>,
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, gene);
let actor_handle = Handle::new(sender);
tokio::spawn(start(actor));
actor_handle
}