use super::manifest_writer;
use super::{RunSummary, sink::ExportSink};
use crate::config::IncrementalCursorMode;
use crate::error::Result;
use crate::plan::{ExtractionStrategy, IncrementalCursorPlan, KeysetPlan, ResolvedRunPlan};
use crate::source::{self, Source};
use crate::state::StateStore;
use crate::types::CursorState;
use crate::{destination, format};
fn keyset_plan(plan: &ResolvedRunPlan) -> &KeysetPlan {
match &plan.strategy {
ExtractionStrategy::Keyset(kp) => kp,
_ => unreachable!("keyset runner called with non-keyset plan"),
}
}
pub(crate) struct KeysetPage {
pub(crate) parts: Vec<super::commit::PartRecord>,
pub(crate) rows: usize,
pub(crate) schema: Option<arrow::datatypes::Schema>,
pub(crate) next_cursor: Option<String>,
pub(crate) column_checksums: std::collections::BTreeMap<String, u64>,
pub(crate) checksum_key_column: Option<String>,
}
pub(crate) fn read_keyset_page(
src: &mut dyn Source,
plan: &ResolvedRunPlan,
key_plan: &IncrementalCursorPlan,
page_size: usize,
cursor: Option<&str>,
dest: &dyn destination::Destination,
part_base: &str,
) -> Result<Option<KeysetPage>> {
read_keyset_page_bounded(
src, plan, key_plan, page_size, cursor, None, dest, part_base,
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn read_keyset_page_bounded(
src: &mut dyn Source,
plan: &ResolvedRunPlan,
key_plan: &IncrementalCursorPlan,
page_size: usize,
cursor: Option<&str>,
upper: Option<&str>,
dest: &dyn destination::Destination,
part_base: &str,
) -> Result<Option<KeysetPage>> {
let cursor_state = cursor.map(|v| CursorState {
export_name: plan.export_name.clone(),
last_cursor_value: Some(v.to_string()),
last_run_at: None,
});
let mut sink = ExportSink::new(plan)?;
src.export(
&source::ExportRequest::unwrapped(&plan.base_query, &plan.tuning, &plan.column_overrides)
.with_incremental(Some(key_plan))
.with_cursor(cursor_state.as_ref())
.with_upper_bound(upper)
.with_page_limit(page_size),
&mut sink,
)?;
if let Some(w) = sink.writer.take() {
w.finish()?;
}
let rows = sink.total_rows;
if rows == 0 {
return Ok(None); }
let schema = sink.dest_schema.as_deref().cloned();
let parts = super::commit::write_sink_parts(
dest,
&mut sink,
plan.validate.then_some(plan.format),
|idx, count| super::commit::part_indexed_name(part_base, idx, count),
)?;
let checksum_key_column = sink.checksum_key_col.and(sink.cursor_column.clone());
Ok(Some(KeysetPage {
parts,
rows,
schema,
next_cursor: sink.effective_cursor(),
column_checksums: std::mem::take(&mut sink.column_checksums),
checksum_key_column,
}))
}
fn percentile_offset(total: i64, i: usize, parts: usize) -> i64 {
total * i as i64 / parts as i64
}
fn highest_range_max(range_maxes: Vec<Option<String>>) -> Option<String> {
range_maxes.into_iter().rev().flatten().next()
}
fn nth_row_clause(st: crate::config::SourceType, off: i64) -> String {
use crate::config::SourceType::*;
match st {
Postgres | Mysql => format!("LIMIT 1 OFFSET {off}"),
Mssql => format!("OFFSET {off} ROWS FETCH NEXT 1 ROWS ONLY"),
Mongo => unreachable!("parallel keyset sampling is a SQL path; Mongo uses $sample"),
}
}
fn sample_key_boundaries(
src: &mut dyn Source,
plan: &ResolvedRunPlan,
key: &str,
parts: usize,
floor: Option<&str>,
ceil: Option<&str>,
) -> Result<Vec<String>> {
let st = plan.source.source_type;
let base = &plan.base_query;
let k = crate::sql::quote_ident(st, key);
let mut preds: Vec<String> = Vec::new();
if let Some(lo) = floor {
preds.push(format!(
"{k} > {}",
crate::source::query::inline_literal(st, lo)
));
}
if let Some(hi) = ceil {
preds.push(format!(
"{k} <= {}",
crate::source::query::inline_literal(st, hi)
));
}
let where_clause = if preds.is_empty() {
String::new()
} else {
format!("WHERE {}", preds.join(" AND "))
};
let total: i64 = src
.query_scalar(&format!(
"SELECT COUNT(*) FROM ({base}) AS _rivet_pk_cnt {where_clause}"
))?
.as_deref()
.and_then(|s| s.trim().parse::<i64>().ok())
.unwrap_or(0);
if total <= 1 {
return Ok(vec![]);
}
let mut bounds: Vec<String> = Vec::with_capacity(parts.saturating_sub(1));
for i in 1..parts {
let off = percentile_offset(total, i, parts);
let nth = nth_row_clause(st, off);
let sql =
format!("SELECT {k} FROM ({base}) AS _rivet_pk {where_clause} ORDER BY {k} {nth}");
if let Some(v) = src.query_scalar(&sql)?
&& bounds.last().map(String::as_str) != Some(v.as_str())
{
bounds.push(v);
}
}
Ok(bounds)
}
fn sanitize_run_id(s: &str) -> String {
s.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' || c == '_' {
c
} else {
'_'
}
})
.collect()
}
#[allow(clippy::type_complexity)]
fn sample_parallel_ranges(
src: &mut dyn Source,
plan: &ResolvedRunPlan,
key: &str,
parallel: usize,
floor: Option<&str>,
ceil: Option<&str>,
) -> Result<Vec<(usize, Option<String>, Option<String>, bool)>> {
let bounds = sample_key_boundaries(src, plan, key, parallel, floor, ceil)?;
let mut ranges = Vec::with_capacity(bounds.len() + 1);
let mut prev: Option<String> = floor.map(str::to_string);
for (i, b) in bounds.iter().enumerate() {
ranges.push((i, prev.clone(), Some(b.clone()), false));
prev = Some(b.clone());
}
let last = ranges.len();
ranges.push((last, prev, ceil.map(str::to_string), false));
Ok(ranges)
}
fn run_keyset_parallel(
src: &mut dyn Source,
plan: &ResolvedRunPlan,
summary: &mut RunSummary,
key_plan: IncrementalCursorPlan,
parallel: usize,
state: Option<&StateStore>,
) -> Result<()> {
use std::sync::Mutex;
use std::sync::atomic::{AtomicI64, Ordering};
let kp = keyset_plan(plan);
let key = kp.key_column.clone();
let page_size = kp.chunk_size;
let checkpoint = kp.checkpoint;
let resume_run_id: Option<String> = if checkpoint {
state
.and_then(|s| s.get_resume_run_id(&plan.export_name).ok())
.flatten()
} else {
None
};
let incremental = kp.incremental;
let (floor, ceil): (Option<String>, Option<String>) = if incremental && resume_run_id.is_none()
{
let anchor = state
.and_then(|s| s.get(&plan.export_name).ok())
.and_then(|c| c.last_cursor_value);
let key_q = crate::sql::quote_ident(plan.source.source_type, &key);
let cur_max = src.query_scalar(&format!(
"SELECT MAX({key_q}) FROM ({}) AS _rivet_pk_max",
plan.base_query
))?;
let advances = match (&anchor, &cur_max) {
(_, None) => false, (None, Some(_)) => true, (Some(a), Some(c)) => key_advances(a, c), };
if !advances {
log::info!(
"export '{}': parallel keyset incremental — no new rows past the anchor, nothing to export",
plan.export_name
);
return Ok(());
}
(anchor, cur_max)
} else {
(None, None)
};
let (floor_r, ceil_r) = (floor.as_deref(), ceil.as_deref());
let ranges: Vec<(usize, Option<String>, Option<String>, bool)> = match (&resume_run_id, state) {
(Some(rid), Some(st)) => {
summary.run_id = rid.clone();
summary.resumed = true;
let rows = st.load_keyset_ranges(&plan.export_name, rid)?;
if rows.is_empty() {
let fresh = sample_parallel_ranges(src, plan, &key, parallel, floor_r, ceil_r)?;
st.persist_keyset_ranges(&plan.export_name, rid, &lo_hi_pairs(&fresh))?;
fresh
} else {
rows.into_iter()
.map(|r| (r.range_index as usize, r.lo, r.hi, r.done))
.collect()
}
}
(None, Some(st)) if checkpoint => {
let fresh = sample_parallel_ranges(src, plan, &key, parallel, floor_r, ceil_r)?;
st.persist_keyset_ranges(&plan.export_name, &summary.run_id, &lo_hi_pairs(&fresh))?;
st.set_resume_run_id(&plan.export_name, &summary.run_id)?;
fresh
}
_ => sample_parallel_ranges(src, plan, &key, parallel, floor_r, ceil_r)?,
};
let anchor_ceiling: Option<String> = ranges.last().and_then(|(_, _, hi, _)| hi.clone());
let anchor_floor: Option<String> = ranges.first().and_then(|(_, lo, _, _)| lo.clone());
let total_ranges = ranges.len();
let pending: Vec<(usize, Option<String>, Option<String>)> = ranges
.into_iter()
.filter(|(_, _, _, done)| !done)
.map(|(idx, lo, hi, _)| (idx, lo, hi))
.collect();
if parallel > 1 && total_ranges == 1 {
log::warn!(
"export '{}': parallel keyset requested {} workers but sampled 0 boundaries — \
running as a SINGLE worker. The key may be a type the boundary probe cannot \
read; data is complete but the parallel speed-up is absent.",
plan.export_name,
parallel
);
}
log::info!(
"export '{}': parallel keyset — {} range(s), {} to run{}, page size {}",
plan.export_name,
total_ranges,
pending.len(),
if resume_run_id.is_some() {
" (resume)"
} else {
""
},
page_size
);
let dest = std::sync::Arc::new(destination::create_destination(&plan.destination)?);
crate::manifest::guard_manifest_mode(&**dest, "batch")?;
let ext = format::create_format(plan.format, plan.compression, plan.compression_level, None)
.file_extension()
.to_string();
let run_tag = sanitize_run_id(&summary.run_id);
let run_id = summary.run_id.clone();
let state_ref = if checkpoint {
state.map(|s| s.state_ref().clone())
} else {
None
};
let fmt_label = plan.format.label();
let cmp_label = plan.compression.label();
let rows = AtomicI64::new(0);
let parts_mx: Mutex<Vec<super::commit::PartRecord>> = Mutex::new(Vec::new());
#[allow(clippy::type_complexity)]
let checksums_mx: Mutex<Vec<(std::collections::BTreeMap<String, u64>, Option<String>)>> =
Mutex::new(Vec::new());
let fingerprint: std::sync::OnceLock<arrow::datatypes::Schema> = std::sync::OnceLock::new();
let range_max: Mutex<Vec<Option<String>>> = Mutex::new(vec![None; total_ranges]);
let errors: Mutex<Vec<String>> = Mutex::new(Vec::new());
std::thread::scope(|scope| {
for (ridx, lo, hi) in pending.iter().cloned() {
let dest = std::sync::Arc::clone(&dest);
let (plan_r, key_plan_r, ext_r, tag_r, key_r) =
(plan, &key_plan, &ext, run_tag.as_str(), key.as_str());
let (rows_r, parts_r, checks_r, fp_r, rmax_r, errs_r) = (
&rows,
&parts_mx,
&checksums_mx,
&fingerprint,
&range_max,
&errors,
);
let (sref_r, rid_r, fmt_r, cmp_r) = (&state_ref, run_id.as_str(), fmt_label, cmp_label);
scope.spawn(move || {
let mut wsrc = match source::create_source(&plan_r.source) {
Ok(s) => s,
Err(e) => {
errs_r
.lock()
.unwrap()
.push(format!("range {ridx}: connect: {e:#}"));
return;
}
};
let mut cursor = lo;
let mut pages = 0usize;
let mut rmax: Option<String> = None;
let mut range_parts: Vec<crate::state::KeysetRangePart> = Vec::new();
let mut local_parts: Vec<super::commit::PartRecord> = Vec::new();
let mut local_checks: Vec<(
std::collections::BTreeMap<String, u64>,
Option<String>,
)> = Vec::new();
loop {
if let Err(e) = crate::test_hook::maybe_error_at_index(
"keyset_parallel_worker",
ridx as i64,
) {
errs_r.lock().unwrap().push(format!("range {ridx}: {e}"));
return;
}
let base = format!(
"{}_{}_pk_w{}_{}.{}",
plan_r.export_name, tag_r, ridx, pages, ext_r
);
let page = match read_keyset_page_bounded(
&mut *wsrc,
plan_r,
key_plan_r,
page_size,
cursor.as_deref(),
hi.as_deref(),
&**dest,
&base,
) {
Ok(p) => p,
Err(e) => {
errs_r
.lock()
.unwrap()
.push(format!("range {ridx}: page {pages}: {e:#}"));
return;
}
};
let Some(page) = page else { break };
rows_r.fetch_add(page.rows as i64, Ordering::Relaxed);
if let Some(sc) = &page.schema {
let _ = fp_r.set(sc.clone());
}
rmax = page.next_cursor.clone().or(rmax);
for p in &page.parts {
range_parts.push(crate::state::KeysetRangePart {
file_name: p.file_name.clone(),
rows: p.rows,
bytes: p.bytes as i64,
});
}
local_parts.extend(page.parts);
local_checks.push((page.column_checksums, page.checksum_key_column));
let last_page = page.rows < page_size;
if !last_page {
match page.next_cursor {
Some(v) => cursor = Some(v),
None => {
errs_r.lock().unwrap().push(format!(
"range {ridx}: could not advance the '{key_r}' cursor at page \
{pages} (NULL or unsupported type)"
));
return;
}
}
}
pages += 1;
if last_page {
break;
}
}
if let Some(sref) = sref_r
&& let Err(e) = crate::state::StateStore::commit_keyset_range_at_ref(
sref,
rid_r,
&plan_r.export_name,
ridx as i64,
&range_parts,
fmt_r,
Some(cmp_r),
)
{
errs_r
.lock()
.unwrap()
.push(format!("range {ridx}: checkpoint commit: {e:#}"));
return;
}
crate::test_hook::maybe_exit_at_index(
"keyset_parallel_range_committed",
ridx as i64,
);
rmax_r.lock().unwrap()[ridx] = rmax;
parts_r.lock().unwrap().extend(local_parts);
checks_r.lock().unwrap().extend(local_checks);
});
}
});
let errs = errors.into_inner().unwrap();
if !errs.is_empty() {
anyhow::bail!(
"export '{}': parallel keyset failed on {} range(s): {}",
plan.export_name,
errs.len(),
errs.join("; ")
);
}
summary.total_rows += rows.into_inner();
if plan.validate {
summary.validated = Some(true);
}
if let Some(sc) = fingerprint.get() {
manifest_writer::record_run_schema_fingerprint(summary, sc);
}
summary.cursor_high = highest_range_max(range_max.into_inner().unwrap());
summary.cursor_low = None;
let file_log_state = if checkpoint { None } else { state };
let parts = parts_mx.into_inner().unwrap();
for (idx, rec) in parts.iter().enumerate() {
super::commit::record_part(
plan,
summary,
file_log_state,
rec,
super::commit::PartKind::Page {
page_index: idx as i64,
},
);
}
if let Some(st) = state
&& checkpoint
{
super::chunked::rehydrate_manifest_parts_from_file_log(st, &run_id, summary)?;
}
let mut acc: std::collections::BTreeMap<String, u64> = std::collections::BTreeMap::new();
let mut ck_key: Option<String> = None;
for (m, k) in checksums_mx.into_inner().unwrap() {
super::commit::accumulate_column_checksums(&mut acc, &m);
if ck_key.is_none() {
ck_key = k;
}
}
super::commit::harvest_column_checksums(summary, acc, ck_key);
log::info!(
"export '{}': parallel keyset complete — {} range(s), {} parts, {} rows",
plan.export_name,
total_ranges,
parts.len(),
summary.total_rows
);
if let (Some(sc), Some(st)) = (fingerprint.get(), state) {
super::schema_drift::check_from_sink_schema(
st,
&plan.export_name,
sc,
plan.schema_drift_policy,
summary,
)?;
}
if incremental && let Some(hi) = &anchor_ceiling {
summary.cursor_high = Some(hi.clone());
summary.cursor_low = anchor_floor.clone();
if let Some(st) = state {
st.update(&plan.export_name, hi)?;
}
}
Ok(())
}
fn key_advances(anchor: &str, candidate: &str) -> bool {
if let (Ok(a), Ok(b)) = (anchor.parse::<i128>(), candidate.parse::<i128>()) {
return b > a;
}
if let (Ok(a), Ok(b)) = (anchor.parse::<f64>(), candidate.parse::<f64>())
&& let Some(ord) = b.partial_cmp(&a)
{
return ord.is_gt();
}
candidate > anchor
}
fn lo_hi_pairs(
ranges: &[(usize, Option<String>, Option<String>, bool)],
) -> Vec<(Option<String>, Option<String>)> {
ranges
.iter()
.map(|(_, lo, hi, _)| (lo.clone(), hi.clone()))
.collect()
}
pub(crate) fn run_keyset(
src: &mut dyn Source,
plan: &ResolvedRunPlan,
summary: &mut RunSummary,
state: Option<&StateStore>,
) -> Result<()> {
let kp = keyset_plan(plan);
let key_plan = IncrementalCursorPlan {
primary_column: kp.key_column.clone(),
fallback_column: None,
mode: IncrementalCursorMode::SingleColumn,
};
if kp.parallel > 1 {
return run_keyset_parallel(src, plan, summary, key_plan, kp.parallel, state);
}
log::info!(
"export '{}': keyset (seek) pagination on '{}', page size {}",
plan.export_name,
kp.key_column,
kp.chunk_size
);
let resume_run_id: Option<String> = if kp.checkpoint {
match state {
Some(st) => st.get_resume_run_id(&plan.export_name)?,
None => None,
}
} else {
None
};
let recovering_crash = resume_run_id.is_some();
summary.resumed = recovering_crash;
let mut last: Option<String> = if kp.checkpoint && (recovering_crash || kp.incremental) {
state
.and_then(|s| s.get(&plan.export_name).ok())
.and_then(|cs| cs.last_cursor_value)
} else {
None
};
summary.cursor_low = last.clone();
if kp.checkpoint
&& let Some(st) = state
{
match &resume_run_id {
Some(rid) => {
summary.run_id = rid.clone();
super::chunked::rehydrate_manifest_parts_from_file_log(st, rid, summary)?;
}
None => {
if !kp.incremental {
st.clear_cursor_value(&plan.export_name)?;
}
st.set_resume_run_id(&plan.export_name, &summary.run_id)?;
}
}
}
crate::test_hook::maybe_panic_at("keyset_after_open_before_first_page");
let mut pages: usize = 0;
let mut drift_schema: Option<arrow::datatypes::Schema> = None;
let mut checksums_acc: std::collections::BTreeMap<String, u64> =
std::collections::BTreeMap::new();
let mut checksum_key_column: Option<String> = None;
let dest = destination::create_destination(&plan.destination)?;
crate::manifest::guard_manifest_mode(dest.as_ref(), "batch")?;
let ext = format::create_format(plan.format, plan.compression, plan.compression_level, None)
.file_extension()
.to_string();
let stamp = chrono::Utc::now().format("%Y%m%d_%H%M%S_%3f").to_string();
loop {
let base = format!("{}_{}_keyset{}.{}", plan.export_name, stamp, pages, ext);
let Some(page) = read_keyset_page(
src,
plan,
&key_plan,
kp.chunk_size,
last.as_deref(),
dest.as_ref(),
&base,
)?
else {
summary.cursor_high = last.clone();
break;
};
if let Some(sc) = &page.schema {
manifest_writer::record_run_schema_fingerprint(summary, sc);
if drift_schema.is_none() {
drift_schema = Some(sc.clone());
}
}
super::commit::accumulate_column_checksums(&mut checksums_acc, &page.column_checksums);
if checksum_key_column.is_none() {
checksum_key_column = page.checksum_key_column.clone();
}
summary.total_rows += page.rows as i64;
if plan.validate {
summary.validated = Some(true);
}
for rec in &page.parts {
super::commit::record_part(
plan,
summary,
state,
rec,
super::commit::PartKind::Page {
page_index: pages as i64,
},
);
}
if kp.checkpoint
&& let (Some(st), Some(v)) = (state, page.next_cursor.as_ref())
{
st.update(&plan.export_name, v)?;
}
crate::test_hook::maybe_panic_at(&format!("after_keyset_page:{pages}"));
log::info!(
"export '{}': keyset page {} — {} rows",
plan.export_name,
pages,
page.rows
);
pages += 1;
if page.rows < kp.chunk_size {
summary.cursor_high = page.next_cursor.clone().or_else(|| last.clone());
break;
}
match page.next_cursor {
Some(v) => last = Some(v),
None => {
summary.offending_value = last.clone();
summary.cursor_high = last.clone();
anyhow::bail!(
"export '{}': keyset could not read the '{}' value from the last row of page {} \
(NULL or unsupported type) — cannot advance safely (last readable key: {}). \
The key must be NOT NULL and one of: integer, float, string, timestamp, date, uuid.",
plan.export_name,
kp.key_column,
pages - 1,
last.as_deref().unwrap_or("<none>"),
);
}
}
}
super::commit::harvest_column_checksums(summary, checksums_acc, checksum_key_column);
if kp.checkpoint
&& !kp.incremental
&& let Some(st) = state
{
st.clear_resume_run_id(&plan.export_name)?;
}
crate::test_hook::maybe_panic_at("keyset_after_data_complete");
log::info!(
"export '{}': keyset complete — {} page(s), {} rows",
plan.export_name,
pages,
summary.total_rows
);
if let (Some(sc), Some(st)) = (&drift_schema, state) {
super::schema_drift::check_from_sink_schema(
st,
&plan.export_name,
sc,
plan.schema_drift_policy,
summary,
)?;
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::SourceType;
#[test]
fn highest_range_max_takes_the_top_populated_range() {
let s = |x: &str| Some(x.to_string());
assert_eq!(highest_range_max(vec![s("k1"), s("k2"), s("k3")]), s("k3"));
assert_eq!(highest_range_max(vec![s("k1"), s("k2"), None]), s("k2"));
assert_eq!(highest_range_max(vec![s("lo"), None, s("hi")]), s("hi"));
assert_eq!(highest_range_max(vec![None, None]), None);
assert_eq!(highest_range_max(vec![]), None);
}
#[test]
fn percentile_offset_partitions_evenly() {
assert_eq!(percentile_offset(1000, 1, 4), 250);
assert_eq!(percentile_offset(1000, 2, 4), 500);
assert_eq!(percentile_offset(1000, 3, 4), 750);
assert_eq!(percentile_offset(999, 1, 3), 333);
assert_eq!(percentile_offset(999, 2, 3), 666);
}
#[test]
fn nth_row_clause_is_per_dialect() {
assert_eq!(
nth_row_clause(SourceType::Postgres, 250),
"LIMIT 1 OFFSET 250"
);
assert_eq!(nth_row_clause(SourceType::Mysql, 250), "LIMIT 1 OFFSET 250");
assert_eq!(
nth_row_clause(SourceType::Mssql, 250),
"OFFSET 250 ROWS FETCH NEXT 1 ROWS ONLY"
);
}
#[test]
fn sanitize_run_id_keeps_safe_chars_and_replaces_the_rest() {
assert_eq!(sanitize_run_id("run-2026_01A9"), "run-2026_01A9");
assert_eq!(sanitize_run_id("a/b c:d.e"), "a_b_c_d_e");
assert_eq!(sanitize_run_id("../etc"), "___etc");
assert_eq!(sanitize_run_id("ABCabc012"), "ABCabc012");
}
#[test]
fn key_advances_is_numeric_not_lexical() {
assert!(key_advances("999", "1000"));
assert!(!key_advances("1000", "999"));
assert!(!key_advances("5", "5")); assert!(key_advances("18446744073709551614", "18446744073709551615"));
assert!(key_advances("1.5", "2.0"));
assert!(!key_advances("2.0", "1.5"));
assert!(key_advances("2026-01-01T00:00:00Z", "2026-01-02T00:00:00Z"));
assert!(!key_advances(
"2026-01-02T00:00:00Z",
"2026-01-01T00:00:00Z"
));
}
#[test]
fn lo_hi_pairs_projects_the_bounds() {
let ranges = vec![
(0usize, None, Some("k0500".to_string()), false),
(
1,
Some("k0500".to_string()),
Some("k1000".to_string()),
false,
),
(2, Some("k1000".to_string()), None, false),
];
assert_eq!(
lo_hi_pairs(&ranges),
vec![
(None, Some("k0500".to_string())),
(Some("k0500".to_string()), Some("k1000".to_string())),
(Some("k1000".to_string()), None),
]
);
}
}