This repository has no description
0

Configure Feed

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

core / bobbin / crates / ingest / examples / smoke.rs
3.1 kB 89 lines
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}