Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion hyperbytedb-cli/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "hyperbytedb-cli"
version = "0.8.3"
version = "0.8.5"
edition = "2024"
description = "Interactive CLI client for HyperbyteDB — InfluxDB v1-compatible REPL, query, write, and admin"
license = "Apache-2.0"
Expand Down
2 changes: 1 addition & 1 deletion hyperbytedb-proxy/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "hyperbytedb-proxy"
version = "0.8.3"
version = "0.8.5"
edition = "2024"
description = "Health-aware HTTP reverse proxy for hyperbytedb. Provides connection draining, request hold-and-retry across rolling restarts, and per-backend health gating"
license = "Apache-2.0"
Expand Down
2 changes: 1 addition & 1 deletion hyperbytedb/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "hyperbytedb"
version = "0.8.3"
version = "0.8.5"
edition = "2024"

[dependencies]
Expand Down
2 changes: 1 addition & 1 deletion hyperbytedb/src/adapters/http/middleware.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
use axum::http::{Response, header};
use uuid::Uuid;

const VERSION: &str = "HyperbyteDB-0.8.2";
const VERSION: &str = "HyperbyteDB-0.8.5";
const BUILD: &str = "OSS";
const REQUEST_ID_HEADER: &str = "Request-Id";
const X_REQUEST_ID_HEADER: &str = "X-Request-Id";
Expand Down
67 changes: 67 additions & 0 deletions hyperbytedb/src/application/ingest_metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,14 @@ use std::sync::Arc;

use parking_lot::RwLock;

use crate::domain::chdb_naming::quote_backticks;
use crate::domain::column_mapping::ColumnMapping;
use crate::domain::field_type::merge_field_type_map;
use crate::domain::point::Point;
use crate::domain::series::series_id_for_point;
use crate::error::HyperbytedbError;
use crate::ports::metadata::{MeasurementMeta, MetadataPort};
use crate::ports::query::QueryPort;

/// Cardinality limits (0 = unlimited for that bound), matching [`crate::config::CardinalityConfig`].
#[derive(Debug, Clone, Copy, Default)]
Expand Down Expand Up @@ -226,6 +229,70 @@ pub async fn backfill_tag_metadata(
.await
}

/// Load series rows from a chDB `_series` table and persist them to RocksDB for
/// `SHOW SERIES` / `SHOW TAG VALUES` (used after MV destination dimension backfill).
pub async fn register_series_from_series_table(
metadata: &Arc<dyn MetadataPort>,
query_port: &Arc<dyn QueryPort>,
db: &str,
rp: &str,
measurement: &str,
dest_meta: &MeasurementMeta,
series_table_quoted: &str,
) -> Result<(), HyperbytedbError> {
let mapping = ColumnMapping::from_measurement_meta(dest_meta);
let tag_keys: Vec<String> = dest_meta.tag_keys.clone();

let select_cols: Vec<String> = if tag_keys.is_empty() {
vec!["series_id".to_string()]
} else {
let mut cols = vec!["series_id".to_string()];
for key in &tag_keys {
cols.push(quote_backticks(&mapping.tag_column_name(key)));
}
cols
};

let sql = format!(
"SELECT {} FROM {series_table_quoted} FINAL FORMAT TabSeparated",
select_cols.join(", ")
);
let raw = query_port.execute_sql(&sql).await?;
if raw.trim().is_empty() {
return Ok(());
}

let mut entries: Vec<(u64, BTreeMap<String, String>)> = Vec::new();
let mut tag_pairs: Vec<(String, String)> = Vec::new();
for line in raw.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let parts: Vec<&str> = line.split('\t').collect();
let sid: u64 = parts[0].parse().map_err(|e| {
HyperbytedbError::Internal(format!("invalid series_id in {series_table_quoted}: {e}"))
})?;
let mut tags = BTreeMap::new();
for (i, key) in tag_keys.iter().enumerate() {
let val = parts.get(i + 1).map(|s| s.to_string()).unwrap_or_default();
tags.insert(key.clone(), val.clone());
if !val.is_empty() {
tag_pairs.push((key.clone(), val));
}
}
entries.push((sid, tags));
}
if entries.is_empty() {
return Ok(());
}

metadata
.register_series_batch(db, rp, measurement, &entries)
.await?;
backfill_tag_metadata(metadata, db, rp, measurement, tag_pairs).await
}

/// Fast-path metadata preparation for columnar batches.
///
/// Works directly from the wire format without expanding to `Vec<Point>`,
Expand Down
29 changes: 25 additions & 4 deletions hyperbytedb/src/application/materialized_view_service.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
use std::sync::Arc;

use crate::application::ingest_metadata::register_series_from_series_table;
use crate::domain::chdb_naming::{
quoted_fact_mv_name, quoted_series_mv_name, quoted_series_table_name, quoted_table_name,
unquoted_fact_mv_name, unquoted_series_mv_name,
Expand Down Expand Up @@ -415,21 +416,41 @@ impl MaterializedViewService {

if backfill_on_create {
self.query_port.execute_sql(&backfill_fact).await?;

let backfill_series = format!("INSERT INTO {dest_series}\n{series_select}");
self.query_port.execute_sql(&backfill_series).await?;
} else {
tracing::info!(
mv = %mv.name,
db = %mv.database,
"skipping materialized view historical backfill (use WITH BACKFILL to enable)"
"skipping materialized view fact backfill (use WITH BACKFILL to enable historical data)"
);
}

let backfill_series = format!("INSERT INTO {dest_series}\n{series_select}");
self.query_port.execute_sql(&backfill_series).await?;

self.metadata
.register_measurement(&dest_db, &dest_rp, &dest_meta)
.await?;

if let Err(e) = register_series_from_series_table(
&self.metadata,
&self.query_port,
&dest_db,
&dest_rp,
&dest_measurement,
&dest_meta,
&dest_series,
)
.await
{
tracing::warn!(
mv = %mv.name,
db = %mv.database,
dest = %dest_measurement,
error = %e,
"failed to persist MV destination series metadata"
);
}

Ok::<_, HyperbytedbError>(())
}
.await;
Expand Down
133 changes: 133 additions & 0 deletions hyperbytedb/tests/compat/ddl_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1010,6 +1010,12 @@ async fn create_materialized_view_without_backfill_leaves_dest_empty_until_new_w
"dest fact table should be empty before any post-create writes"
);

let dest_series_count: u64 = mv_fact_row_count(&ctx, "mvdb_autogen_metrics_nobf_series").await;
assert!(
dest_series_count > 0,
"dest series table should be seeded at create even without fact backfill"
);

let t2 = t + 1_000_000_000;
ctx.write_and_flush("mvdb", &format!("metrics,host=h1 value=5 {t2}"))
.await
Expand Down Expand Up @@ -1042,6 +1048,133 @@ async fn create_materialized_view_without_backfill_leaves_dest_empty_until_new_w
);
}

#[tokio::test]
#[serial(chdb)]
async fn create_materialized_view_without_backfill_seeds_dest_series_metadata() {
use chrono::{TimeZone, Utc};

let ctx = match TestContext::new() {
Ok(c) => c,
Err(_) => {
eprintln!("skipping MV dest series metadata test: chDB not available");
return;
}
};

fn ts(h: u32, m: u32) -> i64 {
Utc.with_ymd_and_hms(2016, 8, 28, h, m, 0)
.unwrap()
.timestamp_nanos_opt()
.unwrap()
}

ctx.metadata.create_database("gameservers").await.unwrap();
ctx.metadata
.create_retention_policy(
"gameservers",
hyperbytedb::domain::database::RetentionPolicy {
name: "default_high".to_string(),
duration: Some(std::time::Duration::from_secs(7 * 24 * 3600)),
shard_group_duration: std::time::Duration::from_secs(3600),
replication_factor: 1,
is_default: false,
},
)
.await
.unwrap();

let minute = ts(8, 0);
let lines = format!(
"server_stats,region_id=us,server_id=s1 cpu=1i {minute}\n\
server_stats,region_id=eu,server_id=s2 cpu=1i {minute}\n\
server_stats,region_id=ap,server_id=s3 cpu=1i {minute}\n\
server_stats,region_id=us,server_id=s4 cpu=1i {minute}\n\
server_stats,region_id=eu,server_id=s5 cpu=1i {minute}",
);
ctx.write_and_flush("gameservers", &lines).await.unwrap();

let create_resp = ctx
.query(
"gameservers",
r#"CREATE MATERIALIZED VIEW "mv_server_stats" ON "gameservers" AS SELECT count("cpu") AS "num_servers" INTO "default_high"."server_stats" FROM "server_stats" GROUP BY time(1m), "region_id""#,
)
.await
.unwrap();
assert!(
create_resp.results[0].error.is_none(),
"create MV failed: {:?}",
create_resp.results[0].error
);

let dest_series_count: u64 =
mv_fact_row_count(&ctx, "gameservers_default_high_server_stats_series").await;
assert_eq!(
dest_series_count, 3,
"dest series should collapse to one row per region_id"
);

let tag_keys = ctx
.query(
"gameservers",
r#"SHOW TAG KEYS FROM "default_high"."server_stats""#,
)
.await
.unwrap();
assert!(
tag_keys.results[0].error.is_none(),
"{:?}",
tag_keys.results[0].error
);
let keys: Vec<&str> = tag_keys.results[0].series.as_ref().unwrap()[0]
.values
.iter()
.filter_map(|row| row.first().and_then(|v| v.as_str()))
.collect();
assert!(keys.contains(&"region_id"));
assert!(!keys.contains(&"server_id"));

let tag_values = ctx
.query(
"gameservers",
r#"SHOW TAG VALUES FROM "default_high"."server_stats" WITH KEY = "region_id""#,
)
.await
.unwrap();
assert!(
tag_values.results[0].error.is_none(),
"{:?}",
tag_values.results[0].error
);
let regions: Vec<&str> = tag_values.results[0].series.as_ref().unwrap()[0]
.values
.iter()
.filter_map(|row| row.get(1).and_then(|v| v.as_str()))
.collect();
assert_eq!(regions.len(), 3);
assert!(regions.contains(&"us"));
assert!(regions.contains(&"eu"));
assert!(regions.contains(&"ap"));

let show_series = ctx
.query(
"gameservers",
r#"SHOW SERIES FROM "default_high"."server_stats""#,
)
.await
.unwrap();
assert!(
show_series.results[0].error.is_none(),
"{:?}",
show_series.results[0].error
);
let series_keys = &show_series.results[0].series.as_ref().unwrap()[0].values;
assert_eq!(
series_keys.len(),
3,
"SHOW SERIES should list collapsed dest series"
);
}

#[tokio::test]
#[serial(chdb)]
async fn create_materialized_view_with_backfill_populates_history() {
Expand Down
Loading