use vantage_core::Result;
use vantage_types::Record;
use ciborium::Value as CborValue;
use crate::dio::{Dio, DioEvent, DioInner};
use crate::ops::{ChangeFlash, FlashKind};
impl Dio {
pub async fn flash(&self, mut flash: ChangeFlash) -> Result<()> {
let Some(id) = flash.id().map(str::to_string) else {
return crate::dio::worker::run_write_through(self, flash).await;
};
let _pending = self.inner.pending_flashes.begin(id.clone());
let pre = self.inner.cache.get_value(&id).await?;
flash.ensure_before(pre.as_ref());
stage_in_cache(&self.inner, &flash, pre.as_ref()).await?;
let _ = self
.inner
.event_bus
.send(DioEvent::WritePending { id: id.clone() });
match crate::dio::worker::run_write_through(self, flash.clone()).await {
Ok(()) => {
reassert_confirmed(&self.inner, &flash).await?;
let _ = self.inner.event_bus.send(DioEvent::RecordChanged { id });
Ok(())
}
Err(err) => {
match &pre {
Some(prev) => self.inner.cache.insert_value(&id, prev).await?,
None => self.inner.cache.delete_value(&id).await?,
}
let _ = self.inner.event_bus.send(DioEvent::WriteReverted {
id,
error: err.to_string(),
});
Err(err)
}
}
}
pub async fn flash_patch(
&self,
id: impl Into<String>,
partial: Record<CborValue>,
) -> Result<()> {
self.flash(ChangeFlash::new(FlashKind::Patch, Some(id.into()), partial))
.await
}
pub async fn flash_insert(
&self,
id: impl Into<String>,
record: Record<CborValue>,
) -> Result<()> {
self.flash(ChangeFlash::insert(id, record)).await
}
pub async fn flash_replace(
&self,
id: impl Into<String>,
record: Record<CborValue>,
) -> Result<()> {
self.flash(ChangeFlash::replace(id, record)).await
}
pub async fn flash_delete(&self, id: impl Into<String>) -> Result<()> {
self.flash(ChangeFlash::delete(id)).await
}
}
async fn reassert_confirmed(inner: &DioInner, flash: &ChangeFlash) -> Result<()> {
let Some(id) = flash.id() else {
return Ok(());
};
match flash.kind() {
FlashKind::Insert | FlashKind::Replace | FlashKind::Patch => {
let mut merged = inner.cache.get_value(id).await?.unwrap_or_default();
for (k, v) in flash.patch() {
merged.insert(k.clone(), v.clone());
}
inner.cache.insert_value(id, &merged).await
}
FlashKind::Delete => inner.cache.delete_value(id).await,
FlashKind::Clear => Ok(()),
}
}
async fn stage_in_cache(
inner: &DioInner,
flash: &ChangeFlash,
pre: Option<&Record<CborValue>>,
) -> Result<()> {
let Some(id) = flash.id() else {
return Ok(());
};
match flash.kind() {
FlashKind::Insert | FlashKind::Replace => inner.cache.insert_value(id, flash.patch()).await,
FlashKind::Patch => {
let mut merged = pre.cloned().unwrap_or_default();
for (k, v) in flash.patch() {
merged.insert(k.clone(), v.clone());
}
inner.cache.insert_value(id, &merged).await
}
FlashKind::Delete => inner.cache.delete_value(id).await,
FlashKind::Clear => Ok(()),
}
}