Skip to content
Merged
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
165 changes: 153 additions & 12 deletions magicblock-ledger/src/store/api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -669,20 +669,46 @@ impl Ledger {
encoded_transaction: &[u8],
status: TransactionStatusMeta,
) -> LedgerResult<()> {
// 1. Write Transaction Status
self.write_transaction_status(
slot,
index,
signature,
writable_keys,
readonly_keys,
status,
)?;
// Publish the transaction and every lookup index atomically. Readers
// must never find a signature before its bytes and status are available.
let status: generated::TransactionStatusMeta = status.into();
let address_count = (writable_keys.len() + readonly_keys.len()) as u64;
let mut batch = self.db.batch();
batch.put_bytes::<cf::Transaction>(
(signature, slot),
encoded_transaction,
);
batch.put_bytes::<cf::TransactionStatus>(
(signature, slot),
&status.encode_to_vec(),
);
batch.put::<cf::SlotSignatures>((slot, index), &signature)?;
for address in writable_keys {
batch.put::<cf::AddressSignatures>(
(*address, slot, index, signature),
&AddressSignatureMeta { writeable: true },
)?;
}
for address in readonly_keys {
batch.put::<cf::AddressSignatures>(
(*address, slot, index, signature),
&AddressSignatureMeta { writeable: false },
)?;
}
self.db.write(batch)?;

// 2. Write Transaction (raw bincode bytes)
self.transaction_cf
.put_bytes((signature, slot), encoded_transaction)?;
// Update cached counts only after the whole write succeeds.
self.transaction_cf.try_increase_entry_counter(1);
self.transaction_status_cf.try_increase_entry_counter(1);
self.slot_signatures_cf.try_increase_entry_counter(1);
self.address_signatures_cf
.try_increase_entry_counter(address_count);
let outcome_count = if status.err.is_none() {
&self.transaction_successful_status_count
} else {
&self.transaction_failed_status_count
};
try_increase_entry_counter(outcome_count, 1);

Ok(())
}
Expand Down Expand Up @@ -834,6 +860,8 @@ impl Ledger {
})
}

// Status-only fixtures are used by the focused index/status tests.
#[cfg(test)]
fn write_transaction_status(
&self,
slot: Slot,
Expand Down Expand Up @@ -1538,6 +1566,119 @@ mod tests {
}
}

#[test]
fn test_signature_index_only_exposes_complete_transactions() {
use std::sync::{atomic::AtomicBool, Barrier};

let ledger_path = get_tmp_ledger_path_auto_delete!();
let store = Arc::new(Ledger::open(ledger_path.path()).unwrap());
let address = Pubkey::new_unique();
let readonly =
(0..90).map(|_| Pubkey::new_unique()).collect::<Vec<_>>();
let keys = [vec![address], readonly.clone()].concat();
let transaction = VersionedTransaction {
signatures: vec![Signature::new_unique()],
message: VersionedMessage::V0(v0::Message {
header: MessageHeader {
num_required_signatures: 1,
num_readonly_unsigned_accounts: readonly.len() as u8,
..Default::default()
},
account_keys: keys.clone(),
..Default::default()
}),
};
let status = TransactionStatusMeta {
pre_balances: vec![1; keys.len()],
post_balances: vec![1; keys.len()],
pre_token_balances: Some(vec![]),
post_token_balances: Some(vec![]),
rewards: Some(vec![]),
..Default::default()
};
let done = AtomicBool::new(false);
let start = Barrier::new(2);
const TRANSACTIONS: u64 = 200;

// Initialize the cached counters before writing.
assert_eq!(store.count_transactions().unwrap(), 0);
assert_eq!(store.count_transaction_successful_status().unwrap(), 0);
assert_eq!(store.count_transaction_failed_status().unwrap(), 0);
assert_eq!(store.count_transaction_status().unwrap(), 0);
assert_eq!(store.count_address_signatures().unwrap(), 0);

std::thread::scope(|scope| {
scope.spawn(|| {
start.wait();
for slot in 1..=TRANSACTIONS {
let mut transaction = transaction.clone();
transaction.signatures[0] = Signature::new_unique();
let encoded = serialize(&transaction).unwrap();
store
.write_transaction(
transaction.signatures[0],
slot,
0,
vec![&address],
readonly.iter().collect(),
&encoded,
status.clone(),
)
.unwrap();
}
done.store(true, Ordering::Release);
});
start.wait();
let mut observed = 0;
loop {
let signatures = store
.get_confirmed_signatures_for_address(
address,
Slot::MAX,
None,
None,
1,
)
.unwrap();
for info in signatures.infos {
let complete = store
.get_complete_transaction(info.signature, Slot::MAX)
.unwrap();
assert!(
complete.is_some(),
"indexed signature {} at slot {} has no transaction",
info.signature,
info.slot,
);
assert_eq!(
complete.unwrap().tx_with_meta.get_status_meta(),
Some(status.clone()),
);
observed += 1;
}
if done.load(Ordering::Acquire) {
break;
}
}
assert!(observed > 0);
});

assert_eq!(store.count_transactions().unwrap(), TRANSACTIONS as i64);
assert_eq!(
store.count_transaction_status().unwrap(),
TRANSACTIONS as i64
);
assert_eq!(
store.count_address_signatures().unwrap(),
(TRANSACTIONS * keys.len() as u64) as i64,
);
assert_eq!(
store.count_transaction_successful_status().unwrap(),
TRANSACTIONS as i64,
);
assert_eq!(store.count_transaction_failed_status().unwrap(), 0);
}

#[test]
fn test_get_complete_transaction_by_signature() {
init_logger!();
Expand Down
Loading