hamelin_analysis 0.23.1

Analysis utilities for Hamelin query language
Documentation
//! Infer input timestamp lower bounds from an output timestamp lower bound.

use super::backward_pass::{compute_input_range_backward, InputRange};
use super::*;

/// Find a conservative source timestamp lower bound for output at or after
/// `output_lower_bound`. The upper bound remains open; this does not limit future data.
///
/// Unlike incremental analysis, this starts with an output requirement, not a
/// range of changed source rows. Window requirements compose backwards through
/// timestamp projections and referenced DEF pipelines. The returned bound also
/// covers intermediate timestamps because an execution range can apply to every
/// source in the query. The caller must still filter the final output.
///
/// Unsupported commands, unbounded windows, and untraceable timestamp lineage
/// return an error when a safe finite lower bound cannot be derived.
pub fn compute_input_lower_bound_for_query(
    statement: &TypedStatement,
    output_lower_bound: Timestamp,
    field: &SimpleIdentifier,
) -> Result<Timestamp, IncrementalAnalysisError> {
    check_no_template_parameters(statement)?;
    if !statement.side_effect.is_none() {
        return Err(IncrementalAnalysisError::DmlNotSupported);
    }
    let mut requirements = HashMap::new();
    let mut earliest = compute_input_range_backward(
        &statement.pipeline,
        &mut requirements,
        InputRange {
            start: output_lower_bound,
            end: None,
        },
        field,
        false,
    )?;
    // DEF dependencies precede their consumers. Merge all consumer requirements
    // before visiting a definition, so shared pipelines are analyzed only once.
    for definition in statement.pipeline_defs.iter().rev() {
        let name = definition
            .name
            .valid_ref()
            .map_err(|e| IncrementalAnalysisError::TreeHadError(e.as_ref().clone()))?;
        if let Some(start) = requirements.remove(name) {
            earliest = earliest.union(compute_input_range_backward(
                &definition.pipeline,
                &mut requirements,
                start,
                field,
                false,
            )?);
        }
    }
    Ok(earliest.start.min(output_lower_bound))
}