polydat_core/iteration/source.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Data sources: typed sequences that drive workload iteration.
5//!
6//! A **source** is a data provider with identity. It knows what it yields
7//! (schema), how much it has (extent), where the consumer is (cursor),
8//! and how to partition across concurrent fibers.
9//!
10//! Sources replace the `cycles` counter as the workload iteration driver.
11//! A Polydat program declares one with the `cursor` keyword
12//! (`cursor q = range(0, N) [over <partition>]`). A host pulls from
13//! sources to drive dispatch; when a source is exhausted, the activation
14//! is done.
15//!
16//! ## Source Types
17//!
18//! - **Range**: `range(0, N)` — a finite sequence of ordinals, served by
19//! [`RangeSourceFactory`]. Replaces `cycles: N`.
20//! - **Extending**: `until_elapsed(base, min_ms[, delta])` and the other
21//! `until_*` constructors, compiled to a `CursorKind::Extending*` and
22//! served by [`ExtendingRangeSourceFactory`].
23//! - **Host-supplied**: any [`DataSourceFactory`] a host implements, such
24//! as a dataset reader whose items carry vectors or metadata; this
25//! crate ships only the range factories.
26//!
27//! ## Crate Sovereignty
28//!
29//! All source API surface lives here in `polydat-core`, re-exported by
30//! `polydat` at `polydat::iteration::source`. Host runtimes and adapters
31//! consume these types but don't define them.
32
33use std::sync::Arc;
34use std::sync::atomic::{AtomicU64, Ordering};
35
36use crate::ast::{PortType, Value};
37
38/// Whether a source item can be reconstructed from its ordinal without
39/// consulting or advancing mutable source state.
40///
41/// The default is deliberately conservative. A forward-only cursor is not
42/// sufficient evidence that an earlier item is safe to render again.
43#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
44pub enum SourceReplayStability {
45 /// Rendering is destructive, stateful, externally mutable, or otherwise
46 /// not proven stable for repeated calls at one ordinal.
47 #[default]
48 Consumptive,
49 /// `render_item(ordinal)` is pure, total, and byte-identical for the
50 /// lifetime of the advertised source generation.
51 StableByOrdinal,
52}
53
54/// A compact description of values which can be generated without rendering
55/// an opaque source item first.
56#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
57pub enum SourceValueForm {
58 /// The source may be replayable, but the item must still be rendered.
59 #[default]
60 Opaque,
61 /// The yielded scalar value is the ordinal itself. A typed consumer may
62 /// apply its declared narrowing or wrapping conversion while vectorizing.
63 Ordinal,
64}
65
66/// Runtime source capability used by offset-stamped batch plans.
67///
68/// `generation` must change whenever rendering the same ordinal may produce a
69/// different value. Activation identity remains executor-owned because a
70/// phase rewind can intentionally replay one unchanged source generation as a
71/// new logical evaluation.
72#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
73pub struct SourceReplayContract {
74 /// Whether rendering an ordinal again gives the same value.
75 pub stability: SourceReplayStability,
76 /// The generation of the source's values; changes when an ordinal may render differently.
77 pub generation: u64,
78 /// What a value is: an ordinal, or opaque.
79 pub value_form: SourceValueForm,
80}
81
82impl SourceReplayContract {
83 /// The contract of a source that is consumed as it advances: no replay.
84 pub const fn consumptive() -> Self {
85 Self {
86 stability: SourceReplayStability::Consumptive,
87 generation: 0,
88 value_form: SourceValueForm::Opaque,
89 }
90 }
91
92 /// The contract of a source stable by ordinal in `generation`.
93 pub const fn stable_ordinal(generation: u64) -> Self {
94 Self {
95 stability: SourceReplayStability::StableByOrdinal,
96 generation,
97 value_form: SourceValueForm::Ordinal,
98 }
99 }
100
101 /// Whether any owned ordinal renders again without advancing the source.
102 pub const fn is_stable_by_ordinal(self) -> bool {
103 matches!(self.stability, SourceReplayStability::StableByOrdinal)
104 }
105
106 /// True when the source ordinal is simultaneously a replay key, scalar
107 /// value, packet number, and low-order lane clock.
108 pub const fn is_perfect_ordinal(self) -> bool {
109 self.is_stable_by_ordinal() && matches!(self.value_form, SourceValueForm::Ordinal)
110 }
111}
112
113/// A single item yielded by a source.
114#[derive(Clone, Debug)]
115pub struct SourceItem {
116 /// Position in the source sequence.
117 pub ordinal: u64,
118 /// Named field values. Empty for range sources (ordinal IS the data).
119 /// For dataset sources: `[("vector", Value::Json(...)), ("metadata", Value::U64(...))]`.
120 pub fields: Vec<(String, Value)>,
121}
122
123impl SourceItem {
124 /// Create a range item (ordinal only, no fields).
125 pub fn ordinal(ordinal: u64) -> Self {
126 Self {
127 ordinal,
128 fields: Vec::new(),
129 }
130 }
131
132 /// Create an item with ordinal and named fields.
133 pub fn with_fields(ordinal: u64, fields: Vec<(String, Value)>) -> Self {
134 Self { ordinal, fields }
135 }
136
137 /// Get a field value by name.
138 pub fn field(&self, name: &str) -> Option<&Value> {
139 self.fields.iter().find(|(n, _)| n == name).map(|(_, v)| v)
140 }
141}
142
143/// Schema describing what a source yields.
144#[derive(Clone, Debug)]
145pub struct SourceSchema {
146 /// Source name as declared in the Polydat graph.
147 pub name: String,
148 /// Field names and types available for projection (e.g., `ordinal: U64`, `vector: Json`).
149 pub projections: Vec<(String, PortType)>,
150 /// Known extent, if finite. `None` for infinite sources.
151 ///
152 /// `None` may also mean the extent is computable but only at runtime
153 /// (e.g., the cursor's `range(...)` bounds depend on iteration-variable
154 /// externs). In that case `extent_outputs` carries the kernel output
155 /// names whose values yield `[start, end)` once externs are bound.
156 pub extent: Option<u64>,
157 /// Optional aux output names `(start, end)` that the runtime can pull
158 /// from the kernel after externs are populated to compute extent. Set
159 /// when `range(...)` bounds are non-literal (e.g., wire-bound dataset
160 /// function calls). The compiler also writes back to `extent` if both
161 /// values fold to constants at compile time.
162 pub extent_outputs: Option<(String, String)>,
163 /// Optional cursor limit clamp, taken from
164 /// `CompileOptions::cursor_limit` (a host's `--limit`-style option).
165 /// Applied after runtime extent evaluation.
166 pub extent_limit: Option<u64>,
167 /// What kind of cursor this is. Default `Range`, the bounded
168 /// `range(start, end)` cursor; the `Extending*` variants declare a
169 /// runtime-extending cursor under a wall-clock, count, or
170 /// predicate policy. A host picks [`ExtendingRangeSourceFactory`]
171 /// for those, and `cursor_over_partitions` treats them as
172 /// open-extent when it resolves an `over` clause.
173 pub cursor_kind: CursorKind,
174 /// Name of the kernel output that carries the
175 /// partition-narrowing source (the `over <expr>` clause on
176 /// the cursor declaration, cursor_partitions.md §7.2).
177 /// `cursor_over_partitions` pulls it and
178 /// resolves it through `cursor_partition::resolve_over` into the
179 /// full partition list: a spec string or a `PartitionSpec`
180 /// resolves against the cursor's extent, a `Partition` or a
181 /// `PartitionList` is re-projected onto it, and `Value::None`
182 /// yields no partitions. A host, such as activation setup,
183 /// selects one partition and writes it with `set_cursor`.
184 ///
185 /// `None` means the cursor was declared without an `over`
186 /// clause; the cursor uses its full declared extent.
187 pub partition_output: Option<String>,
188 /// The partitions the `over` clause denotes, resolved by the
189 /// compiler when the clause is a literal spec and the extent is
190 /// known at build (engines.md §3.5). A host reads them
191 /// without evaluating anything; `cursor_over_partitions` returns
192 /// them without a pull. `None` when the clause or the extent is
193 /// only known at run time, or the cursor has no `over` clause.
194 pub partitions: Option<Vec<crate::iteration::cursor_partition::Partition>>,
195}
196
197/// Cursor-construction discriminator. Set by the Polydat compiler
198/// when it recognises a cursor's constructor expression
199/// (`range(...)`, `until_elapsed(...)`, etc.); read by the
200/// executor at phase setup to instantiate the matching
201/// data-source factory + policy.
202///
203/// Each `Extending*` variant carries the Polydat output names the
204/// runtime should pull at phase setup to resolve the policy's
205/// parameters (same pattern as `extent_outputs` for `range`).
206/// `delta_output` is optional — when `None`, the cursor's base
207/// is used as the extension delta.
208#[derive(Clone, Debug, Default)]
209pub enum CursorKind {
210 /// Bounded `range(start, end)` cursor. Default — keeps
211 /// every existing call site working.
212 #[default]
213 Range,
214 /// `until_elapsed(base, min_ms[, delta])`.
215 /// Extends while elapsed_ms < min_ms.
216 ExtendingTimed {
217 /// The output holding the minimum elapsed milliseconds.
218 min_ms_output: String,
219 /// The output holding the extension step, if one was given.
220 delta_output: Option<String>,
221 },
222 /// `until_passes(base, min_passes[, delta])`.
223 /// Extends while completed passes < min_passes.
224 ExtendingPasses {
225 /// The output holding the minimum number of passes.
226 min_passes_output: String,
227 /// The output holding the extension step, if one was given.
228 delta_output: Option<String>,
229 },
230 /// `until_count(base, min_count[, delta])`.
231 /// Extends while raw consumed count < min_count.
232 ExtendingCount {
233 /// The output holding the minimum consumed count.
234 min_count_output: String,
235 /// The output holding the extension step, if one was given.
236 delta_output: Option<String>,
237 },
238 /// `until_elapsed_and_passes(base, min_ms, min_passes[, delta])`.
239 /// AND: stops when EITHER target is reached. Extends while
240 /// BOTH conditions are still below target.
241 ExtendingElapsedAndPasses {
242 /// The output holding the minimum elapsed milliseconds.
243 min_ms_output: String,
244 /// The output holding the minimum number of passes.
245 min_passes_output: String,
246 /// The output holding the extension step, if one was given.
247 delta_output: Option<String>,
248 },
249 /// `until_elapsed_or_passes(base, min_ms, min_passes[, delta])`.
250 /// OR: stops only when BOTH targets are reached. Extends
251 /// while EITHER condition is still below target.
252 ExtendingElapsedOrPasses {
253 /// The output holding the minimum elapsed milliseconds.
254 min_ms_output: String,
255 /// The output holding the minimum number of passes.
256 min_passes_output: String,
257 /// The output holding the extension step, if one was given.
258 delta_output: Option<String>,
259 },
260}
261
262/// Consumption API for data sources. One instance per fiber.
263///
264/// The interaction model has two phases:
265///
266/// 1. **Reserve** — `reserve(stride)` atomically claims a range of
267/// ordinals via CAS on the shared cursor. This touches shared
268/// state but is instantaneous (one atomic op). Returns `None`
269/// when the source is globally exhausted.
270///
271/// 2. **Render** — the fiber uses the reserved range with its own
272/// Polydat instance to produce field values. No shared state, no
273/// contention between fibers. For range sources, rendering is
274/// trivial (ordinal IS the data). For dataset sources, rendering
275/// reads vectors/metadata from mmap'd storage.
276///
277/// The `next_chunk` convenience method combines both phases. Use
278/// `reserve` directly when the rendering is handled by the
279/// executor's Polydat fiber.
280pub trait DataSource: Send {
281 /// Atomically reserve up to `stride` ordinals from the source.
282 ///
283 /// Returns the half-open range `[start..end)` of reserved
284 /// ordinals, or `None` if the source is exhausted. The range
285 /// may be shorter than `stride` at the tail of the source.
286 ///
287 /// This is the only method that touches shared state (the
288 /// global cursor). It must be lock-free — a single CAS or
289 /// fetch_add.
290 fn reserve(&mut self, stride: usize) -> Option<std::ops::Range<u64>>;
291
292 /// Pull the next item. `None` = source exhausted.
293 fn next(&mut self) -> Option<SourceItem> {
294 let range = self.reserve(1)?;
295 Some(self.render_item(range.start))
296 }
297
298 /// Pull up to `limit` items. Combines reserve + render.
299 /// Returns fewer than `limit` only when the source is globally
300 /// exhausted. Empty vec = exhausted.
301 fn next_chunk(&mut self, limit: usize) -> Vec<SourceItem> {
302 let range = match self.reserve(limit) {
303 Some(r) => r,
304 None => return Vec::new(),
305 };
306 (range.start..range.end)
307 .map(|ordinal| self.render_item(ordinal))
308 .collect()
309 }
310
311 /// Produce a source item for a previously reserved ordinal.
312 ///
313 /// This is the fiber-local rendering step — no shared state.
314 /// For range sources: returns `SourceItem::ordinal(ordinal)`.
315 /// For dataset sources: reads vector/metadata from storage.
316 fn render_item(&self, ordinal: u64) -> SourceItem;
317
318 /// Known extent, if finite.
319 fn extent(&self) -> Option<u64>;
320
321 /// Items consumed so far (for progress reporting).
322 fn consumed(&self) -> u64;
323
324 /// The schema of items this source yields.
325 fn schema(&self) -> &SourceSchema;
326
327 /// Replay/addressability capability for offset-stamped batch execution.
328 /// External source implementations inherit the safe consumptive default
329 /// until they explicitly prove the stronger contract.
330 fn replay_contract(&self) -> SourceReplayContract {
331 SourceReplayContract::consumptive()
332 }
333}
334
335/// Factory that creates per-fiber `DataSource` readers.
336///
337/// Holds shared state (atomic cursor, partition pool) that's
338/// distributed across readers. Each fiber gets its own reader.
339///
340/// ## Dispatch model
341///
342/// The **stride** is the stanza length — the number of source items
343/// a fiber acquires as an atomic unit. One stanza of ops processes
344/// one stride of source items. Strides are inseparable: a fiber
345/// that acquires a stride processes all items before acquiring the
346/// next.
347///
348/// The default implementation (`RangeSourceFactory`) uses a shared
349/// atomic cursor — all fibers pull strides from the same counter,
350/// producing natural monotonic striping. This is correct for range
351/// sources where items are independent ordinals.
352///
353/// For dataset sources with locality benefits (mmap prefetch),
354/// factories can implement partitioned allocation: each fiber gets
355/// a pre-assigned range of strides, and when exhausted, steals
356/// strides from a shared pool. The stride is the minimum unit of
357/// work stealing — a fiber never steals partial stanzas.
358pub trait DataSourceFactory: Send + Sync {
359 /// Create a new reader for a fiber.
360 fn create_reader(&self) -> Box<dyn DataSource>;
361
362 /// Schema for all readers from this factory.
363 fn schema(&self) -> &SourceSchema;
364
365 /// Global items consumed across all readers (for progress reporting).
366 fn global_consumed(&self) -> u64;
367
368 /// Known extent, if finite. Same as schema().extent but avoids clone.
369 fn global_extent(&self) -> Option<u64> {
370 self.schema().extent
371 }
372
373 /// Replay/addressability capability shared by readers from this factory.
374 fn replay_contract(&self) -> SourceReplayContract {
375 SourceReplayContract::consumptive()
376 }
377
378 /// Start a new round over the same ordinal domain, returning
379 /// `true`, or return `false` when the factory cannot rewind.
380 ///
381 /// A factory that rewinds promises that every reservation
382 /// linearized after the call returns, from any reader of the
383 /// factory, belongs to the new round, and that the new round
384 /// hands out each of its ordinals at most once. This is safe
385 /// under any concurrency: readers may reserve while the rewind
386 /// runs, and a reservation that overlaps it belongs to one round
387 /// or the other. The default returns `false` and changes nothing.
388 /// A host that re-runs a source between poll rounds rejects a
389 /// factory that returns `false` rather than completing silently
390 /// after the first round (cursor_partitions.md §8).
391 fn rewind_for_poll(&self) -> bool {
392 false
393 }
394}
395
396// =========================================================================
397// RangeSource: finite sequence of ordinals
398// =========================================================================
399
400/// Factory for range sources. Shared atomic cursor distributes
401/// ordinals across fibers.
402pub struct RangeSourceFactory {
403 cursor: Arc<AtomicU64>,
404 end: u64,
405 schema: SourceSchema,
406}
407
408impl RangeSourceFactory {
409 /// Create a range source from `[start, end)`.
410 pub fn new(start: u64, end: u64) -> Self {
411 Self {
412 cursor: Arc::new(AtomicU64::new(start)),
413 end,
414 schema: SourceSchema {
415 name: "_range".into(),
416 projections: vec![("ordinal".into(), PortType::U64)],
417 extent: Some(end.saturating_sub(start)),
418 extent_outputs: None,
419 extent_limit: None,
420 cursor_kind: CursorKind::Range,
421 partition_output: None,
422 partitions: None,
423 },
424 }
425 }
426
427 /// Create a range source with a named schema.
428 pub fn named(name: &str, start: u64, end: u64) -> Self {
429 let mut factory = Self::new(start, end);
430 factory.schema.name = name.to_string();
431 factory
432 }
433}
434
435impl DataSourceFactory for RangeSourceFactory {
436 fn create_reader(&self) -> Box<dyn DataSource> {
437 Box::new(RangeSource {
438 cursor: self.cursor.clone(),
439 end: self.end,
440 consumed: 0,
441 schema: self.schema.clone(),
442 })
443 }
444
445 fn schema(&self) -> &SourceSchema {
446 &self.schema
447 }
448
449 fn global_consumed(&self) -> u64 {
450 let pos = self.cursor.load(Ordering::Relaxed);
451 let start = self.end.saturating_sub(self.schema.extent.unwrap_or(0));
452 pos.saturating_sub(start)
453 .min(self.schema.extent.unwrap_or(u64::MAX))
454 }
455
456 fn replay_contract(&self) -> SourceReplayContract {
457 SourceReplayContract::stable_ordinal(0)
458 }
459
460 fn rewind_for_poll(&self) -> bool {
461 // The round's whole state is the one cursor atomic, and every
462 // reservation is a `fetch_add` on it. The store and the
463 // reservations share that location's single modification
464 // order, so each reservation reads either the old round's
465 // cursor or a value derived from this store, and no other
466 // memory is published with it. Relaxed is therefore the
467 // cheapest ordering that is correct under any concurrency.
468 let start = self.end.saturating_sub(self.schema.extent.unwrap_or(0));
469 self.cursor.store(start, Ordering::Relaxed);
470 true
471 }
472}
473
474/// Per-fiber range reader. Pulls ordinals from a shared atomic cursor.
475struct RangeSource {
476 cursor: Arc<AtomicU64>,
477 end: u64,
478 consumed: u64,
479 schema: SourceSchema,
480}
481
482impl DataSource for RangeSource {
483 fn reserve(&mut self, stride: usize) -> Option<std::ops::Range<u64>> {
484 let base = self.cursor.fetch_add(stride as u64, Ordering::Relaxed);
485 if base >= self.end {
486 return None;
487 }
488 let actual_end = (base + stride as u64).min(self.end);
489 let count = actual_end - base;
490 self.consumed += count;
491 Some(base..actual_end)
492 }
493
494 fn render_item(&self, ordinal: u64) -> SourceItem {
495 SourceItem::ordinal(ordinal)
496 }
497
498 fn extent(&self) -> Option<u64> {
499 self.schema.extent
500 }
501
502 fn consumed(&self) -> u64 {
503 self.consumed
504 }
505
506 fn schema(&self) -> &SourceSchema {
507 &self.schema
508 }
509
510 fn replay_contract(&self) -> SourceReplayContract {
511 SourceReplayContract::stable_ordinal(0)
512 }
513}
514
515// =========================================================================
516// ExtendingRangeSource: runtime-growable extent
517// =========================================================================
518
519/// Snapshot of cursor state passed to an [`ExtensionPolicy`]
520/// at each end-reach decision point. Policies are pure
521/// predicates over this context — they don't carry their own
522/// clocks or counters.
523#[derive(Clone, Copy, Debug)]
524pub struct ExtensionContext {
525 /// Wall-clock milliseconds since the current round started: the
526 /// source factory's construction (typically phase start) for the
527 /// first round, and the latest `rewind_for_poll()` after that.
528 pub elapsed_ms: u64,
529 /// Global ordinals consumed so far — `cursor.load() - start`.
530 /// Pass count is `consumed / base`.
531 pub consumed: u64,
532 /// The `base` chunk size declared by the cursor. Used by
533 /// policies that reason in passes rather than raw counts.
534 pub base: u64,
535}
536
537impl ExtensionContext {
538 /// Convenience: integer pass count (consumed / base).
539 /// Returns 0 when `base == 0` (degenerate cursor).
540 pub fn passes(&self) -> u64 {
541 self.consumed.checked_div(self.base).unwrap_or(0)
542 }
543}
544
545/// Policy that decides whether to extend an
546/// `ExtendingRangeSource` when its current end is reached.
547///
548/// Implementations are pure predicates over the
549/// [`ExtensionContext`] — no internal state, no side effects.
550/// The source provides elapsed time and consumed counts; the
551/// policy returns `Some(delta)` to grow the extent or `None`
552/// to terminate. The source consults the policy under its round
553/// lock, so calls never overlap; a reader that finds the end
554/// reached after another reader's `None` consults it again.
555pub trait ExtensionPolicy: Send + Sync {
556 /// Decide how to extend (if at all) given the current
557 /// cursor context.
558 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64>;
559}
560
561/// Factory for ExtendingRangeSource. Differs from
562/// [`RangeSourceFactory`] in that `end` is an atomic that the
563/// extension policy may grow over the lifetime of the phase.
564/// `global_extent()` returns the CURRENT end so phase-status
565/// displays reflect any growth honestly.
566pub struct ExtendingRangeSourceFactory {
567 cursor: Arc<AtomicU64>,
568 end: Arc<AtomicU64>,
569 /// Serializes the round's end decisions: a reader's extension
570 /// and a rewind each read and write the cursor and end pair
571 /// under it, so an extension never acts on a pair a rewind has
572 /// half replaced. Reservations below the end do not take it.
573 round: Arc<std::sync::Mutex<()>>,
574 start: u64,
575 /// Per-pass chunk size — also the default extension delta
576 /// when the policy reports "continue". Exposed to the
577 /// policy via `ExtensionContext::base` so pass-count
578 /// predicates work.
579 base: u64,
580 /// Hard upper bound on growth (cursor_partitions.md §7.2). When a cursor is
581 /// narrowed by a partition (`until_elapsed(...) over p`),
582 /// the partition's end ordinal caps the extension: the
583 /// policy keeps making its time / pass / count decisions,
584 /// but the source terminates the moment the partition is
585 /// exhausted, whichever comes first. `None` = unbounded
586 /// (the policy alone decides).
587 max_end: Option<u64>,
588 /// The round's wall-clock baseline, which every reader of the
589 /// factory shares ([`RoundClock`]).
590 clock: RoundClock,
591 policy: Arc<dyn ExtensionPolicy>,
592 schema: SourceSchema,
593}
594
595/// The wall-clock baseline an extension policy's elapsed time is
596/// measured from. The baseline is the start of the current round:
597/// the factory's construction for the first round, and each
598/// `rewind_for_poll()` for the rounds after it. It is written and
599/// read under the factory's round lock, so every reader of the
600/// factory measures from the same instant.
601#[derive(Clone)]
602struct RoundClock {
603 /// The instant the factory was constructed.
604 epoch: std::time::Instant,
605 /// Nanoseconds from `epoch` to the start of the current round.
606 round_start_ns: Arc<AtomicU64>,
607}
608
609impl RoundClock {
610 fn new() -> Self {
611 Self {
612 epoch: std::time::Instant::now(),
613 round_start_ns: Arc::new(AtomicU64::new(0)),
614 }
615 }
616
617 fn now_ns(&self) -> u64 {
618 u64::try_from(self.epoch.elapsed().as_nanos()).unwrap_or(u64::MAX)
619 }
620
621 /// Start a new round at the current instant.
622 fn restart(&self) {
623 self.round_start_ns.store(self.now_ns(), Ordering::Relaxed);
624 }
625
626 /// Whole milliseconds since the current round started.
627 fn elapsed_ms(&self) -> u64 {
628 let start = self.round_start_ns.load(Ordering::Relaxed);
629 self.now_ns().saturating_sub(start) / 1_000_000
630 }
631}
632
633impl ExtendingRangeSourceFactory {
634 /// Build with an initial extent `[start, start + initial_extent)`.
635 /// The extension policy is consulted only when the cursor
636 /// reaches the end — the first stride consumed produces
637 /// indices in the initial range.
638 pub fn new(
639 name: &str,
640 start: u64,
641 initial_extent: u64,
642 policy: Arc<dyn ExtensionPolicy>,
643 ) -> Self {
644 let end = start.saturating_add(initial_extent);
645 Self {
646 cursor: Arc::new(AtomicU64::new(start)),
647 end: Arc::new(AtomicU64::new(end)),
648 round: Arc::new(std::sync::Mutex::new(())),
649 start,
650 base: initial_extent,
651 max_end: None,
652 clock: RoundClock::new(),
653 policy,
654 schema: SourceSchema {
655 name: name.to_string(),
656 projections: vec![("ordinal".into(), PortType::U64)],
657 extent: Some(initial_extent),
658 extent_outputs: None,
659 cursor_kind: CursorKind::Range,
660 extent_limit: None,
661 partition_output: None,
662 partitions: None,
663 },
664 }
665 }
666
667 /// Cap growth at `max_end` (absolute ordinal, exclusive) —
668 /// the partition-narrowing bound (cursor_partitions.md §7.2). The initial
669 /// extent is clamped too, so a base chunk larger than the
670 /// partition never reserves past it.
671 pub fn bounded(mut self, max_end: u64) -> Self {
672 self.max_end = Some(max_end);
673 let clamped = self.end.load(Ordering::Acquire).min(max_end);
674 self.end.store(clamped, Ordering::Release);
675 self
676 }
677}
678
679impl DataSourceFactory for ExtendingRangeSourceFactory {
680 fn create_reader(&self) -> Box<dyn DataSource> {
681 Box::new(ExtendingRangeSource {
682 cursor: self.cursor.clone(),
683 end: self.end.clone(),
684 round: self.round.clone(),
685 policy: self.policy.clone(),
686 start: self.start,
687 base: self.base,
688 max_end: self.max_end,
689 clock: self.clock.clone(),
690 consumed: 0,
691 schema: self.schema.clone(),
692 })
693 }
694
695 fn schema(&self) -> &SourceSchema {
696 &self.schema
697 }
698
699 fn global_consumed(&self) -> u64 {
700 let pos = self.cursor.load(Ordering::Relaxed);
701 pos.saturating_sub(self.start)
702 }
703
704 fn global_extent(&self) -> Option<u64> {
705 // Live extent — readers / status displays see growth.
706 Some(self.end.load(Ordering::Acquire).saturating_sub(self.start))
707 }
708
709 fn rewind_for_poll(&self) -> bool {
710 // A new round starts over: the end at the base chunk, capped by
711 // a partition bound (`max_end`), the cursor at the start, and
712 // the policy's elapsed time at zero, so the round extends as
713 // the first one did. Under the round lock no reader extends
714 // the end or reads the clock meanwhile. The
715 // end is written before the cursor, and the cursor with
716 // release, so a reader whose acquire load sees the new cursor
717 // also sees the new end; a reader that still sees the old
718 // cursor fails its claim or makes it in the old round.
719 let _round = self
720 .round
721 .lock()
722 .unwrap_or_else(std::sync::PoisonError::into_inner);
723 let mut end = self.start.saturating_add(self.base);
724 if let Some(cap) = self.max_end {
725 end = end.min(cap);
726 }
727 self.clock.restart();
728 self.end.store(end, Ordering::Relaxed);
729 self.cursor.store(self.start, Ordering::Release);
730 true
731 }
732
733 fn replay_contract(&self) -> SourceReplayContract {
734 SourceReplayContract::stable_ordinal(0)
735 }
736}
737
738/// Per-fiber reader for an extending range source. A reservation
739/// below the end is one compare-and-swap on the shared cursor and
740/// takes no lock. At the end the reader takes the factory's round
741/// lock, re-reads the cursor and end, and consults the policy only
742/// if the end is still reached, so one reader extends the round and
743/// the others claim from the new end.
744struct ExtendingRangeSource {
745 cursor: Arc<AtomicU64>,
746 end: Arc<AtomicU64>,
747 /// The factory's round lock
748 /// ([`ExtendingRangeSourceFactory::round`]).
749 round: Arc<std::sync::Mutex<()>>,
750 policy: Arc<dyn ExtensionPolicy>,
751 start: u64,
752 base: u64,
753 /// Partition cap — see
754 /// [`ExtendingRangeSourceFactory::bounded`].
755 max_end: Option<u64>,
756 clock: RoundClock,
757 consumed: u64,
758 schema: SourceSchema,
759}
760
761impl DataSource for ExtendingRangeSource {
762 fn reserve(&mut self, stride: usize) -> Option<std::ops::Range<u64>> {
763 loop {
764 let cur = self.cursor.load(Ordering::Acquire);
765 let end = self.end.load(Ordering::Acquire);
766 if cur < end {
767 // Try to claim [cur, min(cur+stride, end)).
768 let target = (cur.saturating_add(stride as u64)).min(end);
769 match self
770 .cursor
771 .compare_exchange(cur, target, Ordering::AcqRel, Ordering::Acquire)
772 {
773 Ok(_) => {
774 let count = target - cur;
775 self.consumed += count;
776 return Some(cur..target);
777 }
778 Err(_) => continue, // raced; retry
779 }
780 }
781 // The cursor has reached the end. Whether the round ends
782 // or grows is decided under the round lock over a fresh
783 // read of the pair, so the decision never acts on a pair a
784 // rewind has half replaced, and readers that reach the end
785 // together consult the policy one at a time: the first
786 // extends, and the others see the new end and claim.
787 let _round = self
788 .round
789 .lock()
790 .unwrap_or_else(std::sync::PoisonError::into_inner);
791 let cur = self.cursor.load(Ordering::Acquire);
792 let end = self.end.load(Ordering::Acquire);
793 if cur < end {
794 continue;
795 }
796 // A partition-bound cursor terminates the moment the
797 // partition is exhausted, whether or not the policy's
798 // time / pass / count target was reached — no policy
799 // consultation past the cap.
800 if let Some(max) = self.max_end
801 && end >= max
802 {
803 return None;
804 }
805 let ctx = ExtensionContext {
806 elapsed_ms: self.clock.elapsed_ms(),
807 consumed: end.saturating_sub(self.start),
808 base: self.base,
809 };
810 match self.policy.next_extension(&ctx) {
811 Some(delta) if delta > 0 => {
812 // Grow the end, clamped at the partition cap when
813 // one is set.
814 let mut new_end = end.saturating_add(delta);
815 if let Some(max) = self.max_end {
816 new_end = new_end.min(max);
817 }
818 if new_end == end {
819 return None;
820 }
821 self.end.store(new_end, Ordering::Release);
822 continue;
823 }
824 _ => return None,
825 }
826 }
827 }
828
829 fn render_item(&self, ordinal: u64) -> SourceItem {
830 SourceItem::ordinal(ordinal)
831 }
832
833 fn extent(&self) -> Option<u64> {
834 // Live extent — same convention as the factory's
835 // `global_extent`: subscribers see the current ceiling
836 // even after it grows.
837 Some(self.end.load(Ordering::Acquire).saturating_sub(self.start))
838 }
839
840 fn consumed(&self) -> u64 {
841 self.consumed
842 }
843
844 fn schema(&self) -> &SourceSchema {
845 &self.schema
846 }
847
848 fn replay_contract(&self) -> SourceReplayContract {
849 SourceReplayContract::stable_ordinal(0)
850 }
851}
852
853// =========================================================================
854// Cursors: provenance-driven cursor targeting
855// =========================================================================
856
857/// A cursor target: a DataSource reader paired with its Polydat input index.
858struct CursorTarget {
859 /// The DataSource reader that provides values for this cursor.
860 reader: Box<dyn DataSource>,
861 /// The Polydat input index where the cursor's ordinal is injected.
862 input_index: usize,
863 /// Source name (for diagnostics).
864 source_name: String,
865 /// Stable process-independent identity used in ordinal packet stamps.
866 stream_id: u64,
867}
868
869/// Ownership token for one range reserved from a perfect ordinal source.
870///
871/// The token contains no reader borrow or rendered payload because the source
872/// contract proves that the scalar value is the ordinal itself. Moving this
873/// value transfers responsibility for every ordinal in `range`; the batch
874/// executor must drain all of them or process the remainder through its scalar
875/// recovery path. The shared source cursor is never rewound.
876#[derive(Debug)]
877pub struct OrdinalBatchLease {
878 source_name: String,
879 stream_id: u64,
880 source_generation: u64,
881 input_index: usize,
882 sequence: u64,
883 range: std::ops::Range<u64>,
884}
885
886impl OrdinalBatchLease {
887 /// The cursor the lease reserves from.
888 pub fn source_name(&self) -> &str {
889 &self.source_name
890 }
891
892 /// The stream the ordinals belong to.
893 pub const fn stream_id(&self) -> u64 {
894 self.stream_id
895 }
896
897 /// The generation of the source's values.
898 pub const fn source_generation(&self) -> u64 {
899 self.source_generation
900 }
901
902 /// The kernel input the cursor's ordinal is written to.
903 pub const fn input_index(&self) -> usize {
904 self.input_index
905 }
906
907 /// The lease's position in reservation order.
908 pub const fn sequence(&self) -> u64 {
909 self.sequence
910 }
911
912 /// The ordinals reserved.
913 pub fn range(&self) -> std::ops::Range<u64> {
914 self.range.clone()
915 }
916
917 /// How many ordinals the lease holds.
918 pub const fn len_u64(&self) -> u64 {
919 self.range.end - self.range.start
920 }
921
922 /// Whether the lease holds no ordinal.
923 pub const fn is_empty(&self) -> bool {
924 self.range.start == self.range.end
925 }
926}
927
928#[derive(Clone, Debug, PartialEq, Eq)]
929/// Why a cursor batch could not be reserved.
930pub enum CursorBatchError {
931 /// No ordinals were asked for.
932 ZeroDemand,
933 /// The reservation names several targets; a batch drives one.
934 RequiresSingleTarget {
935 /// The targets named.
936 targets: usize,
937 },
938 /// The source is not a perfect ordinal stream.
939 SourceNotPerfectOrdinal {
940 /// The source's name.
941 source: String,
942 },
943}
944
945impl std::fmt::Display for CursorBatchError {
946 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
947 match self {
948 Self::ZeroDemand => f.write_str("ordinal batch demand must be greater than zero"),
949 Self::RequiresSingleTarget { targets } => write!(
950 f,
951 "Tier-1 ordinal batching requires exactly one cursor target, found {targets}"
952 ),
953 Self::SourceNotPerfectOrdinal { source } => write!(
954 f,
955 "source '{source}' does not declare the perfect ordinal replay contract"
956 ),
957 }
958 }
959}
960
961impl std::error::Error for CursorBatchError {}
962
963/// Provenance-driven advancer that targets only the cursor nodes
964/// relevant to a specific set of output fields.
965///
966/// Built by a host from the output names it will read, tracing Polydat
967/// provenance from those fields back to root cursor nodes. Only those cursors
968/// advance — unused cursors are left untouched.
969pub struct Cursors {
970 targets: Vec<CursorTarget>,
971 /// Last items read from each target (for injecting into Polydat state).
972 last_items: Vec<Option<SourceItem>>,
973 /// Total advances performed.
974 advances: u64,
975 /// Reservation order for owned ordinal leases from this cursor set.
976 next_batch_sequence: u64,
977}
978
979impl Cursors {
980 /// Build an advancer from a Polydat program, a set of output field names,
981 /// and a map of source name → DataSourceFactory.
982 ///
983 /// Traces provenance: for each field, finds the output node, gets its
984 /// input provenance bitmask, and identifies which inputs are cursor
985 /// ordinals (matching `{source}__ordinal` pattern). Creates a reader
986 /// for each targeted source.
987 pub fn for_fields(
988 program: &crate::kernel::PolydatProgram,
989 field_names: &[&str],
990 source_factories: &std::collections::HashMap<String, Arc<dyn DataSourceFactory>>,
991 ) -> Self {
992 // Collect the union of input provenance for all referenced
993 // fields — exact at any input count (multi-word ProvMask;
994 // the former one-word form over-matched cursor inputs >= 63).
995 let mut combined_provenance = crate::kernel::ProvMask::empty();
996 for name in field_names {
997 if let Some((node_idx, _)) = program.resolve_output(name)
998 && let Some(prov) = program.input_provenance_for(node_idx)
999 {
1000 combined_provenance.union_with(prov);
1001 }
1002 }
1003
1004 // Find cursor inputs in the provenance
1005 let input_names = program.input_names();
1006 let mut targets = Vec::new();
1007 let mut seen_sources = std::collections::HashSet::new();
1008
1009 for (idx, input_name) in input_names.iter().enumerate() {
1010 if !combined_provenance.contains(idx) {
1011 continue;
1012 }
1013
1014 // Check if this input is a source projection ({source}__ordinal)
1015 if let Some(source_name) = input_name.strip_suffix("__ordinal") {
1016 if seen_sources.contains(source_name) {
1017 continue;
1018 }
1019 seen_sources.insert(source_name.to_string());
1020
1021 if let Some(factory) = source_factories.get(source_name) {
1022 let stream_id =
1023 xxhash_rust::xxh3::xxh3_64(format!("{source_name}@{idx}").as_bytes());
1024 targets.push(CursorTarget {
1025 reader: factory.create_reader(),
1026 input_index: idx,
1027 source_name: source_name.to_string(),
1028 stream_id,
1029 });
1030 }
1031 }
1032 }
1033
1034 let target_count = targets.len();
1035 Cursors {
1036 targets,
1037 last_items: vec![None; target_count],
1038 advances: 0,
1039 next_batch_sequence: 0,
1040 }
1041 }
1042
1043 /// Reserve one owned batch from a single perfect ordinal cursor.
1044 ///
1045 /// This is the source-side entry point for Tier-1 SIMD execution. It is
1046 /// intentionally unavailable for multiple cursor targets or sources that
1047 /// merely happen to be monotonic: only the explicit perfect-ordinal
1048 /// contract permits payload-free replay from the returned range.
1049 pub fn reserve_ordinal_batch(
1050 &mut self,
1051 demand: usize,
1052 ) -> Result<Option<OrdinalBatchLease>, CursorBatchError> {
1053 if demand == 0 {
1054 return Err(CursorBatchError::ZeroDemand);
1055 }
1056 if self.targets.len() != 1 {
1057 return Err(CursorBatchError::RequiresSingleTarget {
1058 targets: self.targets.len(),
1059 });
1060 }
1061
1062 let target = &mut self.targets[0];
1063 let contract = target.reader.replay_contract();
1064 if !contract.is_perfect_ordinal() {
1065 return Err(CursorBatchError::SourceNotPerfectOrdinal {
1066 source: target.source_name.clone(),
1067 });
1068 }
1069 let Some(range) = target.reader.reserve(demand) else {
1070 return Ok(None);
1071 };
1072 let claimed = range.end - range.start;
1073 let lease = OrdinalBatchLease {
1074 source_name: target.source_name.clone(),
1075 stream_id: target.stream_id,
1076 source_generation: contract.generation,
1077 input_index: target.input_index,
1078 sequence: self.next_batch_sequence,
1079 range,
1080 };
1081 self.next_batch_sequence = self.next_batch_sequence.wrapping_add(1);
1082 self.advances = self.advances.saturating_add(claimed);
1083 Ok(Some(lease))
1084 }
1085
1086 /// Advance all targeted cursors. Returns `false` if any targeted
1087 /// cursor is exhausted (no more data).
1088 ///
1089 /// After advancing, the new ordinals and field projections are
1090 /// available via `inject_into_state()`.
1091 pub fn advance(&mut self) -> bool {
1092 for (i, target) in self.targets.iter_mut().enumerate() {
1093 match target.reader.next() {
1094 Some(item) => {
1095 self.last_items[i] = Some(item);
1096 }
1097 None => return false, // this cursor exhausted
1098 }
1099 }
1100 self.advances += 1;
1101 true
1102 }
1103
1104 /// Inject the current cursor values into a Polydat state.
1105 ///
1106 /// Sets each cursor's ordinal at its input index. Field projections
1107 /// are not written here; a host reads them from `last_items()`.
1108 pub fn inject_into_state(&self, state: &mut crate::kernel::PolydatState) {
1109 for (i, target) in self.targets.iter().enumerate() {
1110 if let Some(ref item) = self.last_items[i] {
1111 state.set_input(target.input_index, crate::ast::Value::U64(item.ordinal));
1112 // Field projections (e.g. base__vector) are the host's to
1113 // write from `last_items()`; only the ordinal lands here.
1114 }
1115 }
1116 }
1117
1118 /// The last items read from all targets, for a host to project fields from.
1119 pub fn last_items(&self) -> &[Option<SourceItem>] {
1120 &self.last_items
1121 }
1122
1123 /// Known extent of the driving cursor (smallest among targeted
1124 /// cursors with known extent). Used for progress reporting.
1125 pub fn extent(&self) -> Option<u64> {
1126 self.targets.iter().filter_map(|t| t.reader.extent()).min()
1127 }
1128
1129 /// Total advances performed so far.
1130 pub fn consumed(&self) -> u64 {
1131 self.advances
1132 }
1133
1134 /// Number of targeted cursors.
1135 pub fn target_count(&self) -> usize {
1136 self.targets.len()
1137 }
1138
1139 /// Whether this advancer has any targets.
1140 pub fn is_empty(&self) -> bool {
1141 self.targets.is_empty()
1142 }
1143}
1144
1145// =========================================================================
1146// Built-in ExtensionPolicy implementations
1147// =========================================================================
1148//
1149// Each policy is a pure predicate over `ExtensionContext` and
1150// returns `Some(delta)` when the cursor should continue or
1151// `None` when it should stop. The delta is independent of the
1152// stop condition — workloads can extend in chunks unrelated to
1153// the cursor's base size (e.g. base=10000 but delta=1000 to
1154// check the condition more often).
1155
1156/// Extend while `ctx.elapsed_ms < min_ms`, projecting the
1157/// remaining work from the observed rate and committing it in
1158/// one batch.
1159///
1160/// Per call: estimate `repeats = floor((consumed * remaining_ms
1161/// / elapsed_ms) / base)` — the number of base-sized passes
1162/// needed to fill the time remaining at the current rate. Apply
1163/// a 5% under-bias (`repeats * 95 / 100`) so successive
1164/// extensions converge from below the target rather than
1165/// overshooting once. Commit `repeats * base` cycles.
1166///
1167/// The under-bias gives a geometric convergence: each call
1168/// covers ~95% of the time remaining, so 3–4 calls typically
1169/// land within 0.01% of the time budget. The first call (no
1170/// rate signal yet — `elapsed_ms == 0` or `consumed == 0`)
1171/// falls back to one base chunk via the `delta` field.
1172pub struct UntilElapsedPolicy {
1173 /// The elapsed milliseconds to reach.
1174 pub min_ms: u64,
1175 /// The extension step for the first call, before a rate is known.
1176 pub delta: u64,
1177}
1178
1179impl ExtensionPolicy for UntilElapsedPolicy {
1180 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1181 if ctx.elapsed_ms >= self.min_ms {
1182 return None;
1183 }
1184 // No rate signal yet — first end-reach with elapsed time
1185 // below the clock resolution, or a degenerate empty
1186 // cursor. Step by the declared `delta` (which the
1187 // executor sets to `base` when the workload didn't
1188 // supply an explicit `delta` arg) so the next call has
1189 // a rate measurement to project from.
1190 if ctx.elapsed_ms == 0 || ctx.consumed == 0 {
1191 return Some(self.delta.max(1));
1192 }
1193 let remaining_ms = self.min_ms - ctx.elapsed_ms;
1194 // u128 saturating math: consumed * remaining_ms can
1195 // overflow u64 for long-running phases with high
1196 // throughput (e.g. 10^9 ops over 10^4 ms).
1197 let est_remaining =
1198 (ctx.consumed as u128).saturating_mul(remaining_ms as u128) / (ctx.elapsed_ms as u128);
1199 let biased = est_remaining.saturating_mul(95) / 100;
1200 let base = ctx.base.max(1) as u128;
1201 let repeats = biased / base;
1202 if repeats == 0 {
1203 return None;
1204 }
1205 let delta = repeats.saturating_mul(base);
1206 Some(u64::try_from(delta).unwrap_or(u64::MAX))
1207 }
1208}
1209
1210/// Extend by `delta` while `ctx.passes() < min_passes`. Passes
1211/// are whole multiples of `ctx.base`.
1212pub struct UntilPassesPolicy {
1213 /// The passes to reach.
1214 pub min_passes: u64,
1215 /// The extension step.
1216 pub delta: u64,
1217}
1218
1219impl ExtensionPolicy for UntilPassesPolicy {
1220 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1221 if ctx.passes() < self.min_passes {
1222 Some(self.delta)
1223 } else {
1224 None
1225 }
1226 }
1227}
1228
1229/// Extend by `delta` while `ctx.consumed < min_count`.
1230pub struct UntilCountPolicy {
1231 /// The consumed count to reach.
1232 pub min_count: u64,
1233 /// The extension step.
1234 pub delta: u64,
1235}
1236
1237impl ExtensionPolicy for UntilCountPolicy {
1238 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1239 if ctx.consumed < self.min_count {
1240 Some(self.delta)
1241 } else {
1242 None
1243 }
1244 }
1245}
1246
1247/// Continue (extend) while ALL child policies say continue.
1248/// Stops the moment any child policy returns `None`. The
1249/// delta is the minimum of the children's deltas — conservative
1250/// step size keeps any single condition from over-shooting its
1251/// stop point.
1252pub struct AndPolicy {
1253 /// The policies that must all ask to continue.
1254 pub policies: Vec<Arc<dyn ExtensionPolicy>>,
1255}
1256
1257impl ExtensionPolicy for AndPolicy {
1258 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1259 let mut min_delta = u64::MAX;
1260 for p in &self.policies {
1261 let d = p.next_extension(ctx)?;
1262 min_delta = min_delta.min(d);
1263 }
1264 if min_delta == u64::MAX || min_delta == 0 {
1265 None
1266 } else {
1267 Some(min_delta)
1268 }
1269 }
1270}
1271
1272/// Continue (extend) while ANY child policy says continue.
1273/// Stops only when every child policy returns `None`. The
1274/// delta is the maximum of the policies that said continue —
1275/// matches the most aggressive child still pushing forward.
1276pub struct OrPolicy {
1277 /// The policies of which any may ask to continue.
1278 pub policies: Vec<Arc<dyn ExtensionPolicy>>,
1279}
1280
1281impl ExtensionPolicy for OrPolicy {
1282 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1283 let mut max_delta: Option<u64> = None;
1284 for p in &self.policies {
1285 if let Some(d) = p.next_extension(ctx) {
1286 max_delta = Some(max_delta.map(|m| m.max(d)).unwrap_or(d));
1287 }
1288 }
1289 max_delta
1290 }
1291}
1292
1293/// Back-compat alias for the original time-only policy.
1294/// Constructs an [`UntilElapsedPolicy`] with delta = base.
1295pub struct TimeElapsedPolicy {
1296 inner: UntilElapsedPolicy,
1297}
1298
1299impl TimeElapsedPolicy {
1300 /// A policy extending in steps of `base` until `min_ms` have elapsed.
1301 pub fn new(base: u64, min_ms: u64) -> Self {
1302 Self {
1303 inner: UntilElapsedPolicy {
1304 min_ms,
1305 delta: base,
1306 },
1307 }
1308 }
1309}
1310
1311impl ExtensionPolicy for TimeElapsedPolicy {
1312 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1313 self.inner.next_extension(ctx)
1314 }
1315}
1316
1317#[cfg(test)]
1318mod tests {
1319 use super::*;
1320
1321 #[test]
1322 fn range_sources_declare_the_perfect_ordinal_replay_contract() {
1323 let factory = RangeSourceFactory::new(10, 20);
1324 assert!(factory.replay_contract().is_perfect_ordinal());
1325 let reader = factory.create_reader();
1326 assert!(reader.replay_contract().is_perfect_ordinal());
1327 let item = reader.render_item(13);
1328 assert_eq!(item.ordinal, 13);
1329 assert!(item.fields.is_empty());
1330 }
1331
1332 /// Policy that always extends by `base` — stands in for a
1333 /// time/pass policy whose target is far away, so the
1334 /// partition cap is the only thing that can stop growth.
1335 struct AlwaysExtend;
1336 impl ExtensionPolicy for AlwaysExtend {
1337 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1338 Some(ctx.base.max(1))
1339 }
1340 }
1341
1342 #[test]
1343 fn extending_source_bounded_terminates_at_partition_end() {
1344 // `until_*(base, ...) over p` — base-sized chunks
1345 // walk within the partition; the cap stops growth even
1346 // though the policy would keep extending. Partition
1347 // [100, 125) with base 10 → 10 + 10 + 5, then exhausted.
1348 let factory =
1349 ExtendingRangeSourceFactory::new("q", 100, 10, Arc::new(AlwaysExtend)).bounded(125);
1350 let mut reader = factory.create_reader();
1351 let mut total = 0u64;
1352 let mut last_end = 100;
1353 while let Some(r) = reader.reserve(7) {
1354 assert!(r.end <= 125, "reservation past the partition cap: {r:?}");
1355 assert_eq!(r.start, last_end, "contiguous reservations");
1356 last_end = r.end;
1357 total += r.end - r.start;
1358 }
1359 assert_eq!(total, 25, "exactly the partition's cardinality");
1360 assert_eq!(last_end, 125);
1361 }
1362
1363 #[test]
1364 fn extending_source_bounded_clamps_oversized_base() {
1365 // A base chunk larger than the partition never reserves
1366 // past it.
1367 let factory =
1368 ExtendingRangeSourceFactory::new("q", 0, 1000, Arc::new(AlwaysExtend)).bounded(30);
1369 let mut reader = factory.create_reader();
1370 let r = reader.reserve(usize::MAX).unwrap();
1371 assert_eq!(r, 0..30);
1372 assert!(reader.reserve(1).is_none());
1373 }
1374
1375 #[test]
1376 fn extending_source_unbounded_keeps_policy_semantics() {
1377 // Without a cap the policy alone decides — three
1378 // extensions of a terminating policy.
1379 struct NTimes(std::sync::atomic::AtomicU64);
1380 impl ExtensionPolicy for NTimes {
1381 fn next_extension(&self, ctx: &ExtensionContext) -> Option<u64> {
1382 if self.0.fetch_add(1, Ordering::Relaxed) < 3 {
1383 Some(ctx.base)
1384 } else {
1385 None
1386 }
1387 }
1388 }
1389 let factory = ExtendingRangeSourceFactory::new(
1390 "q",
1391 0,
1392 10,
1393 Arc::new(NTimes(std::sync::atomic::AtomicU64::new(0))),
1394 );
1395 let mut reader = factory.create_reader();
1396 let mut total = 0u64;
1397 while let Some(r) = reader.reserve(64) {
1398 total += r.end - r.start;
1399 }
1400 assert_eq!(total, 40, "initial 10 + three 10-ordinal extensions");
1401 }
1402
1403 #[test]
1404 fn range_source_yields_ordinals() {
1405 let factory = RangeSourceFactory::new(0, 5);
1406 let mut reader = factory.create_reader();
1407 assert_eq!(reader.extent(), Some(5));
1408
1409 for i in 0..5 {
1410 let item = reader.next().unwrap();
1411 assert_eq!(item.ordinal, i);
1412 assert!(item.fields.is_empty());
1413 }
1414 assert!(reader.next().is_none());
1415 assert_eq!(reader.consumed(), 5);
1416 }
1417
1418 #[test]
1419 fn range_source_chunk() {
1420 let factory = RangeSourceFactory::new(0, 10);
1421 let mut reader = factory.create_reader();
1422
1423 let chunk = reader.next_chunk(3);
1424 assert_eq!(chunk.len(), 3);
1425 assert_eq!(chunk[0].ordinal, 0);
1426 assert_eq!(chunk[2].ordinal, 2);
1427
1428 let chunk = reader.next_chunk(100);
1429 assert_eq!(chunk.len(), 7); // only 7 remaining
1430 assert_eq!(chunk[0].ordinal, 3);
1431 assert_eq!(chunk[6].ordinal, 9);
1432
1433 let chunk = reader.next_chunk(1);
1434 assert!(chunk.is_empty()); // exhausted
1435 }
1436
1437 #[test]
1438 fn range_source_concurrent_readers() {
1439 let factory = RangeSourceFactory::new(0, 100);
1440 let mut r1 = factory.create_reader();
1441 let mut r2 = factory.create_reader();
1442
1443 // Each reader gets unique ordinals from the shared cursor
1444 let a = r1.next().unwrap().ordinal;
1445 let b = r2.next().unwrap().ordinal;
1446 assert_ne!(a, b);
1447
1448 // Drain both readers
1449 let mut total = 2;
1450 while r1.next().is_some() {
1451 total += 1;
1452 }
1453 while r2.next().is_some() {
1454 total += 1;
1455 }
1456 assert_eq!(total, 100);
1457 }
1458
1459 #[test]
1460 fn source_item_field_access() {
1461 let item = SourceItem::with_fields(
1462 42,
1463 vec![
1464 ("name".into(), Value::Str("test".into())),
1465 ("score".into(), Value::F64(0.95)),
1466 ],
1467 );
1468 assert_eq!(item.ordinal, 42);
1469 assert_eq!(item.field("name"), Some(&Value::Str("test".into())));
1470 assert_eq!(item.field("score"), Some(&Value::F64(0.95)));
1471 assert_eq!(item.field("missing"), None);
1472 }
1473
1474 #[test]
1475 fn range_source_named() {
1476 let factory = RangeSourceFactory::named("users", 0, 1000);
1477 assert_eq!(factory.schema().name, "users");
1478 assert_eq!(factory.schema().extent, Some(1000));
1479 }
1480
1481 // ── ExtendingRangeSource ─────────────────────────────────
1482
1483 /// Trivial extension policy for tests: lets the caller
1484 /// decide how many extensions remain. Each call to
1485 /// `next_extension` decrements the count.
1486 struct FixedExtensions {
1487 delta: u64,
1488 remaining: std::sync::atomic::AtomicU64,
1489 }
1490
1491 impl FixedExtensions {
1492 fn new(delta: u64, times: u64) -> Self {
1493 Self {
1494 delta,
1495 remaining: std::sync::atomic::AtomicU64::new(times),
1496 }
1497 }
1498 }
1499
1500 impl ExtensionPolicy for FixedExtensions {
1501 fn next_extension(&self, _ctx: &ExtensionContext) -> Option<u64> {
1502 let prev = self.remaining.fetch_sub(1, Ordering::Relaxed);
1503 if prev == 0 || prev > i64::MAX as u64 {
1504 self.remaining.store(0, Ordering::Relaxed);
1505 None
1506 } else {
1507 Some(self.delta)
1508 }
1509 }
1510 }
1511
1512 #[test]
1513 fn extending_source_consumes_initial_extent_then_extends() {
1514 // Initial extent = 5; policy extends once by 5 more.
1515 // Total expected: 10 ordinals consumed.
1516 let policy = Arc::new(FixedExtensions::new(5, 1));
1517 let factory = ExtendingRangeSourceFactory::new("ext", 0, 5, policy);
1518 let mut reader = factory.create_reader();
1519 let mut got: Vec<u64> = Vec::new();
1520 while let Some(item) = reader.next() {
1521 got.push(item.ordinal);
1522 if got.len() > 50 {
1523 panic!("runaway extension");
1524 }
1525 }
1526 assert_eq!(got, (0..10).collect::<Vec<u64>>());
1527 assert_eq!(reader.consumed(), 10);
1528 }
1529
1530 #[test]
1531 fn extending_source_zero_extensions_behaves_like_fixed_range() {
1532 let policy = Arc::new(FixedExtensions::new(0, 0));
1533 let factory = ExtendingRangeSourceFactory::new("ext", 0, 3, policy);
1534 let mut reader = factory.create_reader();
1535 let mut got: Vec<u64> = Vec::new();
1536 while let Some(item) = reader.next() {
1537 got.push(item.ordinal);
1538 }
1539 assert_eq!(got, vec![0, 1, 2]);
1540 }
1541
1542 #[test]
1543 fn extending_source_global_extent_grows_after_extension() {
1544 // global_extent must reflect the CURRENT end so phase
1545 // status displays show growth honestly.
1546 let policy = Arc::new(FixedExtensions::new(10, 2));
1547 let factory = ExtendingRangeSourceFactory::new("ext", 0, 5, policy);
1548 assert_eq!(factory.global_extent(), Some(5));
1549 let mut reader = factory.create_reader();
1550 // Drain the first 5 to force an extension.
1551 for _ in 0..5 {
1552 reader.next().unwrap();
1553 }
1554 // Trigger the extension by attempting one more pull.
1555 let _ = reader.next().unwrap();
1556 assert_eq!(
1557 factory.global_extent(),
1558 Some(15),
1559 "extent should grow by the extension delta"
1560 );
1561 }
1562
1563 #[test]
1564 fn extending_source_chunked_reservation_caps_at_current_end() {
1565 // A stride request that crosses the current `end` must
1566 // return a SHORT range up to `end`, not an over-claim.
1567 // Otherwise the next reservation would jump past the
1568 // extended range and lose ordinals.
1569 let policy = Arc::new(FixedExtensions::new(5, 1));
1570 let factory = ExtendingRangeSourceFactory::new("ext", 0, 3, policy);
1571 let mut reader = factory.create_reader();
1572 let range = reader.reserve(10).expect("first reserve");
1573 // Initial end was 3; the reservation must cap there.
1574 assert_eq!(range, 0..3);
1575 // Next reserve triggers the extension and consumes the
1576 // remainder.
1577 let range = reader.reserve(10).expect("second reserve");
1578 assert_eq!(range, 3..8);
1579 }
1580
1581 fn ctx_at(elapsed_ms: u64, consumed: u64, base: u64) -> ExtensionContext {
1582 ExtensionContext {
1583 elapsed_ms,
1584 consumed,
1585 base,
1586 }
1587 }
1588
1589 #[test]
1590 fn until_elapsed_policy_bootstrap_falls_back_to_delta_when_no_rate_signal() {
1591 // First end-reach: consumed or elapsed effectively zero —
1592 // no rate to project from. Policy returns `delta` exactly
1593 // so the next call has a measurement to work with.
1594 let policy = UntilElapsedPolicy {
1595 min_ms: 50,
1596 delta: 7,
1597 };
1598 assert_eq!(policy.next_extension(&ctx_at(0, 0, 0)), Some(7));
1599 assert_eq!(policy.next_extension(&ctx_at(49, 0, 0)), Some(7));
1600 assert_eq!(policy.next_extension(&ctx_at(0, 0, 100)), Some(7));
1601 }
1602
1603 #[test]
1604 fn until_elapsed_policy_stops_at_or_past_min_ms() {
1605 let policy = UntilElapsedPolicy {
1606 min_ms: 50,
1607 delta: 7,
1608 };
1609 assert_eq!(policy.next_extension(&ctx_at(50, 100, 100)), None);
1610 assert_eq!(policy.next_extension(&ctx_at(1000, 100, 100)), None);
1611 }
1612
1613 #[test]
1614 fn until_elapsed_policy_projects_remaining_cycles_with_under_bias() {
1615 // base=100, consumed=100 (one pass), elapsed=200ms,
1616 // min_ms=1000 → remaining=800ms. Rate = 100/200 = 0.5
1617 // cycles/ms; project 0.5 * 800 = 400 cycles. Under-bias
1618 // 5% → 380. Round to multiples of base (100) → 3 repeats
1619 // → 300 cycles.
1620 let policy = UntilElapsedPolicy {
1621 min_ms: 1000,
1622 delta: 100,
1623 };
1624 assert_eq!(
1625 policy.next_extension(&ctx_at(200, 100, 100)),
1626 Some(300),
1627 "expected 3-pass batch (380 under-biased, floored to base multiples)",
1628 );
1629 }
1630
1631 #[test]
1632 fn until_elapsed_policy_returns_none_when_remaining_is_under_one_pass() {
1633 // base=100, consumed=10000 over 990ms, min_ms=1000 →
1634 // remaining=10ms. Project 10000/990*10 ≈ 101 cycles.
1635 // Under-bias → 95. Rounded to base multiples → 0 repeats.
1636 // Policy terminates rather than under-shooting the budget
1637 // with a wasted partial pass.
1638 let policy = UntilElapsedPolicy {
1639 min_ms: 1000,
1640 delta: 100,
1641 };
1642 assert_eq!(policy.next_extension(&ctx_at(990, 10000, 100)), None);
1643 }
1644
1645 #[test]
1646 fn until_elapsed_policy_converges_geometrically() {
1647 // Sanity: simulate the extension loop. Each iteration
1648 // commits ~95% of the remaining time budget; the
1649 // residual is bounded below by the base-pass rounding,
1650 // so convergence is one-base-coarse rather than
1651 // arbitrarily tight.
1652 let policy = UntilElapsedPolicy {
1653 min_ms: 1000,
1654 delta: 10,
1655 };
1656 // Pretend the first pass took 10ms (rate = 1 cycle/ms).
1657 let mut elapsed = 10u64;
1658 let mut consumed = 10u64;
1659 let mut iters = 0;
1660 while let Some(delta) = policy.next_extension(&ctx_at(elapsed, consumed, 10)) {
1661 iters += 1;
1662 assert!(iters < 10, "geometric series should converge fast");
1663 consumed += delta;
1664 // Same rate: 1 cycle/ms.
1665 elapsed += delta;
1666 }
1667 assert!(elapsed <= 1000, "must under-shoot, got elapsed={elapsed}");
1668 // Residual ≤ ~5% (under-bias) + 1 base pass (rounding) ≈ 6% of target.
1669 assert!(
1670 elapsed >= 940,
1671 "must come within ~6% of target, got {elapsed}"
1672 );
1673 }
1674
1675 #[test]
1676 fn until_passes_policy_counts_in_base_multiples() {
1677 let policy = UntilPassesPolicy {
1678 min_passes: 3,
1679 delta: 100,
1680 };
1681 // 0 passes done — extend.
1682 assert_eq!(policy.next_extension(&ctx_at(0, 0, 100)), Some(100));
1683 // 2 passes done (200 consumed @ base=100) — extend.
1684 assert_eq!(policy.next_extension(&ctx_at(0, 200, 100)), Some(100));
1685 // 3 passes done — stop.
1686 assert_eq!(policy.next_extension(&ctx_at(0, 300, 100)), None);
1687 // 4 passes done — stop.
1688 assert_eq!(policy.next_extension(&ctx_at(0, 400, 100)), None);
1689 }
1690
1691 #[test]
1692 fn until_count_policy_uses_raw_consumed() {
1693 let policy = UntilCountPolicy {
1694 min_count: 250,
1695 delta: 50,
1696 };
1697 assert_eq!(policy.next_extension(&ctx_at(0, 0, 100)), Some(50));
1698 assert_eq!(policy.next_extension(&ctx_at(0, 249, 100)), Some(50));
1699 assert_eq!(policy.next_extension(&ctx_at(0, 250, 100)), None);
1700 }
1701
1702 #[test]
1703 fn and_policy_stops_when_any_child_stops() {
1704 // time<5000 AND passes<3.
1705 let time = Arc::new(UntilElapsedPolicy {
1706 min_ms: 5000,
1707 delta: 10,
1708 });
1709 let passes = Arc::new(UntilPassesPolicy {
1710 min_passes: 3,
1711 delta: 20,
1712 });
1713 let and = AndPolicy {
1714 policies: vec![time, passes],
1715 };
1716 // Both still want to continue: delta = min(10, 20) = 10.
1717 assert_eq!(and.next_extension(&ctx_at(0, 0, 100)), Some(10));
1718 // Time done — stop.
1719 assert_eq!(and.next_extension(&ctx_at(5000, 0, 100)), None);
1720 // Passes done — stop.
1721 assert_eq!(and.next_extension(&ctx_at(0, 300, 100)), None);
1722 }
1723
1724 #[test]
1725 fn or_policy_continues_if_any_child_continues() {
1726 let time = Arc::new(UntilElapsedPolicy {
1727 min_ms: 5000,
1728 delta: 10,
1729 });
1730 let passes = Arc::new(UntilPassesPolicy {
1731 min_passes: 3,
1732 delta: 20,
1733 });
1734 let or = OrPolicy {
1735 policies: vec![time, passes],
1736 };
1737 // Both want to continue → delta = max(10, 20) = 20.
1738 assert_eq!(or.next_extension(&ctx_at(0, 0, 100)), Some(20));
1739 // Time done, passes still wants → delta = 20.
1740 assert_eq!(or.next_extension(&ctx_at(5000, 0, 100)), Some(20));
1741 // Passes done, time still wants → delta = 10.
1742 assert_eq!(or.next_extension(&ctx_at(0, 300, 100)), Some(10));
1743 // Both done → stop.
1744 assert_eq!(or.next_extension(&ctx_at(5000, 300, 100)), None);
1745 }
1746
1747 #[test]
1748 fn time_elapsed_policy_compat_alias_still_works() {
1749 let p = TimeElapsedPolicy::new(7, 1);
1750 std::thread::sleep(std::time::Duration::from_millis(20));
1751 assert_eq!(p.next_extension(&ctx_at(20, 0, 0)), None);
1752 }
1753}
1754
1755#[cfg(test)]
1756mod rewind_tests {
1757 use super::*;
1758
1759 fn drain(reader: &mut dyn DataSource) -> Vec<u64> {
1760 let mut out = Vec::new();
1761 while let Some(item) = reader.next() {
1762 out.push(item.ordinal);
1763 }
1764 out
1765 }
1766
1767 /// A factory that rewinds hands out the same ordinal range again:
1768 /// what a host that re-runs a source between rounds relies on.
1769 #[test]
1770 fn a_range_factory_rewinds_to_the_same_ordinals() {
1771 let factory = RangeSourceFactory::new(3, 6);
1772 assert_eq!(drain(factory.create_reader().as_mut()), vec![3, 4, 5]);
1773 assert!(factory.create_reader().next().is_none());
1774 assert!(factory.rewind_for_poll());
1775 assert_eq!(drain(factory.create_reader().as_mut()), vec![3, 4, 5]);
1776 }
1777
1778 /// The extending factory rewinds too: the cursor to the start and
1779 /// the extent to the base chunk, capped by a partition bound.
1780 #[test]
1781 fn an_extending_factory_rewinds_to_its_base_chunk() {
1782 struct Never;
1783 impl ExtensionPolicy for Never {
1784 fn next_extension(&self, _ctx: &ExtensionContext) -> Option<u64> {
1785 None
1786 }
1787 }
1788 let factory = ExtendingRangeSourceFactory::new("rows", 10, 4, Arc::new(Never));
1789 assert_eq!(
1790 drain(factory.create_reader().as_mut()),
1791 vec![10, 11, 12, 13]
1792 );
1793 assert_eq!(factory.global_consumed(), 4);
1794 assert!(factory.rewind_for_poll());
1795 assert_eq!(factory.global_consumed(), 0);
1796 assert_eq!(factory.global_extent(), Some(4));
1797 assert_eq!(
1798 drain(factory.create_reader().as_mut()),
1799 vec![10, 11, 12, 13]
1800 );
1801 let bounded = ExtendingRangeSourceFactory::new("rows", 0, 8, Arc::new(Never)).bounded(5);
1802 assert_eq!(drain(bounded.create_reader().as_mut()), vec![0, 1, 2, 3, 4]);
1803 assert!(bounded.rewind_for_poll());
1804 assert_eq!(drain(bounded.create_reader().as_mut()), vec![0, 1, 2, 3, 4]);
1805 }
1806
1807 /// Readers reserve while the factory rewinds, over and over. Every
1808 /// reservation lies in the round's domain. A reservation that
1809 /// starts after a rewind returns and finishes before the next one
1810 /// begins belongs to that round, and one round's such reservations
1811 /// never overlap. After the readers stop, a rewind and one reader
1812 /// drain the domain exactly once.
1813 fn rewind_under_concurrent_readers(
1814 factory: Arc<dyn DataSourceFactory>,
1815 domain: std::ops::Range<u64>,
1816 stride: usize,
1817 ) {
1818 use std::sync::atomic::{AtomicBool, AtomicU64};
1819 const READERS: usize = 4;
1820 const REWINDS: u64 = 300;
1821 let begun = Arc::new(AtomicU64::new(0));
1822 let returned = Arc::new(AtomicU64::new(0));
1823 let stop = Arc::new(AtomicBool::new(false));
1824 let readers: Vec<_> = (0..READERS)
1825 .map(|_| {
1826 let (factory, domain) = (factory.clone(), domain.clone());
1827 let (begun, returned, stop) = (begun.clone(), returned.clone(), stop.clone());
1828 std::thread::spawn(move || {
1829 let mut reader = factory.create_reader();
1830 let mut settled = Vec::new();
1831 while !stop.load(Ordering::SeqCst) {
1832 let round = returned.load(Ordering::SeqCst);
1833 let quiet = begun.load(Ordering::SeqCst) == round;
1834 let reserved = reader.reserve(stride);
1835 let still = begun.load(Ordering::SeqCst) == round;
1836 match reserved {
1837 Some(r) => {
1838 assert!(
1839 domain.start <= r.start
1840 && r.end <= domain.end
1841 && r.start < r.end,
1842 "{r:?} lies outside the round {domain:?}"
1843 );
1844 if quiet && still {
1845 settled.push((round, r));
1846 }
1847 }
1848 None => std::thread::yield_now(),
1849 }
1850 }
1851 settled
1852 })
1853 })
1854 .collect();
1855 for _ in 0..REWINDS {
1856 std::thread::yield_now();
1857 begun.fetch_add(1, Ordering::SeqCst);
1858 assert!(factory.rewind_for_poll());
1859 returned.fetch_add(1, Ordering::SeqCst);
1860 }
1861 stop.store(true, Ordering::SeqCst);
1862 let mut settled: Vec<(u64, std::ops::Range<u64>)> = readers
1863 .into_iter()
1864 .flat_map(|h| h.join().expect("a reader panicked"))
1865 .collect();
1866 settled.sort_by_key(|(round, r)| (*round, r.start));
1867 for pair in settled.windows(2) {
1868 let ((a_round, a), (b_round, b)) = (&pair[0], &pair[1]);
1869 assert!(
1870 a_round != b_round || a.end <= b.start,
1871 "round {a_round} handed out {a:?} and {b:?}"
1872 );
1873 }
1874 assert!(factory.rewind_for_poll());
1875 let mut reader = factory.create_reader();
1876 let mut all = Vec::new();
1877 while let Some(r) = reader.reserve(stride) {
1878 all.extend(r);
1879 }
1880 assert_eq!(all, domain.collect::<Vec<_>>());
1881 }
1882
1883 #[test]
1884 fn a_range_rewind_is_safe_under_concurrent_readers() {
1885 rewind_under_concurrent_readers(Arc::new(RangeSourceFactory::new(5, 505)), 5..505, 7);
1886 }
1887
1888 #[test]
1889 fn an_extending_rewind_is_safe_under_concurrent_readers() {
1890 let policy = Arc::new(UntilCountPolicy {
1891 min_count: 300,
1892 delta: 100,
1893 });
1894 let factory = ExtendingRangeSourceFactory::new("rows", 5, 100, policy);
1895 rewind_under_concurrent_readers(Arc::new(factory), 5..305, 7);
1896 let bounded = ExtendingRangeSourceFactory::new(
1897 "rows",
1898 5,
1899 100,
1900 Arc::new(UntilCountPolicy {
1901 min_count: 300,
1902 delta: 100,
1903 }),
1904 )
1905 .bounded(250);
1906 rewind_under_concurrent_readers(Arc::new(bounded), 5..250, 7);
1907 }
1908
1909 /// A rewind restarts the policy's elapsed time, so a time-based
1910 /// policy grows the second round past the base chunk just as it
1911 /// grew the first. The first round runs until its time is spent;
1912 /// measured from the factory's construction, the second round
1913 /// would start with no time left and end at the base chunk.
1914 #[test]
1915 fn a_rewind_restarts_the_elapsed_time_of_a_time_based_policy() {
1916 const BASE: u64 = 64;
1917 let factory = ExtendingRangeSourceFactory::new(
1918 "rows",
1919 0,
1920 BASE,
1921 Arc::new(TimeElapsedPolicy::new(BASE, 40)),
1922 );
1923 fn round(reader: &mut dyn DataSource) -> u64 {
1924 let mut claimed = 0u64;
1925 while let Some(r) = reader.reserve(BASE as usize) {
1926 claimed += r.end - r.start;
1927 }
1928 claimed
1929 }
1930 let mut reader = factory.create_reader();
1931 let first = round(reader.as_mut());
1932 assert!(first > BASE, "the first round ends at the base chunk");
1933 assert!(factory.rewind_for_poll());
1934 let second = round(reader.as_mut());
1935 assert!(
1936 second > BASE,
1937 "the second round claimed {second} ordinals, only the base chunk"
1938 );
1939 }
1940}