pub mod bandwidth;
pub(crate) mod cache;
pub mod metrics;
pub mod seed;
pub(crate) mod source;
pub mod testing;
pub mod transport;
pub(crate) mod verify;
pub use bandwidth::BandwidthGovernor;
pub use cache::CacheStats;
pub use metrics::{FetchObserver, FetchReport};
pub use seed::{derive_key, seed_blob, R2Target, SeedOutcome};
pub(crate) use source::BlobSource;
pub use source::{FetchProgress, FetchTier, ProgressSink};
use bytes::Bytes;
use std::{collections::HashMap, path::PathBuf, sync::Arc};
use tokio::sync::RwLock;
use crate::cache::Cache;
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error("invalid hash: {0}")]
InvalidHash(String),
#[error("all sources exhausted for {0}")]
FetchFailed(BlakeHash),
#[error("hash mismatch: expected {expected}, got {actual}")]
HashMismatch {
expected: BlakeHash,
actual: BlakeHash,
},
#[error("io: {0}")]
Io(#[from] std::io::Error),
}
pub type Result<T, E = Error> = std::result::Result<T, E>;
#[derive(Clone, Copy, PartialEq, Eq, Hash)]
pub struct BlakeHash([u8; 32]);
impl BlakeHash {
pub fn hash(data: &[u8]) -> Self {
Self(*blake3::hash(data).as_bytes())
}
pub fn from_bytes(bytes: [u8; 32]) -> Self {
Self(bytes)
}
pub fn from_hex(s: &str) -> Result<Self> {
let raw = hex::decode(s).map_err(|_| Error::InvalidHash(s.to_string()))?;
let arr: [u8; 32] = raw
.try_into()
.map_err(|_| Error::InvalidHash(format!("expected 32 bytes: {s}")))?;
Ok(Self(arr))
}
pub fn to_hex(&self) -> String {
hex::encode(self.0)
}
pub fn as_bytes(&self) -> &[u8; 32] {
&self.0
}
pub fn verify(&self, data: &[u8]) -> bool {
*self == Self::hash(data)
}
}
impl std::fmt::Debug for BlakeHash {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "BlakeHash({}…)", &self.to_hex()[..8])
}
}
impl std::fmt::Display for BlakeHash {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.to_hex())
}
}
impl std::str::FromStr for BlakeHash {
type Err = Error;
fn from_str(s: &str) -> Result<Self> {
Self::from_hex(s)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SeedRole {
Passive,
Participant,
Permanent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum PeerTier {
Cloud,
Camp,
Workstation,
Mobile,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BwCaps {
pub up_mbit: u32,
pub down_mbit: u32,
}
impl BwCaps {
pub fn passive() -> Self {
Self {
up_mbit: 0,
down_mbit: 50,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct BandwidthPolicy {
roles: HashMap<PeerTier, BwCaps>,
}
impl BandwidthPolicy {
pub fn role(mut self, tier: PeerTier, caps: BwCaps) -> Self {
self.roles.insert(tier, caps);
self
}
pub fn caps_for(&self, tier: PeerTier) -> BwCaps {
self.roles
.get(&tier)
.copied()
.unwrap_or_else(BwCaps::passive)
}
}
#[derive(Debug, Clone)]
pub struct Discovery {
pub lan: bool,
pub swarm: bool,
pub relays: Vec<String>,
}
impl Default for Discovery {
fn default() -> Self {
Self {
lan: true,
swarm: true,
relays: vec![],
}
}
}
impl Discovery {
pub fn with_relays(mut self, relays: impl IntoIterator<Item = impl Into<String>>) -> Self {
self.relays = relays.into_iter().map(Into::into).collect();
self
}
pub fn lan_only() -> Self {
Self {
lan: true,
swarm: false,
relays: vec![],
}
}
pub fn none() -> Self {
Self {
lan: false,
swarm: false,
relays: vec![],
}
}
}
pub struct AssetClassConfig {
pub name: &'static str,
pub permanent_seeds: Vec<String>,
pub cdn_fallback: Option<String>,
pub discovery: Discovery,
pub bandwidth: BandwidthPolicy,
pub cache_dir: Option<PathBuf>,
pub cache_budget_bytes: u64,
pub observer: Option<FetchObserver>,
}
impl Default for AssetClassConfig {
fn default() -> Self {
Self {
name: "default",
permanent_seeds: vec![],
cdn_fallback: None,
discovery: Discovery::default(),
bandwidth: BandwidthPolicy::default(),
cache_dir: None,
cache_budget_bytes: 1024 * 1024 * 1024,
observer: None,
}
}
}
struct ClassInner {
config: AssetClassConfig,
cache: Cache,
sources: RwLock<Vec<Arc<dyn BlobSource>>>,
role: RwLock<(SeedRole, PeerTier)>,
governor: BandwidthGovernor,
}
#[derive(Clone)]
pub struct AssetClass(Arc<ClassInner>);
impl AssetClass {
pub async fn register(config: AssetClassConfig) -> Result<Self> {
let cdn: Option<Arc<dyn BlobSource>> = if let Some(url) = &config.cdn_fallback {
match transport::http::HttpFetcher::new(url.as_str()) {
Ok(f) => Some(Arc::new(f)),
Err(e) => {
tracing::warn!(url = %url, "HttpFetcher init failed: {e}");
None
}
}
} else {
None
};
let initial_sources: Vec<Arc<dyn BlobSource>> = cdn.into_iter().collect();
let governor = BandwidthGovernor::new(config.bandwidth.clone());
governor.probe_os();
let cache = Cache::new(config.cache_dir.as_deref(), config.cache_budget_bytes)?;
Ok(Self(Arc::new(ClassInner {
config,
cache,
sources: RwLock::new(initial_sources),
role: RwLock::new((SeedRole::Participant, PeerTier::Workstation)),
governor,
})))
}
pub async fn set_role(&self, role: SeedRole, tier: PeerTier) -> Result<()> {
*self.0.role.write().await = (role, tier);
Ok(())
}
pub fn asset(&self, hash: BlakeHash) -> Asset {
Asset {
class: self.clone(),
hash,
}
}
pub fn name(&self) -> &str {
self.0.config.name
}
pub(crate) fn permanent_seeds(&self) -> &[String] {
&self.0.config.permanent_seeds
}
pub fn governor(&self) -> &BandwidthGovernor {
&self.0.governor
}
pub fn set_battery(&self, on_battery: bool) {
self.0.governor.set_battery(on_battery);
}
pub fn set_metered(&self, metered: bool) {
self.0.governor.set_metered(metered);
}
pub async fn cache_stats(&self) -> CacheStats {
self.0
.cache
.stats(self.0.config.cache_budget_bytes)
.await
}
pub(crate) fn add_source(&self, source: Arc<dyn BlobSource>) {
if let Ok(mut sources) = self.0.sources.try_write() {
sources.push(source);
} else {
tracing::warn!(
class = self.name(),
"add_source: lock contended, source dropped"
);
}
}
pub(crate) async fn fetch_bytes(&self, hash: BlakeHash) -> Result<Bytes> {
self.fetch_bytes_with_progress(hash, None).await
}
pub(crate) async fn fetch_bytes_reported(
&self,
hash: BlakeHash,
sink: Option<ProgressSink>,
) -> (Result<Bytes>, FetchReport) {
let started = std::time::Instant::now();
if let Some(bytes) = self.0.cache.get(&hash).await {
tracing::trace!(%hash, "cache hit");
let report = self.report(
hash,
Some(FetchTier::Cache),
None,
bytes.len() as u64,
started,
);
return (Ok(bytes), report);
}
let sources: Vec<Arc<dyn BlobSource>> = self.0.sources.read().await.clone();
let chain = source::FetchChain::new(sources);
match chain.fetch_with_progress(&hash, sink.as_ref()).await {
Some(hit) => {
if let Err(e) = self.0.cache.put(hash, hit.bytes.clone()).await {
tracing::warn!(%hash, "cache write failed: {e}");
}
let report = self.report(
hash,
Some(hit.tier),
hit.peer_id,
hit.bytes.len() as u64,
started,
);
(Ok(hit.bytes), report)
}
None => {
let report = self.report(hash, None, None, 0, started);
(Err(Error::FetchFailed(hash)), report)
}
}
}
fn report(
&self,
hash: BlakeHash,
tier: Option<FetchTier>,
peer_id: Option<String>,
bytes_served: u64,
started: std::time::Instant,
) -> FetchReport {
let report = FetchReport {
class: self.0.config.name,
hash,
tier,
peer_id,
bytes_served,
duration_ms: started.elapsed().as_millis() as u64,
};
if let Some(observer) = &self.0.config.observer {
observer(&report);
}
report
}
pub(crate) async fn fetch_bytes_with_progress(
&self,
hash: BlakeHash,
sink: Option<ProgressSink>,
) -> Result<Bytes> {
self.fetch_bytes_reported(hash, sink).await.0
}
}
#[derive(Clone)]
pub struct Asset {
class: AssetClass,
hash: BlakeHash,
}
impl Asset {
pub async fn fetch(&self) -> Result<Bytes> {
self.class.fetch_bytes(self.hash).await
}
pub async fn fetch_with_progress(&self, sink: ProgressSink) -> Result<Bytes> {
self.class
.fetch_bytes_with_progress(self.hash, Some(sink))
.await
}
pub async fn fetch_reported(&self) -> (Result<Bytes>, FetchReport) {
self.class.fetch_bytes_reported(self.hash, None).await
}
pub async fn is_cached(&self) -> bool {
self.class.0.cache.contains(&self.hash).await
}
pub fn hash(&self) -> BlakeHash {
self.hash
}
}