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