From f30d2e7cc3fe5f1e5732b264227e48bc895a3c6b Mon Sep 17 00:00:00 2001 From: lashman Date: Wed, 20 May 2026 15:39:20 +0300 Subject: [PATCH] ap peer discovery --- cruciverb-server/src/activitypub/discover.rs | 188 +++++++++++++++++++ 1 file changed, 188 insertions(+) create mode 100644 cruciverb-server/src/activitypub/discover.rs diff --git a/cruciverb-server/src/activitypub/discover.rs b/cruciverb-server/src/activitypub/discover.rs new file mode 100644 index 0000000..b3b9087 --- /dev/null +++ b/cruciverb-server/src/activitypub/discover.rs @@ -0,0 +1,188 @@ +use std::sync::Arc; +use std::time::Duration; +use std::collections::HashSet; + +use crate::state::AppState; +use super::inbox::check_is_cruciverb; + +pub fn start_discovery_worker(state: Arc) { + tokio::spawn(async move { + // wait 2 minutes before first run to let the server stabilize + tokio::time::sleep(Duration::from_secs(120)).await; + loop { + if state.ap.is_some() && !super::is_frozen() { + discover_peers(&state).await; + } + // run every 6 hours + tokio::time::sleep(Duration::from_secs(6 * 3600)).await; + } + }); +} + +async fn discover_peers(state: &Arc) { + let ap = match state.ap.as_ref() { + Some(ap) => ap, + None => return, + }; + + // collect known domains (peers + followers) to avoid re-checking + let known_domains: HashSet = { + let dbi = state.db.lock().await; + let mut domains = HashSet::new(); + domains.insert(ap.domain.clone()); + + if let Ok(mut stmt) = dbi.db.prepare("SELECT domain FROM ap_peers") { + if let Ok(rows) = stmt.query_map([], |row| row.get::<_, String>(0)) { + for r in rows.flatten() { domains.insert(r); } + } + } + if let Ok(mut stmt) = dbi.db.prepare("SELECT domain FROM ap_blocked") { + if let Ok(rows) = stmt.query_map([], |row| row.get::<_, String>(0)) { + for r in rows.flatten() { domains.insert(r); } + } + } + domains + }; + + // search our own outbox for posts with cruciverb tags to find + // domains that have interacted with us (via boosts, likes, etc.) + // and check peer outboxes for cruciverb-tagged posts from unknown domains + + let client = reqwest::Client::builder() + .timeout(Duration::from_secs(10)) + .build() + .unwrap_or_default(); + + // strategy 1: check known peers' followers for other cruciverb instances + let peer_domains: Vec = { + let dbi = state.db.lock().await; + let mut stmt = dbi.db.prepare("SELECT domain, actor_uri FROM ap_peers").unwrap(); + stmt.query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))) + .unwrap().filter_map(|r| r.ok()).map(|(d, _)| d).collect() + }; + + for peer_domain in &peer_domains { + // fetch peer's followers to find other cruciverb instances + let followers_url = format!("https://{}/@cruciverb/followers", peer_domain); + if let Ok(resp) = client.get(&followers_url) + .header("Accept", "application/activity+json") + .send().await + { + if let Ok(body) = resp.json::().await { + if let Some(items) = body.get("orderedItems").and_then(|i| i.as_array()) { + for item in items { + if let Some(uri) = item.as_str() { + let domain = uri.split("//").nth(1) + .and_then(|s| s.split('/').next()) + .unwrap_or(""); + if !domain.is_empty() && !known_domains.contains(domain) { + try_discover_instance(state, domain).await; + } + } + } + } + } + } + } + + // strategy 2: search for the hashtag on known Mastodon/Pleroma instances + // these have public hashtag timeline APIs + let follower_domains: Vec = { + let dbi = state.db.lock().await; + let mut domains = Vec::new(); + if let Ok(mut stmt) = dbi.db.prepare("SELECT actor_uri FROM ap_followers WHERE accepted = 1") { + if let Ok(rows) = stmt.query_map([], |row| row.get::<_, String>(0)) { + for uri in rows.flatten() { + if let Some(domain) = uri.split("//").nth(1).and_then(|s| s.split('/').next()) { + if !domains.contains(&domain.to_string()) { + domains.push(domain.to_string()); + } + } + } + } + } + domains + }; + + for domain in &follower_domains { + // Mastodon API: GET /api/v1/timelines/tag/cruciverb-clue + let tag_url = format!("https://{}/api/v1/timelines/tag/cruciverb-clue?limit=20", domain); + if let Ok(resp) = client.get(&tag_url).send().await { + if let Ok(posts) = resp.json::>().await { + for post in &posts { + // extract the account's domain + if let Some(acct) = post.get("account") + .and_then(|a| a.get("url")) + .and_then(|u| u.as_str()) + { + let post_domain = acct.split("//").nth(1) + .and_then(|s| s.split('/').next()) + .unwrap_or(""); + if !post_domain.is_empty() && !known_domains.contains(post_domain) { + try_discover_instance(state, post_domain).await; + } + } + } + } + } + } +} + +async fn try_discover_instance(state: &Arc, domain: &str) { + if check_is_cruciverb(domain).await { + let ap = match state.ap.as_ref() { + Some(ap) => ap, + None => return, + }; + + let actor_uri = format!("https://{}/@cruciverb", domain); + if crate::security::is_internal_url(&actor_uri) { + tracing::warn!("refusing to discover internal URL: {}", actor_uri); + return; + } + let client = reqwest::Client::builder() + .timeout(Duration::from_secs(10)) + .build() + .unwrap_or_default(); + + // fetch their actor document for inbox and key + if let Ok(resp) = client.get(&actor_uri) + .header("Accept", "application/activity+json") + .send().await + { + if let Ok(actor) = resp.json::().await { + let inbox = actor.get("inbox").and_then(|v| v.as_str()).unwrap_or(""); + let key_pem = actor.get("publicKey") + .and_then(|pk| pk.get("publicKeyPem")) + .and_then(|v| v.as_str()) + .unwrap_or(""); + + if !inbox.is_empty() { + let dbi = state.db.lock().await; + let now = chrono::Utc::now().to_rfc3339(); + dbi.db.execute( + "INSERT OR IGNORE INTO ap_peers (domain, actor_uri, inbox_uri, public_pem, last_seen, created_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?5)", + rusqlite::params![domain, actor_uri, inbox, key_pem, now], + ).ok(); + + // send a Follow to establish mutual connection + let follow = serde_json::json!({ + "@context": "https://www.w3.org/ns/activitystreams", + "id": format!("{}/follow/{}", ap.actor_uri, uuid::Uuid::new_v4()), + "type": "Follow", + "actor": ap.actor_uri, + "object": actor_uri, + }); + let payload = serde_json::to_string(&follow).unwrap(); + dbi.db.execute( + "INSERT INTO ap_delivery_queue (inbox_uri, payload, next_retry, created_at) VALUES (?1, ?2, ?3, ?3)", + rusqlite::params![inbox, payload, now], + ).ok(); + + tracing::info!("discovered cruciverb instance at {} via hashtag search", domain); + } + } + } + } +}