This repository has no description
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 ®istry,
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, ®istry, &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, ®istry, &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, ®istry, &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, ®istry, &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, ®istry, &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}