Skip to content

Commit 9fc96f8

Browse files
committed
fix(ilp): keep the line time out of a declared time-key projection
An ILP batch now projects its inferred time column apart from the tag and field columns. The catalog merge drops it for a timeseries collection whose declared time key names one of its fields, so a steady ingest stream no longer proposes a new descriptor version on every flush. The batch flush also takes its write leases through retry_through_drain, so a flush that meets a descriptor drain waits it out instead of dropping the connection.
1 parent 17c434b commit 9fc96f8

6 files changed

Lines changed: 308 additions & 22 deletions

File tree

‎nodedb/src/control/catalog_entry/persist_collection.rs‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,10 @@ pub async fn persist_collection_replicated(
4141
/// Union ingest-inferred fields into a collection's schema projection and
4242
/// persist the result through the replicated metadata path.
4343
///
44+
/// `time_column` is the column a timeseries ingest inferred for the row
45+
/// time. [`crate::control::security::catalog::merge_inferred_fields`] decides
46+
/// whether the collection's projection takes it.
47+
///
4448
/// Returns `true` when the projection changed and a new descriptor version was
4549
/// proposed, `false` when the collection is absent or already carries every
4650
/// inferred field (the overwhelmingly common case on a steady ingest stream —
@@ -66,6 +70,7 @@ pub async fn merge_collection_fields_replicated(
6670
database_id: DatabaseId,
6771
tenant_id: u64,
6872
name: &str,
73+
time_column: Option<&(String, String)>,
6974
inferred_fields: &[(String, String)],
7075
) -> crate::Result<bool> {
7176
let Some(mut coll) =
@@ -76,7 +81,11 @@ pub async fn merge_collection_fields_replicated(
7681
else {
7782
return Ok(false);
7883
};
79-
if !crate::control::security::catalog::merge_inferred_fields(&mut coll, inferred_fields) {
84+
if !crate::control::security::catalog::merge_inferred_fields(
85+
&mut coll,
86+
time_column,
87+
inferred_fields,
88+
) {
8089
return Ok(false);
8190
}
8291
persist_collection_replicated(state, &coll).await?;

‎nodedb/src/control/security/catalog/collections.rs‎

Lines changed: 125 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,11 @@ use super::types::{COLLECTIONS, StoredCollection, SystemCatalog, catalog_err};
1616
/// appended. Bitemporal collections always expose their reserved BIGINT
1717
/// fields exactly once. Returns `true` when the projection changed.
1818
///
19+
/// `time_column` is the column an ingest inferred for the row time. A
20+
/// timeseries collection whose declared time key names one of its fields
21+
/// stores the row time in that field. Its projection therefore never gains
22+
/// the inferred time column. Every other collection takes it as a field.
23+
///
1924
/// Deliberately pure: a collection descriptor is replicated catalog state, so
2025
/// the merged record has to reach storage through the replicated metadata path
2126
/// (see `catalog_entry::persist_collection`) rather than a local write. Mutating
@@ -24,10 +29,12 @@ use super::types::{COLLECTIONS, StoredCollection, SystemCatalog, catalog_err};
2429
/// replaying that entry after a restart wedges the metadata applier.
2530
pub fn merge_inferred_fields(
2631
collection: &mut StoredCollection,
32+
time_column: Option<&(String, String)>,
2733
inferred_fields: &[(String, String)],
2834
) -> bool {
35+
let time_column = time_column.filter(|_| !declares_time_key_field(collection));
2936
let mut changed = false;
30-
for (field, field_type) in inferred_fields {
37+
for (field, field_type) in time_column.into_iter().chain(inferred_fields) {
3138
// Reserved bitemporal columns are schema-owned. Never let an ingest
3239
// projection supply their type or add a duplicate; normalization below
3340
// owns them entirely.
@@ -75,6 +82,20 @@ pub fn merge_inferred_fields(
7582
changed
7683
}
7784

85+
/// Whether a timeseries collection's declared time key names one of its
86+
/// fields. The Data Plane then builds its memtable from that declaration and
87+
/// writes each row's time into that field.
88+
fn declares_time_key_field(collection: &StoredCollection) -> bool {
89+
let nodedb_types::CollectionType::Columnar(nodedb_types::ColumnarProfile::Timeseries {
90+
time_key,
91+
..
92+
}) = &collection.collection_type
93+
else {
94+
return false;
95+
};
96+
collection.fields.iter().any(|(field, _)| field == time_key)
97+
}
98+
7899
impl SystemCatalog {
79100
/// Store a collection record. A record with no incarnation is refused
80101
/// with [`crate::Error::CollectionUnstamped`].
@@ -514,15 +535,18 @@ mod tests {
514535

515536
assert!(merge_inferred_fields(
516537
&mut coll,
538+
None,
517539
&[("first".to_owned(), "BIGINT".to_owned())]
518540
));
519541
assert!(merge_inferred_fields(
520542
&mut coll,
543+
None,
521544
&[("second".to_owned(), "FLOAT".to_owned())]
522545
));
523546
// A known name never re-types an existing column, and reports no change.
524547
assert!(!merge_inferred_fields(
525548
&mut coll,
549+
None,
526550
&[("first".to_owned(), "BOOLEAN".to_owned())]
527551
));
528552

@@ -536,6 +560,103 @@ mod tests {
536560
);
537561
}
538562

563+
fn ilp_time_column() -> (String, String) {
564+
("timestamp".to_owned(), "TIMESTAMP".to_owned())
565+
}
566+
567+
/// A collection declared with `ts BIGINT TIME_KEY` stores the ILP line
568+
/// time in `ts`. The inferred `timestamp` column must not reach its
569+
/// projection, or every flush proposes a new descriptor version.
570+
#[test]
571+
fn merge_skips_the_inferred_time_column_for_a_declared_time_key() {
572+
let mut coll = make_coll(1, "crash_ilp_ts_bulk");
573+
coll.collection_type = CollectionType::timeseries("ts", "1h");
574+
coll.fields = vec![
575+
("ts".to_owned(), "BIGINT TIME_KEY".to_owned()),
576+
("value".to_owned(), "BIGINT".to_owned()),
577+
];
578+
579+
assert!(
580+
!merge_inferred_fields(
581+
&mut coll,
582+
Some(&ilp_time_column()),
583+
&[("value".to_owned(), "BIGINT".to_owned())]
584+
),
585+
"an ILP batch carrying only declared fields changes nothing"
586+
);
587+
assert!(
588+
merge_inferred_fields(
589+
&mut coll,
590+
Some(&ilp_time_column()),
591+
&[
592+
("host".to_owned(), "VARCHAR".to_owned()),
593+
("value".to_owned(), "BIGINT".to_owned()),
594+
("load".to_owned(), "FLOAT".to_owned()),
595+
]
596+
),
597+
"a new tag and a new field still reach the projection"
598+
);
599+
assert_eq!(
600+
coll.fields,
601+
vec![
602+
("ts".to_owned(), "BIGINT TIME_KEY".to_owned()),
603+
("value".to_owned(), "BIGINT".to_owned()),
604+
("host".to_owned(), "VARCHAR".to_owned()),
605+
("load".to_owned(), "FLOAT".to_owned()),
606+
]
607+
);
608+
}
609+
610+
/// A field literally called `timestamp` is a field, not the line time.
611+
/// It reaches the projection of a collection with a declared time key.
612+
#[test]
613+
fn merge_keeps_a_field_named_timestamp_for_a_declared_time_key() {
614+
let mut coll = make_coll(1, "metrics");
615+
coll.collection_type = CollectionType::timeseries("ts", "1h");
616+
coll.fields = vec![("ts".to_owned(), "TIMESTAMP".to_owned())];
617+
618+
assert!(merge_inferred_fields(
619+
&mut coll,
620+
Some(&ilp_time_column()),
621+
&[("timestamp".to_owned(), "BIGINT".to_owned())]
622+
));
623+
assert_eq!(
624+
coll.fields,
625+
vec![
626+
("ts".to_owned(), "TIMESTAMP".to_owned()),
627+
("timestamp".to_owned(), "BIGINT".to_owned()),
628+
]
629+
);
630+
}
631+
632+
/// With no declared time key among its fields, the Data Plane infers the
633+
/// schema and stores the line time under the inferred name. The
634+
/// projection follows it.
635+
#[test]
636+
fn merge_adds_the_inferred_time_column_without_a_declared_time_key() {
637+
let mut undeclared = make_coll(1, "events");
638+
assert!(merge_inferred_fields(
639+
&mut undeclared,
640+
Some(&ilp_time_column()),
641+
&[("value".to_owned(), "FLOAT".to_owned())]
642+
));
643+
assert_eq!(
644+
undeclared.fields,
645+
vec![ilp_time_column(), ("value".to_owned(), "FLOAT".to_owned())]
646+
);
647+
648+
// A time key absent from the field list resolves no declaration, so
649+
// the Data Plane infers here too.
650+
let mut unresolved = make_coll(1, "cpu");
651+
unresolved.collection_type = CollectionType::timeseries("ts", "1h");
652+
assert!(merge_inferred_fields(
653+
&mut unresolved,
654+
Some(&ilp_time_column()),
655+
&[]
656+
));
657+
assert_eq!(unresolved.fields, vec![ilp_time_column()]);
658+
}
659+
539660
#[test]
540661
fn merge_inferred_fields_adds_bitemporal_reserved_fields_once() {
541662
let mut coll = make_coll(1, "audit");
@@ -544,9 +665,10 @@ mod tests {
544665

545666
assert!(merge_inferred_fields(
546667
&mut coll,
668+
None,
547669
&[("value".to_owned(), "FLOAT".to_owned())]
548670
));
549-
assert!(!merge_inferred_fields(&mut coll, &[]));
671+
assert!(!merge_inferred_fields(&mut coll, None, &[]));
550672
for reserved in [TS_SYSTEM, TS_VALID_FROM, TS_VALID_UNTIL] {
551673
assert_eq!(
552674
coll.fields
@@ -570,6 +692,7 @@ mod tests {
570692

571693
assert!(merge_inferred_fields(
572694
&mut coll,
695+
None,
573696
&[
574697
(TS_VALID_FROM.to_owned(), "VARCHAR".to_owned()),
575698
(TS_VALID_UNTIL.to_owned(), "BOOLEAN".to_owned()),

0 commit comments

Comments
 (0)