Files
jeremyhi 78eefbe4c4 perf: batch schema export requests (#9408)
* perf: batch schema export requests

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: propagate legacy schema export failures

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix: keep credentials out of SQL response errors

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test: use a typo-safe dotted catalog name

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* refactor(cli): address schema export review feedback

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

---------

Signed-off-by: jeremyhi <fengjiachun@gmail.com>
2026-09-30 06:23:33 +00:00

2129 lines
74 KiB
Rust

// Copyright 2023 Greptime Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::path::PathBuf;
use std::sync::Arc;
use common_query::OutputData;
use common_time::Timestamp;
use common_time::range::TimestampRange;
use frontend::instance::Instance;
use operator::statement::export_logical_tables::{LogicalTableExport, LogicalTableExportLimits};
use session::context::QueryContext;
use tests_integration::cluster::GreptimeDbClusterBuilder;
use tests_integration::standalone::GreptimeDbStandaloneBuilder;
use tests_integration::test_util::execute_sql as sql;
use tokio_util::sync::CancellationToken;
async fn table(instance: &Arc<Instance>, name: &str) -> table::TableRef {
instance
.catalog_manager()
.table("greptime", "public", name, Some(&QueryContext::arc()))
.await
.unwrap()
.unwrap()
}
async fn values(instance: &Arc<Instance>, query: &str) -> Vec<Vec<datatypes::value::Value>> {
let output = sql(instance, query).await;
let batches = match output.data {
OutputData::Stream(stream) => common_recordbatch::util::collect_batches(stream)
.await
.unwrap(),
OutputData::RecordBatches(batches) => batches,
_ => panic!("expected query rows"),
};
batches
.iter()
.flat_map(|b| {
let columns = datatypes::vectors::Helper::try_into_vectors(b.columns()).unwrap();
(0..b.num_rows())
.map(|row| columns.iter().map(|col| col.get(row)).collect::<Vec<_>>())
.collect::<Vec<_>>()
})
.collect()
}
async fn create_metric_export_source_tables(
instance: &Arc<Instance>,
physical: &str,
encoding: &str,
) -> ([String; 3], Vec<table::TableRef>, String) {
sql(instance, &format!("CREATE TABLE {physical} (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY) PARTITION ON COLUMNS (host) (host < 'm', host >= 'm') ENGINE=metric WITH (physical_metric_table='', primary_key_encoding='{encoding}')")).await;
for (suffix, extra, key) in [
("cpu.v1", "zone_tag STRING,", ", zone_tag"),
("requests", "service_tag STRING,", ", service_tag"),
("empty", "", ""),
] {
let name = format!("{physical}_{suffix}");
sql(instance, &format!("CREATE TABLE \"{name}\" (host STRING, {extra} val DOUBLE, ts TIMESTAMP TIME INDEX, PRIMARY KEY(host{key})) ENGINE=metric WITH (on_physical_table='{physical}')")).await;
if suffix != "empty" {
sql(instance, &format!("INSERT INTO \"{name}\" (host,val,ts) VALUES ('a',1,1),('a',NULL,2),('z',3,3),('z',4,4)")).await;
}
}
// This live but unselected table models rows outside the routing whitelist.
sql(instance, &format!("CREATE TABLE {physical}_excluded (host STRING, huge_tag STRING, val DOUBLE, ts TIMESTAMP TIME INDEX, PRIMARY KEY(host, huge_tag)) ENGINE=metric WITH (on_physical_table='{physical}')")).await;
sql(
instance,
&format!(
"INSERT INTO {physical}_excluded (host, huge_tag, val, ts) VALUES ('z','ignore',9,2)"
),
)
.await;
let names = ["cpu.v1", "requests", "empty"].map(|suffix| format!("{physical}_{suffix}"));
let tables = vec![
table(instance, &names[0]).await,
table(instance, &names[1]).await,
table(instance, &names[2]).await,
];
let renamed = format!("renamed_{physical}");
sql(
instance,
&format!("ALTER TABLE {physical} RENAME {renamed}"),
)
.await;
(names, tables, renamed)
}
async fn roundtrip(instance: &Arc<Instance>) {
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
for (physical, encoding) in [("phy", "dense"), ("other_phy", "sparse")] {
let (names, tables, renamed) =
create_metric_export_source_tables(instance, physical, encoding).await;
let unit = LogicalTableExport::try_new(table(instance, &renamed).await, &tables).unwrap();
let range =
TimestampRange::new(Timestamp::new_millisecond(2), Timestamp::new_millisecond(4))
.unwrap();
sql(instance, &format!("CREATE TABLE target_{physical} (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY) ENGINE=metric WITH (physical_metric_table='')")).await;
let mut limits = LogicalTableExportLimits::default();
limits.writer.row_group_rows = 1;
for partitions in [1, 2, 4] {
let directory = destination.path().join(format!("{physical}_{partitions}"));
let mut ctx = QueryContext::with("greptime", "public");
ctx.set_extension(
query::datafusion::QUERY_PARALLELISM_HINT,
partitions.to_string(),
);
let summary = instance
.statement_executor()
.export_logical_tables(
&unit,
directory.to_str().unwrap(),
&Default::default(),
Some(&range),
limits,
&CancellationToken::new(),
Arc::new(ctx),
)
.await
.unwrap();
assert_eq!(summary.rows, 4);
assert_eq!(summary.files, 3);
assert_eq!(summary.skipped_rows, 1);
for (index, name) in names.iter().enumerate() {
let restored = format!("restore_{physical}_{partitions}_{index}");
let (extra, key) = [
("zone_tag STRING,", ", zone_tag"),
("service_tag STRING,", ", service_tag"),
("", ""),
][index];
sql(instance, &format!("CREATE TABLE {restored} (host STRING, {extra} val DOUBLE, ts TIMESTAMP TIME INDEX, PRIMARY KEY(host{key})) ENGINE=metric WITH (on_physical_table='target_{physical}')")).await;
sql(
instance,
&format!(
"COPY {restored} FROM '{}/{}.parquet' WITH (FORMAT='parquet')",
directory.display(),
name
),
)
.await;
let expected = values(
instance,
&format!("SELECT * FROM \"{name}\" WHERE ts >= 2 AND ts < 4 ORDER BY host, ts"),
)
.await;
let actual = values(
instance,
&format!("SELECT * FROM {restored} ORDER BY host, ts"),
)
.await;
assert_eq!(actual, expected);
}
}
let cancellation = CancellationToken::new();
cancellation.cancel();
let directory = destination.path().join(format!("cancel_{physical}"));
let result = instance
.statement_executor()
.export_logical_tables(
&unit,
directory.to_str().unwrap(),
&Default::default(),
None,
LogicalTableExportLimits::default(),
&cancellation,
QueryContext::arc(),
)
.await;
assert!(matches!(
result,
Err(operator::error::Error::LogicalTableExportCancelled { .. })
));
assert!(!directory.exists());
}
}
#[tokio::test(flavor = "multi_thread")]
async fn physical_export_standalone_roundtrip() {
common_telemetry::init_default_ut_logging();
let standalone = GreptimeDbStandaloneBuilder::new("physical_export")
.build()
.await;
roundtrip(standalone.fe_instance()).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn physical_export_distributed_roundtrip() {
common_telemetry::init_default_ut_logging();
let cluster = GreptimeDbClusterBuilder::new("physical_export")
.await
.with_datanodes(2)
.with_local_file_access(
common_datasource::object_store::LocalFileAccess::sandboxed(
common_test_util::find_workspace_path("."),
)
.unwrap(),
)
.build(false)
.await;
roundtrip(cluster.fe_instance()).await;
}
fn database_export_request(directory: &std::path::Path) -> table::requests::CopyDatabaseRequest {
table::requests::CopyDatabaseRequest {
catalog_name: "greptime".into(),
schema_name: "public".into(),
location: format!("{}/", directory.display()),
with: [
("format".into(), "parquet".into()),
("parallelism".into(), "2".into()),
]
.into(),
connection: Default::default(),
time_range: Some(
TimestampRange::new(Timestamp::new_millisecond(2), Timestamp::new_millisecond(4))
.unwrap(),
),
}
}
async fn database_export_roundtrip(instance: &Arc<Instance>, parallelism: usize) {
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
let (first_logical_table_names, _, renamed_physical_table) =
create_metric_export_source_tables(instance, "db_a", "dense").await;
let (second_logical_table_names, _, _) =
create_metric_export_source_tables(instance, "db_b", "sparse").await;
sql(
instance,
"CREATE TABLE audit (host STRING, val DOUBLE, ts TIMESTAMP TIME INDEX, PRIMARY KEY(host))",
)
.await;
sql(
instance,
"INSERT INTO audit VALUES ('a',1,1),('z',NULL,2),('z',4,3)",
)
.await;
sql(instance, "CREATE VIEW dashboard AS SELECT * FROM audit").await;
let selected = vec![
first_logical_table_names[0].clone(),
first_logical_table_names[2].clone(),
second_logical_table_names[1].clone(),
"audit".into(),
];
let mut names = selected.clone();
names.extend([renamed_physical_table, "dashboard".into()]);
let mut req = database_export_request(&destination.path().join("data"));
req.with
.insert("parallelism".into(), parallelism.to_string());
let executor = instance.statement_executor();
let captured = executor
.capture_database_export_tables(&req, None, &QueryContext::arc())
.await
.unwrap();
assert_eq!(captured.len(), 9);
let captured = executor
.capture_database_export_tables(&req, Some(&names), &QueryContext::arc())
.await
.unwrap();
assert_eq!(captured.len(), 4);
let plan = executor
.prepare_database_export(req.clone(), captured)
.await
.unwrap();
assert_eq!(plan.job_count_for_test(), 3);
assert_eq!(
std::fs::read_dir(destination.path().join("data"))
.unwrap()
.count(),
0
);
let summary = instance
.export_database_for_test(
req,
Some(&names),
&CancellationToken::new(),
QueryContext::arc(),
)
.await
.unwrap();
assert_eq!(summary.rows, 6);
let expected = selected
.iter()
.map(|name| {
destination
.path()
.join("data")
.join(format!("{name}.parquet"))
})
.collect::<std::collections::BTreeSet<_>>();
assert_eq!(
summary
.output_files
.into_iter()
.map(PathBuf::from)
.collect::<std::collections::BTreeSet<_>>(),
expected
);
assert_eq!(
std::fs::read_dir(destination.path().join("data"))
.unwrap()
.count(),
4
);
sql(instance, "CREATE TABLE restored_phy (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY) ENGINE=metric WITH (physical_metric_table='')").await;
for (index, name) in selected.iter().enumerate() {
let restored = format!("restored_{index}");
if index < 3 {
let (extra, key) = [
("zone_tag STRING,", ", zone_tag"),
("", ""),
("service_tag STRING,", ", service_tag"),
][index];
sql(instance, &format!("CREATE TABLE {restored} (host STRING, {extra} val DOUBLE, ts TIMESTAMP TIME INDEX, PRIMARY KEY(host{key})) ENGINE=metric WITH (on_physical_table='restored_phy')")).await;
} else {
sql(
instance,
&format!("CREATE TABLE {restored} LIKE \"{name}\""),
)
.await;
}
sql(
instance,
&format!(
"COPY {restored} FROM '{}/data/{name}.parquet' WITH (FORMAT='parquet')",
destination.path().display()
),
)
.await;
assert_eq!(
table(instance, name).await.schema().column_schemas(),
table(instance, &restored).await.schema().column_schemas()
);
assert_eq!(
values(
instance,
&format!("SELECT * FROM {restored} ORDER BY host, ts")
)
.await,
values(
instance,
&format!("SELECT * FROM \"{name}\" WHERE ts >= 2 AND ts < 4 ORDER BY host, ts")
)
.await
);
}
let ordinary_req = database_export_request(&destination.path().join("captured"));
let captured = executor
.capture_database_export_tables(
&ordinary_req,
Some(&["audit".into()]),
&QueryContext::arc(),
)
.await
.unwrap();
let plan = executor
.prepare_database_export(ordinary_req, captured)
.await
.unwrap();
sql(instance, "ALTER TABLE audit RENAME original_audit").await;
sql(instance, "CREATE TABLE audit LIKE original_audit").await;
sql(instance, "INSERT INTO audit VALUES ('replacement',99,2)").await;
let result = executor
.export_database(plan, &CancellationToken::new(), QueryContext::arc())
.await
.unwrap();
assert_eq!(result.rows, 2);
sql(
instance,
"CREATE TABLE captured_restore LIKE original_audit",
)
.await;
sql(
instance,
&format!(
"COPY captured_restore FROM '{}/captured/audit.parquet' WITH (FORMAT='parquet')",
destination.path().display()
),
)
.await;
assert_eq!(
values(instance, "SELECT * FROM captured_restore ORDER BY host,ts").await,
values(
instance,
"SELECT * FROM original_audit WHERE ts >= 2 AND ts < 4 ORDER BY host,ts"
)
.await
);
for suffix in ["?attempt=/", "#attempt/"] {
let path = destination.path().join("invalid_destination");
let mut req = database_export_request(&path);
req.location = format!("{}{suffix}", url::Url::from_file_path(&path).unwrap());
req.with.insert("parallelism".into(), "1".into());
let result = instance
.export_database_for_test(
req,
Some(&["audit".into(), "original_audit".into()]),
&CancellationToken::new(),
QueryContext::arc(),
)
.await;
assert!(matches!(
result,
Err(frontend::error::Error::TableOperation {
source: operator::error::Error::InvalidCopyDatabasePath { .. },
..
})
));
assert!(!path.exists());
assert!(
tests_integration::test_util::try_execute_sql(
instance,
&format!(
"COPY DATABASE public TO '{}{suffix}' WITH (FORMAT='parquet')",
url::Url::from_file_path(&path).unwrap()
)
)
.await
.is_err()
);
assert!(!path.exists());
}
let token = CancellationToken::new();
token.cancel();
let req = database_export_request(&destination.path().join("cancelled"));
let result = instance
.export_database_for_test(req, Some(&names), &token, QueryContext::arc())
.await;
assert!(matches!(
result,
Err(frontend::error::Error::TableOperation {
source: operator::error::Error::DatabaseExportCancelled { .. },
..
})
));
assert!(!destination.path().join("cancelled").exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn database_export_standalone_roundtrip() {
let standalone = GreptimeDbStandaloneBuilder::new("database_export")
.build()
.await;
database_export_roundtrip(standalone.fe_instance(), 1).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn database_export_distributed_roundtrip() {
let cluster = GreptimeDbClusterBuilder::new("database_export")
.await
.with_datanodes(2)
.with_local_file_access(
common_datasource::object_store::LocalFileAccess::sandboxed(
common_test_util::find_workspace_path("."),
)
.unwrap(),
)
.build(false)
.await;
database_export_roundtrip(cluster.fe_instance(), 4).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn database_export_rejects_invalid_members_before_output() {
use common_meta::key::table_info::TableInfoKey;
use common_meta::key::table_route::{TableRouteKey, TableRouteValue};
use common_meta::key::{MetadataKey, MetadataValue};
use common_meta::rpc::store::PutRequest;
let standalone = GreptimeDbStandaloneBuilder::new("invalid_database_export")
.build()
.await;
let instance = standalone.fe_instance();
let (_, tables, physical) =
create_metric_export_source_tables(instance, "invalid_phy", "dense").await;
let executor = instance.statement_executor();
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
for name in ["Foo", "foo"] {
sql(
instance,
&format!("CREATE TABLE \"{name}\" (ts TIMESTAMP TIME INDEX)"),
)
.await;
}
let case_tables = vec![table(instance, "Foo").await, table(instance, "foo").await];
assert_ne!(
case_tables[0].table_info().table_id(),
case_tables[1].table_info().table_id()
);
let case_path = destination.path().join("case_aliases");
let mut case_req = database_export_request(&case_path);
case_req.with.insert("parallelism".into(), "1".into());
let result = instance
.export_database_for_test(
case_req.clone(),
Some(&["Foo".into(), "foo".into()]),
&CancellationToken::new(),
QueryContext::arc(),
)
.await;
assert!(matches!(result,
Err(frontend::error::Error::TableOperation {
source: operator::error::Error::InvalidDatabaseExport { reason }, ..
}) if reason == "duplicate output name: foo"));
assert!(!case_path.exists());
case_req.location = "s3://export-bucket/data/".into();
case_req.connection.extend([
("region".into(), "us-east-1".into()),
("access_key_id".into(), "test-key".into()),
("secret_access_key".into(), "test-secret".into()),
]);
let plan = executor
.prepare_database_export(case_req, case_tables)
.await
.unwrap();
assert_eq!(plan.job_count_for_test(), 2);
let req = database_export_request(&destination.path().join("data"));
let key = TableRouteKey::new(tables[2].table_info().table_id()).to_bytes();
let original = standalone
.kv_backend
.get(&key)
.await
.unwrap()
.unwrap()
.value;
for route in [
None,
Some(TableRouteValue::physical(vec![])),
Some(TableRouteValue::logical(u32::MAX)),
] {
if let Some(route) = route {
standalone
.kv_backend
.put(
PutRequest::new()
.with_key(key.clone())
.with_value(route.try_as_raw_value().unwrap()),
)
.await
.unwrap();
} else {
standalone.kv_backend.delete(&key, false).await.unwrap();
}
let result = executor
.prepare_database_export(req.clone(), tables.clone())
.await;
assert!(matches!(
result,
Err(operator::error::Error::InvalidDatabaseExport { .. })
));
assert!(!destination.path().join("data").exists());
}
standalone
.kv_backend
.put(PutRequest::new().with_key(key.clone()).with_value(original))
.await
.unwrap();
let plan = executor
.prepare_database_export(req.clone(), tables.clone())
.await
.unwrap();
// Preparation does not waive PR3's membership revalidation before a scan.
standalone.kv_backend.delete(&key, false).await.unwrap();
let result = executor
.export_database(plan, &CancellationToken::new(), QueryContext::arc())
.await;
assert!(matches!(
result,
Err(operator::error::Error::InvalidLogicalTableExport { .. })
));
assert_eq!(
std::fs::read_dir(destination.path().join("data"))
.unwrap()
.count(),
0
);
let physical_id = table(instance, &physical).await.table_info().table_id();
let route = TableRouteValue::logical(physical_id)
.try_as_raw_value()
.unwrap();
standalone
.kv_backend
.put(PutRequest::new().with_key(key).with_value(route))
.await
.unwrap();
standalone
.kv_backend
.delete(&TableInfoKey::new(physical_id).to_bytes(), false)
.await
.unwrap();
let result = executor
.prepare_database_export(req.clone(), tables.clone())
.await;
assert!(matches!(
result,
Err(operator::error::Error::InvalidDatabaseExport { .. })
));
let duplicate_name = tables[0].table_info().name.clone();
let result = executor
.prepare_database_export(req, vec![tables[0].clone(), tables[0].clone()])
.await;
assert!(matches!(result,
Err(operator::error::Error::InvalidDatabaseExport { reason })
if reason == format!("duplicate output name: {duplicate_name}")));
assert_eq!(
std::fs::read_dir(destination.path().join("data"))
.unwrap()
.count(),
0
);
}
#[tokio::test(flavor = "multi_thread")]
async fn database_export_preserves_valid_table_names() {
let standalone = GreptimeDbStandaloneBuilder::new("database_export_names")
.build()
.await;
let instance = standalone.fe_instance();
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
sql(instance, "CREATE TABLE names_physical (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY) ENGINE=metric WITH (physical_metric_table='')").await;
let mut names = vec!["ordinary#b".to_string(), "metric#b".to_string()];
if !cfg!(windows) {
names.extend(["ordinary:b".to_string(), "metric:b".to_string()]);
}
for name in &names {
let engine = if name.starts_with("metric") {
"ENGINE=metric WITH (on_physical_table='names_physical')"
} else {
""
};
sql(instance, &format!("CREATE TABLE \"{name}\" (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY) {engine}")).await;
sql(
instance,
&format!("INSERT INTO \"{name}\" (ts,val,host) VALUES (1,42,'h')"),
)
.await;
}
sql(
instance,
"CREATE VIEW names_view AS SELECT * FROM \"ordinary#b\"",
)
.await;
sql(instance, "CREATE DATABASE names_restored").await;
for name in &names {
sql(instance, &format!("CREATE TABLE names_restored.\"{name}\" (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY)")).await;
}
for (attempt, file_url, legacy) in [
("plain", false, false),
("url", true, false),
("legacy", true, true),
] {
let directory = destination.path().join(attempt);
let mut req = database_export_request(&directory);
req.time_range = None;
if file_url {
req.location = url::Url::from_directory_path(&directory)
.unwrap()
.to_string();
}
if legacy {
let output = sql(
instance,
&format!(
"COPY DATABASE public TO '{}' WITH (FORMAT='parquet')",
req.location
),
)
.await;
assert!(matches!(output.data, OutputData::AffectedRows(rows) if rows == names.len()));
} else {
let summary = instance
.export_database_for_test(
req.clone(),
None,
&CancellationToken::new(),
QueryContext::arc(),
)
.await
.unwrap();
assert_eq!(summary.rows, names.len());
let expected = names
.iter()
.map(|name| {
let path = directory.join(format!("{name}.parquet"));
if file_url {
url::Url::from_file_path(path).unwrap().to_string()
} else {
path.to_str().unwrap().to_string()
}
})
.collect::<std::collections::BTreeSet<_>>();
let actual = summary
.output_files
.into_iter()
.collect::<std::collections::BTreeSet<_>>();
if file_url {
assert_eq!(actual, expected);
} else {
assert_eq!(
actual
.into_iter()
.map(PathBuf::from)
.collect::<std::collections::BTreeSet<_>>(),
expected
.into_iter()
.map(PathBuf::from)
.collect::<std::collections::BTreeSet<_>>()
);
}
}
for name in &names {
let path = directory.join(format!("{name}.parquet"));
assert!(path.is_file());
}
assert_eq!(std::fs::read_dir(&directory).unwrap().count(), names.len());
let output = sql(
instance,
&format!(
"COPY DATABASE names_restored FROM '{}' WITH (FORMAT='parquet')",
req.location
),
)
.await;
assert!(matches!(output.data, OutputData::AffectedRows(rows) if rows == names.len()));
for name in &names {
assert_eq!(
values(
instance,
&format!("SELECT ts, val, host FROM names_restored.\"{name}\"")
)
.await,
values(instance, &format!("SELECT ts, val, host FROM \"{name}\"")).await,
);
sql(
instance,
&format!("TRUNCATE TABLE names_restored.\"{name}\""),
)
.await;
}
}
}
#[tokio::test(flavor = "multi_thread")]
async fn packed_copy_standalone_heterogeneous_streams() {
use common_datasource::packed_snapshot::{ObjectKind, PackIndex, PackObject, PackTable};
let standalone = GreptimeDbStandaloneBuilder::new("packed_copy")
.build()
.await;
let instance = standalone.fe_instance();
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
let data_directory = destination.path().join("data/public/1");
std::fs::create_dir_all(&data_directory).unwrap();
let directory = data_directory.as_path();
let mut index = PackIndex {
version: 1,
objects: vec![],
tables: vec![],
};
let mut pack = Vec::new();
let mut expected = Vec::new();
let mut creates = Vec::new();
let mut physical_creates = Vec::new();
let mut physical_ids = Vec::new();
// Metric tag schemas vary; the standalone file has a different value type.
for (i, (name, typ, value)) in [
("literal.name", "DOUBLE", Some("42")),
("strings", "DOUBLE", Some("43,'tag'")),
("empty", "DOUBLE", None),
("standalone", "STRING", Some("'standalone'")),
]
.into_iter()
.enumerate()
{
let create = if i < 3 {
let physical = format!(
"CREATE TABLE IF NOT EXISTS physical_{i} (ts TIMESTAMP TIME INDEX, val {typ}) ENGINE=metric WITH(physical_metric_table='')"
);
sql(instance, &physical).await;
physical_creates.push(physical);
physical_ids.push(
table(instance, &format!("physical_{i}"))
.await
.table_info()
.ident
.table_id,
);
let tag = if i == 1 {
", host STRING PRIMARY KEY"
} else {
""
};
format!(
"CREATE TABLE \"{name}\" (ts TIMESTAMP TIME INDEX, val {typ}{tag}) ENGINE=metric WITH(on_physical_table='physical_{i}')"
)
} else {
format!("CREATE TABLE \"{name}\" (ts TIMESTAMP TIME INDEX, val {typ})")
};
sql(instance, &create).await;
creates.push(create.clone());
if let Some(value) = value {
let columns = if i == 1 { "ts,val,host" } else { "ts,val" };
sql(
instance,
&format!("INSERT INTO \"{name}\" ({columns}) VALUES (1,{value})"),
)
.await;
}
expected.push(values(instance, &format!("SELECT * FROM \"{name}\"")).await);
let path = directory.join(format!("table-{i}.parquet"));
sql(
instance,
&format!(
"COPY \"{name}\" TO '{}' WITH (FORMAT='parquet')",
path.display()
),
)
.await;
let bytes = std::fs::read(&path).unwrap();
let (object, offset) = if i == 3 {
index.objects.push(PackObject {
path: format!("table-{i}.parquet"),
kind: ObjectKind::Parquet,
length: bytes.len() as u64,
});
(format!("table-{i}.parquet"), 0)
} else {
let offset = pack.len() as u64;
pack.extend_from_slice(&bytes);
std::fs::remove_file(&path).unwrap();
("pack-0.bin".into(), offset)
};
index.tables.push(PackTable {
table_name: name.into(),
object,
offset,
length: bytes.len() as u64,
row_count: u64::from(value.is_some()),
});
sql(instance, &format!("DROP TABLE \"{name}\"")).await;
sql(instance, &create).await;
}
index.objects.push(PackObject {
path: "pack-0.bin".into(),
kind: ObjectKind::Pack,
length: pack.len() as u64,
});
std::fs::write(directory.join("pack-0.bin"), pack).unwrap();
std::fs::write(
directory.join("pack-index.json"),
serde_json::to_vec(&index).unwrap(),
)
.unwrap();
struct DenyPackedTable(std::sync::atomic::AtomicBool);
impl auth::PermissionChecker for DenyPackedTable {
fn check_permission(
&self,
_: auth::UserInfoRef,
_: auth::PermissionReq,
) -> auth::error::Result<auth::PermissionResp> {
Ok(auth::PermissionResp::Allow)
}
fn check_permission_with_table_targets(
&self,
_: auth::UserInfoRef,
req: auth::PermissionReq,
targets: auth::PermissionTableTargets,
) -> auth::error::Result<auth::PermissionResp> {
let denied = self.0.load(std::sync::atomic::Ordering::SeqCst)
&& !req.is_readonly()
&& matches!(targets, auth::PermissionTableTargets::Resolved(ref tables) if tables.iter().any(|t| t.table == "strings"));
Ok(if denied {
auth::PermissionResp::Reject
} else {
auth::PermissionResp::Allow
})
}
}
let checker = Arc::new(DenyPackedTable(std::sync::atomic::AtomicBool::new(true)));
instance
.plugins()
.insert::<auth::PermissionCheckerRef>(checker.clone());
let statement = format!(
"COPY DATABASE public FROM '{}/' WITH (FORMAT='parquet', metric_data_layout='packed', parallelism=2)",
directory.display()
);
let denied = servers::query_handler::sql::SqlQueryHandler::do_query(
instance.as_ref(),
&statement,
QueryContext::arc(),
)
.await;
assert!(denied.into_iter().all(|r| r.is_err()));
checker.0.store(false, std::sync::atomic::Ordering::SeqCst);
for entry in &index.tables {
assert!(
values(instance, &format!("SELECT * FROM \"{}\"", entry.table_name))
.await
.is_empty()
);
}
let output = sql(instance, &format!("COPY DATABASE public FROM '{}/' WITH (FORMAT='parquet', metric_data_layout='packed', parallelism=2)", directory.display())).await;
assert!(matches!(output.data, OutputData::AffectedRows(3)));
for (entry, expected) in index.tables.iter().zip(&expected) {
assert_eq!(
values(instance, &format!("SELECT * FROM \"{}\"", entry.table_name)).await,
*expected
);
}
if let Ok(bucket) = std::env::var("GT_S3_BUCKET")
&& !bucket.is_empty()
{
let prefix = format!(
"packed-reader/{}/",
destination.path().file_name().unwrap().to_str().unwrap()
);
let mut connection = std::collections::HashMap::from([
(
"access_key_id".into(),
std::env::var("GT_S3_ACCESS_KEY_ID").unwrap(),
),
(
"secret_access_key".into(),
std::env::var("GT_S3_ACCESS_KEY").unwrap(),
),
("region".into(), std::env::var("GT_S3_REGION").unwrap()),
]);
if let Ok(endpoint) = std::env::var("GT_S3_ENDPOINT_URL") {
connection.insert("endpoint".into(), endpoint);
}
let location = format!("s3://{bucket}/{prefix}");
let store = common_datasource::object_store::build_backend(
&location,
&connection,
&common_datasource::object_store::LocalFileAccess::Disabled,
)
.await
.unwrap();
for path in index
.objects
.iter()
.map(|o| o.path.as_str())
.chain(["pack-index.json"])
{
store
.write(path, std::fs::read(directory.join(path)).unwrap())
.await
.unwrap();
}
for (entry, create) in index.tables.iter().zip(&creates) {
sql(instance, &format!("DROP TABLE \"{}\"", entry.table_name)).await;
sql(instance, create).await;
}
let connection_sql = connection
.iter()
.map(|(k, v)| format!("{k}='{}'", v.replace('\'', "''")))
.collect::<Vec<_>>()
.join(",");
sql(instance, &format!("COPY DATABASE public FROM '{location}' WITH (FORMAT='parquet', metric_data_layout='packed', parallelism=2) CONNECTION ({connection_sql})")).await;
for (entry, expected) in index.tables.iter().zip(&expected) {
assert_eq!(
values(instance, &format!("SELECT * FROM \"{}\"", entry.table_name)).await,
*expected
);
}
for path in index
.objects
.iter()
.map(|o| o.path.as_str())
.chain(["pack-index.json"])
{
store.delete(path).await.unwrap();
}
}
// Run the real CLI importer through HTTP against the same generated fixture.
use clap::Parser;
use cli::export_v2::manifest::{ChunkMeta, DataFormat, Manifest, TimeRange};
let mut manifest = Manifest::new_full(
"greptime".into(),
vec!["public".into()],
TimeRange::unbounded(),
DataFormat::Parquet,
);
manifest.version = 2;
manifest.data_layout = Some("metric-parquet-packs".into());
let mut chunk = ChunkMeta::new(1, TimeRange::unbounded());
chunk.mark_completed(
index
.objects
.iter()
.map(|o| format!("data/public/1/{}", o.path))
.chain(["data/public/1/pack-index.json".into()])
.collect(),
None,
);
manifest.chunks.push(chunk);
let root = destination.path();
std::fs::create_dir_all(root.join("schema/ddl")).unwrap();
std::fs::write(
root.join("manifest.json"),
serde_json::to_vec(&manifest).unwrap(),
)
.unwrap();
std::fs::write(
root.join("schema/ddl/public.sql"),
physical_creates
.iter()
.chain(&creates)
.cloned()
.collect::<Vec<_>>()
.join(";\n")
+ ";",
)
.unwrap();
for entry in &index.tables {
sql(instance, &format!("DROP TABLE \"{}\"", entry.table_name)).await;
}
for i in 0..physical_ids.len() {
sql(instance, &format!("DROP TABLE physical_{i}")).await;
}
let server = servers::http::HttpServerBuilder::new(Default::default())
.with_sql_handler(instance.clone())
.with_user_provider(Arc::new(
auth::static_user_provider_from_option("static_user_provider:cmd:user=password")
.unwrap(),
))
.build();
let client =
servers::http::test_helpers::TestClient::new(server.build(server.make_app()).unwrap())
.await;
let state_path = root.join("import-state.json");
let snapshot_uri = url::Url::from_file_path(root).unwrap();
let command = cli::import_v2::ImportV2Command::parse_from([
"import-v2",
"--addr",
client.base_url().trim_start_matches("http://"),
"--from",
snapshot_uri.as_str(),
"--state-path",
state_path.to_str().unwrap(),
"--auth-basic",
"user:password",
"--no-proxy",
"--progress",
"never",
]);
let importer = command.build().await.unwrap();
std::fs::write(
root.join("schema/schemas.json"),
serde_json::to_vec(
&serde_json::json!([{"catalog":"greptime","name":"public","options":{}}]),
)
.unwrap(),
)
.unwrap();
#[derive(clap::Parser)]
struct VerifyArgs {
#[command(subcommand)]
command: cli::export_v2::ExportV2Command,
}
let verifier =
VerifyArgs::parse_from(["export-v2", "verify", "--snapshot", snapshot_uri.as_str()])
.command
.build()
.await
.unwrap();
verifier.do_work().await.unwrap();
let original = serde_json::to_value(&manifest).unwrap();
let initial_tables = values(instance, "SHOW TABLES").await;
let pack_path = directory.join("pack-0.bin");
let pack_bytes = std::fs::read(&pack_path).unwrap();
for case in [
"schema_only",
"missing_chunks",
"missing_files",
"manifest_checksum",
"chunk_checksum",
"duplicate",
"extra",
"missing_object",
"truncated_object",
] {
let mut invalid = original.clone();
let expected_error = match case {
"schema_only" => {
invalid["schema_only"] = true.into();
"schema_only"
}
"missing_chunks" => {
invalid.as_object_mut().unwrap().remove("chunks");
"chunks array"
}
"missing_files" => {
invalid["chunks"][0]
.as_object_mut()
.unwrap()
.remove("files");
"files array"
}
"manifest_checksum" => {
invalid["checksum"] = "unsupported".into();
"checksums"
}
"chunk_checksum" => {
invalid["chunks"][0]["checksum"] = "unsupported".into();
"checksums"
}
"duplicate" => {
let file = invalid["chunks"][0]["files"][0].clone();
invalid["chunks"][0]["files"]
.as_array_mut()
.unwrap()
.push(file);
"duplicate"
}
"extra" => {
std::fs::write(directory.join("extra.bin"), b"extra").unwrap();
invalid["chunks"][0]["files"]
.as_array_mut()
.unwrap()
.push("data/public/1/extra.bin".into());
"inventory"
}
"missing_object" => {
std::fs::remove_file(&pack_path).unwrap();
"expected length"
}
"truncated_object" => {
std::fs::write(&pack_path, &pack_bytes[..pack_bytes.len() - 1]).unwrap();
"expected length"
}
_ => unreachable!(),
};
std::fs::write(
root.join("manifest.json"),
serde_json::to_vec(&invalid).unwrap(),
)
.unwrap();
for resumed in [false, true] {
if resumed {
std::fs::write(&state_path, b"state must not be read or changed").unwrap();
}
let error = importer.do_work().await.unwrap_err();
assert!(
format!("{error:?}").contains(expected_error),
"{case}: {error:?}"
);
assert_eq!(
values(instance, "SHOW TABLES").await,
initial_tables,
"{case}"
);
if resumed {
assert_eq!(
std::fs::read(&state_path).unwrap(),
b"state must not be read or changed"
);
std::fs::remove_file(&state_path).unwrap();
} else {
assert!(!state_path.exists(), "{case}");
}
}
assert!(verifier.do_work().await.is_err(), "{case}");
std::fs::write(&pack_path, &pack_bytes).unwrap();
if case == "extra" {
std::fs::remove_file(directory.join("extra.bin")).unwrap();
}
}
std::fs::write(
root.join("manifest.json"),
serde_json::to_vec(&manifest).unwrap(),
)
.unwrap();
verifier.do_work().await.unwrap();
importer.do_work().await.unwrap();
assert!(!state_path.exists());
for (i, old_id) in physical_ids.iter().enumerate() {
assert_ne!(
*old_id,
table(instance, &format!("physical_{i}"))
.await
.table_info()
.ident
.table_id
);
}
for (entry, expected) in index.tables.iter().zip(&expected) {
assert_eq!(
values(instance, &format!("SELECT * FROM \"{}\"", entry.table_name)).await,
*expected
);
}
for entry in &index.tables {
sql(instance, &format!("DROP TABLE \"{}\"", entry.table_name)).await;
}
index.tables.last_mut().unwrap().row_count += 1;
std::fs::write(
directory.join("pack-index.json"),
serde_json::to_vec(&index).unwrap(),
)
.unwrap();
let failure = importer.do_work().await.unwrap_err();
assert!(
format!("{failure:?}").contains("Chunk 1 import failed"),
"{failure:?}"
);
let failed: serde_json::Value =
serde_json::from_slice(&std::fs::read(&state_path).unwrap()).unwrap();
assert_eq!(failed["ddl_completed"], true);
assert_eq!(failed["tasks"][0]["status"], "failed");
let saved_state = std::fs::read(&state_path).unwrap();
std::fs::write(&pack_path, &pack_bytes[..pack_bytes.len() - 1]).unwrap();
let preflight_error = importer.do_work().await.unwrap_err();
assert!(
format!("{preflight_error:?}").contains("expected length"),
"{preflight_error:?}"
);
assert_eq!(std::fs::read(&state_path).unwrap(), saved_state);
std::fs::write(&pack_path, &pack_bytes).unwrap();
let failed_copy = servers::query_handler::sql::SqlQueryHandler::do_query(
instance.as_ref(),
&statement,
QueryContext::arc(),
)
.await
.remove(0)
.err()
.unwrap();
assert!(
format!("{failed_copy:?}").contains("row_count"),
"{failed_copy:?}"
);
// A failed chunk may have inserted rows; clear them before explicit replay.
for (entry, create) in index.tables.iter().zip(&creates) {
sql(instance, &format!("DROP TABLE \"{}\"", entry.table_name)).await;
sql(instance, create).await;
}
index.tables.last_mut().unwrap().row_count -= 1;
std::fs::write(
directory.join("pack-index.json"),
serde_json::to_vec(&index).unwrap(),
)
.unwrap();
importer.do_work().await.unwrap();
assert!(!state_path.exists());
for (entry, expected) in index.tables.iter().zip(&expected) {
assert_eq!(
values(instance, &format!("SELECT * FROM \"{}\"", entry.table_name)).await,
*expected
);
}
// The same importer still restores the original per-table manifest layout.
let packed_bytes = std::fs::read(directory.join("pack-0.bin")).unwrap();
let mut files = Vec::new();
for entry in &index.tables {
let name = format!("{}.parquet", entry.table_name);
if entry.object == "pack-0.bin" {
std::fs::write(
directory.join(&name),
&packed_bytes[entry.offset as usize..(entry.offset + entry.length) as usize],
)
.unwrap();
} else {
std::fs::rename(directory.join(&entry.object), directory.join(&name)).unwrap();
}
files.push(format!("data/public/1/{name}"));
sql(instance, &format!("DROP TABLE \"{}\"", entry.table_name)).await;
}
manifest.version = 1;
manifest.data_layout = None;
manifest.chunks[0].files = files;
std::fs::write(
root.join("manifest.json"),
serde_json::to_vec(&manifest).unwrap(),
)
.unwrap();
importer.do_work().await.unwrap();
for (entry, expected) in index.tables.iter().zip(&expected) {
assert_eq!(
values(instance, &format!("SELECT * FROM \"{}\"", entry.table_name)).await,
*expected
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn packed_copy_large_column_in_middle_stream() {
use common_datasource::packed_snapshot::{ObjectKind, PackIndex, PackObject, PackTable};
use datafusion::parquet::arrow::ArrowWriter;
use datafusion::parquet::basic::Compression;
use datafusion::parquet::file::properties::WriterProperties;
use datatypes::arrow::array::{
ArrayRef, BooleanArray, Int64Array, StringArray, TimestampMillisecondArray,
};
use datatypes::arrow::datatypes::{Field, Schema};
use datatypes::arrow::record_batch::RecordBatch;
let standalone = GreptimeDbStandaloneBuilder::new("packed_large_column")
.build()
.await;
let instance = standalone.fe_instance();
let directory = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
let large = "x".repeat(20 * 1024 * 1024);
let arrays: [ArrayRef; 3] = [
Arc::new(Int64Array::from(vec![42])),
Arc::new(StringArray::from(vec![large.as_str()])),
Arc::new(BooleanArray::from(vec![true])),
];
let mut packed = Vec::new();
let mut index = PackIndex {
version: 1,
objects: vec![],
tables: vec![],
};
for (i, (array, typ)) in arrays
.into_iter()
.zip(["BIGINT", "STRING", "BOOLEAN"])
.enumerate()
{
sql(
instance,
&format!("CREATE TABLE restored_{i} (ts TIMESTAMP TIME INDEX, val {typ})"),
)
.await;
let timestamps: ArrayRef = Arc::new(TimestampMillisecondArray::from(vec![1]));
let schema = Arc::new(Schema::new(vec![
Field::new("ts", timestamps.data_type().clone(), false),
Field::new("val", array.data_type().clone(), false),
]));
let batch = RecordBatch::try_new(schema.clone(), vec![timestamps, array]).unwrap();
let properties = WriterProperties::builder()
.set_compression(Compression::UNCOMPRESSED)
.set_dictionary_enabled(false)
.build();
let mut bytes = Vec::new();
let mut writer = ArrowWriter::try_new(&mut bytes, schema, Some(properties)).unwrap();
writer.write(&batch).unwrap();
let metadata = writer.close().unwrap();
if i == 1 {
assert!(metadata.row_group(0).column(1).compressed_size() > 16 * 1024 * 1024);
assert!(!packed.is_empty());
}
index.tables.push(PackTable {
table_name: format!("restored_{i}"),
object: "pack-0.bin".into(),
offset: packed.len() as u64,
length: bytes.len() as u64,
row_count: 1,
});
packed.extend(bytes);
}
index.objects.push(PackObject {
path: "pack-0.bin".into(),
kind: ObjectKind::Pack,
length: packed.len() as u64,
});
std::fs::write(directory.path().join("pack-0.bin"), packed).unwrap();
std::fs::write(
directory.path().join("pack-index.json"),
serde_json::to_vec(&index).unwrap(),
)
.unwrap();
let output = sql(instance, &format!("COPY DATABASE public FROM '{}/' WITH (FORMAT='parquet', metric_data_layout='packed', parallelism=2)", directory.path().display())).await;
assert!(matches!(output.data, OutputData::AffectedRows(3)));
for (i, expected) in [
datatypes::value::Value::Int64(42),
datatypes::value::Value::String(large.into()),
datatypes::value::Value::Boolean(true),
]
.into_iter()
.enumerate()
{
assert_eq!(
values(instance, &format!("SELECT val FROM restored_{i}")).await,
vec![vec![expected]]
);
}
}
#[derive(clap::Parser)]
struct ExportDataCli {
#[command(subcommand)]
command: cli::DataCommand,
}
async fn run_data_cli(args: &[&str]) -> Result<(), common_error::ext::BoxedError> {
use clap::Parser;
ExportDataCli::try_parse_from(std::iter::once("greptime-data").chain(args.iter().copied()))
.unwrap()
.command
.build()
.await?
.do_work()
.await
}
async fn export_http(instance: Arc<Instance>) -> (String, tokio::task::JoinHandle<()>) {
let server = servers::http::HttpServerBuilder::new(servers::http::HttpOptions::default())
.with_sql_handler(instance)
.build();
let app = server.build(server.make_app()).unwrap();
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap().to_string();
let task = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
(addr, task)
}
struct FailSecondChunk {
fail: std::sync::atomic::AtomicBool,
copies: std::sync::atomic::AtomicUsize,
}
impl servers::interceptor::SqlQueryInterceptor for FailSecondChunk {
type Error = frontend::error::Error;
fn pre_execute(
&self,
statement: Option<&sql::statements::statement::Statement>,
_plan: Option<&datafusion_expr::LogicalPlan>,
_ctx: session::context::QueryContextRef,
) -> Result<(), Self::Error> {
use sql::statements::copy::{Copy, CopyDatabase};
use sql::statements::statement::Statement;
if let Some(Statement::Copy(Copy::CopyDatabase(CopyDatabase::To(arg)))) = statement {
self.copies
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if arg.location.ends_with("/z_later/2/")
&& self.fail.swap(false, std::sync::atomic::Ordering::SeqCst)
{
return Err(operator::error::InvalidDatabaseExportSnafu {
reason: "injected terminated COPY failure",
}
.build()
.into());
}
}
Ok(())
}
}
#[tokio::test]
async fn metric_export_v2_cli_resume_roundtrip() {
metric_export_v2_cli_roundtrip(false, false).await;
}
#[tokio::test]
async fn metric_export_v2_cli_s3_resume_roundtrip() {
// The CLI takes an explicit `--s3-endpoint`, so this only runs against an
// S3-compatible endpoint (MinIO in CI), not against bucket+region alone.
if std::env::var("GT_S3_ENDPOINT_URL").is_ok_and(|e| !e.is_empty())
&& std::env::var("GT_S3_BUCKET").is_ok_and(|b| !b.is_empty())
{
metric_export_v2_cli_roundtrip(true, false).await;
}
}
#[tokio::test(flavor = "multi_thread")]
async fn packed_export_v2_cli_local_roundtrip() {
metric_export_v2_cli_roundtrip(false, true).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn packed_export_v2_cli_s3_roundtrip() {
if std::env::var("GT_S3_ENDPOINT_URL").is_ok_and(|e| !e.is_empty())
&& std::env::var("GT_S3_BUCKET").is_ok_and(|b| !b.is_empty())
{
metric_export_v2_cli_roundtrip(true, true).await;
}
}
async fn exported_table_path(
store: &object_store::ObjectStore,
packed: bool,
chunk: u32,
name: &str,
) -> String {
let prefix = format!("data/public/{chunk}/");
if !packed {
return format!("{prefix}{name}.parquet");
}
let index: common_datasource::packed_snapshot::PackIndex = serde_json::from_slice(
&store
.read(&format!("{prefix}pack-index.json"))
.await
.unwrap()
.to_bytes(),
)
.unwrap();
format!(
"{prefix}{}",
index
.tables
.iter()
.find(|t| t.table_name == name)
.unwrap()
.object
)
}
async fn metric_export_v2_cli_roundtrip(s3: bool, packed: bool) {
use servers::interceptor::SqlQueryInterceptorRef;
let plugins = common_base::Plugins::new();
let faults = Arc::new(FailSecondChunk {
fail: std::sync::atomic::AtomicBool::new(true),
copies: std::sync::atomic::AtomicUsize::new(0),
});
plugins.insert::<SqlQueryInterceptorRef<frontend::error::Error>>(faults.clone());
let source = GreptimeDbStandaloneBuilder::new("metric_v2_source")
.with_experimental_metric_export()
.with_plugin(plugins)
.build()
.await;
let instance = source.fe_instance();
let mut names = Vec::new();
for physical in ["v2_a", "v2_b"] {
let (logical, _, renamed) =
create_metric_export_source_tables(instance, physical, "dense").await;
sql(
instance,
&format!("ALTER TABLE {renamed} RENAME {physical}"),
)
.await;
names.extend(logical);
names.push(format!("{physical}_excluded"));
}
sql(
instance,
"CREATE TABLE audit (host STRING PRIMARY KEY, val DOUBLE, ts TIMESTAMP TIME INDEX)",
)
.await;
sql(instance, "INSERT INTO audit VALUES ('a',1,1),('z',NULL,3)").await;
sql(instance, "CREATE VIEW dashboard AS SELECT * FROM audit").await;
let special_name = "audit.dashboard";
let quoted_special = special_name.replace('"', "\"\"");
sql(
instance,
&format!("CREATE VIEW \"{quoted_special}\" AS SELECT * FROM audit"),
)
.await;
if packed && !s3 {
for id in 0..128 {
let name = format!("batch_{id:03}");
sql(instance, &format!("CREATE TABLE {name} (host STRING PRIMARY KEY, val DOUBLE, ts TIMESTAMP TIME INDEX) ENGINE=metric WITH (on_physical_table='v2_a')")).await;
names.push(name);
}
}
sql(instance, "CREATE DATABASE z_later").await;
if packed {
sql(instance, "CREATE DATABASE empty_schema").await;
}
sql(
instance,
"CREATE TABLE z_later.events (ts TIMESTAMP TIME INDEX, val BIGINT)",
)
.await;
sql(instance, "INSERT INTO z_later.events VALUES (1,10),(3,20)").await;
names.push("audit".into());
if s3 {
sql(
instance,
"CREATE TABLE bulk (host STRING PRIMARY KEY, payload STRING, ts TIMESTAMP TIME INDEX)",
)
.await;
for batch in 0..64 {
let rows = (0..64)
.map(|row| {
let payload = (0..128)
.map(|_| uuid::Uuid::new_v4().simple().to_string())
.collect::<String>();
format!("('{}','{payload}',1)", batch * 64 + row)
})
.collect::<Vec<_>>()
.join(",");
sql(instance, &format!("INSERT INTO bulk VALUES {rows}")).await;
}
names.push("bulk".into());
}
let (addr, server) = export_http(instance.clone()).await;
if packed && !s3 {
let client = cli::DatabaseClient::new(
addr.clone(),
"greptime".into(),
None,
std::time::Duration::from_secs(60),
None,
true,
);
let error = client.sql_in_public(
"SHOW CREATE TABLE audit; SHOW CREATE TABLE missing_middle; SHOW CREATE VIEW dashboard"
).await.unwrap_err();
assert!(error.to_string().contains("SQL request failed"), "{error}");
}
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
for experimental in [false, true] {
for layout in ["packed", "invalid"] {
if experimental && layout == "packed" {
continue;
}
let path = destination
.path()
.join(format!("unsupported-{experimental}-{layout}"));
let uri = url::Url::from_directory_path(&path).unwrap();
let statement = format!(
"COPY DATABASE public TO '{uri}' WITH (FORMAT='parquet', experimental_metric_export='{experimental}', metric_data_layout='{layout}')"
);
let result = servers::query_handler::sql::SqlQueryHandler::do_query(
instance.as_ref(),
&statement,
QueryContext::arc(),
)
.await
.remove(0);
let error = result.expect_err("export must reject unsupported layouts");
assert!(
format!("{error:?}").contains("metric_data_layout"),
"{error:?}"
);
assert!(!path.exists());
}
}
if packed && !s3 {
for window in [None, Some("3ms")] {
let path = destination.path().join("empty-range");
let uri = url::Url::from_directory_path(&path).unwrap();
let mut args = vec![
"export-v2",
"create",
"--addr",
&addr,
"--to",
uri.as_str(),
"--schemas",
"public",
"--experimental-metric-export",
"--metric-data-layout",
"packed",
"--no-proxy",
"--start-time",
"1970-01-01T00:00:00Z",
"--end-time",
"1970-01-01T00:00:00Z",
];
if let Some(window) = window {
args.extend(["--chunk-time-window", window]);
}
let error = run_data_cli(&args).await.unwrap_err();
assert!(
error
.to_string()
.contains("Packed export requires --start-time to be earlier than --end-time"),
"{error}"
);
assert!(!path.exists());
}
}
let (uri, store, storage_args) = if s3 {
let endpoint = std::env::var("GT_S3_ENDPOINT_URL").unwrap();
let bucket = std::env::var("GT_S3_BUCKET").unwrap();
let region = std::env::var("GT_S3_REGION").unwrap();
let key = std::env::var("GT_S3_ACCESS_KEY_ID").unwrap();
let secret = std::env::var("GT_S3_ACCESS_KEY").unwrap();
let root = format!("pr04b-{}", uuid::Uuid::new_v4());
let store = object_store::ObjectStore::new(
object_store::services::S3::default()
.endpoint(&endpoint)
.bucket(&bucket)
.root(&root)
.region(&region)
.access_key_id(&key)
.secret_access_key(&secret),
)
.unwrap();
(
format!("s3://{bucket}/{root}"),
store,
vec![
"--s3".into(),
"--s3-endpoint".into(),
endpoint,
"--s3-region".into(),
region,
"--s3-access-key-id".into(),
key,
"--s3-secret-access-key".into(),
secret,
],
)
} else {
(
url::Url::from_file_path(destination.path())
.unwrap()
.to_string(),
object_store::ObjectStore::new(
object_store::services::Fs::default().root(destination.path().to_str().unwrap()),
)
.unwrap(),
Vec::<String>::new(),
)
};
let mut args = vec![
"export-v2",
"create",
"--addr",
&addr,
"--to",
&uri,
"--schemas",
if packed {
"public,z_later,empty_schema"
} else {
"public,z_later"
},
"--experimental-metric-export",
"--no-proxy",
"--start-time",
"1970-01-01T00:00:00Z",
"--end-time",
"1970-01-01T00:00:00.006Z",
"--chunk-time-window",
"3ms",
"--progress",
"never",
];
args.extend(storage_args.iter().map(String::as_str));
if packed {
args.extend(["--metric-data-layout", "packed"]);
}
assert!(run_data_cli(&args).await.is_err());
let before: cli::export_v2::manifest::Manifest =
serde_json::from_slice(&store.read("manifest.json").await.unwrap().to_vec()).unwrap();
assert_eq!(
before.chunks[0].status,
cli::export_v2::manifest::ChunkStatus::Completed
);
assert_eq!(
before.chunks[1].status,
cli::export_v2::manifest::ChunkStatus::Failed
);
if s3 {
assert!(
store
.stat(&exported_table_path(&store, packed, 1, "bulk").await)
.await
.unwrap()
.content_length()
> 5 * 1024 * 1024
);
}
let mut preserved = Vec::new();
for path in &before.chunks[0].files {
preserved.push((path.clone(), store.read(path).await.unwrap().to_vec()));
}
assert!(
store
.exists(&exported_table_path(&store, packed, 2, "audit").await)
.await
.unwrap()
);
store
.write("data/public/2/unknown.txt", "keep")
.await
.unwrap();
let copies = faults.copies.load(std::sync::atomic::Ordering::SeqCst);
assert!(run_data_cli(&args).await.is_err());
assert_eq!(
faults.copies.load(std::sync::atomic::Ordering::SeqCst),
copies
);
assert_eq!(
store
.read("data/public/2/unknown.txt")
.await
.unwrap()
.to_vec(),
b"keep"
);
store.delete("data/public/2/unknown.txt").await.unwrap();
run_data_cli(&args).await.unwrap();
let after: cli::export_v2::manifest::Manifest =
serde_json::from_slice(&store.read("manifest.json").await.unwrap().to_vec()).unwrap();
assert!(after.is_complete());
assert_eq!(after.version, if packed { 2 } else { 1 });
assert_eq!(before.data_layout, after.data_layout);
if packed {
use common_datasource::packed_snapshot::{ObjectKind, PackIndex};
for chunk in &after.chunks {
let empty: PackIndex = serde_json::from_slice(
&store
.read(&format!("data/empty_schema/{}/pack-index.json", chunk.id))
.await
.unwrap()
.to_bytes(),
)
.unwrap();
empty.validate_membership([]).unwrap();
let prefix = format!("data/public/{}/", chunk.id);
let index: PackIndex = serde_json::from_slice(
&store
.read(&format!("{prefix}pack-index.json"))
.await
.unwrap()
.to_bytes(),
)
.unwrap();
index
.validate_membership(names.iter().map(String::as_str))
.unwrap();
assert_eq!(
index
.objects
.iter()
.filter(|o| o.kind == ObjectKind::Pack)
.count(),
1
);
assert!(index.tables.iter().filter(|t| t.row_count == 0).count() >= 2);
let mut expected = index
.objects
.iter()
.map(|o| format!("{prefix}{}", o.path))
.collect::<Vec<_>>();
expected.push(format!("{prefix}pack-index.json"));
expected.sort();
assert_eq!(
chunk
.files
.iter()
.filter(|f| f.starts_with(&prefix))
.cloned()
.collect::<Vec<_>>(),
expected
);
}
}
assert_eq!(
serde_json::to_value(&before.chunks[0]).unwrap(),
serde_json::to_value(&after.chunks[0]).unwrap()
);
for (path, bytes) in preserved {
assert_eq!(bytes, store.read(&path).await.unwrap().to_vec());
}
let mut verify = vec!["export-v2", "verify", "--snapshot", &uri];
verify.extend(storage_args.iter().map(String::as_str));
run_data_cli(&verify).await.unwrap();
let target = GreptimeDbStandaloneBuilder::new("metric_v2_target")
.build()
.await;
for id in 0..24 {
sql(
target.fe_instance(),
&format!("CREATE TABLE occupied_{id} (ts TIMESTAMP TIME INDEX)"),
)
.await;
}
let (target_addr, target_server) = export_http(target.fe_instance().clone()).await;
let state = destination.path().join("restore-state.json");
let mut import = vec![
"import-v2",
"--addr",
&target_addr,
"--from",
&uri,
"--no-proxy",
"--state-path",
state.to_str().unwrap(),
"--progress",
"never",
];
import.extend(storage_args.iter().map(String::as_str));
run_data_cli(&import).await.unwrap();
for name in names {
let source_table = table(instance, &name).await;
let target_table = table(target.fe_instance(), &name).await;
assert_ne!(
source_table.table_info().table_id(),
target_table.table_info().table_id()
);
assert_eq!(
source_table.schema().column_schemas(),
target_table.schema().column_schemas()
);
let query = format!("SELECT * FROM \"{name}\" ORDER BY ts,host");
assert_eq!(
values(instance, &query).await,
values(target.fe_instance(), &query).await
);
}
assert_eq!(
values(instance, "SELECT * FROM z_later.events ORDER BY ts").await,
values(
target.fe_instance(),
"SELECT * FROM z_later.events ORDER BY ts"
)
.await
);
for view in ["dashboard", special_name] {
let query = format!(
"SELECT * FROM \"{}\" ORDER BY ts,host",
view.replace('"', "\"\"")
);
assert_eq!(
values(instance, &query).await,
values(target.fe_instance(), &query).await
);
assert_eq!(
table(instance, view).await.schema().column_schemas(),
table(target.fe_instance(), view)
.await
.schema()
.column_schemas()
);
}
if packed {
let schema_uri = format!("{uri}/schema-only");
let mut schema_args = vec![
"export-v2",
"create",
"--addr",
&addr,
"--to",
&schema_uri,
"--schemas",
"public",
"--schema-only",
"--experimental-metric-export",
"--metric-data-layout",
"packed",
"--no-proxy",
"--progress",
"never",
];
schema_args.extend(storage_args.iter().map(String::as_str));
run_data_cli(&schema_args).await.unwrap();
let schema_manifest: cli::export_v2::manifest::Manifest = serde_json::from_slice(
&store
.read("schema-only/manifest.json")
.await
.unwrap()
.to_bytes(),
)
.unwrap();
assert!(schema_manifest.schema_only && schema_manifest.chunks.is_empty());
assert_eq!(schema_manifest.version, 1);
assert!(schema_manifest.data_layout.is_none());
}
server.abort();
target_server.abort();
if s3 {
store.delete_with("/").recursive(true).await.unwrap();
}
}
#[tokio::test]
async fn metric_export_v2_disabled_and_legacy_cli() {
let source = GreptimeDbStandaloneBuilder::new("metric_v2_disabled")
.build()
.await;
let (addr, server) = export_http(source.fe_instance().clone()).await;
let destination = tempfile::tempdir_in(common_test_util::find_workspace_path(".")).unwrap();
let uri = url::Url::from_file_path(destination.path())
.unwrap()
.to_string();
let manifest = destination.path().join("manifest.json");
std::fs::write(&manifest, b"preserve-before-capability-check").unwrap();
let args = [
"export-v2",
"create",
"--addr",
&addr,
"--to",
&uri,
"--force",
"--experimental-metric-export",
"--no-proxy",
"--progress",
"never",
];
assert!(run_data_cli(&args).await.is_err());
assert_eq!(
std::fs::read(&manifest).unwrap(),
b"preserve-before-capability-check"
);
std::fs::remove_file(&manifest).unwrap();
sql(
source.fe_instance(),
"CREATE TABLE ordinary (ts TIMESTAMP TIME INDEX, val BIGINT)",
)
.await;
sql(source.fe_instance(), "INSERT INTO ordinary VALUES (1,2)").await;
run_data_cli(&[
"export-v2",
"create",
"--addr",
&addr,
"--to",
&uri,
"--schemas",
"public",
"--no-proxy",
"--progress",
"never",
])
.await
.unwrap();
run_data_cli(&["export-v2", "verify", "--snapshot", &uri])
.await
.unwrap();
server.abort();
}
#[tokio::test]
async fn metric_export_v2_refuses_missing_or_malformed_capability_before_force() {
use axum::http::StatusCode;
use serde_json::json;
let column = json!({"name": "EXPERIMENTAL_METRIC_EXPORT", "data_type": "String"});
let records = |rows, columns| {
json!({"records": {
"schema": {"column_schemas": columns},
"rows": rows
}})
};
let valid = records(json!([["true"]]), json!([column.clone()]));
let mut responses = [
(StatusCode::BAD_REQUEST, json!([])),
(StatusCode::OK, json!([])),
(StatusCode::OK, json!([[true]])),
(StatusCode::OK, json!([["true", "extra"]])),
(StatusCode::OK, json!([["true"], ["true"]])),
]
.map(|(status, rows)| (status, json!([records(rows, json!([column.clone()]))])))
.to_vec();
responses.extend([
(
StatusCode::OK,
json!([
valid.clone(),
records(json!([["false"]]), json!([column.clone()]))
]),
),
(
StatusCode::OK,
json!([records(json!([["true"]]), json!([]))]),
),
(
StatusCode::OK,
json!([records(json!([["true"]]), json!([column.clone(), column]))]),
),
(
StatusCode::OK,
json!([records(
json!([["true"]]),
json!([{"name": "EXPERIMENTAL_METRIC_EXPORT", "data_type": "Boolean"}])
)]),
),
(
StatusCode::OK,
json!([records(
json!([["true"]]),
json!([{"name": "OTHER", "data_type": "String"}])
)]),
),
]);
for (status, output) in responses {
let app =
axum::Router::new().route(
"/v1/sql",
axum::routing::post(
move |axum::Form(form): axum::Form<
std::collections::HashMap<String, String>,
>| async move {
assert_eq!(form["sql"], "SHOW VARIABLES experimental_metric_export");
(
status,
axum::Json(json!({
"execution_time_ms": 0,
"output": output
})),
)
},
),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap().to_string();
let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
let destination = tempfile::tempdir().unwrap();
let snapshot = destination.path().join("snapshot");
let uri = url::Url::from_file_path(&snapshot).unwrap().to_string();
let args = [
"export-v2",
"create",
"--addr",
&addr,
"--to",
&uri,
"--experimental-metric-export",
"--force",
"--no-proxy",
"--progress",
"never",
];
let error = run_data_cli(&args).await.unwrap_err();
if status == StatusCode::OK {
assert!(
error.to_string().contains("Metric export requires"),
"{error}"
);
}
assert!(!snapshot.exists());
std::fs::create_dir(&snapshot).unwrap();
let manifest = snapshot.join("manifest.json");
std::fs::write(&manifest, b"preserve").unwrap();
assert!(run_data_cli(&args).await.is_err());
assert_eq!(std::fs::read(manifest).unwrap(), b"preserve");
server.abort();
}
}