This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / bobbin / crates / knot-ingest / src / orchestrator.rs
17 kB 463 lines
1use std::collections::{HashMap, HashSet}; 2use std::sync::{Arc, Mutex}; 3use std::time::Duration; 4 5use bobbin_edge_index::EdgeStore; 6use bobbin_knot_proxy::KnotHost; 7use bobbin_runtime::{Clock, WsTransport}; 8use bobbin_types::knot_acl::{KnotHostKey, host_to_knot_did}; 9use chrono::{DateTime, Utc}; 10use futures::StreamExt; 11use jacquard_common::DefaultStr; 12use jacquard_common::types::did::Did; 13use tokio_util::sync::CancellationToken; 14 15use crate::client::{AclListing, Completeness, KnotClient, knot_endpoint}; 16use crate::gate::CapabilityGate; 17use crate::registry::KnotRegistry; 18use crate::roster::{AclOp, Cursor, Roster}; 19use crate::stream::{StreamConfig, run_stream}; 20 21const POLL_INTERVAL: Duration = Duration::from_secs(30); 22const RECONCILE_INTERVAL: Duration = Duration::from_secs(300); 23const PROBE_CONCURRENCY: usize = 16; 24 25pub struct Orchestrator { 26 pub client: Arc<KnotClient>, 27 pub gate: Arc<CapabilityGate>, 28 pub registry: Arc<KnotRegistry>, 29 pub store: Arc<EdgeStore>, 30 pub ws: Arc<dyn WsTransport>, 31 pub clock: Arc<dyn Clock>, 32 pub dev: bool, 33 pub allow_private: bool, 34 pub cancel: CancellationToken, 35} 36 37impl Orchestrator { 38 pub async fn run(self) { 39 let mut subscribed: HashMap<KnotHostKey, CancellationToken> = HashMap::new(); 40 let mut unspawnable: HashSet<KnotHostKey> = HashSet::new(); 41 loop { 42 if self.cancel.is_cancelled() { 43 break; 44 } 45 self.discover(&mut subscribed, &mut unspawnable).await; 46 tokio::select! { 47 _ = self.cancel.cancelled() => break, 48 _ = self.clock.sleep(POLL_INTERVAL) => {} 49 } 50 } 51 } 52 53 async fn discover( 54 &self, 55 subscribed: &mut HashMap<KnotHostKey, CancellationToken>, 56 unspawnable: &mut HashSet<KnotHostKey>, 57 ) { 58 let candidates: Vec<KnotHostKey> = self 59 .registry 60 .hosts() 61 .into_iter() 62 .filter(|host| !subscribed.contains_key(host) && !unspawnable.contains(host)) 63 .collect(); 64 let approved: Vec<KnotHostKey> = futures::stream::iter(candidates) 65 .map(|host| async move { self.gate.has_knot_acl(&host).await.then_some(host) }) 66 .buffer_unordered(PROBE_CONCURRENCY) 67 .filter_map(|approved| async move { approved }) 68 .collect() 69 .await; 70 approved 71 .into_iter() 72 .for_each(|host| match self.spawn(&host) { 73 Some(token) => { 74 subscribed.insert(host, token); 75 } 76 None => { 77 tracing::warn!(host = %host, "skipping unspawnable knot endpoint"); 78 unspawnable.insert(host); 79 } 80 }); 81 } 82 83 fn spawn(&self, host: &KnotHostKey) -> Option<CancellationToken> { 84 let endpoint = knot_endpoint(host.as_str(), self.dev, self.allow_private).ok()?; 85 let knot_did = host_to_knot_did(host.as_str())?; 86 let roster = Arc::new(Mutex::new(Roster::new( 87 self.store.clone(), 88 knot_did, 89 self.registry.clone(), 90 host.clone(), 91 ))); 92 let token = self.cancel.child_token(); 93 94 let stream_cfg = StreamConfig { 95 ws: self.ws.clone(), 96 clock: self.clock.clone(), 97 cancel: token.clone(), 98 }; 99 let stream_roster = roster.clone(); 100 let stream_endpoint = endpoint.clone(); 101 tokio::spawn(async move { 102 run_stream(&stream_cfg, &stream_endpoint, &stream_roster, 0).await; 103 }); 104 105 let client = self.client.clone(); 106 let registry = self.registry.clone(); 107 let clock = self.clock.clone(); 108 let host_owned = host.clone(); 109 let reconcile_token = token.clone(); 110 tokio::spawn(async move { 111 reconcile_loop( 112 &client, 113 &endpoint, 114 &host_owned, 115 &registry, 116 &roster, 117 &*clock, 118 &reconcile_token, 119 ) 120 .await; 121 }); 122 123 Some(token) 124 } 125} 126 127async fn reconcile_loop( 128 client: &KnotClient, 129 endpoint: &KnotHost, 130 host: &KnotHostKey, 131 registry: &KnotRegistry, 132 roster: &Mutex<Roster>, 133 clock: &dyn Clock, 134 cancel: &CancellationToken, 135) { 136 loop { 137 if cancel.is_cancelled() { 138 return; 139 } 140 reconcile_once(client, endpoint, host, registry, roster).await; 141 tokio::select! { 142 _ = cancel.cancelled() => return, 143 _ = clock.sleep(RECONCILE_INTERVAL) => {} 144 } 145 } 146} 147 148async fn reconcile_once( 149 client: &KnotClient, 150 endpoint: &KnotHost, 151 host: &KnotHostKey, 152 registry: &KnotRegistry, 153 roster: &Mutex<Roster>, 154) { 155 let horizon = roster.lock().unwrap().max_cursor(); 156 157 match client.list_members(endpoint).await { 158 Ok(AclListing { 159 entries, 160 completeness, 161 }) => { 162 let present: HashSet<Did<DefaultStr>> = 163 entries.iter().map(|entry| entry.subject.clone()).collect(); 164 let mut guard = roster.lock().unwrap(); 165 entries.into_iter().for_each(|entry| { 166 guard.apply_member(AclOp::Add, entry.subject, Cursor(nanos(entry.created_at))); 167 }); 168 if completeness == Completeness::Complete { 169 guard.reap_members(&present, horizon); 170 } 171 } 172 Err(err) => tracing::warn!(host = %host, error = %err, "knot member reconcile failed"), 173 } 174 175 futures::stream::iter(registry.repos(host)) 176 .for_each(|repo| async move { 177 match client.list_collaborators(endpoint, &repo).await { 178 Ok(AclListing { 179 entries, 180 completeness, 181 }) => { 182 let present: HashSet<Did<DefaultStr>> = 183 entries.iter().map(|entry| entry.subject.clone()).collect(); 184 let mut guard = roster.lock().unwrap(); 185 entries.into_iter().for_each(|entry| { 186 guard.apply_collaborator( 187 AclOp::Add, 188 repo.clone(), 189 entry.subject, 190 Cursor(nanos(entry.created_at)), 191 ); 192 }); 193 if completeness == Completeness::Complete { 194 guard.reap_collaborators(&repo, &present, horizon); 195 } 196 } 197 Err(err) => { 198 tracing::warn!(host = %host, repo = %repo.as_ref(), error = %err, "knot collaborator reconcile failed") 199 } 200 } 201 }) 202 .await; 203 204 roster.lock().unwrap().purge_legacy(); 205} 206 207fn nanos(timestamp: DateTime<Utc>) -> i64 { 208 timestamp.timestamp_nanos_opt().unwrap_or(0) 209} 210 211#[cfg(test)] 212mod tests { 213 use super::*; 214 use crate::client::authority; 215 use bobbin_runtime::{ReqwestHttp, RuntimeHasher}; 216 use bobbin_types::edges::Edge; 217 use bobbin_types::ids::{EdgeKey, SubjectRef, nsid_static}; 218 use jacquard_common::types::string::AtUri; 219 use serde_json::json; 220 use wiremock::matchers::{method, path, query_param}; 221 use wiremock::{Mock, MockServer, ResponseTemplate}; 222 223 fn did(s: &str) -> Did<DefaultStr> { 224 Did::new_owned(s).unwrap() 225 } 226 227 fn member_count(store: &EdgeStore, subject: &str) -> u64 { 228 store.count(&EdgeKey::new( 229 nsid_static("sh.tangled.knot.member"), 230 SubjectRef::Did(did(subject)), 231 )) 232 } 233 234 fn collaborator_count(store: &EdgeStore, repo: &str) -> u64 { 235 store.count(&EdgeKey::new( 236 nsid_static("sh.tangled.repo.collaborator"), 237 SubjectRef::Did(did(repo)), 238 )) 239 } 240 241 #[tokio::test] 242 async fn reconcile_backfills_members_and_collaborators() { 243 let server = MockServer::start().await; 244 let endpoint = KnotHost::parse(&server.uri()).unwrap(); 245 let host = KnotHostKey::new(&authority(&endpoint)); 246 247 Mock::given(method("GET")) 248 .and(path("/xrpc/sh.tangled.knot.listMembers")) 249 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 250 "items": [ 251 {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}, 252 {"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} 253 ] 254 }))) 255 .mount(&server) 256 .await; 257 Mock::given(method("GET")) 258 .and(path("/xrpc/sh.tangled.repo.listCollaborators")) 259 .and(query_param("subject", "did:plc:scallop")) 260 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 261 "items": [ 262 {"subject": "did:plc:olaren", "addedBy": "did:plc:boltless", "createdAt": "2026-06-03T00:00:00Z"} 263 ] 264 }))) 265 .mount(&server) 266 .await; 267 268 let registry = Arc::new(KnotRegistry::new()); 269 registry.observe_repo(&host, did("did:plc:scallop")); 270 271 let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); 272 let knot_did = host_to_knot_did(host.as_str()).unwrap(); 273 let roster = Mutex::new(Roster::new( 274 store.clone(), 275 knot_did, 276 registry.clone(), 277 host.clone(), 278 )); 279 let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); 280 281 reconcile_once(&client, &endpoint, &host, &registry, &roster).await; 282 283 assert_eq!(member_count(&store, "did:plc:boltless"), 1); 284 assert_eq!(member_count(&store, "did:plc:akshay"), 1); 285 assert_eq!(collaborator_count(&store, "did:plc:scallop"), 1); 286 } 287 288 #[tokio::test] 289 async fn reconcile_reaps_departed_member() { 290 let server = MockServer::start().await; 291 let endpoint = KnotHost::parse(&server.uri()).unwrap(); 292 let host = KnotHostKey::new(&authority(&endpoint)); 293 294 let registry = Arc::new(KnotRegistry::new()); 295 let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); 296 let knot_did = host_to_knot_did(host.as_str()).unwrap(); 297 let roster = Mutex::new(Roster::new( 298 store.clone(), 299 knot_did, 300 registry.clone(), 301 host.clone(), 302 )); 303 let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); 304 305 Mock::given(method("GET")) 306 .and(path("/xrpc/sh.tangled.knot.listMembers")) 307 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 308 "items": [ 309 {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}, 310 {"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} 311 ] 312 }))) 313 .mount(&server) 314 .await; 315 reconcile_once(&client, &endpoint, &host, &registry, &roster).await; 316 assert_eq!(member_count(&store, "did:plc:boltless"), 1); 317 assert_eq!(member_count(&store, "did:plc:akshay"), 1); 318 319 server.reset().await; 320 Mock::given(method("GET")) 321 .and(path("/xrpc/sh.tangled.knot.listMembers")) 322 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 323 "items": [ 324 {"subject": "did:plc:akshay", "addedBy": "did:plc:akshay", "createdAt": "2026-06-02T00:00:00Z"} 325 ] 326 }))) 327 .mount(&server) 328 .await; 329 reconcile_once(&client, &endpoint, &host, &registry, &roster).await; 330 331 assert_eq!( 332 member_count(&store, "did:plc:boltless"), 333 0, 334 "a member dropped from the authoritative snapshot is reaped on reconcile" 335 ); 336 assert_eq!(member_count(&store, "did:plc:akshay"), 1); 337 } 338 339 #[tokio::test] 340 async fn reconcile_skips_reap_when_member_list_truncated() { 341 let server = MockServer::start().await; 342 let endpoint = KnotHost::parse(&server.uri()).unwrap(); 343 let host = KnotHostKey::new(&authority(&endpoint)); 344 345 Mock::given(method("GET")) 346 .and(path("/xrpc/sh.tangled.knot.listMembers")) 347 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 348 "items": [{"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"}], 349 "cursor": "more" 350 }))) 351 .mount(&server) 352 .await; 353 354 let registry = Arc::new(KnotRegistry::new()); 355 let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); 356 let knot_did = host_to_knot_did(host.as_str()).unwrap(); 357 let roster = Mutex::new(Roster::new( 358 store.clone(), 359 knot_did, 360 registry.clone(), 361 host.clone(), 362 )); 363 let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); 364 365 let stayed = did("did:plc:akshay"); 366 roster 367 .lock() 368 .unwrap() 369 .apply_member(AclOp::Add, stayed, Cursor(1_000_000)); 370 assert_eq!(member_count(&store, "did:plc:akshay"), 1); 371 372 reconcile_once(&client, &endpoint, &host, &registry, &roster).await; 373 374 assert_eq!( 375 member_count(&store, "did:plc:akshay"), 376 1, 377 "a truncated member snapshot must not reap members it could not enumerate" 378 ); 379 assert_eq!(member_count(&store, "did:plc:boltless"), 1); 380 } 381 382 fn seed_legacy_edge(store: &EdgeStore, kind: &'static str, subject: &str, source: &str) { 383 let source = AtUri::new_owned(source).unwrap(); 384 store.upsert_source( 385 &source, 386 vec![Edge { 387 kind: nsid_static(kind), 388 subject: SubjectRef::Did(did(subject)), 389 source: source.clone(), 390 sort_micros: 0, 391 }], 392 ); 393 } 394 395 #[tokio::test] 396 async fn reconcile_purges_seeded_legacy_acl() { 397 let server = MockServer::start().await; 398 let endpoint = KnotHost::parse(&server.uri()).unwrap(); 399 let host = KnotHostKey::new(&authority(&endpoint)); 400 401 Mock::given(method("GET")) 402 .and(path("/xrpc/sh.tangled.knot.listMembers")) 403 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 404 "items": [ 405 {"subject": "did:plc:boltless", "addedBy": "did:plc:akshay", "createdAt": "2026-06-01T00:00:00Z"} 406 ] 407 }))) 408 .mount(&server) 409 .await; 410 Mock::given(method("GET")) 411 .and(path("/xrpc/sh.tangled.repo.listCollaborators")) 412 .and(query_param("subject", "did:plc:scallop")) 413 .respond_with(ResponseTemplate::new(200).set_body_json(json!({ 414 "items": [ 415 {"subject": "did:plc:olaren", "addedBy": "did:plc:boltless", "createdAt": "2026-06-03T00:00:00Z"} 416 ] 417 }))) 418 .mount(&server) 419 .await; 420 421 let registry = Arc::new(KnotRegistry::new()); 422 let repo = did("did:plc:scallop"); 423 registry.observe_repo(&host, repo.clone()); 424 425 let store = Arc::new(EdgeStore::new(RuntimeHasher::default())); 426 let member_source = "at://did:plc:akshay/sh.tangled.knot.member/r1"; 427 seed_legacy_edge( 428 &store, 429 "sh.tangled.knot.member", 430 "did:plc:boltless", 431 member_source, 432 ); 433 registry.note_legacy_member(AtUri::new_owned(member_source).unwrap(), &host); 434 seed_legacy_edge( 435 &store, 436 "sh.tangled.repo.collaborator", 437 "did:plc:scallop", 438 "at://did:plc:akshay/sh.tangled.repo.collaborator/r2", 439 ); 440 441 let knot_did = host_to_knot_did(host.as_str()).unwrap(); 442 let roster = Mutex::new(Roster::new( 443 store.clone(), 444 knot_did, 445 registry.clone(), 446 host.clone(), 447 )); 448 let client = KnotClient::new(ReqwestHttp::shared(reqwest::Client::new())); 449 450 reconcile_once(&client, &endpoint, &host, &registry, &roster).await; 451 452 assert_eq!( 453 member_count(&store, "did:plc:boltless"), 454 1, 455 "stale legacy member purged, knot-owned member kept" 456 ); 457 assert_eq!( 458 collaborator_count(&store, "did:plc:scallop"), 459 1, 460 "stale legacy collaborator purged, knot-owned collaborator kept" 461 ); 462 } 463}