This repository has no description
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}