This repository has no description
1use std::sync::Arc;
2use std::time::Duration;
3
4use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex};
5use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run};
6use bobbin_record_lru::{NoopRecordStore, RecordStore};
7use bobbin_resolver::IdentityResolver;
8use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs};
9use bobbin_types::search::NoopSearchSink;
10use futures::stream::{self, StreamExt};
11use tokio_util::sync::CancellationToken;
12use url::Url;
13
14const MIN_EVENTS_PER_RUN: u64 = 100;
15
16#[tokio::main(flavor = "multi_thread")]
17async fn main() {
18 tracing_subscriber::fmt::try_init().ok();
19
20 let endpoint = std::env::args()
21 .nth(1)
22 .unwrap_or_else(|| "ws://127.0.0.1:3010".to_owned());
23 let ticks: u32 = std::env::args()
24 .nth(2)
25 .and_then(|s| s.parse().ok())
26 .unwrap_or(6);
27
28 let url = Url::parse(&endpoint).expect("valid hydrant base url");
29 let hasher = RuntimeHasher::from_entropy(&OsEntropy);
30 let store = Arc::new(EdgeStore::new(hasher.clone()));
31 let coverage = Arc::new(CoverageWatch::new());
32
33 let cfg = IngestConfig::new(url);
34 let cancel = CancellationToken::new();
35 let runtime = IngestRuntime {
36 store: store.clone(),
37 issue_states: Arc::new(StateIndex::new(hasher.clone())),
38 pull_statuses: Arc::new(StateIndex::new(hasher.clone())),
39 coverage: coverage.clone(),
40 search: Arc::new(NoopSearchSink),
41 records: Arc::new(NoopRecordStore) as Arc<dyn RecordStore>,
42 resolver: Arc::new(RepoIdResolver::detached(hasher.clone())),
43 identity: Arc::new(IdentityResolver::detached(hasher)),
44 clock: Arc::new(SystemClock::new()),
45 entropy: Arc::new(OsEntropy),
46 ws: TungsteniteWs::shared(),
47 cancel: cancel.clone(),
48 disconnects: None,
49 warming_shadow: None,
50 warming_buffer: None,
51 knot_registry: None,
52 knot_gate: None,
53 };
54 let task = tokio::spawn(async move {
55 let _ = run(cfg, runtime).await;
56 });
57
58 stream::iter(0..ticks)
59 .for_each(|i| {
60 let store = store.clone();
61 let coverage = coverage.clone();
62 async move {
63 tokio::time::sleep(Duration::from_secs(5)).await;
64 let snap = coverage.snapshot();
65 println!(
66 "[{i:02}] coverage={snap:?} edge_keys={} sources={}",
67 store.key_count(),
68 store.source_count()
69 );
70 }
71 })
72 .await;
73
74 cancel.cancel();
75 let _ = task.await;
76
77 let final_snap = coverage.snapshot();
78 let events = final_snap.events_processed();
79 println!(
80 "done: events={events} edge_keys={} sources={} ready={}",
81 store.key_count(),
82 store.source_count(),
83 final_snap.is_ready(),
84 );
85 assert!(
86 events >= MIN_EVENTS_PER_RUN,
87 "smoke received {events} events from {endpoint} across {ticks} five-second ticks, below threshold of {MIN_EVENTS_PER_RUN}; hydrant may be silent or unreachable",
88 );
89}