use crate::errors::AppError;
const NO_LIMIT: i64 = -1;
pub(in crate::commands::enrich) fn limit_param(limit: Option<usize>) -> i64 {
limit.map_or(NO_LIMIT, |n| i64::try_from(n).unwrap_or(i64::MAX))
}
pub(in crate::commands::enrich) fn limit_clause(index: usize) -> String {
format!("LIMIT ?{index}")
}
pub(super) fn placeholder_list(start: usize, count: usize) -> String {
(start..start + count)
.map(|i| format!("?{i}"))
.collect::<Vec<_>>()
.join(", ")
}
pub(in crate::commands::enrich) fn page_take(page_size: usize, remaining: Option<usize>) -> usize {
let page = page_size.max(1);
match remaining {
Some(r) => r.min(page),
None => page,
}
}
pub(in crate::commands::enrich) fn keyset_collect<T, F>(
global_limit: Option<usize>,
page_size: usize,
mut fetch_page: F,
) -> Result<Vec<T>, AppError>
where
F: FnMut(i64, usize) -> Result<Vec<(i64, T)>, AppError>,
{
let mut out = Vec::new();
keyset_for_each(global_limit, page_size, &mut fetch_page, |page| {
out.extend(page);
Ok(())
})?;
Ok(out)
}
pub(in crate::commands::enrich) fn keyset_for_each<T, F, G>(
global_limit: Option<usize>,
page_size: usize,
fetch_page: &mut F,
on_page: G,
) -> Result<usize, AppError>
where
F: FnMut(i64, usize) -> Result<Vec<(i64, T)>, AppError>,
G: FnMut(Vec<T>) -> Result<(), AppError>,
{
keyset_for_each_selected(
global_limit,
page_size,
&mut |after, want| {
Ok(fetch_page(after, want)?
.into_iter()
.map(|(id, value)| (id, Some(value)))
.collect())
},
on_page,
)
}
pub(in crate::commands::enrich) fn keyset_for_each_selected<T, F, G>(
global_limit: Option<usize>,
page_size: usize,
fetch_page: &mut F,
mut on_page: G,
) -> Result<usize, AppError>
where
F: FnMut(i64, usize) -> Result<Vec<(i64, Option<T>)>, AppError>,
G: FnMut(Vec<T>) -> Result<(), AppError>,
{
let mut after_id: i64 = 0;
let mut remaining = global_limit;
let mut total = 0usize;
loop {
let want = page_take(page_size, remaining);
if want == 0 {
break;
}
let page = fetch_page(after_id, want)?;
if page.is_empty() {
break;
}
let scanned = page.len();
after_id = page.last().map(|(id, _)| *id).unwrap_or(after_id);
let values: Vec<T> = page.into_iter().filter_map(|(_, v)| v).collect();
let delivered = values.len();
total = total.saturating_add(delivered);
if delivered > 0 {
on_page(values)?;
}
if let Some(r) = remaining.as_mut() {
*r = r.saturating_sub(delivered);
}
if scanned < want {
break;
}
}
Ok(total)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn page_take_honours_remaining_budget() {
assert_eq!(page_take(512, None), 512);
assert_eq!(page_take(512, Some(10)), 10);
assert_eq!(page_take(512, Some(0)), 0);
assert_eq!(page_take(0, None), 1);
}
#[test]
fn keyset_collect_pages_until_empty() {
let pages: Vec<Vec<(i64, String)>> = vec![
vec![(1, "a".into()), (2, "b".into())],
vec![(3, "c".into())],
vec![],
];
let mut idx = 0;
let out = keyset_collect(None, 2, |after, want| {
assert!(want <= 2);
let page = pages.get(idx).cloned().unwrap_or_default();
idx += 1;
if let Some((id, _)) = page.first() {
assert!(*id > after || after == 0);
}
Ok(page)
})
.unwrap();
assert_eq!(out, vec!["a", "b", "c"]);
}
#[test]
fn keyset_for_each_drops_pages_and_respects_limit() {
let mut seen = Vec::new();
let total = keyset_for_each(
Some(3),
2,
&mut |after, want| {
let start = after + 1;
let page: Vec<(i64, i64)> = (start..start + want as i64)
.map(|id| (id, id * 10))
.collect();
Ok(page)
},
|page| {
seen.push(page);
Ok(())
},
)
.unwrap();
assert_eq!(total, 3);
assert_eq!(seen, vec![vec![10, 20], vec![30]]);
}
}