Skip to content
Draft
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
918 changes: 436 additions & 482 deletions Cargo.lock

Large diffs are not rendered by default.

32 changes: 16 additions & 16 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -20,22 +20,22 @@ license = "MIT OR Apache-2.0"
alto-chain = { version = "2026.9.1", path = "chain" }
alto-client = { version = "2026.9.1", path = "client" }
alto-types = { version = "2026.9.1", path = "types" }
commonware-actor = "2026.9.0"
commonware-broadcast = "2026.9.0"
commonware-codec = "2026.9.0"
commonware-consensus = "2026.9.0"
commonware-cryptography = "2026.9.0"
commonware-deployer = { version = "2026.9.0", default-features = false }
commonware-formatting = "2026.9.0"
commonware-macros = "2026.9.0"
commonware-p2p = "2026.9.0"
commonware-resolver = "2026.9.0"
commonware-runtime = "2026.9.0"
commonware-storage = "2026.9.0"
commonware-stream = "2026.9.0"
commonware-utils = "2026.9.0"
commonware-math = "2026.9.0"
commonware-parallel = "2026.9.0"
commonware-actor = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-broadcast = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-codec = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-consensus = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-cryptography = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-deployer = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4", default-features = false }
commonware-formatting = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-macros = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-p2p = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-resolver = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-runtime = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-storage = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-stream = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-utils = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-math = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
commonware-parallel = { git = "https://github.com/commonwarexyz/monorepo", rev = "b35162940f119b93a65ecee8874987c3d8a1beb4" }
mimalloc = "0.1.52"
thiserror = "2.0.12"
bytes = "1.7.1"
Expand Down
41 changes: 30 additions & 11 deletions chain/src/application.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ use alto_types::{Block, Context, Scheme};
use commonware_actor::Feedback;
use commonware_consensus::{
marshal::{ancestry::Ancestry, Update},
Application as ConsensusApplication, Heightable, Reporter,
Application as ConsensusApplication, HandoffPolicy, Heightable, Reporter,
};
use commonware_cryptography::Digestible;
use commonware_runtime::{Clock, Metrics, Spawner, Storage};
Expand All @@ -28,16 +28,18 @@ pub struct Application<S: Scheme> {
backfiller: Option<indexer::Producer>,
delay_ms: u64,
block_size: usize,
handoff_policy: HandoffPolicy,
_scheme: PhantomData<S>,
}

impl<S: Scheme> Application<S> {
pub fn new(delay_ms: u64, block_size: u32) -> Self {
pub fn new(delay_ms: u64, block_size: u32, handoff_policy: HandoffPolicy) -> Self {
Self {
backfiller: None,
delay_ms,
block_size: usize::try_from(block_size)
.expect("configured block size is unsupported on this platform"),
handoff_policy,
_scheme: PhantomData,
}
}
Expand All @@ -58,6 +60,10 @@ where
type Block = Block;
type Input = ();

fn handoff_policy(&self, _context: &Self::Context) -> HandoffPolicy {
self.handoff_policy
}

async fn propose(
&mut self,
(mut runtime_context, context): (E, Self::Context),
Expand Down Expand Up @@ -174,12 +180,15 @@ mod tests {
use commonware_consensus::{
marshal::ancestry,
types::{Height, Round, View},
HandoffPublication,
};
use commonware_cryptography::{ed25519, sha256, Digest as _, Hasher, Sha256, Signer};
use commonware_runtime::{deterministic, Runner as _, Supervisor as _};
use std::sync::Arc;

const DELAY_MS: u64 = 10;
const PREPARE_HANDOFF: HandoffPolicy =
HandoffPolicy::Prepare(HandoffPublication::AfterCertification);

fn test_context(view: u64, parent: (View, sha256::Digest)) -> Context {
Context {
Expand Down Expand Up @@ -211,11 +220,21 @@ mod tests {
.expect("expected proposal")
}

#[test]
fn pipelines_handoffs_through_regular_proposer() {
let application = Application::<VrfScheme>::new(DELAY_MS, 0, PREPARE_HANDOFF);
let context = test_context(1, (View::zero(), sha256::Digest::EMPTY));
let policy =
ConsensusApplication::<deterministic::Context>::handoff_policy(&application, &context);

assert_eq!(policy, PREPARE_HANDOFF);
}

#[test]
fn verify_waits_until_future_block_enters_skew_window() {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let mut application = Application::new(DELAY_MS, 0);
let mut application = Application::new(DELAY_MS, 0, PREPARE_HANDOFF);

let now = context.current().epoch_millis();
let parent = Block::new(
Expand Down Expand Up @@ -262,7 +281,7 @@ mod tests {

// Timestamp validity does not depend on the local proposal delay.
for delay_ms in [0, DELAY_MS] {
let mut application = Application::new(delay_ms, 0);
let mut application = Application::new(delay_ms, 0, PREPARE_HANDOFF);
for (timestamp, valid) in [(now, false), (now + 1, true), (now + 2, true)] {
let block = Block::new(
test_context(2, (View::new(1), parent.digest())),
Expand All @@ -288,7 +307,7 @@ mod tests {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let block_size = 4;
let mut application = Application::new(DELAY_MS, block_size);
let mut application = Application::new(DELAY_MS, block_size, PREPARE_HANDOFF);

let now = context.current().epoch_millis();
let parent = Block::new(
Expand Down Expand Up @@ -328,7 +347,7 @@ mod tests {
fn verify_returns_immediately_for_mature_block_timestamp() {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let mut application = Application::new(DELAY_MS, 0);
let mut application = Application::new(DELAY_MS, 0, PREPARE_HANDOFF);

context.sleep(Duration::from_millis(10)).await;
let now = context.current().epoch_millis();
Expand Down Expand Up @@ -359,7 +378,7 @@ mod tests {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let delay_ms = 37;
let mut application = Application::new(delay_ms, 0);
let mut application = Application::new(delay_ms, 0, PREPARE_HANDOFF);

let now = context.current().epoch_millis();
let parent = Block::new(
Expand Down Expand Up @@ -389,7 +408,7 @@ mod tests {
for parent_timestamp in [0, 10, 15, MAX_BLOCK_TIMESTAMP_MS] {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let mut application = Application::new(0, 0);
let mut application = Application::new(0, 0, PREPARE_HANDOFF);

// Cover past, equal, and future parents, including the protocol's maximum.
context.sleep(Duration::from_millis(10)).await;
Expand Down Expand Up @@ -431,7 +450,7 @@ mod tests {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let block_size = 128;
let mut application = Application::new(DELAY_MS, block_size);
let mut application = Application::new(DELAY_MS, block_size, PREPARE_HANDOFF);

let now = context.current().epoch_millis();
let parent = Block::new(
Expand All @@ -458,7 +477,7 @@ mod tests {
fn verify_rejects_timestamp_above_maximum() {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let mut application = Application::new(DELAY_MS, 0);
let mut application = Application::new(DELAY_MS, 0, PREPARE_HANDOFF);

let now = context.current().epoch_millis();
let parent = Block::new(
Expand Down Expand Up @@ -489,7 +508,7 @@ mod tests {
fn propose_panics_when_parent_timestamp_is_maximum() {
let runner = deterministic::Runner::default();
runner.start(|context| async move {
let mut application = Application::new(DELAY_MS, 0);
let mut application = Application::new(DELAY_MS, 0, PREPARE_HANDOFF);

// Adding the proposal delay to a parent at the timestamp limit exceeds that limit.
let parent = Block::new(
Expand Down
51 changes: 42 additions & 9 deletions chain/src/engine.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
use crate::{
application::Application,
indexer::{self, Client},
HandoffMode,
};
use alto_types::{Activity, Block, Finalization, Scheme, EPOCH, EPOCH_LENGTH};
use commonware_broadcast::buffered;
Expand Down Expand Up @@ -39,6 +40,7 @@ use governor::clock::Clock as GClock;
use rand::{CryptoRng, Rng};
use std::{
num::{NonZero, NonZeroUsize},
sync::Arc,
time::{Duration, Instant},
};
use tracing::{error, info, warn};
Expand Down Expand Up @@ -69,13 +71,18 @@ const STABLE_LEADER_STALL_TIMEOUT: Duration = Duration::from_secs(12);
/// Round-robin leader election used with native standard certificates.
pub type StableElector = RoundRobin<Sha256>;

/// Builds stable leader election with a bounded optimistic view window.
/// Builds standard-certificate round-robin election.
pub fn stable_elector(term_length: NonZero<u32>, optimistic_views: u64) -> StableElector {
RoundRobin::default().with_term(
TermLength::new(term_length),
STABLE_LEADER_STALL_TIMEOUT,
ViewDelta::new(optimistic_views),
)
let elector = RoundRobin::default();
if term_length.get() == 1 {
elector
} else {
elector.with_term(
TermLength::new(term_length),
STABLE_LEADER_STALL_TIMEOUT,
ViewDelta::new(optimistic_views),
)
}
}

/// Configuration for the [Engine].
Expand All @@ -98,6 +105,7 @@ pub struct Config<
pub mailbox_size: usize,
pub deque_size: usize,
pub block_size: u32,
pub handoff_mode: HandoffMode,
/// Minimum interval from the parent's timestamp in milliseconds. Zero disables pacing.
pub proposal_delay_ms: u64,

Expand Down Expand Up @@ -138,7 +146,7 @@ where
Standard<Block>,
ConstantProvider<CS, Epoch>,
immutable::Archive<E, Digest, Finalization<CS>>,
immutable::Archive<E, Digest, Block>,
immutable::Archive<E, Digest, Arc<Block>>,
FixedEpocher,
S,
>,
Expand Down Expand Up @@ -290,7 +298,7 @@ where
.get()
.saturating_mul(SYNCER_ACTIVITY_TIMEOUT_MULTIPLIER),
),
start: marshal::Start::Genesis(genesis),
start: marshal::Start::Genesis(genesis.into()),
prunable_items_per_section: PRUNABLE_ITEMS_PER_SECTION,
replay_buffer: REPLAY_BUFFER,
key_write_buffer: WRITE_BUFFER,
Expand All @@ -307,7 +315,11 @@ where
// Create the application and, when an indexer is configured, a backfill
// queue of finalized digests so block uploads can resume after
// restarts.
let mut app = Application::new(proposal_delay_ms, cfg.block_size);
let mut app = Application::new(
proposal_delay_ms,
cfg.block_size,
cfg.handoff_mode.application_policy(),
);
let (pusher, consumer) = if let Some(indexer) = cfg.indexer {
let queue = queue::shared::init(
context.child("queue"),
Expand Down Expand Up @@ -515,4 +527,25 @@ mod tests {
assert_eq!(terms.stall_timeout(), Some(STABLE_LEADER_STALL_TIMEOUT));
assert_eq!(terms.optimistic_views(), ViewDelta::new(37));
}

#[test]
fn stable_elector_configures_single_view_terms() {
use commonware_consensus::types::{Round, View};

let Fixture { schemes, .. } =
standard::fixture::<MinSig, _>(&mut StdRng::seed_from_u64(1), NAMESPACE, 4);
let elector = elector::Config::<StandardScheme>::build(
stable_elector(NZU32!(1), 37),
schemes[0].participants(),
);
let terms = elector::Elector::terms(&elector);
assert_eq!(terms.length(), TermLength::ONE);
assert_eq!(terms.stall_timeout(), None);
assert_eq!(terms.optimistic_views(), ViewDelta::zero());
assert!(elector::Elector::elect_without_certificate(
&elector,
Round::new(Epoch::new(0), View::new(2))
)
.is_some());
}
}
3 changes: 2 additions & 1 deletion chain/src/indexer/backfiller/consumer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ use commonware_runtime::{
};
use commonware_storage::queue;
use commonware_utils::futures::{OptionFuture, Pool};
use std::{num::NonZeroUsize, time::Duration};
use std::{num::NonZeroUsize, sync::Arc, time::Duration};
use tracing::{debug, warn};

/// Final outcome for one backfill queue entry.
Expand Down Expand Up @@ -248,6 +248,7 @@ where
NextBlock::Ready(block) => return Some(*block),
NextBlock::FetchFromMarshal => {
if let Some(block) = marshal.get_block(Identifier::Digest(digest)).await {
let block = Arc::unwrap_or_clone(block);
uploads.lock().cache_block(block.clone());
return Some(block);
}
Expand Down
4 changes: 2 additions & 2 deletions chain/src/indexer/backfiller/state.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use alto_types::Block;
use bytes::{Buf, BufMut};
use commonware_codec::{self, FixedSize, Read, Write};
use bytes::BufMut;
use commonware_codec::{self, Buf, FixedSize, Read, Write};
use commonware_cryptography::{sha256::Digest, Digestible};
use commonware_utils::{sync::Mutex, PrioritySet};
use std::{collections::BTreeMap, sync::Arc};
Expand Down
Loading
Loading