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.4 kB 196 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( 120 hasher.clone(), 121 bobbin_resolver::DEFAULT_IDENTITY_CACHE_ENTRIES, 122 )), 123 clock: clock.clone(), 124 entropy: entropy.clone(), 125 ws: mem_ws, 126 cancel: cancel.clone(), 127 disconnects: Some(disconnects.clone()), 128 warming_shadow: Some(warming_shadow.clone()), 129 warming_buffer: warming_buffer_enabled.then(|| warming_buffer.clone()), 130 knot_registry: None, 131 knot_gate: None, 132 }; 133 let ingest_config = IngestConfig { 134 hydrant_base, 135 start_cursor: bobbin_edge_index::HydrantCursor::new(0), 136 parallelism, 137 }; 138 let mut ingest_handle = tokio::spawn(async move { 139 let _ = run_ingest(ingest_config, ingest_runtime).await; 140 }); 141 142 let script = hooks.script; 143 let initial_report = tokio::select! { 144 biased; 145 _ = clock.sleep(max_virtual_runtime) => SimReport { 146 workload: workload_name, 147 seed, 148 outcome: SimOutcome::TimedOut, 149 virtual_runtime: max_virtual_runtime, 150 virtual_clock_end: clock.now_unix_micros(), 151 events_processed: coverage.snapshot().events_processed(), 152 last_cursor: coverage.snapshot().last_cursor().raw(), 153 edge_count: store.key_count() as u64, 154 resolver_hits: 0, 155 resolver_misses: 0, 156 consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), 157 disconnect_count: disconnects.count(), 158 last_disconnect: disconnects.snapshot(), 159 warming_shadow: warming_shadow.snapshot(), 160 warming_buffer: warming_buffer.snapshot(), 161 failure_reason: Some(format!( 162 "max_virtual_runtime {max_virtual_runtime:?} exhausted", 163 )), 164 }, 165 r = script => r, 166 }; 167 168 cancel.cancel(); 169 let drain_deadline = clock.sleep(Duration::from_secs(30)); 170 tokio::pin!(drain_deadline); 171 tokio::select! { 172 _ = &mut drain_deadline => { 173 ingest_handle.abort(); 174 let _ = ingest_handle.await; 175 } 176 res = &mut ingest_handle => { 177 let _ = res; 178 } 179 } 180 181 let stats = resolver.stats(); 182 SimReport { 183 edge_count: store.key_count() as u64, 184 events_processed: coverage.snapshot().events_processed(), 185 last_cursor: coverage.snapshot().last_cursor().raw(), 186 resolver_hits: stats.hits, 187 resolver_misses: stats.miss_count(), 188 consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), 189 disconnect_count: disconnects.count(), 190 last_disconnect: disconnects.snapshot(), 191 warming_shadow: warming_shadow.snapshot(), 192 warming_buffer: warming_buffer.snapshot(), 193 ..initial_report 194 } 195 } 196}