vector_core/community/v2/
streamauth.rs1use nostr_sdk::prelude::FinalizeEvent;
24use std::collections::HashMap;
25use std::sync::atomic::{AtomicBool, Ordering};
26use std::sync::{LazyLock, Mutex};
27
28use nostr_sdk::prelude::{Client, ClientMessage, ClientNotification, EventBuilder, Keys, RelayUrl, StreamExt};
29
30use super::community::CommunityV2;
31
32static REGISTRY: LazyLock<Mutex<HashMap<[u8; 32], Keys>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
34static RESPONDER_RUNNING: AtomicBool = AtomicBool::new(false);
36static CHALLENGES: LazyLock<Mutex<HashMap<RelayUrl, String>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
43
44pub fn register(keys: impl IntoIterator<Item = Keys>) -> usize {
46 let mut reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
47 let mut added = 0;
48 for k in keys {
49 if reg.insert(k.public_key().to_bytes(), k).is_none() {
50 added += 1;
51 }
52 }
53 added
54}
55
56pub fn register_community(c: &CommunityV2) -> usize {
61 let mut keys: Vec<Keys> = vec![
62 super::derive::control_group_key(&c.community_root, c.id(), c.root_epoch).keys().clone(),
63 super::derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch).keys().clone(),
64 super::derive::dissolved_group_key(c.id()).keys().clone(),
65 ];
66 for ch in &c.channels {
67 if ch.private && ch.key.is_none() {
68 continue;
69 }
70 let (secret, epoch) = c.channel_secret(ch);
71 keys.push(super::derive::channel_group_key(&secret, &ch.id, epoch).keys().clone());
72 }
73 for pk_keys in rekey_plane_keys(c) {
75 keys.push(pk_keys);
76 }
77 register(keys)
78}
79
80fn rekey_plane_keys(c: &CommunityV2) -> Vec<Keys> {
92 use crate::community::Epoch;
93 let mut out = vec![super::derive::base_rekey_group_key(&c.community_root, c.id(), Epoch(c.root_epoch.0.saturating_add(1))).keys().clone()];
94 let cid_hex = crate::simd::hex::bytes_to_hex_32(&c.id().0);
95 let roots = super::service::channel_rekey_addressing_roots(c.community_root, &cid_hex);
96 for ch in &c.channels {
97 if ch.private {
98 for root in &roots {
99 out.push(super::derive::channel_rekey_group_key(root, &ch.id, Epoch(ch.epoch.0.saturating_add(1))).keys().clone());
100 }
101 }
102 }
103 out
104}
105
106fn remember_challenge(relay: &RelayUrl, challenge: &str) -> bool {
112 CHALLENGES
113 .lock()
114 .unwrap_or_else(|e| e.into_inner())
115 .insert(relay.clone(), challenge.to_string())
116 .as_deref()
117 != Some(challenge)
118}
119
120fn remembered_challenge(relay: &RelayUrl) -> Option<String> {
122 CHALLENGES.lock().unwrap_or_else(|e| e.into_inner()).get(relay).cloned()
123}
124
125static RESUB_AT: LazyLock<Mutex<HashMap<RelayUrl, std::time::Instant>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
128const RESUB_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(30);
129
130fn resub_cooldown_elapsed(relay: &RelayUrl) -> bool {
133 let mut map = RESUB_AT.lock().unwrap_or_else(|e| e.into_inner());
134 let now = std::time::Instant::now();
135 match map.get(relay) {
136 Some(at) if now.duration_since(*at) < RESUB_COOLDOWN => false,
137 _ => {
138 map.insert(relay.clone(), now);
139 true
140 }
141 }
142}
143
144pub fn clear() {
147 REGISTRY.lock().unwrap_or_else(|e| e.into_inner()).clear();
148 CHALLENGES.lock().unwrap_or_else(|e| e.into_inner()).clear();
149 RESUB_AT.lock().unwrap_or_else(|e| e.into_inner()).clear();
150 RESPONDER_RUNNING.store(false, Ordering::SeqCst);
151}
152
153pub fn is_empty() -> bool {
155 REGISTRY.lock().unwrap_or_else(|e| e.into_inner()).is_empty()
156}
157
158fn sign_all(challenge: &str, relay: &RelayUrl) -> Vec<nostr_sdk::prelude::Event> {
161 let reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
162 reg.values()
163 .filter_map(|keys| EventBuilder::auth(challenge, relay.clone()).finalize(keys).ok())
164 .collect()
165}
166
167async fn authenticate_streams(client: &Client, relay: &RelayUrl, challenge: &str) -> usize {
171 let events = sign_all(challenge, relay);
172 let Ok(Some(r)) = client.relay(relay.clone()).await else {
173 return 0;
174 };
175 let mut sent = 0;
176 for ev in events {
177 if r.send_msg(ClientMessage::auth(ev)).await.is_ok() {
178 sent += 1;
179 }
180 }
181 sent
182}
183
184pub fn ensure_responder(client: &Client) {
194 if RESPONDER_RUNNING.swap(true, Ordering::SeqCst) {
195 return; }
197 let client = client.clone();
198 let session = crate::state::SessionGuard::capture();
199 tokio::spawn(async move {
200 let mut notifications = client.notifications();
201 while let Some(n) = notifications.next().await {
202 if !session.is_valid() {
203 break;
204 }
205 if let ClientNotification::Message { relay_url, message } = n {
206 if let nostr_sdk::prelude::RelayMessage::Auth { challenge } = *message {
207 let challenge = challenge.into_owned();
208 let gate_key = format!(
213 "auth_gate:{}",
214 crate::inbox_relays::normalize_relay_url(relay_url.as_str())
215 );
216 if crate::db::get_sql_setting(gate_key.clone()).ok().flatten().is_none() {
217 let _ = crate::db::set_sql_setting(gate_key, "1".to_string());
218 }
219 let fresh_connection = remember_challenge(&relay_url, &challenge);
220 if !is_empty() {
221 authenticate_streams(&client, &relay_url, &challenge).await;
222 if fresh_connection && resub_cooldown_elapsed(&relay_url) {
228 super::realtime::resubscribe_relay(&client, &relay_url).await;
229 }
230 }
231 }
232 }
233 }
234 RESPONDER_RUNNING.store(false, Ordering::SeqCst);
235 });
236}
237
238pub fn prime(community: &CommunityV2) {
242 register_community(community);
243 if let Some(client) = crate::state::nostr_client() {
244 ensure_responder(&client);
245 }
246}
247
248pub async fn prime_auth(client: &Client, relays: &[String]) {
261 if is_empty() || relays.is_empty() {
262 return;
263 }
264 ensure_responder(client);
265 let urls: Vec<RelayUrl> = relays.iter().filter_map(|r| RelayUrl::parse(r).ok()).collect();
266 if urls.is_empty() {
267 return;
268 }
269 for url in &urls {
270 if let Some(challenge) = remembered_challenge(url) {
271 authenticate_streams(client, url, &challenge).await;
272 }
273 }
274 let authors: Vec<nostr_sdk::prelude::PublicKey> = {
275 let reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
276 reg.keys().filter_map(|pk| nostr_sdk::prelude::PublicKey::from_slice(pk).ok()).collect()
277 };
278 if authors.is_empty() {
279 return;
280 }
281 let filter = nostr_sdk::prelude::Filter::new()
282 .kind(nostr_sdk::prelude::Kind::Custom(super::stream::KIND_WRAP))
283 .authors(authors)
284 .limit(1);
285 let _ = tokio::time::timeout(std::time::Duration::from_secs(8), client
287 .fetch_events(nostr_sdk::prelude::ReqTarget::manual(
288 urls.into_iter().map(|u| (u, vec![filter.clone()])),
289 ))
290 .timeout(std::time::Duration::from_secs(6))).await;
291}
292
293#[cfg(test)]
294mod tests {
295 use super::*;
296 use crate::community::v2::control::{genesis, CommunityMetadata};
297 use crate::community::v2::community::{ChannelV2, CommunityV2};
298 use crate::community::{ChannelId, Epoch};
299 use nostr_sdk::prelude::Keys;
300
301 fn community_with_channel_mix() -> CommunityV2 {
304 let owner = Keys::generate();
305 let g = genesis(&owner, CommunityMetadata { name: "auth-test".into(), ..Default::default() }, 1_000).unwrap();
306 let mut c = CommunityV2::from_genesis(&g, "auth-test", None, vec!["wss://gated.example".into()], 0);
307 c.channels.push(ChannelV2 { id: ChannelId([2u8; 32]), name: "keyed-private".into(), private: true, key: Some([7u8; 32]), epoch: Epoch(3), voice: None, meta_custom: None, meta_extra: Default::default() });
308 c.channels.push(ChannelV2 { id: ChannelId([3u8; 32]), name: "keyless-private".into(), private: true, key: None, epoch: Epoch(0), voice: None, meta_custom: None, meta_extra: Default::default() });
309 c
310 }
311
312 fn registered(pk: &nostr_sdk::prelude::PublicKey) -> bool {
313 REGISTRY.lock().unwrap_or_else(|e| e.into_inner()).contains_key(&pk.to_bytes())
314 }
315
316 #[test]
320 fn register_community_covers_planes_and_skips_keyless() {
321 let c = community_with_channel_mix();
322 let added = register_community(&c);
323 assert!(added >= 6, "control+guestbook+dissolved+public chat+keyed chat+rekeys = at least 6 new keys, got {added}");
324
325 use super::super::derive;
326 assert!(registered(&derive::control_group_key(&c.community_root, c.id(), c.root_epoch).pk()));
327 assert!(registered(&derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch).pk()));
328 assert!(registered(&derive::dissolved_group_key(c.id()).pk()));
329 let public = &c.channels[0];
331 let (secret, epoch) = c.channel_secret(public);
332 assert!(registered(&derive::channel_group_key(&secret, &public.id, epoch).pk()));
333 let keyed = &c.channels[1];
335 assert!(registered(&derive::channel_group_key(&[7u8; 32], &keyed.id, Epoch(3)).pk()));
336 let keyless = &c.channels[2];
339 assert!(!registered(&derive::channel_group_key(&c.community_root, &keyless.id, Epoch(0)).pk()));
340 assert!(registered(&derive::base_rekey_group_key(&c.community_root, c.id(), Epoch(c.root_epoch.0 + 1)).pk()));
342 assert!(registered(&derive::channel_rekey_group_key(&c.community_root, &keyed.id, Epoch(4)).pk()));
343
344 assert_eq!(register_community(&c), 0);
346 }
347
348 #[test]
352 fn late_registered_keys_sign_against_the_remembered_challenge() {
353 let relay = RelayUrl::parse("wss://late-keys.example").unwrap();
354 let early = Keys::generate();
355 register([early.clone()]);
356 assert!(remember_challenge(&relay, "challenge-1"), "first sighting is a fresh connection");
358 let late = Keys::generate();
360 register([late.clone()]);
361 let challenge = remembered_challenge(&relay).expect("challenge was remembered");
363 let events = sign_all(&challenge, &relay);
364 let signers: Vec<_> = events.iter().map(|e| e.pubkey).collect();
365 assert!(signers.contains(&early.public_key()));
366 assert!(signers.contains(&late.public_key()));
367 }
368
369 #[test]
372 fn signed_auth_events_are_nip42_shaped() {
373 let relay = RelayUrl::parse("wss://shape.example").unwrap();
374 let key = Keys::generate();
375 register([key.clone()]);
376 let events = sign_all("shape-challenge", &relay);
377 let ev = events.iter().find(|e| e.pubkey == key.public_key()).expect("signed by the registered key");
378 assert_eq!(ev.kind, nostr_sdk::prelude::Kind::Authentication);
379 assert!(ev.verify().is_ok(), "signature + id must verify");
380 let tag_values: Vec<String> = ev.tags.iter().filter_map(|t| t.content().map(String::from)).collect();
381 assert!(tag_values.iter().any(|v| v == "shape-challenge"), "carries the challenge tag");
382 }
383
384 #[test]
387 fn challenge_value_change_detects_a_new_connection() {
388 let relay = RelayUrl::parse("wss://conn-detect.example").unwrap();
389 assert!(remember_challenge(&relay, "c1"), "first sighting");
390 assert!(!remember_challenge(&relay, "c1"), "same value = same connection");
391 assert!(remember_challenge(&relay, "c2"), "new value = reconnected");
392 assert_eq!(remembered_challenge(&relay).as_deref(), Some("c2"), "memory holds the newest");
393 }
394
395 #[test]
398 fn resubscribe_cooldown_bounds_the_reaction() {
399 let relay = RelayUrl::parse("wss://cooldown.example").unwrap();
400 assert!(resub_cooldown_elapsed(&relay), "first trigger passes");
401 assert!(!resub_cooldown_elapsed(&relay), "immediate repeat is suppressed");
402 let other = RelayUrl::parse("wss://cooldown-other.example").unwrap();
403 assert!(resub_cooldown_elapsed(&other), "cooldown is per-relay");
404 }
405}