use std::collections::HashSet;
use std::time::Duration;
use nostr_sdk::prelude::*;
const NEG_CAP_TTL_SECS: u64 = 24 * 3600;
fn cap_key(relay_url: &str) -> String {
format!("neg_cap:{}", relay_url.trim_end_matches('/'))
}
fn now_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
pub fn neg_supported_cached(relay_url: &str) -> Option<bool> {
let raw = crate::db::get_sql_setting(cap_key(relay_url)).ok()??;
let (supported, checked_at) = parse_cap_entry(&raw)?;
(now_secs().saturating_sub(checked_at) < NEG_CAP_TTL_SECS).then_some(supported)
}
pub fn record_neg_support(relay_url: &str, supported: bool) {
let _ = crate::db::set_sql_setting(
cap_key(relay_url),
format!("{}:{}", u8::from(supported), now_secs()),
);
}
fn parse_cap_entry(raw: &str) -> Option<(bool, u64)> {
let (flag, ts) = raw.split_once(':')?;
let supported = match flag {
"1" => true,
"0" => false,
_ => return None,
};
Some((supported, ts.parse().ok()?))
}
pub async fn wait_connected(relay: &Relay, allowance: Duration) -> bool {
let deadline = tokio::time::Instant::now() + allowance;
loop {
match relay.status() {
RelayStatus::Connected => return true,
RelayStatus::Terminated | RelayStatus::Banned => return false,
_ => {}
}
if tokio::time::Instant::now() >= deadline {
return false;
}
tokio::time::sleep(Duration::from_millis(150)).await;
}
}
fn cursor_key(relay_url: &str) -> String {
format!("neg_cursor:{}", relay_url.trim_end_matches('/'))
}
pub fn reconcile_cursor(relay_url: &str) -> Option<u64> {
crate::db::get_sql_setting(cursor_key(relay_url)).ok()??.parse().ok()
}
pub fn advance_reconcile_cursor(relay_url: &str, anchor_secs: u64, session: &crate::state::SessionGuard) {
if !session.is_valid() {
return;
}
let _ = crate::db::advance_u64_setting(cursor_key(relay_url), anchor_secs);
}
pub fn is_transient_sync_error(err: &str) -> bool {
err == "timeout"
|| err.contains("not connected")
|| err.contains("transport dispatcher")
|| err.contains("lagged")
}
pub fn classify_neg_sync_error(err: &str, relay_was_connected: bool) -> Option<bool> {
if err.contains("negentropy not supported")
|| err.contains("unknown negentropy error")
|| (err.contains("negentropy") && err.contains("protocol version"))
{
return Some(false);
}
if err == "timeout" && relay_was_connected {
return Some(false);
}
None
}
pub async fn reconcile_missing(
filter: Filter,
local_items: Vec<(EventId, Timestamp)>,
timeout: Duration,
) -> Result<HashSet<EventId>, String> {
use futures_util::stream::{FuturesUnordered, StreamExt};
let client = crate::state::nostr_client().ok_or("Nostr client not initialized")?;
let opts = SyncOptions::new()
.direction(SyncDirection::Down)
.initial_timeout(timeout)
.dry_run();
let relay_map = client.relays().await;
let trusted = crate::state::active_trusted_relays().await;
let relays: Vec<(String, Relay)> = trusted.iter().filter_map(|url| {
if neg_supported_cached(url) == Some(false) {
crate::log_debug!("[Negentropy] {} skipped (cached: no NIP-77)", url);
return None;
}
let normalized = url.trim_end_matches('/');
relay_map.iter()
.find(|(u, _)| u.as_str().trim_end_matches('/') == normalized)
.map(|(_, r)| (url.to_string(), r.clone()))
}).collect();
drop(relay_map);
if relays.is_empty() {
crate::log_warn!("[Negentropy] No trusted relays available for reconciliation");
return Ok(HashSet::new());
}
let connect_allowance = crate::relay_request_timeout(Duration::from_secs(3)).min(timeout);
let mut futs = FuturesUnordered::new();
for (url, relay) in &relays {
let url = url.clone();
let relay = relay.clone();
let f = filter.clone();
let items = local_items.clone();
let o = opts.clone();
futs.push(async move {
if !wait_connected(&relay, connect_allowance).await {
return (url, None, false);
}
let r = tokio::time::timeout(timeout, relay.sync(f).items(items).opts(o)).await;
let connected = relay.status() == RelayStatus::Connected;
(url, Some(r), connected)
});
}
let session = crate::state::SessionGuard::capture();
let mut missing: HashSet<EventId> = HashSet::new();
while let Some((url, result, connected)) = futs.next().await {
let Some(result) = result else {
crate::log_debug!("[Negentropy] {} skipped: not connected", url);
continue;
};
match result {
Ok(Ok(recon)) => {
let n = recon.remote.len();
missing.extend(recon.remote);
crate::log_debug!("[Negentropy] {} reconciled: {} missing", url, n);
if session.is_valid() {
record_neg_support(&url, true);
}
}
Ok(Err(e)) => {
crate::log_warn!("[Negentropy] {} failed: {}", url, e);
if session.is_valid()
&& classify_neg_sync_error(&e.to_string(), connected) == Some(false)
{
crate::log_info!("[Negentropy] {} marked no-NIP-77 for 24h", url);
record_neg_support(&url, false);
}
}
Err(_) => crate::log_warn!("[Negentropy] {} timed out", url),
}
}
Ok(missing)
}
#[cfg(test)]
mod cap_tests {
use super::*;
#[test]
fn classify_detects_deterministic_refusals_regardless_of_connection() {
for err in [
"negentropy not supported",
"unknown negentropy error",
"negentropy: unsupported protocol version",
] {
assert_eq!(classify_neg_sync_error(err, true), Some(false), "{err}");
assert_eq!(classify_neg_sync_error(err, false), Some(false), "{err}");
}
}
#[test]
fn classify_timeout_only_counts_when_connected() {
assert_eq!(classify_neg_sync_error("timeout", true), Some(false));
assert_eq!(classify_neg_sync_error("timeout", false), None);
}
#[test]
fn classify_ignores_unrelated_errors() {
for err in [
"auth-required: we can't serve DMs to unauthenticated users",
"relay message too large: size=200000, max_size=131072",
"not connected",
"timeout exceeded", ] {
assert_eq!(classify_neg_sync_error(err, true), None, "{err}");
}
}
#[test]
fn transient_errors_never_skip_the_archive() {
assert!(is_transient_sync_error("timeout"));
assert!(is_transient_sync_error("relay not connected"));
assert!(is_transient_sync_error("not connected"));
assert!(is_transient_sync_error("can't send message to the transport dispatcher"));
assert!(is_transient_sync_error("channel lagged by 558"));
assert!(!is_transient_sync_error("negentropy not supported"));
assert!(!is_transient_sync_error("unknown negentropy error"));
assert!(!is_transient_sync_error("blocked: sync too big"));
}
#[test]
fn cap_entry_parses_and_rejects() {
assert_eq!(parse_cap_entry("1:1753900000"), Some((true, 1753900000)));
assert_eq!(parse_cap_entry("0:42"), Some((false, 42)));
assert_eq!(parse_cap_entry("2:42"), None);
assert_eq!(parse_cap_entry("1:"), None);
assert_eq!(parse_cap_entry("1"), None);
assert_eq!(parse_cap_entry("nonsense"), None);
assert_eq!(parse_cap_entry(""), None);
}
#[test]
fn cap_key_normalizes_trailing_slash() {
assert_eq!(cap_key("wss://r.example/"), cap_key("wss://r.example"));
assert_eq!(cursor_key("wss://r.example/"), cursor_key("wss://r.example"));
}
}