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>> {
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_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,
}))
}
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,
};
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();
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
};
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 {
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 {
break;
}
match page.next_cursor {
Some(v) => last = Some(v),
None => anyhow::bail!(
"export '{}': keyset could not read the '{}' value from the last row of page {} \
(NULL or unsupported type) — cannot advance safely. The key must be NOT NULL and \
one of: integer, float, string, timestamp, date, uuid.",
plan.export_name,
kp.key_column,
pages - 1
),
}
}
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(())
}