vector_core/
negentropy.rs1use std::collections::HashSet;
8use std::time::Duration;
9
10use nostr_sdk::prelude::*;
11
12const NEG_CAP_TTL_SECS: u64 = 24 * 3600;
24
25fn cap_key(relay_url: &str) -> String {
26 format!("neg_cap:{}", relay_url.trim_end_matches('/'))
27}
28
29fn now_secs() -> u64 {
30 std::time::SystemTime::now()
31 .duration_since(std::time::UNIX_EPOCH)
32 .map(|d| d.as_secs())
33 .unwrap_or(0)
34}
35
36pub fn neg_supported_cached(relay_url: &str) -> Option<bool> {
39 let raw = crate::db::get_sql_setting(cap_key(relay_url)).ok()??;
40 let (supported, checked_at) = parse_cap_entry(&raw)?;
41 (now_secs().saturating_sub(checked_at) < NEG_CAP_TTL_SECS).then_some(supported)
42}
43
44pub fn record_neg_support(relay_url: &str, supported: bool) {
47 let _ = crate::db::set_sql_setting(
48 cap_key(relay_url),
49 format!("{}:{}", u8::from(supported), now_secs()),
50 );
51}
52
53fn parse_cap_entry(raw: &str) -> Option<(bool, u64)> {
54 let (flag, ts) = raw.split_once(':')?;
55 let supported = match flag {
56 "1" => true,
57 "0" => false,
58 _ => return None,
59 };
60 Some((supported, ts.parse().ok()?))
61}
62
63pub async fn wait_connected(relay: &Relay, allowance: Duration) -> bool {
75 let deadline = tokio::time::Instant::now() + allowance;
76 loop {
77 match relay.status() {
78 RelayStatus::Connected => return true,
79 RelayStatus::Terminated | RelayStatus::Banned => return false,
80 _ => {}
81 }
82 if tokio::time::Instant::now() >= deadline {
83 return false;
84 }
85 tokio::time::sleep(Duration::from_millis(150)).await;
86 }
87}
88
89fn cursor_key(relay_url: &str) -> String {
103 format!("neg_cursor:{}", relay_url.trim_end_matches('/'))
104}
105
106pub fn reconcile_cursor(relay_url: &str) -> Option<u64> {
108 crate::db::get_sql_setting(cursor_key(relay_url)).ok()??.parse().ok()
109}
110
111pub fn advance_reconcile_cursor(relay_url: &str, anchor_secs: u64) {
117 let _ = crate::db::advance_u64_setting(cursor_key(relay_url), anchor_secs);
118}
119
120pub fn is_transient_sync_error(err: &str) -> bool {
126 err == "timeout"
127 || err.contains("not connected")
128 || err.contains("transport dispatcher")
129 || err.contains("lagged")
132}
133
134pub fn classify_neg_sync_error(err: &str, relay_was_connected: bool) -> Option<bool> {
135 if err.contains("negentropy not supported")
136 || err.contains("unknown negentropy error")
137 || (err.contains("negentropy") && err.contains("protocol version"))
138 {
139 return Some(false);
140 }
141 if err == "timeout" && relay_was_connected {
142 return Some(false);
143 }
144 None
145}
146
147pub async fn reconcile_missing(
159 filter: Filter,
160 local_items: Vec<(EventId, Timestamp)>,
161 timeout: Duration,
162) -> Result<HashSet<EventId>, String> {
163 crate::db::scoped(async move {
164 use futures_util::stream::{FuturesUnordered, StreamExt};
165
166 let client = crate::state::nostr_client().ok_or("Nostr client not initialized")?;
167
168 let opts = SyncOptions::new()
169 .direction(SyncDirection::Down)
170 .initial_timeout(timeout)
171 .dry_run();
172
173 let relay_map = client.relays().await;
176 let trusted = crate::state::active_trusted_relays().await;
177 let relays: Vec<(String, Relay)> = trusted.iter().filter_map(|url| {
178 if neg_supported_cached(url) == Some(false) {
179 crate::log_debug!("[Negentropy] {} skipped (cached: no NIP-77)", url);
180 return None;
181 }
182 let normalized = url.trim_end_matches('/');
183 relay_map.iter()
184 .find(|(u, _)| u.as_str().trim_end_matches('/') == normalized)
185 .map(|(_, r)| (url.to_string(), r.clone()))
186 }).collect();
187 drop(relay_map);
188
189 if relays.is_empty() {
190 crate::log_warn!("[Negentropy] No trusted relays available for reconciliation");
191 return Ok(HashSet::new());
192 }
193
194 let connect_allowance = crate::relay_request_timeout(Duration::from_secs(3)).min(timeout);
195 let mut futs = FuturesUnordered::new();
196 for (url, relay) in &relays {
197 let url = url.clone();
198 let relay = relay.clone();
199 let f = filter.clone();
200 let items = local_items.clone();
201 let o = opts.clone();
202 futs.push(async move {
203 if !wait_connected(&relay, connect_allowance).await {
204 return (url, None, false);
205 }
206 let r = tokio::time::timeout(timeout, relay.sync(f).items(items).opts(o)).await;
207 let connected = relay.status() == RelayStatus::Connected;
208 (url, Some(r), connected)
209 });
210 }
211
212 let session = crate::db::current_session();
213 let mut missing: HashSet<EventId> = HashSet::new();
214 while let Some((url, result, connected)) = futs.next().await {
215 let Some(result) = result else {
216 crate::log_debug!("[Negentropy] {} skipped: not connected", url);
217 continue;
218 };
219 match result {
220 Ok(Ok(recon)) => {
221 let n = recon.remote.len();
222 missing.extend(recon.remote);
223 crate::log_debug!("[Negentropy] {} reconciled: {} missing", url, n);
224 record_neg_support(&url, true);
225 }
226 Ok(Err(e)) => {
227 crate::log_warn!("[Negentropy] {} failed: {}", url, e);
228 if session.is_live()
229 && classify_neg_sync_error(&e.to_string(), connected) == Some(false)
230 {
231 crate::log_info!("[Negentropy] {} marked no-NIP-77 for 24h", url);
232 record_neg_support(&url, false);
233 }
234 }
235 Err(_) => crate::log_warn!("[Negentropy] {} timed out", url),
236 }
237 }
238
239 Ok(missing)
240 })
241 .await
242}
243
244#[cfg(test)]
245mod cap_tests {
246 use super::*;
247
248 #[test]
249 fn classify_detects_deterministic_refusals_regardless_of_connection() {
250 for err in [
251 "negentropy not supported",
252 "unknown negentropy error",
253 "negentropy: unsupported protocol version",
254 ] {
255 assert_eq!(classify_neg_sync_error(err, true), Some(false), "{err}");
256 assert_eq!(classify_neg_sync_error(err, false), Some(false), "{err}");
257 }
258 }
259
260 #[test]
261 fn classify_timeout_only_counts_when_connected() {
262 assert_eq!(classify_neg_sync_error("timeout", true), Some(false));
263 assert_eq!(classify_neg_sync_error("timeout", false), None);
264 }
265
266 #[test]
267 fn classify_ignores_unrelated_errors() {
268 for err in [
269 "auth-required: we can't serve DMs to unauthenticated users",
270 "relay message too large: size=200000, max_size=131072",
271 "not connected",
272 "timeout exceeded", ] {
274 assert_eq!(classify_neg_sync_error(err, true), None, "{err}");
275 }
276 }
277
278 #[test]
279 fn transient_errors_never_skip_the_archive() {
280 assert!(is_transient_sync_error("timeout"));
281 assert!(is_transient_sync_error("relay not connected"));
282 assert!(is_transient_sync_error("not connected"));
283 assert!(is_transient_sync_error("can't send message to the transport dispatcher"));
284 assert!(is_transient_sync_error("channel lagged by 558"));
285 assert!(!is_transient_sync_error("negentropy not supported"));
286 assert!(!is_transient_sync_error("unknown negentropy error"));
287 assert!(!is_transient_sync_error("blocked: sync too big"));
288 }
289
290 #[test]
291 fn cap_entry_parses_and_rejects() {
292 assert_eq!(parse_cap_entry("1:1753900000"), Some((true, 1753900000)));
293 assert_eq!(parse_cap_entry("0:42"), Some((false, 42)));
294 assert_eq!(parse_cap_entry("2:42"), None);
295 assert_eq!(parse_cap_entry("1:"), None);
296 assert_eq!(parse_cap_entry("1"), None);
297 assert_eq!(parse_cap_entry("nonsense"), None);
298 assert_eq!(parse_cap_entry(""), None);
299 }
300
301 #[test]
302 fn cap_key_normalizes_trailing_slash() {
303 assert_eq!(cap_key("wss://r.example/"), cap_key("wss://r.example"));
304 assert_eq!(cursor_key("wss://r.example/"), cursor_key("wss://r.example"));
305 }
306
307}