1use crate::encode::Encoder;
34use crate::metrics::LaneCounters;
35use crate::plan::EventPlan;
36use spate_core::checkpoint::{AckIssuer, AckRef};
37use spate_core::error::SourceError;
38use spate_core::record::{PartitionId, RawPayload};
39use spate_core::source::{LaneId, PayloadBatch, SourceLane};
40use std::io::Write;
41use std::ops::Range;
42use std::sync::Arc;
43use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
44use std::time::{Duration, Instant};
45
46#[derive(Debug)]
53pub(crate) struct Shared {
54 pub(crate) exhausted: Box<[AtomicBool]>,
56 pub(crate) paused: Box<[AtomicBool]>,
58 pub(crate) remaining: Box<[AtomicU64]>,
61 pub(crate) open: Box<[AtomicU64]>,
63}
64
65impl Shared {
66 pub(crate) fn new(partitions: usize, budgets: Option<&[u64]>) -> Shared {
69 let flags = || (0..partitions).map(|_| AtomicBool::new(false)).collect();
70 Shared {
71 exhausted: flags(),
72 paused: flags(),
73 remaining: (0..partitions)
74 .map(|i| AtomicU64::new(budgets.map_or(0, |b| b[i])))
75 .collect(),
76 open: (0..partitions).map(|_| AtomicU64::new(0)).collect(),
77 }
78 }
79}
80
81#[derive(Debug)]
84struct Item {
85 key: Range<usize>,
86 value: Range<usize>,
87 offset: i64,
88 timestamp_ms: i64,
89}
90
91pub(crate) struct LaneParts {
93 pub(crate) id: LaneId,
94 pub(crate) index: usize,
95 pub(crate) issuer: AckIssuer,
96 pub(crate) plan: EventPlan,
97 pub(crate) encoder: Arc<Encoder>,
98 pub(crate) counters: Option<LaneCounters>,
99 pub(crate) shared: Arc<Shared>,
100 pub(crate) budget: u64,
101 pub(crate) tick_interval: Duration,
102 pub(crate) events_per_tick: usize,
103}
104
105#[derive(Debug)]
107pub struct DatagenLane {
108 id: LaneId,
109 partition: PartitionId,
110 index: usize,
111 issuer: AckIssuer,
112 plan: EventPlan,
113 encoder: Arc<Encoder>,
114 counters: Option<LaneCounters>,
115 shared: Arc<Shared>,
116 budget: u64,
118 emitted: u64,
119 next_offset: i64,
120 tick_interval: Duration,
121 events_per_tick: usize,
122 next_tick: Instant,
123 tick_budget: usize,
126 arena: Vec<u8>,
127 items: Vec<Item>,
128}
129
130impl DatagenLane {
131 pub(crate) fn new(parts: LaneParts) -> DatagenLane {
132 DatagenLane {
133 id: parts.id,
134 partition: PartitionId(parts.index as u32),
135 index: parts.index,
136 issuer: parts.issuer,
137 plan: parts.plan,
138 encoder: parts.encoder,
139 counters: parts.counters,
140 shared: parts.shared,
141 budget: parts.budget,
142 emitted: 0,
143 next_offset: 0,
144 tick_interval: parts.tick_interval,
145 events_per_tick: parts.events_per_tick,
146 next_tick: Instant::now(),
149 tick_budget: 0,
150 arena: Vec::new(),
151 items: Vec::new(),
152 }
153 }
154
155 fn rate_gate(&mut self, timeout: Duration) -> Option<usize> {
159 if self.tick_interval.is_zero() {
160 return Some(usize::MAX);
161 }
162 if self.tick_budget > 0 {
166 return Some(self.tick_budget);
167 }
168 let now = Instant::now();
169 if now < self.next_tick {
170 park(Duration::min(self.next_tick - now, timeout));
171 return None;
172 }
173 let mut next = self.next_tick.checked_add(self.tick_interval);
174 if next.is_some_and(|at| at <= now) {
175 if let Some(counters) = &self.counters {
178 counters.tick_overruns.increment(1);
179 }
180 next = now.checked_add(self.tick_interval);
181 }
182 let Some(next) = next else {
186 park(timeout);
187 return None;
188 };
189 self.next_tick = next;
190 if let Some(counters) = &self.counters {
191 counters.ticks.increment(1);
192 }
193 self.tick_budget = self.events_per_tick;
194 Some(self.tick_budget)
195 }
196
197 fn fill(&mut self, count: usize) -> Result<(), SourceError> {
199 self.arena.clear();
200 self.items.clear();
201 let mut generated = [0u64; 3];
202 for _ in 0..count {
203 let (event, timestamp_ms) = self.plan.next();
204
205 let key_start = self.arena.len();
206 let _ = write!(self.arena, "{}", event.order_id());
209 let key = key_start..self.arena.len();
210
211 let value_start = self.arena.len();
212 self.encoder.encode(&event, &mut self.arena)?;
213 let value = value_start..self.arena.len();
214
215 generated[crate::metrics::kind(&event)] += 1;
216 self.items.push(Item {
217 key,
218 value,
219 offset: self.next_offset,
220 timestamp_ms,
221 });
222 self.next_offset += 1;
223 self.emitted += 1;
224 }
225 if let Some(counters) = &self.counters {
226 counters.add_generated(generated);
227 }
228 self.shared.remaining[self.index]
230 .store(self.budget.saturating_sub(self.emitted), Ordering::Release);
231 self.shared.open[self.index].store(self.plan.open_orders(), Ordering::Release);
232 Ok(())
233 }
234}
235
236impl SourceLane for DatagenLane {
237 type Batch<'a> = DatagenBatch<'a>;
238
239 fn id(&self) -> LaneId {
240 self.id
241 }
242
243 fn partition(&self) -> PartitionId {
244 self.partition
245 }
246
247 fn poll(
248 &mut self,
249 max_records: usize,
250 timeout: Duration,
251 ) -> Result<Option<DatagenBatch<'_>>, SourceError> {
252 if self.emitted >= self.budget {
253 self.shared.exhausted[self.index].store(true, Ordering::Release);
258 park(timeout);
259 return Ok(None);
260 }
261 if self.shared.paused[self.index].load(Ordering::Acquire) || max_records == 0 {
262 park(timeout);
263 return Ok(None);
264 }
265 let Some(quota) = self.rate_gate(timeout) else {
268 return Ok(None);
269 };
270
271 let count = quota
272 .min(max_records)
273 .min(usize::try_from(self.budget - self.emitted).unwrap_or(usize::MAX));
274 self.fill(count)?;
275 self.tick_budget = self.tick_budget.saturating_sub(count);
276
277 let last_offset = self.next_offset - 1;
278 Ok(Some(DatagenBatch {
279 arena: &self.arena,
280 items: &self.items,
281 next: 0,
282 partition: self.partition,
283 ack: self.issuer.issue(self.partition, last_offset),
284 }))
285 }
286}
287
288#[derive(Debug)]
290pub struct DatagenBatch<'a> {
291 arena: &'a [u8],
292 items: &'a [Item],
293 next: usize,
294 partition: PartitionId,
295 ack: AckRef,
296}
297
298impl<'a> PayloadBatch<'a> for DatagenBatch<'a> {
299 fn next_payload(&mut self) -> Option<RawPayload<'a>> {
300 let item = self.items.get(self.next)?;
301 self.next += 1;
302 Some(RawPayload {
303 bytes: &self.arena[item.value.clone()],
304 key: Some(&self.arena[item.key.clone()]),
305 partition: self.partition,
306 offset: item.offset,
307 timestamp_ms: item.timestamp_ms,
308 })
309 }
310
311 fn ack(&self) -> &AckRef {
312 &self.ack
313 }
314}
315
316pub(crate) fn park(how_long: Duration) {
319 if !how_long.is_zero() {
320 std::thread::sleep(how_long);
321 }
322}