use crate::error::DrandError;
use crate::traits::StorageEngine;
use drand_verify::{G2PubkeyRfc, Pubkey};
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tracing::warn;
use web_time::Duration;
#[cfg(not(target_arch = "wasm32"))]
use hickory_resolver::config::*;
const MAX_STALE_ROUNDS_FOR_HEARTBEAT: u64 = 200;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RawKyn {
#[serde(alias = "round")]
pub kyn: u64,
pub randomness: String,
#[serde(default)]
pub signature: String,
#[serde(default)]
pub is_from_cache: bool,
#[serde(default)]
pub is_unavailable: bool,
}
impl RawKyn {
pub fn unavailable() -> Self {
Self {
kyn: 0,
randomness: String::new(),
signature: String::new(),
is_from_cache: false,
is_unavailable: true,
}
}
pub fn is_usable_for_registration(&self) -> bool {
!self.is_unavailable && !self.is_from_cache
}
pub fn is_usable_for_heartbeat(&self, current_live_kyn: u64) -> bool {
if self.is_unavailable {
return false;
}
if !self.is_from_cache {
return true;
}
let staleness = current_live_kyn.saturating_sub(self.kyn);
staleness <= MAX_STALE_ROUNDS_FOR_HEARTBEAT
}
pub fn verify(&self) -> bool {
if self.is_unavailable {
return true;
}
if crate::config::is_dev_mode() {
return true;
}
let pubkey_bytes: [u8; 96] = match hex::decode(crate::constants::DRAND_PUBLIC_KEY)
.ok()
.and_then(|b| b.try_into().ok())
{
Some(b) => b,
None => return false,
};
let pk = match G2PubkeyRfc::from_fixed(pubkey_bytes) {
Ok(p) => p,
Err(_) => return false,
};
let sig_bytes = match hex::decode(&self.signature) {
Ok(b) => b,
Err(_) => return false,
};
if !pk.verify(self.kyn, &[], &sig_bytes).unwrap_or(false) {
return false;
}
use sha2::{Digest, Sha256};
let expected = Sha256::digest(&sig_bytes);
match hex::decode(&self.randomness) {
Ok(r) => r.as_slice() == expected.as_slice(),
Err(_) => false,
}
}
}
pub struct DrandClient {
http: reqwest::Client,
storage: Option<Arc<dyn StorageEngine>>,
endpoints: Vec<String>,
drand_domain: Vec<String>,
#[cfg(not(target_arch = "wasm32"))]
resolver: hickory_resolver::TokioAsyncResolver,
}
impl DrandClient {
pub fn new(storage: Option<Arc<dyn StorageEngine>>) -> Self {
let config = crate::config::KineticConfig::load();
Self {
#[cfg(not(target_arch = "wasm32"))]
http: reqwest::Client::builder()
.redirect(reqwest::redirect::Policy::none())
.build()
.unwrap_or_default(),
#[cfg(target_arch = "wasm32")]
http: reqwest::Client::new(),
storage,
endpoints: config.drand.endpoints,
drand_domain: config.drand.drand_domain,
#[cfg(not(target_arch = "wasm32"))]
resolver: hickory_resolver::TokioAsyncResolver::tokio(
ResolverConfig::default(),
ResolverOpts::default(),
),
}
}
pub async fn fetch_latest(&self) -> Result<RawKyn, DrandError> {
if crate::config::is_dev_mode() {
return self.load_cached_kyn();
}
let mut endpoints = self.endpoints.clone();
#[cfg(not(target_arch = "wasm32"))]
{
let mut injected_count = 0;
for domain in &self.drand_domain {
if injected_count >= 5 {
break;
}
if let Ok(txt_lookup) = self.resolver.txt_lookup(domain.as_str()).await {
for txt in txt_lookup.iter() {
if injected_count >= 5 {
break;
}
let url_str = txt.to_string();
let url_str = url_str.trim_matches('"').to_string();
if url_str.starts_with("https://") {
endpoints.push(url_str);
injected_count += 1;
}
}
}
}
}
let mut last_error = None;
for endpoint in &endpoints {
match self.fetch_with_backoff(endpoint).await {
Ok(mut kyn) => {
if !kyn.verify() {
warn!(
"Drand endpoint {} returned a cryptographically invalid kyn!",
endpoint
);
last_error = Some(DrandError::InvalidSignature);
continue;
}
let now = web_time::SystemTime::now()
.duration_since(web_time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let estimated_kyn = (now
.saturating_sub(crate::constants::DRAND_GENESIS_TIME))
/ crate::constants::DRAND_PERIOD;
let age = estimated_kyn.saturating_sub(kyn.kyn);
if age > MAX_STALE_ROUNDS_FOR_HEARTBEAT {
warn!(
"Drand endpoint {} returned an unacceptably stale kyn (kyn {}, expected ~{}).",
endpoint, kyn.kyn, estimated_kyn
);
last_error = Some(DrandError::StaleKyn {
expected: estimated_kyn,
got: kyn.kyn,
});
continue;
}
kyn.is_from_cache = false;
kyn.is_unavailable = false;
let _ = self.cache_kyn(&kyn);
return Ok(kyn);
}
Err(e) => {
warn!("Drand endpoint {} unreachable: {}", endpoint, e);
last_error = Some(e);
}
}
}
warn!("All Drand endpoints unreachable — falling back to cached kyn");
self.load_cached_kyn()
.map_err(|_| last_error.unwrap_or(DrandError::AllEndpointsFailed))
}
async fn fetch_with_backoff(&self, url: &str) -> Result<RawKyn, DrandError> {
let mut delay = Duration::from_millis(500);
let max_attempts = 3;
for attempt in 0..max_attempts {
match self
.http
.get(url)
.timeout(Duration::from_secs(5))
.send()
.await
{
Ok(mut resp) if resp.status().is_success() => {
#[cfg(target_arch = "wasm32")]
{
let bytes = resp
.bytes()
.await
.map_err(|e| DrandError::Network(e.to_string()))?;
if bytes.len() > 64 * 1024 {
return Err(DrandError::Network(
"Drand response exceeded 64 KB limit".to_string(),
));
}
return Ok(serde_json::from_slice::<RawKyn>(&bytes)?);
}
#[cfg(not(target_arch = "wasm32"))]
{
let mut body = bytes::BytesMut::new();
while let Some(chunk) = resp
.chunk()
.await
.map_err(|e| DrandError::Network(e.to_string()))?
{
body.extend_from_slice(&chunk);
if body.len() > 64 * 1024 {
return Err(DrandError::Network(
"Drand response exceeded 64 KB limit".to_string(),
));
}
}
return Ok(serde_json::from_slice::<RawKyn>(&body)?);
}
}
Ok(_resp) if attempt < max_attempts - 1 => {
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(delay).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(delay).await;
delay *= 2;
}
Ok(resp) => {
return Err(DrandError::HttpError(resp.status().as_u16()));
}
Err(_) if attempt < max_attempts - 1 => {
#[cfg(not(target_arch = "wasm32"))]
tokio::time::sleep(delay).await;
#[cfg(target_arch = "wasm32")]
gloo_timers::future::sleep(delay).await;
delay *= 2; }
Err(e) => return Err(DrandError::Network(e.to_string())),
}
}
Err(DrandError::AllEndpointsFailed)
}
pub fn cache_kyn(&self, kyn: &RawKyn) -> Result<(), DrandError> {
if let Some(storage) = &self.storage {
let bytes = serde_json::to_vec(kyn)?;
storage.put(crate::constants::DB_PREFIX_LAST_DRAND, &bytes)?;
}
Ok(())
}
pub fn load_cached_kyn(&self) -> Result<RawKyn, DrandError> {
if let Some(storage) = &self.storage {
if let Ok(Some(bytes)) = storage.get(crate::constants::DB_PREFIX_LAST_DRAND) {
if let Ok(mut kyn) = serde_json::from_slice::<RawKyn>(&bytes) {
kyn.is_from_cache = true;
return Ok(kyn);
}
}
}
if crate::config::is_dev_mode() {
tracing::warn!("DEV MODE: Returning mock drand kyn because cache is empty.");
return Ok(RawKyn {
kyn: 5000000,
randomness: "mock_randomness".to_string(),
signature: String::new(),
is_from_cache: true,
is_unavailable: false,
});
}
Err(DrandError::NoCachedKyn)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_valid_quicknet_kyn_verification() {
let kyn = RawKyn {
kyn: 30290678,
randomness: "bd5f53ad61578f2566860e3792d01513b817e34c7de92f4781aa76b53ddef0ea".to_string(),
signature: "ac8313d3ad1f95fe1b380ab6124aade0d4de5919fd60dc846746025ac9aa9d3c434b9dc94c0b75c4efd81aec9e2ef0b9".to_string(),
is_from_cache: false,
is_unavailable: false,
};
assert!(
kyn.verify(),
"Valid Quicknet kyn failed BLS verification"
);
}
#[test]
fn test_invalid_quicknet_kyn_verification() {
let kyn = RawKyn {
kyn: 30290678,
randomness: "bd5f53ad61578f2566860e3792d01513b817e34c7de92f4781aa76b53ddef0ea".to_string(),
signature: "bc8313d3ad1f95fe1b380ab6124aade0d4de5919fd60dc846746025ac9aa9d3c434b9dc94c0b75c4efd81aec9e2ef0b9".to_string(), is_from_cache: false,
is_unavailable: false,
};
if crate::config::is_dev_mode() {
assert!(kyn.verify(), "Dev mode should always pass verification");
} else {
assert!(
!kyn.verify(),
"Invalid Quicknet kyn incorrectly passed BLS verification"
);
}
}
}