From a6278769cdcfce197ff0b3910f56819c447166ed Mon Sep 17 00:00:00 2001 From: "octoaide[bot]" <204759324+octoaide[bot]@users.noreply.github.com> Date: Thu, 3 Sep 2026 13:17:03 -0700 Subject: [PATCH 1/2] Implement signed-chronological iteration and stric --- CHANGELOG.md | 6 + src/event.rs | 424 +++++++++++++++++++++++++++++++++++++++++++-------- 2 files changed, 370 insertions(+), 60 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8997f4fc..d3be6f43 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ Versioning](https://semver.org/spec/v2.0.0.html). ## [Unreleased] +### Fixed + +- Corrected `EventDb` iteration to use signed chronological key order across + the Unix epoch, and made `remove_before` consistently retain events exactly + at its cutoff while deleting all earlier events. + ### Added - Added `EventDb::remove_by_sensors` to delete events whose sensor exactly diff --git a/src/event.rs b/src/event.rs index 5aa150f7..134b4bfa 100644 --- a/src/event.rs +++ b/src/event.rs @@ -156,6 +156,52 @@ use super::{ }; const EVENT_DELETION_BATCH_SIZE: usize = 1000; +const FIRST_NON_NEGATIVE_EVENT_KEY: [u8; 16] = 0_i128.to_be_bytes(); +const FIRST_NEGATIVE_EVENT_KEY: [u8; 16] = i128::MIN.to_be_bytes(); + +type PhysicalEventIterator<'i> = rocksdb::DBIteratorWithThreadMode< + 'i, + rocksdb::OptimisticTransactionDB, +>; + +#[derive(Clone, Copy)] +enum EventKeyRegion { + NonNegative, + Negative, +} + +fn signed_event_key(key_bytes: &[u8]) -> Option { + let key_bytes: [u8; 16] = key_bytes.try_into().ok()?; + Some(i128::from_be_bytes(key_bytes)) +} + +fn is_signed_negative(key_bytes: &[u8]) -> bool { + signed_event_key(key_bytes).is_some_and(i128::is_negative) +} + +fn event_region_read_options(region: EventKeyRegion) -> rocksdb::ReadOptions { + let mut readopts = rocksdb::ReadOptions::default(); + match region { + // The physical database starts with non-negative i128 keys. The first + // negative key is therefore the exclusive end of this region. + EventKeyRegion::NonNegative => { + readopts.set_iterate_upper_bound(FIRST_NEGATIVE_EVENT_KEY); + } + // Negative i128 keys occupy the remainder of the physical key space. + EventKeyRegion::Negative => { + readopts.set_iterate_lower_bound(FIRST_NEGATIVE_EVENT_KEY); + } + } + readopts +} + +fn physical_event_iterator<'i>( + db: &'i rocksdb::OptimisticTransactionDB, + region: EventKeyRegion, + mode: IteratorMode<'_>, +) -> PhysicalEventIterator<'i> { + db.iterator_opt(mode, event_region_read_options(region)) +} // event kind const DNS_COVERT_CHANNEL: &str = "DNS Covert Channel"; @@ -3181,20 +3227,72 @@ impl<'a> EventDb<'a> { } } - /// Creates an iterator over key-value pairs, starting from `key`. + /// Creates an iterator over key-value pairs, starting inclusively from `key`. + /// + /// Keys are interpreted as signed `i128` values without changing their + /// stored bytes. Forward iteration yields the nearest key greater than or + /// equal to `key` first; reverse iteration yields the nearest key less than + /// or equal to `key` first. Both directions follow signed chronological + /// order across the Unix epoch. #[must_use] pub fn iter_from(&self, key: i128, direction: Direction) -> EventIterator<'_> { - let iter = self - .inner - .iterator(IteratorMode::From(&key.to_be_bytes(), direction)); - EventIterator { inner: iter } + let key_bytes = key.to_be_bytes(); + match (is_signed_negative(&key_bytes), direction) { + (true, Direction::Forward) => EventIterator::new( + self.inner, + EventKeyRegion::Negative, + IteratorMode::From(&key_bytes, Direction::Forward), + Some(EventKeyRegion::NonNegative), + Direction::Forward, + ), + (true, Direction::Reverse) => EventIterator::new( + self.inner, + EventKeyRegion::Negative, + IteratorMode::From(&key_bytes, Direction::Reverse), + None, + Direction::Reverse, + ), + (false, Direction::Forward) => EventIterator::new( + self.inner, + EventKeyRegion::NonNegative, + IteratorMode::From(&key_bytes, Direction::Forward), + None, + Direction::Forward, + ), + (false, Direction::Reverse) => EventIterator::new( + self.inner, + EventKeyRegion::NonNegative, + IteratorMode::From(&key_bytes, Direction::Reverse), + Some(EventKeyRegion::Negative), + Direction::Reverse, + ), + } } - /// Creates an iterator over key-value pairs for the entire events. + /// Creates an iterator over all events in ascending signed `i128` + /// chronological key order. #[must_use] pub fn iter_forward(&self) -> EventIterator<'_> { - let iter = self.inner.iterator(IteratorMode::Start); - EventIterator { inner: iter } + EventIterator::new( + self.inner, + EventKeyRegion::Negative, + IteratorMode::From(&FIRST_NEGATIVE_EVENT_KEY, Direction::Forward), + Some(EventKeyRegion::NonNegative), + Direction::Forward, + ) + } + + /// Creates an iterator over all events in descending signed `i128` + /// chronological key order. + #[must_use] + pub fn iter_reverse(&self) -> EventIterator<'_> { + EventIterator::new( + self.inner, + EventKeyRegion::NonNegative, + IteratorMode::End, + Some(EventKeyRegion::Negative), + Direction::Reverse, + ) } #[cfg(test)] @@ -3299,10 +3397,10 @@ impl<'a> EventDb<'a> { /// Removes all events whose timestamp is strictly before `before`. /// - /// Events are stored with an i128 key whose upper 64 bits encode the - /// timestamp in nanoseconds. This method iterates from the beginning - /// of the event database and deletes every entry whose timestamp is - /// earlier than `before`, using batched writes for efficiency. + /// Event keys are interpreted in signed `i128` chronological order. The + /// upper 64 bits encode signed epoch nanoseconds. An event exactly at + /// `before` is retained, and cutoffs outside the supported `i64` + /// epoch-nanosecond range delete none or all events as appropriate. /// /// Returns the number of events deleted. /// @@ -3310,55 +3408,76 @@ impl<'a> EventDb<'a> { /// /// Returns an error if a database operation fails. pub fn remove_before(&self, before: Timestamp) -> Result { - let cutoff_nanos = match timestamp::to_i64_nanos(before) { - Ok(nanos) => nanos, - Err(timestamp::TimestampError::OutOfI64Range(nanos)) => { - if nanos >= 0 { - i64::MAX // far-future cutoff → delete everything - } else { - i64::MIN // far-past cutoff → delete nothing - } - } - Err(timestamp::TimestampError::Invalid(err)) => return Err(err.into()), - }; - let mut deleted: u64 = 0; + let cutoff_nanos = before.as_nanosecond(); + if cutoff_nanos <= i128::from(i64::MIN) { + return Ok(0); + } - loop { - let iter = self.inner.iterator(IteratorMode::Start); - let mut batch = rocksdb::WriteBatchWithTransaction::::default(); - let mut batch_count = 0; - - for item in iter { - let (k, _v) = item.context("cannot read from event database")?; - let key_bytes: [u8; 16] = match k.as_ref().try_into() { - Ok(b) => b, - Err(_) => continue, - }; - let key = i128::from_be_bytes(key_bytes); - let ts = (key >> 64) as i64; + if cutoff_nanos > i128::from(i64::MAX) { + let negative = self.delete_event_region(EventKeyRegion::Negative, None)?; + let non_negative = self.delete_event_region(EventKeyRegion::NonNegative, None)?; + return Ok(negative + non_negative); + } - if ts >= cutoff_nanos { - break; - } + let cutoff_nanos = i64::try_from(cutoff_nanos) + .context("cutoff within the checked i64 epoch-nanosecond range")?; + let cutoff_key = (i128::from(cutoff_nanos) << 64).to_be_bytes(); + if cutoff_nanos < 0 { + self.delete_event_region(EventKeyRegion::Negative, Some(cutoff_key)) + } else { + let negative = self.delete_event_region(EventKeyRegion::Negative, None)?; + let non_negative = if cutoff_key == FIRST_NON_NEGATIVE_EVENT_KEY { + 0 + } else { + self.delete_event_region(EventKeyRegion::NonNegative, Some(cutoff_key))? + }; + Ok(negative + non_negative) + } + } - batch.delete(&k); - batch_count += 1; + fn delete_event_region( + &self, + region: EventKeyRegion, + upper_bound: Option<[u8; 16]>, + ) -> Result { + let mut readopts = event_region_read_options(region); + if let Some(upper_bound) = upper_bound { + readopts.set_iterate_upper_bound(upper_bound); + } + let mode = match region { + EventKeyRegion::NonNegative => IteratorMode::Start, + EventKeyRegion::Negative => { + IteratorMode::From(&FIRST_NEGATIVE_EVENT_KEY, Direction::Forward) + } + }; + let iter = self.inner.iterator_opt(mode, readopts); + let mut batch = rocksdb::WriteBatchWithTransaction::::default(); + let mut batch_count = 0_usize; + let mut deleted = 0_u64; - if batch_count >= EVENT_DELETION_BATCH_SIZE { - break; - } + for item in iter { + let (key, _value) = item.context("cannot read from event database")?; + if signed_event_key(&key).is_none() { + continue; } + batch.delete(&key); + batch_count += 1; - if batch_count == 0 { - break; + if batch_count == EVENT_DELETION_BATCH_SIZE { + self.inner + .write(std::mem::take(&mut batch)) + .context("failed to delete expired events")?; + deleted += u64::try_from(batch_count).expect("batch size fits in u64"); + batch_count = 0; } + } + if batch_count > 0 { self.inner .write(batch) .context("failed to delete expired events")?; - deleted += batch_count as u64; + deleted += u64::try_from(batch_count).expect("batch size fits in u64"); } - Ok(deleted) } @@ -3439,10 +3558,39 @@ impl<'a> EventDb<'a> { #[allow(clippy::module_name_repetitions)] pub struct EventIterator<'i> { - inner: rocksdb::DBIteratorWithThreadMode< - 'i, - rocksdb::OptimisticTransactionDB, - >, + db: &'i rocksdb::OptimisticTransactionDB, + inner: PhysicalEventIterator<'i>, + next_region: Option, + direction: Direction, +} + +impl<'i> EventIterator<'i> { + fn new( + db: &'i rocksdb::OptimisticTransactionDB, + region: EventKeyRegion, + mode: IteratorMode<'_>, + next_region: Option, + direction: Direction, + ) -> Self { + Self { + db, + inner: physical_event_iterator(db, region, mode), + next_region, + direction, + } + } + + fn advance_region(&mut self) -> bool { + let Some(region) = self.next_region.take() else { + return false; + }; + let mode = match self.direction { + Direction::Forward => IteratorMode::Start, + Direction::Reverse => IteratorMode::End, + }; + self.inner = physical_event_iterator(self.db, region, mode); + true + } } #[allow(clippy::module_name_repetitions)] @@ -3471,7 +3619,11 @@ impl Iterator for EventIterator<'_> { fn next(&mut self) -> Option { let (key, kind, time, v) = loop { - let (k, v) = self.inner.next().transpose().ok().flatten()?; + let (k, v) = match self.inner.next() { + Some(Ok(item)) => item, + None if self.advance_region() => continue, + Some(Err(_)) | None => return None, + }; let key: [u8; 16] = if let Ok(key) = k.as_ref().try_into() { key @@ -3959,7 +4111,7 @@ mod tests { BlocklistRadiusFields, BlocklistRdp, BlocklistRdpFields, BlocklistSmb, BlocklistSmbFields, BlocklistSmtp, BlocklistSmtpFields, BlocklistSsh, BlocklistSshFields, BlocklistTls, BlocklistTlsFields, CryptocurrencyMiningPool, - CryptocurrencyMiningPoolFields, DceRpcContext, DgaFields, DnsCovertChannel, + CryptocurrencyMiningPoolFields, DceRpcContext, DgaFields, Direction, DnsCovertChannel, DnsEventFields, DomainGenerationAlgorithm, Event, EventFilter, EventKind, EventMessage, ExternalDdos, ExternalDdosFields, ExtraThreat, ExtraThreatFields, FtpBruteForce, FtpBruteForceFields, FtpEventFields, FtpPlainText, HttpEventFields, HttpThreat, @@ -4078,6 +4230,21 @@ mod tests { event.sensor } + fn message_at_nanos(nanos: i64) -> EventMessage { + let mut message = example_message( + EventKind::DnsCovertChannel, + EventCategory::CommandAndControl, + ); + message.time = timestamp::from_i64_nanos(nanos).expect("i64 nanoseconds fit Timestamp"); + message + } + + fn key_timestamp(key: i128) -> i64 { + (key >> 64) + .try_into() + .expect("the upper half of an event key is an i64 timestamp") + } + #[test] fn event_db_put() { let (_permit, store) = setup_store(); @@ -4107,6 +4274,110 @@ mod tests { assert!(iter.next().is_none()); } + #[test] + fn event_iterators_use_signed_chronological_order() { + let (_permit, store) = setup_store(); + let db = store.events(); + let timestamps = [i64::MIN, -1, 0, 1, i64::MAX]; + + for nanos in [0, i64::MAX, -1, i64::MIN, 1] { + db.put(&message_at_nanos(nanos)).unwrap(); + } + + let forward: Vec<_> = db + .iter_forward() + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(forward, timestamps); + + let reverse: Vec<_> = db + .iter_reverse() + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(reverse, timestamps.into_iter().rev().collect::>()); + } + + #[test] + fn iter_from_is_inclusive_and_crosses_only_the_later_sign_region() { + let (_permit, store) = setup_store(); + let db = store.events(); + let mut keys = HashMap::new(); + for nanos in [-5, -1, 0, 4] { + keys.insert(nanos, db.put(&message_at_nanos(nanos)).unwrap()); + } + + let from_negative: Vec<_> = db + .iter_from( + *keys.get(&-5).expect("inserted key is present"), + Direction::Forward, + ) + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(from_negative, [-5, -1, 0, 4]); + + let from_non_negative: Vec<_> = db + .iter_from( + *keys.get(&4).expect("inserted key is present"), + Direction::Reverse, + ) + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(from_non_negative, [4, 0, -1, -5]); + + let forward_within_non_negative: Vec<_> = db + .iter_from( + *keys.get(&0).expect("inserted key is present"), + Direction::Forward, + ) + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(forward_within_non_negative, [0, 4]); + + let reverse_within_negative: Vec<_> = db + .iter_from( + *keys.get(&-1).expect("inserted key is present"), + Direction::Reverse, + ) + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(reverse_within_negative, [-1, -5]); + + let missing_negative = i128::from(-3_i64) << 64; + let forward: Vec<_> = db + .iter_from(missing_negative, Direction::Forward) + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(forward, [-1, 0, 4]); + + let missing_positive = i128::from(2_i64) << 64; + let reverse: Vec<_> = db + .iter_from(missing_positive, Direction::Reverse) + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + assert_eq!(reverse, [0, -1, -5]); + } + + #[test] + fn iterator_preserves_key_order_with_equal_timestamps() { + let (_permit, store) = setup_store(); + let db = store.events(); + let message = message_at_nanos(-1); + let mut inserted = Vec::new(); + + for _ in 0..3 { + inserted.push(db.put(&message).unwrap()); + } + inserted.sort_unstable(); + + let iterated: Vec<_> = db.iter_forward().map(|item| item.unwrap().0).collect(); + assert_eq!(iterated, inserted); + assert!( + iterated + .iter() + .all(|key| key_timestamp(*key) == -1 && key.to_be_bytes().len() == 16) + ); + } + #[test] fn event_db_put_resolves_country_codes() { let orig_addr = IpAddr::V4(Ipv4Addr::LOCALHOST); @@ -8066,6 +8337,40 @@ mod tests { assert_eq!(deleted, 0); } + #[test] + fn remove_before_obeys_strict_signed_timestamp_boundaries() { + let cases = [ + (i64::MIN, Vec::new()), + (-1, vec![i64::MIN]), + (0, vec![i64::MIN, -1]), + (1, vec![i64::MIN, -1, 0]), + (i64::MAX, vec![i64::MIN, -1, 0, 1]), + ]; + + for (cutoff, expected_deleted) in cases { + let (_permit, store) = setup_store(); + let db = store.events(); + for nanos in [i64::MAX, 0, i64::MIN, 1, -1] { + db.put(&message_at_nanos(nanos)).unwrap(); + } + + assert_eq!( + db.remove_before(timestamp::from_i64_nanos(cutoff).unwrap()) + .unwrap(), + u64::try_from(expected_deleted.len()).unwrap() + ); + let remaining: Vec<_> = db + .iter_forward() + .map(|item| key_timestamp(item.unwrap().0)) + .collect(); + let expected_remaining: Vec<_> = [i64::MIN, -1, 0, 1, i64::MAX] + .into_iter() + .filter(|nanos| !expected_deleted.contains(nanos)) + .collect(); + assert_eq!(remaining, expected_remaining); + } + } + #[test] fn remove_before_exact_cutoff_is_not_deleted() { let (_permit, store) = setup_store(); @@ -8211,10 +8516,9 @@ mod tests { let total: usize = 1_500; for i in 0..total { - let time = - base_time + chrono::Duration::seconds(i64::try_from(i).expect("small value")); + let nanos = i64::try_from(i).expect("small value") - 750; let msg = EventMessage { - time: msg_time(time), + time: timestamp::from_i64_nanos(nanos).expect("small nanosecond value is valid"), kind: EventKind::DnsCovertChannel, fields: fields.clone(), }; @@ -8223,8 +8527,8 @@ mod tests { assert_eq!(db.iter_forward().count(), total); - // Cutoff well after all events. - let cutoff = msg_time(Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).unwrap()); + // The cutoff spans both physical sign regions. + let cutoff = timestamp::from_i64_nanos(1_000).unwrap(); let deleted = db.remove_before(cutoff).unwrap(); assert_eq!(deleted, u64::try_from(total).unwrap()); assert_eq!(db.iter_forward().count(), 0); From d008b19a6b53a7c680155c133e5d2a2b0ccbfc48 Mon Sep 17 00:00:00 2001 From: "octoaide[bot]" <204759324+octoaide[bot]@users.noreply.github.com> Date: Fri, 4 Sep 2026 09:28:55 -0700 Subject: [PATCH 2/2] Address PR feedback --- src/event.rs | 106 ++++++++++++++++++++++++++++++++++----------------- 1 file changed, 72 insertions(+), 34 deletions(-) diff --git a/src/event.rs b/src/event.rs index 134b4bfa..650d52e3 100644 --- a/src/event.rs +++ b/src/event.rs @@ -159,10 +159,9 @@ const EVENT_DELETION_BATCH_SIZE: usize = 1000; const FIRST_NON_NEGATIVE_EVENT_KEY: [u8; 16] = 0_i128.to_be_bytes(); const FIRST_NEGATIVE_EVENT_KEY: [u8; 16] = i128::MIN.to_be_bytes(); -type PhysicalEventIterator<'i> = rocksdb::DBIteratorWithThreadMode< - 'i, - rocksdb::OptimisticTransactionDB, ->; +type PhysicalEventDb = rocksdb::OptimisticTransactionDB; +type PhysicalEventIterator<'i> = rocksdb::DBIteratorWithThreadMode<'i, PhysicalEventDb>; +type PhysicalEventSnapshot<'i> = rocksdb::SnapshotWithThreadMode<'i, PhysicalEventDb>; #[derive(Clone, Copy)] enum EventKeyRegion { @@ -196,11 +195,11 @@ fn event_region_read_options(region: EventKeyRegion) -> rocksdb::ReadOptions { } fn physical_event_iterator<'i>( - db: &'i rocksdb::OptimisticTransactionDB, + snapshot: &PhysicalEventSnapshot<'i>, region: EventKeyRegion, mode: IteratorMode<'_>, ) -> PhysicalEventIterator<'i> { - db.iterator_opt(mode, event_region_read_options(region)) + snapshot.iterator_opt(mode, event_region_read_options(region)) } // event kind @@ -3558,8 +3557,9 @@ impl<'a> EventDb<'a> { #[allow(clippy::module_name_repetitions)] pub struct EventIterator<'i> { - db: &'i rocksdb::OptimisticTransactionDB, + // The iterator must be dropped before the snapshot it uses. inner: PhysicalEventIterator<'i>, + snapshot: PhysicalEventSnapshot<'i>, next_region: Option, direction: Direction, } @@ -3572,9 +3572,11 @@ impl<'i> EventIterator<'i> { next_region: Option, direction: Direction, ) -> Self { + let snapshot = db.snapshot(); + let inner = physical_event_iterator(&snapshot, region, mode); Self { - db, - inner: physical_event_iterator(db, region, mode), + inner, + snapshot, next_region, direction, } @@ -3588,7 +3590,7 @@ impl<'i> EventIterator<'i> { Direction::Forward => IteratorMode::Start, Direction::Reverse => IteratorMode::End, }; - self.inner = physical_event_iterator(self.db, region, mode); + self.inner = physical_event_iterator(&self.snapshot, region, mode); true } } @@ -4361,23 +4363,56 @@ mod tests { fn iterator_preserves_key_order_with_equal_timestamps() { let (_permit, store) = setup_store(); let db = store.events(); - let message = message_at_nanos(-1); + let mut dns_message = example_message( + EventKind::DnsCovertChannel, + EventCategory::CommandAndControl, + ); + dns_message.time = timestamp::from_i64_nanos(-1).unwrap(); + let mut locky_message = example_message(EventKind::LockyRansomware, EventCategory::Impact); + locky_message.time = dns_message.time; let mut inserted = Vec::new(); - for _ in 0..3 { - inserted.push(db.put(&message).unwrap()); + for message in [&locky_message, &dns_message, &dns_message] { + inserted.push(db.put(message).unwrap()); } inserted.sort_unstable(); - let iterated: Vec<_> = db.iter_forward().map(|item| item.unwrap().0).collect(); - assert_eq!(iterated, inserted); + let forward: Vec<_> = db.iter_forward().map(|item| item.unwrap().0).collect(); + assert_eq!(forward, inserted); + let reverse: Vec<_> = db.iter_reverse().map(|item| item.unwrap().0).collect(); + assert_eq!(reverse, inserted.into_iter().rev().collect::>()); assert!( - iterated + forward .iter() .all(|key| key_timestamp(*key) == -1 && key.to_be_bytes().len() == 16) ); } + #[test] + fn event_iterator_uses_one_snapshot_across_sign_regions() { + let (_permit, store) = setup_store(); + let db = store.events(); + let negative_key = db.put(&message_at_nanos(-1)).unwrap(); + let non_negative_key = db.put(&message_at_nanos(1)).unwrap(); + let mut iter = db.iter_forward(); + + assert_eq!(iter.next().unwrap().unwrap().0, negative_key); + + let added_key = db.put(&message_at_nanos(2)).unwrap(); + db.inner.delete(non_negative_key.to_be_bytes()).unwrap(); + + assert_eq!( + iter.map(|item| item.unwrap().0).collect::>(), + [non_negative_key] + ); + assert_eq!( + db.iter_forward() + .map(|item| item.unwrap().0) + .collect::>(), + [negative_key, added_key] + ); + } + #[test] fn event_db_put_resolves_country_codes() { let orig_addr = IpAddr::V4(Ipv4Addr::LOCALHOST); @@ -8458,11 +8493,8 @@ mod tests { let (_permit, store) = setup_store(); let db = store.events(); - let msg = example_message( - EventKind::DnsCovertChannel, - EventCategory::CommandAndControl, - ); - db.put(&msg).unwrap(); + db.put(&message_at_nanos(0)).unwrap(); + db.put(&message_at_nanos(i64::MAX)).unwrap(); // A cutoff so far in the future that nanoseconds overflow (after 2262). let far_future = @@ -8472,7 +8504,7 @@ mod tests { "cutoff should overflow nanosecond representation" ); let deleted = db.remove_before(far_future).unwrap(); - assert_eq!(deleted, 1); + assert_eq!(deleted, 2); assert_eq!(db.iter_forward().count(), 0); } @@ -8481,8 +8513,8 @@ mod tests { let (_permit, store) = setup_store(); let db = store.events(); - // Insert more than BATCH_SIZE (1000) events so deletion spans - // multiple batches. + // Insert more than BATCH_SIZE (1000) events in both physical regions + // so each region has a full-batch flush and a remainder batch. let base_time = Utc.with_ymd_and_hms(2020, 1, 1, 0, 0, 0).unwrap(); let fields = bincode::serialize(&DnsEventFields { sensor: "s1".to_string(), @@ -8514,21 +8546,27 @@ mod tests { }) .unwrap(); - let total: usize = 1_500; - for i in 0..total { - let nanos = i64::try_from(i).expect("small value") - 750; - let msg = EventMessage { - time: timestamp::from_i64_nanos(nanos).expect("small nanosecond value is valid"), - kind: EventKind::DnsCovertChannel, - fields: fields.clone(), - }; - db.put(&msg).unwrap(); + let events_per_region = super::EVENT_DELETION_BATCH_SIZE + 1; + for i in 0..events_per_region { + let offset = i64::try_from(i).expect("small value"); + for nanos in [-offset - 1, offset] { + let msg = EventMessage { + time: timestamp::from_i64_nanos(nanos) + .expect("small nanosecond value is valid"), + kind: EventKind::DnsCovertChannel, + fields: fields.clone(), + }; + db.put(&msg).unwrap(); + } } + let total = events_per_region * 2; assert_eq!(db.iter_forward().count(), total); // The cutoff spans both physical sign regions. - let cutoff = timestamp::from_i64_nanos(1_000).unwrap(); + let cutoff = + timestamp::from_i64_nanos(i64::try_from(events_per_region).expect("small value")) + .unwrap(); let deleted = db.remove_before(cutoff).unwrap(); assert_eq!(deleted, u64::try_from(total).unwrap()); assert_eq!(db.iter_forward().count(), 0);