Skip to main content

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}