Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions crates/consensus-db/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,15 @@ impl DbMetrics {
pub fn observe_delete_time(&self, duration: Duration) {
self.delete_time.observe(duration.as_secs_f64());
}

/// Total number of write operations recorded so far.
///
/// Exposed for tests that assert each committed write transaction is counted
/// exactly once.
#[cfg(test)]
pub(crate) fn write_count(&self) -> u64 {
self.write_count.get()
}
}

impl Default for DbMetrics {
Expand Down
123 changes: 116 additions & 7 deletions crates/consensus-db/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -475,12 +475,16 @@ impl Db {
blocks.insert(height, block_bytes)?;
}

self.insert_certificate(
let certificate_bytes = self.insert_certificate(
&tx,
decided_block.certificate,
CommitCertificateType::Minimal,
Some(proposer),
)?;
#[allow(clippy::arithmetic_side_effects)]
{
write_bytes += certificate_bytes;
}

tx.commit()?;

Expand Down Expand Up @@ -546,7 +550,8 @@ impl Db {
}
}

self.insert_certificate(
let start = Instant::now();
let write_bytes = self.insert_certificate(
&tx,
certificate,
CommitCertificateType::Extended,
Expand All @@ -555,18 +560,28 @@ impl Db {

tx.commit()?;

// Measure through the commit: for redb the commit is where the durable
// write (fsync) cost lives, and `insert_decided_block` records its write
// the same way, so both paths feed `write_time` with the same scope.
self.update_write_metrics(write_bytes, start.elapsed());

Ok(())
}

/// Encode and insert `certificate` into the certificates table within the
/// caller's write transaction, returning the number of bytes written.
///
/// This intentionally does not record write metrics: the caller owns and
/// commits the transaction, so it records a single write observation once
/// the commit succeeds. That keeps a decided block (block + certificate
/// committed together) counted as one write instead of two.
fn insert_certificate(
&self,
tx: &WriteTransaction,
certificate: CommitCertificate<ArcContext>,
certificate_type: CommitCertificateType,
proposer: Option<Address>,
) -> Result<(), StoreError> {
let start = Instant::now();

) -> Result<usize, StoreError> {
let height = certificate.height;

let stored = StoredCommitCertificate {
Expand All @@ -582,9 +597,8 @@ impl Db {
let mut certificates = tx.open_table(CERTIFICATES_TABLE)?;
certificates.insert(height, encoded_certificate)?;
}
self.update_write_metrics(write_bytes, start.elapsed());

Ok(())
Ok(write_bytes)
}

/// Store misbehavior evidence for a given height.
Expand Down Expand Up @@ -1863,6 +1877,45 @@ mod tests {
assert_eq!(retrieved.execution_payload, retrieved_payload);
}

#[tokio::test]
async fn store_decided_block_counts_a_single_write() {
// Regression for #142: insert_decided_block observed the write metrics
// twice — once for the block and once inside insert_certificate — which
// double-counted the certificate write in the write_count and write_time
// metrics. A decided block is a single committed transaction, so it must
// be counted exactly once.
let dir = tempdir().unwrap();
let metrics = DbMetrics::default();
let store = Store::open(
dir.path().join("db"),
metrics.clone(),
DbUpgrade::Skip,
TEST_CACHE_SIZE,
)
.await
.unwrap();

let height = Height::new(1);
let round = Round::new(0);
let payload = arbitrary_payload();
let block_hash = payload.payload_inner.payload_inner.block_hash;
let value_id = ValueId::new(block_hash);
let cert = CommitCertificate::<ArcContext>::new(height, round, value_id, vec![]);
let proposer = Address::new([0u8; 20]);

let writes_before = metrics.write_count();
store
.store_decided_block(cert, payload, proposer)
.await
.unwrap();

assert_eq!(
metrics.write_count(),
writes_before + 1,
"a decided block is one committed transaction and must be counted once"
);
}

#[tokio::test]
async fn test_store_extended_certificate() {
use malachitebft_core_types::{NilOrVal, SignedMessage};
Expand Down Expand Up @@ -1929,6 +1982,62 @@ mod tests {
assert_eq!(retrieved.certificate.commit_signatures.len(), 4);
}

#[tokio::test]
async fn extend_certificate_counts_a_single_write() {
use malachitebft_core_types::{NilOrVal, SignedMessage};

// extend_certificate is the other path whose metric ownership this change
// touches (it now records the write itself instead of relying on
// insert_certificate). A certificate extension is one committed
// transaction, so it must be counted exactly once.
let dir = tempdir().unwrap();
let metrics = DbMetrics::default();
let store = Store::open(
dir.path().join("db"),
metrics.clone(),
DbUpgrade::Skip,
TEST_CACHE_SIZE,
)
.await
.unwrap();

let height = Height::new(1);
let round = Round::new(0);
let payload = arbitrary_payload();
let block_hash = payload.payload_inner.payload_inner.block_hash;
let value_id = ValueId::new(block_hash);

let signature = Signature::from_bytes([0xab; 64]);
let vote =
Vote::new_precommit(height, round, NilOrVal::Val(value_id), Address::new([1u8; 20]));
let cert = CommitCertificate::<ArcContext>::new(
height,
round,
value_id,
vec![SignedMessage::new(vote, signature)],
);

store
.store_decided_block(cert, payload, Address::new([0u8; 20]))
.await
.unwrap();

let mut stored = store.get_certificate(Some(height)).await.unwrap().unwrap();
stored.certificate.commit_signatures.push(CommitSignature::new(
Address::new([4u8; 20]),
Signature::from_bytes([0xcd; 64]),
));

let writes_before = metrics.write_count();
store.extend_certificate(stored.certificate).await.unwrap();

assert_eq!(
metrics.write_count(),
writes_before + 1,
"extending a certificate is one committed transaction and must be counted once"
);
}

#[tokio::test]
async fn test_extend_certificate_without_existing_fails() {
let store = create_store().await;
Expand Down