use crate::actor::respond_or_log_error;
use crate::actor::Actor;
use crate::actor::Handle;
use crate::message::Envelope;
use crate::message::Message;
use crate::message::MtHint;
use crate::message::NvError;
use crate::message::NvResult;
use crate::nvtime::OffsetDateTimeWrapper;
use async_trait::async_trait;
use serde_json::from_str;
use sqlx::Row;
use sqlx::SqlitePool;
use std::collections::HashMap;
use std::fmt;
use std::fs::File;
use std::path::Path;
use time::OffsetDateTime;
use tokio::sync::mpsc;
use tokio::sync::oneshot::Sender;
pub type StoreResult<T> = Result<T, StoreError>;
#[derive(Debug, Clone)]
pub struct StoreError {
pub reason: String,
}
impl fmt::Display for StoreError {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(f, "bad actor state: {}", self.reason)
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
enum StreamOption {
Close,
LeaveOpen,
}
pub struct StoreActor {
pub receiver: mpsc::Receiver<Envelope<f64>>,
pub dbconn: Option<SqlitePool>,
pub namespace: String,
pub disable_duplicate_detection: bool,
}
async fn insert_gene_mapping(
dbconn: &SqlitePool,
path: &String,
text: &String,
) -> Result<(), sqlx::error::Error> {
match sqlx::query("INSERT INTO gene_mappings (path, text) VALUES (?,?)")
.bind(path)
.bind(text)
.execute(dbconn)
.await
{
Ok(_) => Ok(()),
Err(e) => {
log::warn!("persisting gene mapping for {} failed: {:?}", path, e);
Err(e)
}
}
}
async fn insert_update(
dbconn: &SqlitePool,
path: &String,
datetime: OffsetDateTime,
sequence: OffsetDateTime,
values: HashMap<i32, f64>,
) -> Result<(), sqlx::error::Error> {
let dt_wrapper = OffsetDateTimeWrapper::new(datetime);
let sequence_wrapper = OffsetDateTimeWrapper::new(sequence);
match sqlx::query(
"INSERT INTO updates (path, timestamp, sequence, values_str) VALUES (?,?,?,?)",
)
.bind(path.clone())
.bind(dt_wrapper.datetime_num)
.bind(sequence_wrapper.datetime_num)
.bind(
serde_json::to_string(&values)
.map_err(|e| {
log::error!("cannot serialize values: {e:?}");
})
.ok(),
)
.execute(dbconn)
.await
{
Ok(_) => Ok(()),
Err(e) => {
log::warn!("jrnling for {} failed: {:?}", path, e);
Err(e)
}
}
}
async fn get_jrnl(dbconn: &SqlitePool, path: &str) -> StoreResult<Vec<Message<f64>>> {
match get_values(path, dbconn).await {
Ok(v) => Ok(v),
Err(e) => {
log::error!("cannot load update jrnl from db: {e:?}");
Err(StoreError {
reason: format!("cannot load jrnl from db: {e:?}"),
})
}
}
}
async fn get_mappings(dbconn: &SqlitePool, path: &str) -> StoreResult<Vec<Message<f64>>> {
match get_mappings_for_ns(path, dbconn).await {
Ok(v) => Ok(v),
Err(e) => {
log::error!("cannot load mappings from db: {e:?}");
Err(StoreError {
reason: format!("cannot load from db: {e:?}"),
})
}
}
}
async fn stream_message(
stream_to: &Option<mpsc::Sender<Message<f64>>>,
message: Message<f64>,
stream_option: StreamOption,
) {
if let Some(stream_to) = stream_to {
match stream_to.send(message).await {
Ok(_) => (),
Err(err) => {
log::error!("Can not integrate from helper: {}", err);
}
}
if stream_option == StreamOption::Close {
stream_to.closed().await;
};
} else {
log::trace!("no stream available for {message}");
}
}
async fn handle_gene_mapping(
path: String,
text: String,
dbconn: &SqlitePool,
respond_to: Option<Sender<NvResult<Message<f64>>>>,
) {
match insert_gene_mapping(dbconn, &path, &text).await {
Ok(_) => {
log::debug!("gene_mapping '{path}' -> '{text}' persisted");
respond_or_log_error(respond_to, Ok(Message::EndOfStream {}));
}
Err(e) => respond_or_log_error(
respond_to,
Err(NvError {
reason: e.to_string(),
}),
),
}
}
async fn handle_update(
path: String,
datetime: OffsetDateTime,
sequence: OffsetDateTime,
values: HashMap<i32, f64>,
disable_duplicate_detection: bool,
dbconn: &SqlitePool,
respond_to: Option<Sender<NvResult<Message<f64>>>>,
) {
let dt = if disable_duplicate_detection {
sequence
} else {
datetime
};
match insert_update(dbconn, &path, dt, sequence, values).await {
Ok(_) => respond_or_log_error(respond_to, Ok(Message::EndOfStream {})),
Err(e) => respond_or_log_error(
respond_to,
Err(NvError {
reason: e.to_string(),
}),
),
}
}
async fn handle_load_cmd(
path: String,
dbconn: &SqlitePool,
stream_to: Option<mpsc::Sender<Message<f64>>>,
) {
match get_jrnl(dbconn, &path).await {
Ok(rows) => {
for message in rows {
stream_message(&stream_to, message, StreamOption::LeaveOpen).await;
}
}
Err(e) => {
log::error!("cannot load jrnl: {path} {e:?}");
}
};
stream_message(&stream_to, Message::EndOfStream {}, StreamOption::Close).await;
}
async fn handle_gene_mapping_load_cmd(
path: String,
dbconn: &SqlitePool,
stream_to: Option<mpsc::Sender<Message<f64>>>,
) {
match get_mappings(dbconn, &path).await {
Ok(rows) => {
for message in rows {
stream_message(&stream_to, message, StreamOption::LeaveOpen).await;
}
}
Err(e) => {
log::error!("cannot load gene mapping jrnl: {path} {e:?}");
}
};
stream_message(&stream_to, Message::EndOfStream {}, StreamOption::Close).await;
}
#[async_trait]
impl Actor for StoreActor {
async fn handle_envelope(&mut self, envelope: Envelope<f64>) {
if let Some(dbconn) = &self.dbconn {
let Envelope {
message,
respond_to,
stream_to,
datetime: sequence,
..
} = envelope;
match message {
Message::Update {
path,
datetime,
values,
} => {
handle_update(
path,
datetime,
sequence,
values,
self.disable_duplicate_detection,
dbconn,
respond_to,
)
.await;
}
Message::LoadCmd { path, hint } if hint == MtHint::GeneMapping => {
handle_gene_mapping_load_cmd(path, dbconn, stream_to).await;
}
Message::LoadCmd { path, hint } if hint == MtHint::Update => {
handle_load_cmd(path, dbconn, stream_to).await;
}
Message::Content { path, text, hint }
if path.is_some() && hint == MtHint::GeneMapping =>
{
match path {
Some(path) => {
handle_gene_mapping(path, text, dbconn, respond_to).await;
}
_ => {
log::error!("path not set");
}
}
}
m => log::warn!("Unexpected: {m}"),
}
} else {
log::error!("DB not configured");
}
}
async fn start(&mut self) {}
async fn stop(&self) {
if let Some(c) = &self.dbconn {
c.close().await;
}
}
}
async fn get_mappings_for_ns(
path: &str,
dbconn: &SqlitePool,
) -> Result<Vec<Message<f64>>, sqlx::error::Error> {
log::debug!("loading mappings for ns {path}");
sqlx::query("SELECT path, text FROM gene_mappings;")
.bind(path)
.try_map(|row: sqlx::sqlite::SqliteRow| {
let path = match row.try_get(0) {
Ok(p) => p,
Err(e) => {
log::error!("cannot read path");
return Err(sqlx::Error::Decode(Box::new(e)));
}
};
let text = match row.try_get(1) {
Ok(p) => p,
Err(e) => {
log::error!("cannot read text");
return Err(sqlx::Error::Decode(Box::new(e)));
}
};
Ok(Message::Content {
path: Some(path),
text,
hint: MtHint::GeneMapping,
})
})
.fetch_all(dbconn)
.await
}
async fn get_values(
path: &str,
dbconn: &SqlitePool,
) -> Result<Vec<Message<f64>>, sqlx::error::Error> {
sqlx::query("SELECT timestamp, values_str FROM updates WHERE path = ?")
.bind(path)
.try_map(|row: sqlx::sqlite::SqliteRow| {
let date_parsed_num = match from_str(row.get(0)) {
Ok(val) => val,
Err(e) => return Err(sqlx::Error::Decode(Box::new(e))),
};
let date_parsed = OffsetDateTimeWrapper {
datetime_num: date_parsed_num,
};
let values = match row.try_get(1) {
Ok(val_str) => match from_str(val_str) {
Ok(val) => val,
Err(e) => return Err(sqlx::Error::Decode(Box::new(e))),
},
Err(e) => return Err(sqlx::Error::Decode(Box::new(e))),
};
let dt = match date_parsed.to_ts() {
Ok(dt) => dt,
Err(e) => {
log::error!("can not parse date - using 'now': {e}");
OffsetDateTime::now_utc()
}
};
Ok(Message::Update {
path: String::from(path),
datetime: dt,
values,
})
})
.fetch_all(dbconn)
.await
}
impl StoreActor {
const fn new(
receiver: mpsc::Receiver<Envelope<f64>>,
dbconn: Option<SqlitePool>,
namespace: String,
disable_duplicate_detection: bool,
) -> Self {
Self {
receiver,
dbconn,
namespace,
disable_duplicate_detection,
}
}
}
async fn define_gene_mapping_table_if_not_exist(
db_url: &str,
dbconn: &SqlitePool,
) -> StoreResult<()> {
let rows = sqlx::query("PRAGMA journal_mode;")
.fetch_all(dbconn)
.await
.map_err(|e| StoreError {
reason: format!("Failed to fetch journal_mode: {e}"),
})?;
let journal_mode: String = rows[0].get("journal_mode");
log::info!("connected to db in journal_mode for mappings: {journal_mode}");
sqlx::query(
"CREATE TABLE IF NOT EXISTS gene_mappings (
path TEXT NOT NULL,
text TEXT NOT NULL,
PRIMARY KEY (path)
)",
)
.execute(dbconn)
.await
.map_err(|e| StoreError {
reason: format!("Failed to create file {db_url}: {e}"),
})?;
Ok(())
}
async fn define_updates_table_if_not_exist(db_url: &str, dbconn: &SqlitePool) -> StoreResult<()> {
let rows = sqlx::query("PRAGMA journal_mode;")
.fetch_all(dbconn)
.await
.map_err(|e| StoreError {
reason: format!("Failed to fetch journal_mode: {e}"),
})?;
let journal_mode: String = rows[0].get("journal_mode");
log::info!("connected to db in journal_mode: {journal_mode}");
sqlx::query(
"CREATE TABLE IF NOT EXISTS updates (
path TEXT NOT NULL,
timestamp TEXT NOT NULL,
sequence TEXT NOT NULL,
values_str TEXT NOT NULL,
PRIMARY KEY (path, timestamp)
)",
)
.execute(dbconn)
.await
.map_err(|e| StoreError {
reason: format!("Failed to create file {db_url}: {e}"),
})?;
Ok(())
}
async fn enable_wal(db_url: &str, dbconn: &SqlitePool) -> StoreResult<()> {
match sqlx::query("PRAGMA journal_mode = WAL;")
.execute(dbconn)
.await
{
Ok(_) => Ok(()),
Err(e) => Err(StoreError {
reason: format!("Failed to create file {db_url}: {e}"),
}),
}
}
async fn init_db(namespace: String, write_ahead_logging: bool) -> StoreResult<SqlitePool> {
let db_url_string: String = format!("{namespace}.db");
let db_url: &str = &db_url_string;
let db_path = Path::new(db_url);
if !db_path.exists() {
match File::create(db_url) {
Ok(_) => log::debug!("File {} has been created", db_url),
Err(e) => {
return Err(StoreError {
reason: format!("Failed to create file {db_url}: {e}"),
});
}
}
}
match SqlitePool::connect(db_url).await {
Ok(dbconn) => {
if write_ahead_logging {
match enable_wal(db_url, &dbconn).await {
Ok(_) => {}
Err(e) => return Err(e),
}
}
match define_updates_table_if_not_exist(db_url, &dbconn).await {
Ok(_) => match define_gene_mapping_table_if_not_exist(db_url, &dbconn).await {
Ok(_) => Ok(dbconn),
Err(e) => Err(e),
},
Err(e) => Err(e),
}
}
Err(e) => {
log::error!("cannot connect to db: {e:?}");
return Err(StoreError {
reason: format!("{e:?}"),
});
}
}
}
#[must_use]
pub fn new(
bufsz: usize,
namespace: String,
write_ahead_logging: bool,
disable_duplicate_detection: bool,
) -> Handle {
async fn start(mut actor: StoreActor, namespace: String, write_ahead_logging: bool) {
let dbconn = init_db(namespace, write_ahead_logging)
.await
.map_err(|e| {
log::error!("cannot get dbconn: {e:?}");
})
.ok();
actor.dbconn = dbconn;
while let Some(envelope) = actor.receiver.recv().await {
actor.handle_envelope(envelope).await;
}
actor.stop().await;
}
let (sender, receiver) = mpsc::channel(bufsz);
let actor = StoreActor::new(
receiver,
None,
namespace.clone(),
disable_duplicate_detection,
);
let actor_handle = Handle::new(sender);
tokio::spawn(start(actor, namespace, write_ahead_logging));
actor_handle
}