use async_trait::async_trait;
use redb::{ReadableDatabase, ReadableTable, TableDefinition};
use crate::authority::{
AuthorityError, AuthorityId, AuthorityState, AuthorityStore, Drawn, Revocation,
StandingAuthority,
};
use crate::core::{EffectKey, Spend, StoreError, Timestamp};
use super::redb::{RedbStore, be, begin_write};
const TERMS: TableDefinition<(&str, &str), &str> = TableDefinition::new("authority_terms");
type BalanceRow<'a> = (u64, u64, u32, bool, i64, &'a str);
const BALANCE: TableDefinition<(&str, &str), BalanceRow<'static>> =
TableDefinition::new("authority_balance");
type ReceiptRow = (u64, u64, u64, u64, u32);
const RECEIPTS: TableDefinition<(&str, &str, &str), ReceiptRow> =
TableDefinition::new("authority_receipts");
#[async_trait]
impl AuthorityStore for RedbStore {
async fn issue(&self, authority: &StandingAuthority) -> Result<(), AuthorityError> {
authority.validate()?;
let tenant = self.tenant_name();
let id = authority.id.0.clone();
let terms = String::from_utf8(
crate::core::canon::to_bytes(authority)
.map_err(|e| AuthorityError::Unavailable(e.to_string()))?,
)
.map_err(|e| AuthorityError::Unavailable(e.to_string()))?;
let conflict: Result<(), ()> = self
.with_db(move |db| {
let w = begin_write(db)?;
let outcome = {
let mut table = w.open_table(TERMS).map_err(|e| be(&e))?;
let existing = table
.get((tenant.as_str(), id.as_str()))
.map_err(|e| be(&e))?
.map(|v| v.value().to_owned());
match existing {
Some(existing) if existing == terms => Ok(()),
Some(_) => Err(()),
None => {
table
.insert((tenant.as_str(), id.as_str()), terms.as_str())
.map_err(|e| be(&e))?;
let mut balance = w.open_table(BALANCE).map_err(|e| be(&e))?;
balance
.insert(
(tenant.as_str(), id.as_str()),
(0u64, 0u64, 0u32, false, 0i64, ""),
)
.map_err(|e| be(&e))?;
Ok(())
}
}
};
w.commit().map_err(|e| be(&e))?;
Ok(outcome)
})
.await?;
conflict.map_err(|()| AuthorityError::AlreadyIssued(authority.id.clone()))
}
async fn draw(
&self,
id: &AuthorityId,
key: EffectKey,
amount: Spend,
at: Timestamp,
) -> Result<Drawn, AuthorityError> {
let tenant = self.tenant_name();
let name = id.0.clone();
let key = key.to_hex();
let now = at.unix_timestamp();
let outcome: Result<Drawn, AuthorityError> = self
.with_db(move |db| {
let w = begin_write(db)?;
let decided = draw_in(&w, &tenant, &name, &key, amount, now);
w.commit().map_err(|e| be(&e))?;
Ok(decided)
})
.await?;
outcome
}
async fn revoke(
&self,
id: &AuthorityId,
reason: &str,
at: Timestamp,
) -> Result<(), AuthorityError> {
let tenant = self.tenant_name();
let name = id.0.clone();
let reason = reason.to_owned();
let now = at.unix_timestamp();
let found: Result<(), ()> = self
.with_db(move |db| {
let w = begin_write(db)?;
let outcome = {
let mut balance = w.open_table(BALANCE).map_err(|e| be(&e))?;
let row = balance
.get((tenant.as_str(), name.as_str()))
.map_err(|e| be(&e))?
.map(|v| {
let (tokens, minor, draws, revoked, _, _) = v.value();
(tokens, minor, draws, revoked)
});
match row {
None => Err(()),
Some((_, _, _, true)) => Ok(()),
Some((tokens, minor, draws, false)) => {
balance
.insert(
(tenant.as_str(), name.as_str()),
(tokens, minor, draws, true, now, reason.as_str()),
)
.map_err(|e| be(&e))?;
Ok(())
}
}
};
w.commit().map_err(|e| be(&e))?;
Ok(outcome)
})
.await?;
found.map_err(|()| AuthorityError::Unknown(id.clone()))
}
async fn state(&self, id: &AuthorityId) -> Result<Option<AuthorityState>, StoreError> {
let tenant = self.tenant_name();
let name = id.0.clone();
self.with_db(move |db| {
let r = db.begin_read().map_err(|e| be(&e))?;
let (Ok(terms), Ok(balance)) = (r.open_table(TERMS), r.open_table(BALANCE)) else {
return Ok(None);
};
let Some(raw) = terms
.get((tenant.as_str(), name.as_str()))
.map_err(|e| be(&e))?
else {
return Ok(None);
};
let authority: StandingAuthority = serde_json::from_str(raw.value())
.map_err(|e| StoreError::Backend(format!("stored authority is unreadable: {e}")))?;
let (tokens, minor, draws, revoked, revoked_at, reason) = balance
.get((tenant.as_str(), name.as_str()))
.map_err(|e| be(&e))?
.map_or((0, 0, 0, false, 0, String::new()), |v| {
let (t, m, d, rv, at, why) = v.value();
(t, m, d, rv, at, why.to_owned())
});
Ok(Some(AuthorityState {
authority,
drawn: Spend {
tokens,
minor_units: minor,
},
draws,
revoked: revoked.then(|| Revocation {
at: Timestamp::from_unix_timestamp(revoked_at).unwrap_or(Timestamp::UNIX_EPOCH),
reason,
}),
}))
})
.await
}
}
fn draw_in(
w: &redb::WriteTransaction,
tenant: &str,
name: &str,
key: &str,
amount: Spend,
now: i64,
) -> Result<Drawn, AuthorityError> {
let id = || AuthorityId::new(name);
let mut receipts = w.open_table(RECEIPTS).map_err(|e| be(&e))?;
if let Some(prior) = receipts
.get((tenant, name, key))
.map_err(|e| be(&e))?
.map(|v| v.value())
{
let (tokens, minor, rem_tokens, rem_minor, draws) = prior;
return Ok(Drawn {
authority: id(),
amount: Spend {
tokens,
minor_units: minor,
},
remaining: Spend {
tokens: rem_tokens,
minor_units: rem_minor,
},
draws,
});
}
let terms_table = w.open_table(TERMS).map_err(|e| be(&e))?;
let Some(raw) = terms_table.get((tenant, name)).map_err(|e| be(&e))? else {
return Err(AuthorityError::Unknown(id()));
};
let authority: StandingAuthority = serde_json::from_str(raw.value())
.map_err(|e| AuthorityError::Unavailable(format!("stored authority is unreadable: {e}")))?;
drop(raw);
let mut balance = w.open_table(BALANCE).map_err(|e| be(&e))?;
let (drawn_tokens, drawn_minor, draws, revoked, reason) = balance
.get((tenant, name))
.map_err(|e| be(&e))?
.map_or((0, 0, 0, false, String::new()), |v| {
let (t, m, d, rv, _, why) = v.value();
(t, m, d, rv, why.to_owned())
});
let remaining = crate::authority::permits(
&authority,
amount,
Spend {
tokens: drawn_tokens,
minor_units: drawn_minor,
},
draws,
revoked.then_some(reason.as_str()),
now,
)?;
let now_drawn = (
drawn_tokens.saturating_add(amount.tokens),
drawn_minor.saturating_add(amount.minor_units),
draws + 1,
);
balance
.insert(
(tenant, name),
(now_drawn.0, now_drawn.1, now_drawn.2, false, 0i64, ""),
)
.map_err(|e| be(&e))?;
let left = Spend {
tokens: remaining.tokens.saturating_sub(amount.tokens),
minor_units: remaining.minor_units.saturating_sub(amount.minor_units),
};
receipts
.insert(
(tenant, name, key),
(
amount.tokens,
amount.minor_units,
left.tokens,
left.minor_units,
now_drawn.2,
),
)
.map_err(|e| be(&e))?;
Ok(Drawn {
authority: id(),
amount,
remaining: left,
draws: now_drawn.2,
})
}