use crate::time::NanoTime;
#[derive(Clone, Copy, Debug)]
pub struct TimeWindow {
lo: NanoTime,
hi: NanoTime,
}
impl TimeWindow {
pub fn clamp(t0: NanoTime, t1: NanoTime, start: NanoTime, end: NanoTime) -> Self {
Self {
lo: t0.max(start),
hi: t1.min(end),
}
}
pub fn contains(&self, time: NanoTime) -> bool {
self.lo <= time && time < self.hi
}
}
pub struct WindowFilter {
label: &'static str,
window: TimeWindow,
dropped: usize,
}
impl WindowFilter {
pub fn new(label: &'static str, window: TimeWindow) -> Self {
Self {
label,
window,
dropped: 0,
}
}
pub fn keep(&mut self, time: NanoTime) -> bool {
if self.window.contains(time) {
true
} else {
self.dropped += 1;
false
}
}
pub fn finish(self) {
if self.dropped > 0 {
log::warn!(
"{}: dropped {} row(s) outside the requested window [{:?}, {:?}); \
the source returned data beyond the range it was asked for",
self.label,
self.dropped,
self.window.lo,
self.window.hi,
);
}
}
}
#[cfg(any(feature = "kdb", feature = "postgres"))]
pub(crate) fn compute_validated_time_slices(
adapter: &str,
start_time: NanoTime,
end_time: anyhow::Result<NanoTime>,
period: std::time::Duration,
) -> anyhow::Result<Vec<((NanoTime, NanoTime), i32, usize)>> {
if period.is_zero() {
anyhow::bail!("{adapter}: period must be greater than zero");
}
if start_time == NanoTime::ZERO {
anyhow::bail!(
"{adapter}: start_time is NanoTime::ZERO; \
use RunMode::HistoricalFrom with an explicit start time"
);
}
let end_time = match end_time {
Ok(t) if t == NanoTime::MAX => anyhow::bail!(
"{adapter} requires RunFor::Duration; \
RunFor::Forever would generate an unbounded number of slices"
),
Ok(t) => t,
Err(_) => anyhow::bail!(
"{adapter} requires RunFor::Duration; \
RunFor::Cycles does not provide an end time"
),
};
Ok(compute_time_slices(start_time, end_time, period))
}
#[cfg(any(feature = "kdb", feature = "postgres"))]
pub(crate) fn compute_time_slices(
start_time: NanoTime,
end_time: NanoTime,
period: std::time::Duration,
) -> Vec<((NanoTime, NanoTime), i32, usize)> {
const DAY_NANOS: i64 = 86_400_000_000_000;
let period_nanos = period.as_nanos() as i64;
let start_kdb = start_time.to_kdb_timestamp();
let end_kdb = end_time.to_kdb_timestamp();
let start_day = start_kdb.div_euclid(DAY_NANOS);
let end_day = (end_kdb - 1).div_euclid(DAY_NANOS);
let mut result = Vec::new();
for day in start_day..=end_day {
let kdb_date = day as i32;
let midnight_kdb = day * DAY_NANOS;
let next_midnight_kdb = midnight_kdb + DAY_NANOS;
let mut iteration = if day == start_day {
((start_kdb - midnight_kdb) / period_nanos) as usize
} else {
0
};
loop {
let t0 = midnight_kdb + iteration as i64 * period_nanos;
if day == end_day && t0 >= end_kdb {
break;
}
let natural_t1 = t0 + period_nanos;
let t1 = natural_t1.min(next_midnight_kdb);
result.push((
(
NanoTime::from_kdb_timestamp(t0),
NanoTime::from_kdb_timestamp(t1),
),
kdb_date,
iteration,
));
if natural_t1 >= next_midnight_kdb {
break;
}
iteration += 1;
}
}
result
}
#[cfg(all(test, any(feature = "kdb", feature = "postgres")))]
mod tests {
use super::*;
fn kdb_epoch() -> NanoTime {
NanoTime::from_kdb_timestamp(0)
}
const DAY_NANOS: u64 = 86_400_000_000_000;
#[test]
fn test_compute_time_slices_no_stub() {
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(8 * 3600);
let start = epoch;
let end = NanoTime::new(u64::from(epoch) + DAY_NANOS - 1);
let slices = compute_time_slices(start, end, period);
assert_eq!(slices.len(), 3, "expected 3 slices for 8h period");
for &(_, date, _) in &slices {
assert_eq!(date, 0);
}
let period_nanos = period.as_nanos() as u64;
let (t0_0, t1_0) = slices[0].0;
assert_eq!(u64::from(t0_0), u64::from(epoch));
assert_eq!(u64::from(t1_0), u64::from(epoch) + period_nanos);
let (t0_1, t1_1) = slices[1].0;
assert_eq!(u64::from(t0_1), u64::from(epoch) + period_nanos);
assert_eq!(u64::from(t1_1), u64::from(epoch) + 2 * period_nanos);
let (t0_2, t1_2) = slices[2].0;
assert_eq!(u64::from(t0_2), u64::from(epoch) + 2 * period_nanos);
assert_eq!(u64::from(t1_2), u64::from(epoch) + DAY_NANOS);
assert_eq!(slices[0].2, 0);
assert_eq!(slices[1].2, 1);
assert_eq!(slices[2].2, 2);
}
#[test]
fn test_compute_time_slices_with_stub() {
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(5 * 3600);
let start = epoch;
let end = NanoTime::new(u64::from(epoch) + DAY_NANOS - 1);
let slices = compute_time_slices(start, end, period);
assert_eq!(slices.len(), 5, "expected 4 full + 1 stub = 5 slices");
let period_nanos = period.as_nanos() as u64;
let (stub_t0, stub_t1) = slices[4].0;
assert_eq!(u64::from(stub_t0), u64::from(epoch) + 4 * period_nanos);
assert_eq!(u64::from(stub_t1), u64::from(epoch) + DAY_NANOS);
for i in 0..4 {
let t1_i = u64::from(slices[i].0.1);
let t0_next = u64::from(slices[i + 1].0.0);
assert_eq!(t1_i, t0_next, "boundary mismatch at slice {i}/{}", i + 1);
}
}
#[test]
fn test_compute_time_slices_two_days() {
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(12 * 3600); let start = epoch;
let end = NanoTime::new(u64::from(epoch) + 2 * DAY_NANOS - 1);
let slices = compute_time_slices(start, end, period);
assert_eq!(slices.len(), 4);
assert_eq!(slices[0].1, 0);
assert_eq!(slices[0].2, 0);
assert_eq!(slices[1].1, 0);
assert_eq!(slices[1].2, 1);
assert_eq!(slices[2].1, 1);
assert_eq!(slices[2].2, 0); assert_eq!(slices[3].1, 1);
assert_eq!(slices[3].2, 1);
}
#[test]
fn test_compute_time_slices_mid_day_start() {
const SECS_23_59_30: i64 = 86_370 * 1_000_000_000;
const SECS_23_59_59: i64 = 86_399 * 1_000_000_000;
const DAY_NANOS: i64 = 86_400_000_000_000;
let start = NanoTime::from_kdb_timestamp(SECS_23_59_30);
let end = NanoTime::from_kdb_timestamp(SECS_23_59_59);
let next_midnight = NanoTime::from_kdb_timestamp(DAY_NANOS);
let slices = compute_time_slices(start, end, std::time::Duration::from_secs(60));
assert_eq!(slices.len(), 1, "60s: expected 1 slice");
let (t0, t1) = slices[0].0;
assert_eq!(
t0,
NanoTime::from_kdb_timestamp(86_340 * 1_000_000_000),
"60s: t0 should be 23:59:00"
);
assert_eq!(t1, next_midnight, "60s: t1 should be midnight");
assert_eq!(slices[0].2, 1439, "60s: iteration should be 1439");
let slices = compute_time_slices(start, end, std::time::Duration::from_secs(30));
assert_eq!(slices.len(), 1, "30s: expected 1 slice");
let (t0, t1) = slices[0].0;
assert_eq!(t0, start, "30s: t0 should be 23:59:30");
assert_eq!(t1, next_midnight, "30s: t1 should be midnight");
assert_eq!(slices[0].2, 2879, "30s: iteration should be 2879");
let slices = compute_time_slices(start, end, std::time::Duration::from_secs(10));
assert_eq!(slices.len(), 3, "10s: expected 3 slices");
assert_eq!(slices[0].0.0, start, "10s: first t0 should be 23:59:30");
assert_eq!(
slices[0].0.1,
NanoTime::from_kdb_timestamp(86_380 * 1_000_000_000),
"10s: first t1 should be 23:59:40"
);
assert_eq!(
slices[1].0.0,
NanoTime::from_kdb_timestamp(86_380 * 1_000_000_000),
"10s: second t0 should be 23:59:40"
);
assert_eq!(
slices[2].0.1, next_midnight,
"10s: last t1 should be midnight"
);
assert_eq!(slices[0].2, 8637, "10s: first iteration should be 8637");
}
#[test]
fn test_compute_time_slices_exact_midnight_boundary() {
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(8 * 3600);
let end = NanoTime::new(u64::from(epoch) + DAY_NANOS);
let slices = compute_time_slices(epoch, end, period);
assert_eq!(
slices.len(),
3,
"exact midnight end should yield 3 slices (day 0 only)"
);
for &(_, date, _) in &slices {
assert_eq!(date, 0, "all slices should be on day 0");
}
}
#[test]
fn test_compute_time_slices_end_past_midnight() {
const HOUR_NANOS: u64 = 3_600_000_000_000;
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(3600);
let start = epoch;
let end = NanoTime::new(u64::from(epoch) + DAY_NANOS + 30 * 60 * 1_000_000_000);
let slices = compute_time_slices(start, end, period);
assert_eq!(slices.len(), 25, "expected 25 slices");
let last = slices.last().unwrap();
assert_eq!(last.1, 1, "last slice should be on kdb_date 1");
assert_eq!(last.2, 0, "last slice should be iteration 0 of day 1");
assert_eq!(
last.0.0,
NanoTime::new(u64::from(epoch) + DAY_NANOS),
"last slice t0 should be day-1 midnight"
);
assert_eq!(
last.0.1,
NanoTime::new(u64::from(epoch) + DAY_NANOS + HOUR_NANOS),
"last slice t1 should be day-1 01:00"
);
}
#[test]
fn test_compute_time_slices_cross_midnight() {
const HOUR_NANOS: u64 = 3_600_000_000_000;
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(3600);
let start = NanoTime::new(u64::from(epoch) + 23 * HOUR_NANOS);
let end = NanoTime::new(u64::from(epoch) + DAY_NANOS + 30 * 60 * 1_000_000_000);
let slices = compute_time_slices(start, end, period);
assert_eq!(slices.len(), 2, "expected 2 slices");
assert_eq!(slices[0].1, 0);
assert_eq!(slices[0].2, 23);
assert_eq!(slices[0].0.0, start);
assert_eq!(slices[0].0.1, NanoTime::new(u64::from(epoch) + DAY_NANOS));
assert_eq!(slices[1].1, 1);
assert_eq!(slices[1].2, 0);
assert_eq!(slices[1].0.0, NanoTime::new(u64::from(epoch) + DAY_NANOS));
assert_eq!(
slices[1].0.1,
NanoTime::new(u64::from(epoch) + DAY_NANOS + HOUR_NANOS)
);
}
#[test]
fn test_validated_rejects_zero_period() {
let err = compute_validated_time_slices(
"test_adapter",
kdb_epoch(),
Ok(NanoTime::new(u64::from(kdb_epoch()) + DAY_NANOS)),
std::time::Duration::ZERO,
)
.unwrap_err();
assert!(err.to_string().contains("period must be greater than zero"));
assert!(err.to_string().contains("test_adapter"));
}
#[test]
fn test_validated_rejects_zero_start() {
let err = compute_validated_time_slices(
"test_adapter",
NanoTime::ZERO,
Ok(NanoTime::new(DAY_NANOS)),
std::time::Duration::from_secs(3600),
)
.unwrap_err();
assert!(err.to_string().contains("start_time is NanoTime::ZERO"));
}
#[test]
fn test_validated_rejects_forever_and_cycles() {
let err = compute_validated_time_slices(
"test_adapter",
kdb_epoch(),
Ok(NanoTime::MAX),
std::time::Duration::from_secs(3600),
)
.unwrap_err();
assert!(err.to_string().contains("RunFor::Forever"));
let err = compute_validated_time_slices(
"test_adapter",
kdb_epoch(),
Err(anyhow::anyhow!("end_time not available for RunFor::Cycles")),
std::time::Duration::from_secs(3600),
)
.unwrap_err();
assert!(err.to_string().contains("RunFor::Cycles"));
}
#[test]
fn test_validated_passes_through_to_compute() {
let slices = compute_validated_time_slices(
"test_adapter",
kdb_epoch(),
Ok(NanoTime::new(u64::from(kdb_epoch()) + DAY_NANOS)),
std::time::Duration::from_secs(8 * 3600),
)
.unwrap();
assert_eq!(slices.len(), 3);
}
#[test]
fn test_compute_time_slices_end_on_period_boundary_day1() {
const HOUR_NANOS: u64 = 3_600_000_000_000;
let epoch = kdb_epoch();
let period = std::time::Duration::from_secs(2 * 3600); let start = epoch;
let end = NanoTime::new(u64::from(epoch) + DAY_NANOS + 2 * HOUR_NANOS);
let slices = compute_time_slices(start, end, period);
assert_eq!(slices.len(), 13, "expected 13 slices");
assert_eq!(slices.last().unwrap().1, 1);
assert_eq!(slices.last().unwrap().2, 0);
}
}