This repository has no description
0

Configure Feed

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

core / bobbin / crates / bobbin-sim / src / runtime.rs
7.3 kB 193 lines
1use std::num::NonZeroUsize; 2use std::sync::Arc; 3use std::sync::atomic::{AtomicU64, Ordering}; 4use std::time::Duration; 5 6use bobbin_edge_index::{CoverageWatch, EdgeStore}; 7use bobbin_ingest::{ 8 DEFAULT_INGEST_PARALLELISM, DisconnectSink, IngestConfig, IngestRuntime, RepoIdResolver, 9 WarmingBuffer, WarmingShadowBuffer, run as run_ingest, 10}; 11use bobbin_record_lru::NoopRecordStore; 12use bobbin_resolver::IdentityResolver; 13use bobbin_runtime::{ 14 Clock, DEFAULT_MEM_WS_CAPACITY, MemHttpTransport, MemWsTransport, RuntimeHasher, SeededEntropy, 15 SimClock, UnixMicros, 16}; 17use bobbin_slingshot_client::SlingshotClient; 18use bobbin_types::search::NoopSearchSink; 19use tokio_util::sync::CancellationToken; 20use url::Url; 21 22use crate::report::{SimOutcome, SimReport}; 23use crate::workload::{Workload, WorkloadCtx}; 24 25#[derive(Clone, Debug)] 26pub struct SimConfig { 27 pub seed: u64, 28 pub base_unix: UnixMicros, 29 pub max_virtual_runtime: Duration, 30 pub parallelism: NonZeroUsize, 31 pub hydrant_base: Url, 32 pub slingshot_base: Url, 33 pub mem_ws_capacity: usize, 34 pub warming_buffer_enabled: bool, 35} 36 37impl SimConfig { 38 pub fn new(seed: u64) -> Self { 39 Self { 40 seed, 41 base_unix: UnixMicros::new(1_700_000_000_000_000), 42 max_virtual_runtime: Duration::from_secs(60), 43 parallelism: DEFAULT_INGEST_PARALLELISM, 44 hydrant_base: Url::parse("ws://hydrant.sim/").unwrap(), 45 slingshot_base: Url::parse("http://slingshot.sim/").unwrap(), 46 mem_ws_capacity: DEFAULT_MEM_WS_CAPACITY, 47 warming_buffer_enabled: true, 48 } 49 } 50} 51 52pub struct Sim { 53 config: SimConfig, 54 workload: Box<dyn Workload>, 55} 56 57impl Sim { 58 pub fn new(config: SimConfig, workload: Box<dyn Workload>) -> Self { 59 Self { config, workload } 60 } 61 62 pub async fn run(self) -> SimReport { 63 let SimConfig { 64 seed, 65 base_unix, 66 max_virtual_runtime, 67 parallelism, 68 hydrant_base, 69 slingshot_base, 70 mem_ws_capacity, 71 warming_buffer_enabled, 72 } = self.config; 73 74 let entropy = Arc::new(SeededEntropy::new(seed)); 75 let hasher = RuntimeHasher::from_entropy(&*entropy); 76 let clock: Arc<dyn Clock> = Arc::new(SimClock::at(base_unix)); 77 78 let store = Arc::new(EdgeStore::new(hasher.clone())); 79 let coverage = Arc::new(CoverageWatch::new()); 80 let records = Arc::new(NoopRecordStore); 81 let cancel = CancellationToken::new(); 82 let consumer_too_slow_count = Arc::new(AtomicU64::new(0)); 83 let disconnects = Arc::new(DisconnectSink::new()); 84 let warming_shadow = Arc::new(WarmingShadowBuffer::new(hasher.clone())); 85 let warming_buffer = Arc::new(WarmingBuffer::new(hasher.clone())); 86 87 let workload_name = self.workload.name(); 88 let ctx = WorkloadCtx { 89 seed, 90 clock: clock.clone(), 91 entropy: entropy.clone(), 92 store: store.clone(), 93 coverage: coverage.clone(), 94 records: records.clone(), 95 cancel: cancel.clone(), 96 consumer_too_slow_count: consumer_too_slow_count.clone(), 97 }; 98 let hooks = self.workload.build(ctx); 99 100 let slingshot_http = MemHttpTransport::shared(hooks.slingshot.clone(), clock.clone()); 101 let slingshot_client = SlingshotClient::new(slingshot_base, slingshot_http) 102 .expect("slingshot base url is valid"); 103 let resolver = Arc::new(RepoIdResolver::with_slingshot( 104 slingshot_client, 105 clock.clone(), 106 hasher.clone(), 107 )); 108 109 let mem_ws = MemWsTransport::shared_with_capacity(hooks.hydrant.clone(), mem_ws_capacity); 110 111 let ingest_runtime: IngestRuntime<NoopSearchSink> = IngestRuntime { 112 store: store.clone(), 113 issue_states: Arc::new(bobbin_edge_index::StateIndex::new(hasher.clone())), 114 pull_statuses: Arc::new(bobbin_edge_index::StateIndex::new(hasher.clone())), 115 coverage: coverage.clone(), 116 search: Arc::new(NoopSearchSink), 117 records: records.clone() as Arc<dyn bobbin_record_lru::RecordStore>, 118 resolver: resolver.clone(), 119 identity: Arc::new(IdentityResolver::detached(hasher.clone())), 120 clock: clock.clone(), 121 entropy: entropy.clone(), 122 ws: mem_ws, 123 cancel: cancel.clone(), 124 disconnects: Some(disconnects.clone()), 125 warming_shadow: Some(warming_shadow.clone()), 126 warming_buffer: warming_buffer_enabled.then(|| warming_buffer.clone()), 127 knot_registry: None, 128 knot_gate: None, 129 }; 130 let ingest_config = IngestConfig { 131 hydrant_base, 132 start_cursor: bobbin_edge_index::HydrantCursor::new(0), 133 parallelism, 134 }; 135 let mut ingest_handle = tokio::spawn(async move { 136 let _ = run_ingest(ingest_config, ingest_runtime).await; 137 }); 138 139 let script = hooks.script; 140 let initial_report = tokio::select! { 141 biased; 142 _ = clock.sleep(max_virtual_runtime) => SimReport { 143 workload: workload_name, 144 seed, 145 outcome: SimOutcome::TimedOut, 146 virtual_runtime: max_virtual_runtime, 147 virtual_clock_end: clock.now_unix_micros(), 148 events_processed: coverage.snapshot().events_processed(), 149 last_cursor: coverage.snapshot().last_cursor().raw(), 150 edge_count: store.key_count() as u64, 151 resolver_hits: 0, 152 resolver_misses: 0, 153 consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), 154 disconnect_count: disconnects.count(), 155 last_disconnect: disconnects.snapshot(), 156 warming_shadow: warming_shadow.snapshot(), 157 warming_buffer: warming_buffer.snapshot(), 158 failure_reason: Some(format!( 159 "max_virtual_runtime {max_virtual_runtime:?} exhausted", 160 )), 161 }, 162 r = script => r, 163 }; 164 165 cancel.cancel(); 166 let drain_deadline = clock.sleep(Duration::from_secs(30)); 167 tokio::pin!(drain_deadline); 168 tokio::select! { 169 _ = &mut drain_deadline => { 170 ingest_handle.abort(); 171 let _ = ingest_handle.await; 172 } 173 res = &mut ingest_handle => { 174 let _ = res; 175 } 176 } 177 178 let stats = resolver.stats(); 179 SimReport { 180 edge_count: store.key_count() as u64, 181 events_processed: coverage.snapshot().events_processed(), 182 last_cursor: coverage.snapshot().last_cursor().raw(), 183 resolver_hits: stats.hits, 184 resolver_misses: stats.miss_count(), 185 consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), 186 disconnect_count: disconnects.count(), 187 last_disconnect: disconnects.snapshot(), 188 warming_shadow: warming_shadow.snapshot(), 189 warming_buffer: warming_buffer.snapshot(), 190 ..initial_report 191 } 192 } 193}