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, ClientAuthentication, ClientMessage, ClientNotification, 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::guestbook_group_key(&c.community_root, c.id(), c.root_epoch).keys().clone(),
63 super::derive::dissolved_group_key(c.id()).keys().clone(),
64 ];
65 if let Some(k) = super::control::ControlPlane::of(c).signer_keys() {
73 keys.push(k.clone());
74 }
75 for ch in &c.channels {
76 if ch.private && ch.key.is_none() {
77 continue;
78 }
79 let (secret, epoch) = c.channel_secret(ch);
80 keys.push(super::derive::channel_group_key(&secret, &ch.id, epoch).keys().clone());
81 }
82 for pk_keys in rekey_plane_keys(c) {
84 keys.push(pk_keys);
85 }
86 register(keys)
87}
88
89fn rekey_plane_keys(c: &CommunityV2) -> Vec<Keys> {
101 use crate::community::Epoch;
102 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()];
103 let cid_hex = crate::simd::hex::bytes_to_hex_32(&c.id().0);
104 let roots = super::service::channel_rekey_addressing_roots(c.community_root, &cid_hex);
105 for ch in &c.channels {
106 if ch.private {
107 for root in &roots {
108 out.push(super::derive::channel_rekey_group_key(root, &ch.id, Epoch(ch.epoch.0.saturating_add(1))).keys().clone());
109 }
110 }
111 }
112 out
113}
114
115fn remember_challenge(relay: &RelayUrl, challenge: &str) -> bool {
121 CHALLENGES
122 .lock()
123 .unwrap_or_else(|e| e.into_inner())
124 .insert(relay.clone(), challenge.to_string())
125 .as_deref()
126 != Some(challenge)
127}
128
129fn remembered_challenge(relay: &RelayUrl) -> Option<String> {
131 CHALLENGES.lock().unwrap_or_else(|e| e.into_inner()).get(relay).cloned()
132}
133
134static RESUB_AT: LazyLock<Mutex<HashMap<RelayUrl, std::time::Instant>>> = LazyLock::new(|| Mutex::new(HashMap::new()));
137const RESUB_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(30);
138
139fn resub_cooldown_elapsed(relay: &RelayUrl) -> bool {
142 let mut map = RESUB_AT.lock().unwrap_or_else(|e| e.into_inner());
143 let now = std::time::Instant::now();
144 match map.get(relay) {
145 Some(at) if now.duration_since(*at) < RESUB_COOLDOWN => false,
146 _ => {
147 map.insert(relay.clone(), now);
148 true
149 }
150 }
151}
152
153pub fn clear() {
156 REGISTRY.lock().unwrap_or_else(|e| e.into_inner()).clear();
157 CHALLENGES.lock().unwrap_or_else(|e| e.into_inner()).clear();
158 RESUB_AT.lock().unwrap_or_else(|e| e.into_inner()).clear();
159 RESPONDER_RUNNING.store(false, Ordering::SeqCst);
160}
161
162pub fn is_empty() -> bool {
164 REGISTRY.lock().unwrap_or_else(|e| e.into_inner()).is_empty()
165}
166
167fn sign_all(challenge: &str, relay: &RelayUrl) -> Vec<nostr_sdk::prelude::Event> {
170 let reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
171 reg.values()
172 .filter_map(|keys| ClientAuthentication::new(challenge, relay.clone()).finalize(keys).ok())
173 .collect()
174}
175
176async fn authenticate_streams(client: &Client, relay: &RelayUrl, challenge: &str) -> usize {
180 let events = sign_all(challenge, relay);
181 let Ok(Some(r)) = client.relay(relay.clone()).await else {
182 return 0;
183 };
184 let mut sent = 0;
185 for ev in events {
186 if r.send_msg(ClientMessage::auth(ev)).await.is_ok() {
187 sent += 1;
188 }
189 }
190 sent
191}
192
193pub fn ensure_responder(client: &Client) {
203 if RESPONDER_RUNNING.swap(true, Ordering::SeqCst) {
204 return; }
206 let client = client.clone();
207 crate::db::spawn_bound(async move {
208 let mut notifications = client.notifications();
209 while let Some(n) = notifications.next().await {
210 if let ClientNotification::Message { relay_url, message } = n {
211 if let nostr_sdk::prelude::RelayMessage::Auth { challenge } = *message {
212 let challenge = challenge.into_owned();
213 let gate_key = format!(
218 "auth_gate:{}",
219 crate::inbox_relays::normalize_relay_url(relay_url.as_str())
220 );
221 if crate::db::get_sql_setting(gate_key.clone()).ok().flatten().is_none() {
222 let _ = crate::db::set_sql_setting(gate_key, "1".to_string());
223 }
224 let fresh_connection = remember_challenge(&relay_url, &challenge);
225 if !is_empty() {
226 authenticate_streams(&client, &relay_url, &challenge).await;
227 if fresh_connection && resub_cooldown_elapsed(&relay_url) {
233 super::realtime::resubscribe_relay(&client, &relay_url).await;
234 }
235 }
236 }
237 }
238 }
239 RESPONDER_RUNNING.store(false, Ordering::SeqCst);
240 });
241}
242
243pub fn prime(community: &CommunityV2) {
247 register_community(community);
248 if let Some(client) = crate::state::nostr_client() {
249 ensure_responder(&client);
250 }
251}
252
253pub async fn prime_auth(client: &Client, relays: &[String]) {
266 if is_empty() || relays.is_empty() {
267 return;
268 }
269 ensure_responder(client);
270 let urls: Vec<RelayUrl> = relays.iter().filter_map(|r| RelayUrl::parse(r).ok()).collect();
271 if urls.is_empty() {
272 return;
273 }
274 for url in &urls {
275 if let Some(challenge) = remembered_challenge(url) {
276 authenticate_streams(client, url, &challenge).await;
277 }
278 }
279 let authors: Vec<nostr_sdk::prelude::PublicKey> = {
280 let reg = REGISTRY.lock().unwrap_or_else(|e| e.into_inner());
281 reg.keys().filter_map(|pk| nostr_sdk::prelude::PublicKey::from_slice(pk).ok()).collect()
282 };
283 if authors.is_empty() {
284 return;
285 }
286 let filter = nostr_sdk::prelude::Filter::new()
287 .kind(nostr_sdk::prelude::Kind::Custom(super::stream::KIND_WRAP))
288 .authors(authors)
289 .limit(1);
290 let _ = tokio::time::timeout(std::time::Duration::from_secs(8), client
292 .fetch_events(nostr_sdk::prelude::ReqTarget::manual(
293 urls.into_iter().map(|u| (u, vec![filter.clone()])),
294 ))
295 .timeout(std::time::Duration::from_secs(6))).await;
296}
297
298#[cfg(test)]
299mod tests {
300 use super::*;
301 use crate::community::v2::control::{genesis, CommunityMetadata};
302 use crate::community::v2::community::{ChannelV2, CommunityV2};
303 use crate::community::{ChannelId, Epoch};
304 use nostr_sdk::prelude::Keys;
305
306 fn community_with_channel_mix() -> CommunityV2 {
309 let owner = Keys::generate();
310 let g = genesis(&owner, CommunityMetadata { name: "auth-test".into(), ..Default::default() }, 1_000).unwrap();
311 let mut c = CommunityV2::from_genesis(&g, "auth-test", None, vec!["wss://gated.example".into()], 0);
312 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() });
313 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() });
314 c
315 }
316
317 fn registered(pk: &nostr_sdk::prelude::PublicKey) -> bool {
318 REGISTRY.lock().unwrap_or_else(|e| e.into_inner()).contains_key(&pk.to_bytes())
319 }
320
321 #[test]
325 fn register_community_covers_planes_and_skips_keyless() {
326 let c = community_with_channel_mix();
327 let added = register_community(&c);
328 assert!(added >= 6, "control+guestbook+dissolved+public chat+keyed chat+rekeys = at least 6 new keys, got {added}");
329
330 use super::super::derive;
331 assert!(registered(&super::super::control::ControlPlane::of(&c).pk()));
334 assert!(!registered(&derive::control_group_key(&c.community_root, c.id(), c.root_epoch).pk()));
335 assert!(registered(&derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch).pk()));
336 assert!(registered(&derive::dissolved_group_key(c.id()).pk()));
337 let public = &c.channels[0];
339 let (secret, epoch) = c.channel_secret(public);
340 assert!(registered(&derive::channel_group_key(&secret, &public.id, epoch).pk()));
341 let keyed = &c.channels[1];
343 assert!(registered(&derive::channel_group_key(&[7u8; 32], &keyed.id, Epoch(3)).pk()));
344 let keyless = &c.channels[2];
347 assert!(!registered(&derive::channel_group_key(&c.community_root, &keyless.id, Epoch(0)).pk()));
348 assert!(registered(&derive::base_rekey_group_key(&c.community_root, c.id(), Epoch(c.root_epoch.0 + 1)).pk()));
350 assert!(registered(&derive::channel_rekey_group_key(&c.community_root, &keyed.id, Epoch(4)).pk()));
351
352 assert_eq!(register_community(&c), 0);
354 }
355
356 #[test]
360 fn late_registered_keys_sign_against_the_remembered_challenge() {
361 let relay = RelayUrl::parse("wss://late-keys.example").unwrap();
362 let early = Keys::generate();
363 register([early.clone()]);
364 assert!(remember_challenge(&relay, "challenge-1"), "first sighting is a fresh connection");
366 let late = Keys::generate();
368 register([late.clone()]);
369 let challenge = remembered_challenge(&relay).expect("challenge was remembered");
371 let events = sign_all(&challenge, &relay);
372 let signers: Vec<_> = events.iter().map(|e| e.pubkey).collect();
373 assert!(signers.contains(&early.public_key()));
374 assert!(signers.contains(&late.public_key()));
375 }
376
377 #[test]
380 fn signed_auth_events_are_nip42_shaped() {
381 let relay = RelayUrl::parse("wss://shape.example").unwrap();
382 let key = Keys::generate();
383 register([key.clone()]);
384 let events = sign_all("shape-challenge", &relay);
385 let ev = events.iter().find(|e| e.pubkey == key.public_key()).expect("signed by the registered key");
386 assert_eq!(ev.kind, nostr_sdk::prelude::Kind::Authentication);
387 assert!(ev.verify().is_ok(), "signature + id must verify");
388 let tag_values: Vec<String> = ev.tags.iter().filter_map(|t| t.content().map(String::from)).collect();
389 assert!(tag_values.iter().any(|v| v == "shape-challenge"), "carries the challenge tag");
390 }
391
392 #[test]
395 fn challenge_value_change_detects_a_new_connection() {
396 let relay = RelayUrl::parse("wss://conn-detect.example").unwrap();
397 assert!(remember_challenge(&relay, "c1"), "first sighting");
398 assert!(!remember_challenge(&relay, "c1"), "same value = same connection");
399 assert!(remember_challenge(&relay, "c2"), "new value = reconnected");
400 assert_eq!(remembered_challenge(&relay).as_deref(), Some("c2"), "memory holds the newest");
401 }
402
403 #[test]
406 fn resubscribe_cooldown_bounds_the_reaction() {
407 let relay = RelayUrl::parse("wss://cooldown.example").unwrap();
408 assert!(resub_cooldown_elapsed(&relay), "first trigger passes");
409 assert!(!resub_cooldown_elapsed(&relay), "immediate repeat is suppressed");
410 let other = RelayUrl::parse("wss://cooldown-other.example").unwrap();
411 assert!(resub_cooldown_elapsed(&other), "cooldown is per-relay");
412 }
413}