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(
44 hasher,
45 bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES,
46 )),
47 clock: Arc::new(SystemClock::new()),
48 entropy: Arc::new(OsEntropy),
49 ws: TungsteniteWs::shared(),
50 cancel: cancel.clone(),
51 disconnects: None,
52 warming_shadow: None,
53 warming_buffer: None,
54 knot_registry: None,
55 knot_gate: None,
56 };
57 let task = tokio::spawn(async move {
58 let _ = run(cfg, runtime).await;
59 });
60
61 stream::iter(0..ticks)
62 .for_each(|i| {
63 let store = store.clone();
64 let coverage = coverage.clone();
65 async move {
66 tokio::time::sleep(Duration::from_secs(5)).await;
67 let snap = coverage.snapshot();
68 println!(
69 "[{i:02}] coverage={snap:?} edge_keys={} sources={}",
70 store.key_count(),
71 store.source_count()
72 );
73 }
74 })
75 .await;
76
77 cancel.cancel();
78 let _ = task.await;
79
80 let final_snap = coverage.snapshot();
81 let events = final_snap.events_processed();
82 println!(
83 "done: events={events} edge_keys={} sources={} ready={}",
84 store.key_count(),
85 store.source_count(),
86 final_snap.is_ready(),
87 );
88 assert!(
89 events >= MIN_EVENTS_PER_RUN,
90 "smoke received {events} events from {endpoint} across {ticks} five-second ticks, below threshold of {MIN_EVENTS_PER_RUN}; hydrant may be silent or unreachable",
91 );
92}