use std::fmt;
use std::marker::PhantomData;
use async_trait::async_trait;
use bytes::Bytes;
use destream::de;
use freqfs::FileSave;
use futures::future::TryFutureExt;
use futures::join;
use log::debug;
use safecast::{AsType, TryCastFrom, TryCastInto};
use tc_collection::{Collection, CollectionBase};
use tc_error::*;
use tc_scalar::Scalar;
use tc_transact::hash::{AsyncHash, Output, Sha256};
use tc_transact::public::{Route, StateInstance};
use tc_transact::{fs, Replicate};
use tc_transact::{Gateway, IntoView, Transact, Transaction, TxnId};
use tc_value::{Link, Value};
use tcgeneric::{Map, Tuple};
use super::data::{ChainBlock, History};
use super::{CacheBlock, ChainInstance, Recover, HISTORY};
pub struct BlockChain<State, Txn, FE, T> {
history: History<State, Txn, FE>,
subject: T,
}
impl<State, Txn, FE, T> Clone for BlockChain<State, Txn, FE, T>
where
T: Clone,
{
fn clone(&self) -> Self {
Self {
history: self.history.clone(),
subject: self.subject.clone(),
}
}
}
impl<State, Txn, FE, T> BlockChain<State, Txn, FE, T> {
fn new(subject: T, history: History<State, Txn, FE>) -> Self {
Self { subject, history }
}
pub fn history(&self) -> &History<State, Txn, FE> {
&self.history
}
}
impl<State, T> ChainInstance<State, T> for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance,
State::FE: CacheBlock + for<'a> fs::FileSave<'a>,
T: Route<State> + fmt::Debug,
Collection<State::Txn, State::FE>: TryCastFrom<State>,
Scalar: TryCastFrom<State>,
{
fn append_delete(&self, txn_id: TxnId, key: Value) -> TCResult<()> {
self.history.append_delete(txn_id, key)
}
fn append_put(&self, txn: State::Txn, key: Value, value: State) -> TCResult<()> {
self.history.append_put(txn, key, value)
}
fn subject(&self) -> &T {
&self.subject
}
}
#[async_trait]
impl<State> Replicate<State::Txn>
for BlockChain<State, State::Txn, State::FE, CollectionBase<State::Txn, State::FE>>
where
State: StateInstance,
State::FE: CacheBlock,
State: From<Collection<State::Txn, State::FE>> + From<Scalar>,
Collection<State::Txn, State::FE>: TryCastFrom<State>,
CollectionBase<State::Txn, State::FE>: Route<State>,
Scalar: TryCastFrom<State>,
Self: TryCastFrom<State>,
{
async fn replicate(&self, txn: &State::Txn, mut source: Link) -> TCResult<Output<Sha256>> {
let attr = source
.path_mut()
.pop()
.ok_or_else(|| bad_request!("invalid replica link: {source}"))?;
let chain = txn.get(source, attr).await?;
let chain: Self = chain.try_cast_into(|s| {
bad_request!("expected to replicate a chain of blocks but found {:?}", s,)
})?;
self.history()
.replicate(txn, self.subject(), chain.history().clone())
.await?;
AsyncHash::hash(self, *txn.id()).await
}
}
#[async_trait]
impl<State, T> fs::Persist<State::FE> for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance,
State::FE: CacheBlock + for<'a> fs::FileSave<'a>,
T: Route<State> + fs::Persist<State::FE, Txn = State::Txn> + fmt::Debug,
{
type Txn = State::Txn;
type Schema = T::Schema;
async fn create(
txn_id: TxnId,
schema: Self::Schema,
store: fs::Dir<State::FE>,
) -> TCResult<Self> {
debug!("BlockChain::create");
let subject = T::create(txn_id, schema.clone(), store).await?;
let mut dir = subject.dir().try_write_owned()?;
let history = {
let dir = dir.get_or_create_dir(HISTORY.to_string())?;
fs::Dir::load(txn_id, dir).await?
};
let history = History::create(txn_id, (), history.into()).await?;
Ok(BlockChain::new(subject, history))
}
async fn load(
txn_id: TxnId,
schema: Self::Schema,
store: fs::Dir<State::FE>,
) -> TCResult<Self> {
debug!("BlockChain::load {}", std::any::type_name::<T>());
let subject = T::load(txn_id, schema.clone(), store).await?;
let mut dir = subject.dir().write_owned().await;
let history = {
let dir = dir.get_or_create_dir(HISTORY.to_string())?;
fs::Dir::load(txn_id, dir).await?
};
let history = History::load(txn_id, (), history.into()).await?;
Ok(BlockChain::new(subject, history))
}
fn dir(&self) -> tc_transact::fs::Inner<State::FE> {
self.subject.dir()
}
}
#[async_trait]
impl<State, T> Recover<State::FE> for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance + From<Collection<State::Txn, State::FE>> + From<Scalar>,
State::FE: CacheBlock + for<'a> fs::FileSave<'a>,
T: Route<State> + fmt::Debug,
Collection<State::Txn, State::FE>: TryCastFrom<State>,
Scalar: TryCastFrom<State>,
{
type Txn = State::Txn;
async fn recover(&self, txn: &State::Txn) -> TCResult<()> {
let write_ahead_log = self.history.read_log().await?;
for (past_txn_id, mutations) in &write_ahead_log.mutations {
super::data::replay_all(
&self.subject,
past_txn_id,
mutations,
txn,
self.history.store(),
)
.await?;
}
Ok(())
}
}
#[async_trait]
impl<State, T> fs::CopyFrom<State::FE, Self> for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance,
State::FE: CacheBlock + for<'a> fs::FileSave<'a>,
T: Route<State> + fs::Persist<State::FE, Txn = State::Txn> + fmt::Debug,
{
async fn copy_from(
_txn: &State::Txn,
_store: fs::Dir<State::FE>,
_instance: Self,
) -> TCResult<Self> {
Err(not_implemented!("BlockChain::copy_from"))
}
}
#[async_trait]
impl<State, T> AsyncHash for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance,
State::FE: AsType<ChainBlock> + for<'a> fs::FileSave<'a>,
T: Send + Sync,
{
async fn hash(&self, txn_id: TxnId) -> TCResult<Output<Sha256>> {
self.history.hash(txn_id).await
}
}
#[async_trait]
impl<State, T> Transact for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance,
State::FE: AsType<ChainBlock> + for<'a> fs::FileSave<'a>,
T: Transact + Send + Sync,
{
type Commit = T::Commit;
async fn commit(&self, txn_id: TxnId) -> Self::Commit {
debug!("BlockChain::commit");
self.history.write_ahead(txn_id).await;
let guard = self.subject.commit(txn_id).await;
self.history.commit(txn_id).await;
guard
}
async fn rollback(&self, txn_id: &TxnId) {
join!(self.subject.rollback(txn_id), self.history.rollback(txn_id));
}
async fn finalize(&self, txn_id: &TxnId) {
join!(self.subject.finalize(txn_id), self.history.finalize(txn_id));
}
}
#[async_trait]
impl<State, T> de::FromStream for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance
+ de::FromStream<Context = State::Txn>
+ From<Collection<State::Txn, State::FE>>
+ From<Scalar>,
State::FE: CacheBlock + for<'a> fs::FileSave<'a>,
T: Route<State> + de::FromStream<Context = State::Txn> + fmt::Debug,
(Bytes, Map<Tuple<State>>): TryCastFrom<State>,
Collection<State::Txn, State::FE>: TryCastFrom<State>,
Scalar: TryCastFrom<State>,
Value: TryCastFrom<State>,
(Value,): TryCastFrom<State>,
(Value, State): TryCastFrom<State>,
{
type Context = State::Txn;
async fn from_stream<D: de::Decoder>(
txn: State::Txn,
decoder: &mut D,
) -> Result<Self, D::Error> {
decoder.decode_seq(ChainVisitor::new(txn)).await
}
}
#[async_trait]
impl<'en, State, T> IntoView<'en, State::FE> for BlockChain<State, State::Txn, State::FE, T>
where
State: StateInstance,
State::FE: CacheBlock + for<'a> FileSave<'a>,
T: IntoView<'en, State::FE, Txn = State::Txn> + Send + Sync,
{
type Txn = State::Txn;
type View = (T::View, super::data::HistoryView<'en>);
async fn into_view(self, txn: Self::Txn) -> TCResult<Self::View> {
let history = self.history.into_view(txn.clone()).await?;
let subject = self.subject.into_view(txn).await?;
Ok((subject, history))
}
}
impl<State, Txn, FE, T> fmt::Debug for BlockChain<State, Txn, FE, T> {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(f, "BlockChain<{}>", std::any::type_name::<T>())
}
}
struct ChainVisitor<State, Txn, T> {
txn: Txn,
phantom: PhantomData<(State, T)>,
}
impl<State, Txn, T> ChainVisitor<State, Txn, T> {
fn new(txn: Txn) -> Self {
Self {
txn,
phantom: PhantomData,
}
}
}
#[async_trait]
impl<State, T> de::Visitor for ChainVisitor<State, State::Txn, T>
where
State: StateInstance
+ de::FromStream<Context = State::Txn>
+ From<Collection<State::Txn, State::FE>>
+ From<Scalar>,
State::FE: CacheBlock + for<'a> fs::FileSave<'a>,
T: Route<State> + de::FromStream<Context = State::Txn> + fmt::Debug,
(Bytes, Map<Tuple<State>>): TryCastFrom<State>,
Collection<State::Txn, State::FE>: TryCastFrom<State>,
Scalar: TryCastFrom<State>,
Value: TryCastFrom<State>,
(Value,): TryCastFrom<State>,
(Value, State): TryCastFrom<State>,
{
type Value = BlockChain<State, State::Txn, State::FE, T>;
fn expecting() -> &'static str {
"a BlockChain"
}
async fn visit_seq<A: de::SeqAccess>(self, mut seq: A) -> Result<Self::Value, A::Error> {
let subject = seq.next_element::<T>(self.txn.clone()).await?;
let subject = subject.ok_or_else(|| de::Error::invalid_length(0, "a BlockChain schema"))?;
let txn = self.txn.subcontext(HISTORY);
let history = seq
.next_element::<History<State, State::Txn, State::FE>>(txn)
.await?;
let history =
history.ok_or_else(|| de::Error::invalid_length(1, "a BlockChain history"))?;
let write_ahead_log = history.read_log().map_err(de::Error::custom).await?;
for (past_txn_id, mutations) in &write_ahead_log.mutations {
super::data::replay_all(&subject, past_txn_id, mutations, &self.txn, history.store())
.map_err(de::Error::custom)
.await?;
}
Ok(BlockChain::new(subject, history))
}
}