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}