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