use crate::error::{ConsensusError, Result};
use crate::validator::ValidatorSet;
use crate::voter::VOTE_FORMAT_VERSION;
use dashmap::DashMap;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::sync::Arc;
use tenzro_crypto::bls::{BlsKeyPair, BlsPublicKey, BlsSignature};
use tenzro_storage::{CF_AUDIT, CommitteeShape, KvStore};
use tenzro_types::primitives::{Address, Hash};
use tenzro_types::transaction::SignedTransaction;
mod bls_aggregate_serde {
use serde::{Deserialize, Deserializer, Serializer};
pub fn serialize<S: Serializer>(bytes: &[u8; 96], ser: S) -> Result<S::Ok, S::Error> {
ser.serialize_bytes(bytes)
}
pub fn deserialize<'de, D: Deserializer<'de>>(de: D) -> Result<[u8; 96], D::Error> {
let v: Vec<u8> = Vec::<u8>::deserialize(de)?;
if v.len() != 96 {
return Err(serde::de::Error::custom(format!(
"batch-cert bls_aggregate must be exactly 96 bytes, got {}",
v.len()
)));
}
let mut out = [0u8; 96];
out.copy_from_slice(&v);
Ok(out)
}
}
pub const ERASURE_ACTIVATION_N: usize = 100;
pub fn erasure_activation_threshold(n: usize) -> bool {
n >= ERASURE_ACTIVATION_N
}
const BATCH_ID_TAG: &[u8] = b"TENZRO_BATCH_ID:";
const BATCH_ACK_TAG: &[u8] = b"TENZRO_BATCH_ACK:";
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Batch {
pub id: Hash,
pub producer: Address,
pub sequence: u64,
pub transactions: Vec<SignedTransaction>,
}
impl Batch {
pub fn new(producer: Address, sequence: u64, transactions: Vec<SignedTransaction>) -> Self {
let id = Self::compute_id(&producer, sequence, &transactions);
Self {
id,
producer,
sequence,
transactions,
}
}
pub fn compute_id(
producer: &Address,
sequence: u64,
transactions: &[SignedTransaction],
) -> Hash {
let mut h = Sha256::new();
h.update(BATCH_ID_TAG);
h.update(producer.as_bytes());
h.update(sequence.to_le_bytes());
h.update((transactions.len() as u64).to_le_bytes());
for tx in transactions {
h.update(tx.transaction.hash().as_bytes());
}
let mut out = [0u8; 32];
out.copy_from_slice(&h.finalize());
Hash::new(out)
}
pub fn verify_id(&self) -> bool {
Self::compute_id(&self.producer, self.sequence, &self.transactions) == self.id
}
pub fn ack_payload(id: &Hash) -> Vec<u8> {
let mut payload = Vec::with_capacity(BATCH_ACK_TAG.len() + 1 + 32);
payload.extend_from_slice(BATCH_ACK_TAG);
payload.push(VOTE_FORMAT_VERSION);
payload.extend_from_slice(id.as_bytes());
payload
}
pub fn to_bytes(&self) -> Result<Vec<u8>> {
serde_json::to_vec(self)
.map_err(|e| ConsensusError::Internal(format!("batch serialize failed: {e}")))
}
pub fn from_bytes(bytes: &[u8]) -> Result<Self> {
let batch: Batch = serde_json::from_slice(bytes)
.map_err(|e| ConsensusError::Internal(format!("batch deserialize failed: {e}")))?;
if !batch.verify_id() {
return Err(ConsensusError::Internal(
"batch content id does not match bodies".to_string(),
));
}
Ok(batch)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BatchAck {
pub batch_id: Hash,
pub validator: Address,
pub bls_signature: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct BatchAvailabilityCertificate {
pub batch_id: Hash,
pub voting_power: u128,
#[serde(with = "bls_aggregate_serde")]
pub bls_aggregate: [u8; 96],
pub signer_bitmap: Vec<u8>,
}
impl BatchAvailabilityCertificate {
pub fn form(batch_id: Hash, acks: &[BatchAck], validator_set: &ValidatorSet) -> Result<Self> {
let active = validator_set.active_validators();
let n = active.len();
let bitmap_bytes = n.div_ceil(8);
let mut signer_bitmap = vec![0u8; bitmap_bytes];
let mut bls_sigs: Vec<BlsSignature> = Vec::new();
let mut voting_power: u128 = 0;
let normalized = validator_set.normalized_weights();
let mut signed_power: u128 = 0;
let mut seen = vec![false; n];
for ack in acks {
if ack.batch_id != batch_id {
continue;
}
let Some(idx) = validator_set.index_of(&ack.validator) else {
continue;
};
if seen[idx] {
continue;
}
let sig = BlsSignature::from_bytes(&ack.bls_signature).map_err(|e| {
ConsensusError::InvalidSignature(format!(
"batch ack from {} carries malformed BLS signature: {e}",
ack.validator
))
})?;
let pk = BlsPublicKey::from_bytes(&active[idx].bls_public_key).map_err(|e| {
ConsensusError::InvalidSignature(format!(
"validator {} has malformed BLS key: {e}",
ack.validator
))
})?;
let payload = Batch::ack_payload(&batch_id);
let ok = sig.verify(&pk, &payload).map_err(|e| {
ConsensusError::InvalidSignature(format!("batch ack verify raised: {e}"))
})?;
if !ok {
continue;
}
seen[idx] = true;
signer_bitmap[idx / 8] |= 1 << (idx % 8);
bls_sigs.push(sig);
voting_power = voting_power.saturating_add(active[idx].voting_power());
if let Some(w) = normalized.get(idx) {
signed_power = signed_power.saturating_add(*w);
}
}
let quorum_power = validator_set.quorum_voting_power();
if signed_power < quorum_power {
return Err(ConsensusError::InsufficientVotes {
got: signed_power.min(u64::MAX as u128) as u64,
need: quorum_power.min(u64::MAX as u128) as u64,
});
}
let agg = tenzro_crypto::bls::aggregate_signatures(&bls_sigs).map_err(|e| {
ConsensusError::InvalidSignature(format!("batch BLS aggregation failed: {e}"))
})?;
Ok(Self {
batch_id,
voting_power,
bls_aggregate: agg.to_bytes(),
signer_bitmap,
})
}
pub fn verify(&self, validator_set: &ValidatorSet) -> Result<()> {
let active = validator_set.active_validators();
let n = active.len();
let expected_bitmap_bytes = n.div_ceil(8);
if self.signer_bitmap.len() != expected_bitmap_bytes {
return Err(ConsensusError::InvalidSignature(format!(
"batch-cert signer_bitmap length {} does not match expected {} for {} validators",
self.signer_bitmap.len(),
expected_bitmap_bytes,
n
)));
}
let normalized = validator_set.normalized_weights();
let mut signer_pks: Vec<BlsPublicKey> = Vec::new();
let mut total_voting_power: u128 = 0;
let mut signed_power: u128 = 0;
#[allow(clippy::needless_range_loop)]
for bit_index in 0..(expected_bitmap_bytes * 8) {
let byte = self.signer_bitmap[bit_index / 8];
if (byte >> (bit_index % 8)) & 1 == 0 {
continue;
}
if bit_index >= n {
return Err(ConsensusError::InvalidSignature(format!(
"batch-cert bitmap bit {bit_index} set but active set has {n} validators"
)));
}
let validator = &active[bit_index];
let pk = BlsPublicKey::from_bytes(&validator.bls_public_key).map_err(|e| {
ConsensusError::InvalidSignature(format!(
"batch-cert signer {} has malformed BLS key: {e}",
validator.address
))
})?;
signer_pks.push(pk);
total_voting_power = total_voting_power.saturating_add(validator.voting_power());
if let Some(w) = normalized.get(bit_index) {
signed_power = signed_power.saturating_add(*w);
}
}
if signer_pks.is_empty() {
return Err(ConsensusError::InvalidSignature(
"batch-cert bitmap empty — no signers".to_string(),
));
}
let quorum_power = validator_set.quorum_voting_power();
if signed_power < quorum_power {
return Err(ConsensusError::InvalidSignature(format!(
"batch-cert carries {signed_power} stake-weight from {} signers, below quorum {quorum_power}",
signer_pks.len()
)));
}
if self.voting_power != total_voting_power {
return Err(ConsensusError::InvalidSignature(format!(
"batch-cert claims voting_power={} but bitmap-tallied power is {total_voting_power}",
self.voting_power
)));
}
let agg_pk = tenzro_crypto::bls::aggregate_public_keys(&signer_pks).map_err(|e| {
ConsensusError::InvalidSignature(format!(
"batch-cert aggregate public-key reconstruction failed: {e}"
))
})?;
let agg_pk_single = BlsPublicKey::from_bytes(&agg_pk.to_bytes()).map_err(|e| {
ConsensusError::InvalidSignature(format!(
"batch-cert aggregate public-key round-trip failed: {e}"
))
})?;
let agg_sig = BlsSignature::from_bytes(&self.bls_aggregate).map_err(|e| {
ConsensusError::InvalidSignature(format!(
"batch-cert bls_aggregate is not a valid signature: {e}"
))
})?;
let payload = Batch::ack_payload(&self.batch_id);
let ok = agg_sig.verify(&agg_pk_single, &payload).map_err(|e| {
ConsensusError::InvalidSignature(format!("batch-cert aggregate verify raised: {e}"))
})?;
if !ok {
return Err(ConsensusError::InvalidSignature(
"batch-cert aggregate verification rejected the signature".to_string(),
));
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ErasurePlan {
pub shape: Option<CommitteeShape>,
}
impl ErasurePlan {
pub fn for_set_size(n: usize) -> Self {
if !erasure_activation_threshold(n) {
return Self { shape: None };
}
let shape = CommitteeShape::from_committee_size(n).ok();
Self { shape }
}
pub fn is_coded(&self) -> bool {
self.shape.is_some()
}
}
pub struct BatchCertStore {
bls_key: Arc<BlsKeyPair>,
address: Address,
sequence: parking_lot::Mutex<u64>,
certs: DashMap<Hash, BatchAvailabilityCertificate>,
bodies: DashMap<Hash, Batch>,
pending_acks: DashMap<Hash, Vec<BatchAck>>,
storage: Option<Arc<dyn KvStore>>,
}
impl BatchCertStore {
fn cert_key(id: &Hash) -> Vec<u8> {
let mut k = b"batch_cert/".to_vec();
k.extend_from_slice(id.as_bytes());
k
}
fn body_key(id: &Hash) -> Vec<u8> {
let mut k = b"batch_body/".to_vec();
k.extend_from_slice(id.as_bytes());
k
}
pub fn new(bls_key: Arc<BlsKeyPair>, address: Address) -> Self {
Self {
bls_key,
address,
sequence: parking_lot::Mutex::new(0),
certs: DashMap::new(),
bodies: DashMap::new(),
pending_acks: DashMap::new(),
storage: None,
}
}
pub fn with_storage(
bls_key: Arc<BlsKeyPair>,
address: Address,
storage: Arc<dyn KvStore>,
) -> Self {
let store = Self {
bls_key,
address,
sequence: parking_lot::Mutex::new(0),
certs: DashMap::new(),
bodies: DashMap::new(),
pending_acks: DashMap::new(),
storage: Some(storage),
};
store.hydrate();
store
}
fn hydrate(&self) {
let Some(ref storage) = self.storage else {
return;
};
if let Ok(keys) = storage.get_keys_with_prefix(CF_AUDIT, b"batch_cert/") {
for k in keys {
if let Ok(Some(v)) = storage.get(CF_AUDIT, &k)
&& let Ok(cert) = serde_json::from_slice::<BatchAvailabilityCertificate>(&v)
{
self.certs.insert(cert.batch_id, cert);
}
}
}
if let Ok(keys) = storage.get_keys_with_prefix(CF_AUDIT, b"batch_body/") {
for k in keys {
if let Ok(Some(v)) = storage.get(CF_AUDIT, &k)
&& let Ok(batch) = Batch::from_bytes(&v)
{
self.bodies.insert(batch.id, batch);
}
}
}
}
pub fn produce(&self, transactions: Vec<SignedTransaction>) -> Batch {
let seq = {
let mut s = self.sequence.lock();
let cur = *s;
*s += 1;
cur
};
let batch = Batch::new(self.address, seq, transactions);
self.store_body(batch.clone());
batch
}
pub fn encode_for_broadcast(
&self,
batch: &Batch,
validator_count: usize,
) -> Result<Option<tenzro_storage::da::redstuff::EncodedBlob>> {
let plan = ErasurePlan::for_set_size(validator_count);
let Some(shape) = plan.shape else {
return Ok(None);
};
let bytes = batch
.to_bytes()
.map_err(|e| ConsensusError::Internal(format!("batch encode: {e}")))?;
let encoded = tenzro_storage::da::redstuff::encode(&bytes, shape).map_err(|e| {
ConsensusError::Internal(format!("erasure encode (n={validator_count}): {e}"))
})?;
Ok(Some(encoded))
}
pub fn sign_ack(&self, batch_id: &Hash) -> BatchAck {
let payload = Batch::ack_payload(batch_id);
let sig = self.bls_key.sign(&payload);
BatchAck {
batch_id: *batch_id,
validator: self.address,
bls_signature: sig.to_bytes().to_vec(),
}
}
pub fn record_ack(
&self,
batch_id: Hash,
ack: BatchAck,
validator_set: &ValidatorSet,
) -> Result<Option<BatchAvailabilityCertificate>> {
if ack.batch_id != batch_id {
return Ok(None);
}
if self.certs.contains_key(&batch_id) {
return Ok(None);
}
{
let mut entry = self.pending_acks.entry(batch_id).or_default();
if entry.iter().any(|a| a.validator == ack.validator) {
return Ok(None);
}
entry.push(ack);
}
let acks: Vec<BatchAck> = match self.pending_acks.get(&batch_id) {
Some(v) => v.clone(),
None => return Ok(None),
};
match BatchAvailabilityCertificate::form(batch_id, &acks, validator_set) {
Ok(cert) => {
self.install_cert(cert.clone(), validator_set)?;
self.pending_acks.remove(&batch_id);
Ok(Some(cert))
}
Err(ConsensusError::InsufficientVotes { .. }) => Ok(None),
Err(e) => Err(e),
}
}
pub fn store_body(&self, batch: Batch) {
if let Some(ref storage) = self.storage
&& let Ok(bytes) = batch.to_bytes()
{
let _ = storage.put(CF_AUDIT, &Self::body_key(&batch.id), &bytes);
}
self.bodies.insert(batch.id, batch);
}
pub fn install_cert(
&self,
cert: BatchAvailabilityCertificate,
validator_set: &ValidatorSet,
) -> Result<()> {
cert.verify(validator_set)?;
if let Some(ref storage) = self.storage
&& let Ok(bytes) = serde_json::to_vec(&cert)
{
let _ = storage.put(CF_AUDIT, &Self::cert_key(&cert.batch_id), &bytes);
}
self.certs.insert(cert.batch_id, cert);
Ok(())
}
pub fn get_cert(&self, id: &Hash) -> Option<BatchAvailabilityCertificate> {
self.certs.get(id).map(|c| c.clone())
}
pub fn get_body(&self, id: &Hash) -> Option<Batch> {
self.bodies.get(id).map(|b| b.clone())
}
pub fn has_body(&self, id: &Hash) -> bool {
self.bodies.contains_key(id)
}
pub fn evict(&self, id: &Hash) {
self.certs.remove(id);
self.bodies.remove(id);
self.pending_acks.remove(id);
if let Some(ref storage) = self.storage {
let _ = storage.delete(CF_AUDIT, &Self::cert_key(id));
let _ = storage.delete(CF_AUDIT, &Self::body_key(id));
}
}
pub fn cert_count(&self) -> usize {
self.certs.len()
}
pub fn evict_finalized(&self, finalized_tx_hashes: &std::collections::HashSet<Hash>) -> usize {
let to_evict: Vec<Hash> = self
.certs
.iter()
.filter_map(|entry| {
let id = *entry.key();
let body = self.bodies.get(&id)?;
let all_finalized = body
.transactions
.iter()
.all(|tx| finalized_tx_hashes.contains(&tx.transaction.hash()));
if all_finalized { Some(id) } else { None }
})
.collect();
for id in &to_evict {
self.evict(id);
}
to_evict.len()
}
pub fn certified_prefix(&self) -> Vec<Hash> {
let mut ordered: Vec<(Address, u64, Hash)> = self
.certs
.iter()
.filter_map(|entry| {
let id = *entry.key();
self.bodies.get(&id).map(|b| (b.producer, b.sequence, id))
})
.collect();
ordered.sort_by(|a, b| a.0.0.cmp(&b.0.0).then(a.1.cmp(&b.1)));
let certified: std::collections::HashSet<Hash> =
self.certs.iter().map(|e| *e.key()).collect();
let local: Vec<Hash> = ordered.into_iter().map(|(_, _, id)| id).collect();
agree_prefix(&local, &certified)
}
}
pub fn agree_prefix(local: &[Hash], certified: &std::collections::HashSet<Hash>) -> Vec<Hash> {
let mut prefix = Vec::new();
for id in local {
if certified.contains(id) {
prefix.push(*id);
} else {
break;
}
}
prefix
}
#[cfg(test)]
mod tests {
use super::*;
use crate::validator::{ValidatorInfo, ValidatorSet};
use std::collections::HashSet;
use tenzro_crypto::pq::MlDsaSigningKey;
use tenzro_crypto::{KeyPair, KeyType};
fn addr(byte: u8) -> Address {
Address::new([byte; 32])
}
fn make_validator(byte: u8, stake: u128) -> (ValidatorInfo, Arc<BlsKeyPair>) {
let kp = KeyPair::generate(KeyType::Ed25519).unwrap();
let pq = MlDsaSigningKey::generate();
let bls = BlsKeyPair::generate().unwrap();
let info = ValidatorInfo::new(
addr(byte),
kp.public_key().clone(),
pq.verifying_key_bytes().to_vec(),
bls.public_key().to_bytes().to_vec(),
stake,
);
(info, Arc::new(bls))
}
#[test]
fn batch_id_is_content_bound() {
let b1 = Batch::new(addr(1), 0, vec![]);
let b2 = Batch::new(addr(1), 1, vec![]);
assert_ne!(b1.id, b2.id, "sequence must disambiguate empty batches");
assert!(b1.verify_id());
assert!(b2.verify_id());
}
#[test]
fn erasure_plan_is_pure_function_of_n() {
assert!(!ErasurePlan::for_set_size(4).is_coded());
assert!(!ErasurePlan::for_set_size(50).is_coded());
assert!(!ErasurePlan::for_set_size(99).is_coded());
assert!(ErasurePlan::for_set_size(100).is_coded());
assert!(ErasurePlan::for_set_size(1000).is_coded());
let plan = ErasurePlan::for_set_size(1000);
let shape = plan.shape.unwrap();
assert_eq!(shape.f, 333);
assert_eq!(shape.quorum(), 667);
}
#[test]
fn availability_cert_forms_and_verifies_at_quorum() {
let mut infos = Vec::new();
let mut keys = Vec::new();
for i in 0..4u8 {
let (info, bls) = make_validator(i + 1, 1000);
infos.push(info);
keys.push(bls);
}
let vset = ValidatorSet::new(0, infos.clone()).unwrap();
let batch = Batch::new(addr(1), 0, vec![]);
let mut acks = Vec::new();
for i in 0..3usize {
let payload = Batch::ack_payload(&batch.id);
let sig = keys[i].sign(&payload);
acks.push(BatchAck {
batch_id: batch.id,
validator: infos[i].address,
bls_signature: sig.to_bytes().to_vec(),
});
}
let cert = BatchAvailabilityCertificate::form(batch.id, &acks, &vset).unwrap();
assert!(cert.verify(&vset).is_ok());
}
#[test]
fn availability_cert_rejects_below_quorum() {
let mut infos = Vec::new();
let mut keys = Vec::new();
for i in 0..4u8 {
let (info, bls) = make_validator(i + 1, 1000);
infos.push(info);
keys.push(bls);
}
let vset = ValidatorSet::new(0, infos.clone()).unwrap();
let batch = Batch::new(addr(1), 0, vec![]);
let mut acks = Vec::new();
for i in 0..2usize {
let payload = Batch::ack_payload(&batch.id);
let sig = keys[i].sign(&payload);
acks.push(BatchAck {
batch_id: batch.id,
validator: infos[i].address,
bls_signature: sig.to_bytes().to_vec(),
});
}
assert!(BatchAvailabilityCertificate::form(batch.id, &acks, &vset).is_err());
}
#[test]
fn prefix_stops_at_first_uncertified() {
let ids: Vec<Hash> = (0..5u8).map(|i| Hash::new([i; 32])).collect();
let mut certified = HashSet::new();
certified.insert(ids[0]);
certified.insert(ids[1]);
certified.insert(ids[3]); let prefix = agree_prefix(&ids, &certified);
assert_eq!(prefix, vec![ids[0], ids[1]]);
}
#[test]
fn record_ack_forms_cert_at_quorum() {
let mut infos = Vec::new();
let mut keys = Vec::new();
for i in 0..4u8 {
let (info, bls) = make_validator(i + 1, 1000);
infos.push(info);
keys.push(bls);
}
let vset = ValidatorSet::new(0, infos.clone()).unwrap();
let store = BatchCertStore::new(keys[0].clone(), infos[0].address);
let batch = store.produce(vec![]);
let ack0 = store.sign_ack(&batch.id);
assert!(store.record_ack(batch.id, ack0, &vset).unwrap().is_none());
let ack1 = BatchAck {
batch_id: batch.id,
validator: infos[1].address,
bls_signature: keys[1]
.sign(&Batch::ack_payload(&batch.id))
.to_bytes()
.to_vec(),
};
assert!(store.record_ack(batch.id, ack1, &vset).unwrap().is_none());
let ack2 = BatchAck {
batch_id: batch.id,
validator: infos[2].address,
bls_signature: keys[2]
.sign(&Batch::ack_payload(&batch.id))
.to_bytes()
.to_vec(),
};
let cert = store.record_ack(batch.id, ack2, &vset).unwrap();
assert!(cert.is_some(), "third ack should reach quorum");
assert!(store.get_cert(&batch.id).is_some());
assert!(cert.unwrap().verify(&vset).is_ok());
}
#[test]
fn record_ack_dedups_by_validator() {
let mut infos = Vec::new();
let mut keys = Vec::new();
for i in 0..4u8 {
let (info, bls) = make_validator(i + 1, 1000);
infos.push(info);
keys.push(bls);
}
let vset = ValidatorSet::new(0, infos.clone()).unwrap();
let store = BatchCertStore::new(keys[0].clone(), infos[0].address);
let batch = store.produce(vec![]);
let ack = store.sign_ack(&batch.id);
assert!(
store
.record_ack(batch.id, ack.clone(), &vset)
.unwrap()
.is_none()
);
assert!(store.record_ack(batch.id, ack, &vset).unwrap().is_none());
assert!(store.get_cert(&batch.id).is_none());
}
#[test]
fn store_produce_and_fetch_body() {
let (_info, bls) = make_validator(9, 1000);
let store = BatchCertStore::new(bls, addr(9));
let batch = store.produce(vec![]);
assert!(store.has_body(&batch.id));
let fetched = store.get_body(&batch.id).unwrap();
assert_eq!(fetched.id, batch.id);
store.evict(&batch.id);
assert!(!store.has_body(&batch.id));
}
}