use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::fmt;
use time::OffsetDateTime;
use tokio::sync::mpsc;
use tokio::sync::oneshot;
pub type NvResult<T> = Result<T, NvError>;
#[derive(Debug, Clone)]
pub struct NvError {
pub reason: String,
}
impl fmt::Display for NvError {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(f, "{}", self.reason)
}
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PathQuery {
pub path: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct Observations {
pub datetime: String,
pub values: HashMap<i32, f64>,
pub path: String,
}
#[derive(Debug)]
pub struct Envelope<T> {
pub message: Message<T>,
pub respond_to: Option<oneshot::Sender<NvResult<Message<T>>>>,
pub datetime: OffsetDateTime,
pub stream_to: Option<mpsc::Sender<Message<T>>>,
pub stream_from: Option<mpsc::Receiver<Message<T>>>,
}
#[derive(Debug, Clone)]
pub enum MtHint {
Update,
Query,
}
impl fmt::Display for MtHint {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
let display_text = match self {
Self::Query => "query",
Self::Update => "update",
};
write!(f, "{display_text}")
}
}
#[derive(Debug, Clone)]
pub enum Message<T> {
Query {
path: String,
},
Update {
datetime: OffsetDateTime,
path: String,
values: HashMap<i32, T>,
},
StateReport {
datetime: OffsetDateTime,
path: String,
values: HashMap<i32, T>,
},
EndOfStream {},
InitCmd {},
LoadCmd {
path: String,
},
ReadAllCmd {},
TextMsg {
text: String,
hint: MtHint,
},
}
impl<T> fmt::Display for Envelope<T> {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
let display_text = "env TODO";
write!(f, "{display_text}")
}
}
impl<T> fmt::Display for Message<T> {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
let display_text = match self {
Self::TextMsg { text, hint } => format!("[TextMsg {hint}:{text}]"),
Self::LoadCmd { path } => format!("[LoadCmd {path}]"),
Self::ReadAllCmd {} => "[ReadAllCmd]".to_string(),
Self::InitCmd {} => "[InitCmd]".to_string(),
Self::EndOfStream {} => "[EndOfStream]".to_string(),
Self::StateReport { .. } => "[StateReport]".to_string(),
Self::Update { .. } => "[Update]".to_string(),
Self::Query { .. } => "[Query]".to_string(),
};
write!(f, "{display_text}")
}
}
impl<T> Default for Envelope<T> {
fn default() -> Self {
Self {
message: Message::ReadAllCmd {},
respond_to: None,
datetime: OffsetDateTime::now_utc(),
stream_to: None,
stream_from: None,
}
}
}
struct LifeCycleBuilder<T> {
load_from: Option<mpsc::Receiver<Message<T>>>,
send_to: Option<mpsc::Sender<Message<T>>>,
send_to_path: Option<String>,
respond_to: Option<oneshot::Sender<NvResult<Message<T>>>>,
}
impl<T> LifeCycleBuilder<T> {
const fn new() -> Self {
Self {
load_from: None,
send_to: None,
send_to_path: None,
respond_to: None,
}
}
#[allow(clippy::missing_const_for_fn)]
fn with_respond_to(mut self, respond_to: oneshot::Sender<NvResult<Message<T>>>) -> Self {
self.respond_to = Some(respond_to);
self
}
#[allow(clippy::missing_const_for_fn)]
fn with_load_from(mut self, load_from: mpsc::Receiver<Message<T>>) -> Self {
self.load_from = Some(load_from);
self
}
#[allow(clippy::missing_const_for_fn)]
fn with_send_to(mut self, send_to: mpsc::Sender<Message<T>>, send_to_path: String) -> Self {
self.send_to = Some(send_to);
self.send_to_path = Some(send_to_path);
self
}
fn build(self) -> (Envelope<T>, Envelope<T>) {
(
Envelope {
datetime: OffsetDateTime::now_utc(),
respond_to: self.respond_to,
stream_from: self.load_from,
stream_to: None,
message: Message::InitCmd {},
},
Envelope {
datetime: OffsetDateTime::now_utc(),
respond_to: None,
stream_from: None,
stream_to: self.send_to,
message: Message::LoadCmd {
path: self.send_to_path.unwrap_or_default(),
},
},
)
}
}
#[must_use]
pub fn create_init_lifecycle<T>(
path: String,
bufsz: usize,
respond_to: oneshot::Sender<NvResult<Message<T>>>,
) -> (Envelope<T>, Envelope<T>) {
let (tx, rx) = mpsc::channel(bufsz);
let builder = LifeCycleBuilder::new()
.with_load_from(rx)
.with_send_to(tx, path)
.with_respond_to(respond_to);
builder.build()
}