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
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions crates/bin/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,9 @@ tracing.workspace = true
[build-dependencies]
build-info-build.workspace = true

[dev-dependencies]
tempfile = "3"

[lints]
workspace = true

Expand Down
73 changes: 73 additions & 0 deletions crates/bin/src/cluster.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
use std::{
error::Error,
fs::{File, OpenOptions},
io,
path::Path,
};

use silver_common::Enr;
use silver_config::ClusterConfig;
use silver_control::cluster::{AttestationClusterConfig, ClusterStorageConfig};

pub struct ClusterStartup {
pub config: AttestationClusterConfig,
_journal_lock: File,
}

impl ClusterStartup {
pub fn new(
config: &ClusterConfig,
local: &Enr,
directory: &Path,
) -> Result<Self, Box<dyn Error>> {
let node_id = config
.nodes
.iter()
.find(|(_, enr)| enr.node_id() == local.node_id())
.map(|(id, _)| *id)
.ok_or("no local node configured in cluster config")?;
let journal_lock = Self::lock(directory)?;
let path = directory.join("attestation-cluster.wal");
let storage = if config.bootstrap {
ClusterStorageConfig::Create(path)
} else {
ClusterStorageConfig::Open(path)
};
Ok(Self {
config: AttestationClusterConfig::new(
node_id,
config.nodes.keys().copied().collect(),
storage,
),
_journal_lock: journal_lock,
})
}

fn lock(directory: &Path) -> io::Result<File> {
// Retain the lock file: removing it could let processes lock different inodes.
let file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(directory.join("attestation-cluster.lock"))?;
file.try_lock().map_err(|error| {
io::Error::other(format!("cannot lock attestation cluster journal: {error}"))
})?;
Ok(file)
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn journal_lock_is_exclusive_until_the_process_guard_is_dropped() {
let directory = tempfile::tempdir().unwrap();
let first = ClusterStartup::lock(directory.path()).unwrap();
assert!(ClusterStartup::lock(directory.path()).is_err());
drop(first);
let _second = ClusterStartup::lock(directory.path()).unwrap();
}
}
25 changes: 12 additions & 13 deletions crates/bin/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ use std::{
error::Error,
io,
net::IpAddr,
path::Path,
str::FromStr,
sync::Arc,
time::{Duration, Instant, SystemTime, UNIX_EPOCH},
Expand All @@ -28,14 +29,18 @@ use silver_common::{
tracing::initialise_tracing_log,
};
use silver_config::Config;
use silver_control::{Controller, cluster::AttestationClusterConfig, sync_engine::SyncEngine};
use silver_control::{Controller, sync_engine::SyncEngine};
use silver_discovery::{DiscV5, Discovery};
use silver_gossip::GossipHandler;
use silver_httpcore::Bind;
use silver_network::{ClusterNodes, Context, NetworkTile, P2p};
use silver_peer::PeerManager;
use silver_storage::{latest_local_checkpoint, tile::StorageTile};

use crate::cluster::ClusterStartup;

mod cluster;

#[cfg(not(feature = "alloc-profile"))]
#[global_allocator]
static GLOBAL: MiMalloc = MiMalloc;
Expand Down Expand Up @@ -134,19 +139,13 @@ fn main() -> Result<(), Box<dyn Error>> {
local_enr.set_attnets(subnets, keypair.secret_key())?;

// Cluster configuration
let (cluster_config, cluster_nodes) = config
let cluster_startup = config
.cluster_config()
.map(|c| {
let voters = c.nodes.clone();
let node_id = match voters.iter().find(|(_, enr)| enr.node_id() == local_enr.node_id())
{
Some((id, _)) => *id,
None => return Err("no local node configured in cluster config"),
};
Ok((AttestationClusterConfig::new(node_id, voters.keys().copied().collect()), voters))
.map(|cluster| {
ClusterStartup::new(cluster, &local_enr, Path::new(config.data_storage_dir()))
})
.transpose()?
.unzip();
.transpose()?;
let cluster_nodes = config.cluster_config().map(|cluster| cluster.nodes.clone());

let discv5_addr = config.discovery_bind_addr().expect("no discovery port");
let p2p_addr = config.p2p_bind_addr().expect("no p2p port");
Expand Down Expand Up @@ -268,7 +267,7 @@ fn main() -> Result<(), Box<dyn Error>> {
control_rpc_producer,
tcaches,
cluster_outbound_producer,
cluster_config,
cluster_startup.as_ref().map(|cluster| cluster.config.clone()),
SyncEngine::new(
config.syncing_config(),
booting_from_local_checkpoint,
Expand Down
20 changes: 20 additions & 0 deletions crates/config/src/cluster_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,4 +8,24 @@ pub struct ClusterConfig {
/// All cluster nodes, including the local node, keyed by cluster node id.
#[serde(default)]
pub nodes: HashMap<u64, Enr>,
/// Create a new journal exclusively. Disable on restart; recovery never
/// recreates missing journals. Never use this to replace a lost voter
/// journal. The journal lives at
/// `data_storage_dir/attestation-cluster.wal`; its parent directory must
/// exist.
#[serde(default)]
pub bootstrap: bool,
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn journal_creation_requires_explicit_bootstrap() {
let config: ClusterConfig = serde_yml::from_str("nodes: {}").unwrap();
assert!(!config.bootstrap);
let config: ClusterConfig = serde_yml::from_str("nodes: {}\nbootstrap: true").unwrap();
assert!(config.bootstrap);
}
}
1 change: 1 addition & 0 deletions crates/control/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ buffa.workspace = true
bytes.workspace = true
flux.workspace = true
raft.workspace = true
rand.workspace = true
silver_chain_spec.workspace = true
silver_common.workspace = true
silver_gossip.workspace = true
Expand Down
4 changes: 3 additions & 1 deletion crates/control/src/cluster/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ mod command;
mod generated;
mod lock_store;
mod node;
mod persistence;
#[cfg(target_os = "linux")]
mod storage;
mod wire;
Expand All @@ -19,6 +20,7 @@ pub use node::{
AttestationCluster, AttestationClusterConfig, AttestationDecision, ClusterError, ClusterEvent,
ProposalId, ProposeError,
};
pub use persistence::{ClusterStorageConfig, RecoveredStorage};
#[cfg(target_os = "linux")]
pub use storage::{ClusterStorage, ClusterStorageEvent, RecoveredStorage, StorageIdentity};
pub use storage::{ClusterStorage, ClusterStorageEvent, StorageIdentity};
pub(crate) use wire::{decode_message, encode_message};
Loading
Loading