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.2 kB 92 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( 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}