use crate::actor::Actor;
use crate::actor::ActorHandle;
use crate::message::ActorError;
use crate::message::Message;
use crate::message::MessageEnvelope;
use crate::message::Observations;
use async_trait::async_trait;
use time::format_description::well_known::Iso8601;
use time::OffsetDateTime;
use tokio::sync::mpsc;
extern crate serde;
extern crate serde_json;
pub struct JsonDecoder {
pub receiver: mpsc::Receiver<MessageEnvelope>,
pub output: ActorHandle,
}
fn extract_values_from_json(text: &str) -> Result<Observations, String> {
let observations: Observations = match serde_json::from_str(text) {
Ok(o) => o,
Err(e) => return Err(e.to_string()),
};
Ok(observations)
}
fn extract_datetime(datetime_str: &str) -> OffsetDateTime {
match OffsetDateTime::parse(datetime_str, &Iso8601::DEFAULT) {
Ok(d) => d,
Err(e) => {
log::warn!("can not parse datetime {} due to: {}", datetime_str, e);
OffsetDateTime::now_utc()
}
}
}
#[async_trait]
impl Actor for JsonDecoder {
async fn stop(&self) {}
async fn handle_envelope(&mut self, envelope: MessageEnvelope) {
let MessageEnvelope {
message,
respond_to,
datetime,
..
} = envelope;
match &message {
Message::PrintOneCmd { text } => match extract_values_from_json(text) {
Ok(observations) => {
log::trace!("json parsed");
let msg = Message::Update {
path: String::from(&observations.path),
datetime: extract_datetime(&observations.datetime),
values: observations.values,
};
let senv = MessageEnvelope {
message: msg,
respond_to, datetime,
..Default::default()
};
self.output.send(senv).await.expect("cannot send");
}
Err(error) => {
log::warn!("json parse error: {}", error);
if let Some(respond_to) = respond_to {
let etxt = format!("json parse error: {error}");
respond_to
.send(Err(ActorError { reason: etxt }))
.expect("can not return error");
}
}
},
m => {
let senv = MessageEnvelope {
message: m.clone(),
respond_to,
..Default::default()
};
self.output.send(senv).await.expect("cannot send");
}
}
}
}
impl JsonDecoder {
fn new(receiver: mpsc::Receiver<MessageEnvelope>, output: ActorHandle) -> Self {
JsonDecoder { receiver, output }
}
}
pub fn new(bufsz: usize, output: ActorHandle) -> ActorHandle {
async fn start(mut actor: JsonDecoder) {
while let Some(envelope) = actor.receiver.recv().await {
actor.handle_envelope(envelope).await;
}
}
let (sender, receiver) = mpsc::channel(bufsz);
let actor = JsonDecoder::new(receiver, output);
let actor_handle = ActorHandle::new(sender);
tokio::spawn(start(actor));
actor_handle
}