mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-29 08:27:06 +00:00
test(table-catalog): generate DuckDB REST attach SQL (#6749)
This commit is contained in:
@@ -17,7 +17,7 @@ DEFAULT_SPARK_VERSION = "3.5.4"
|
||||
DEFAULT_ICEBERG_VERSION = "1.7.1"
|
||||
DEFAULT_SCALA_VERSION = "2.12"
|
||||
DEFAULT_TRINO_VERSION = "477"
|
||||
DEFAULT_DUCKDB_VERSION = "1.3.2"
|
||||
DEFAULT_DUCKDB_VERSION = "1.5.5"
|
||||
DEFAULT_SNOWFLAKE_CLIENT_VERSION = "operator-recorded"
|
||||
DEFAULT_DATABEND_VERSION = "operator-recorded"
|
||||
DEFAULT_TRINO_SERVER = "http://127.0.0.1:8080"
|
||||
@@ -155,12 +155,13 @@ def engine_compatibility_matrix() -> list[dict[str, Any]]:
|
||||
},
|
||||
{
|
||||
"client": "DuckDB Iceberg",
|
||||
"status": "manual-live-read-probe",
|
||||
"status": "manual-live-harness",
|
||||
"entrypoint": "scripts/table-catalog/engine_compatibility.py --print-live-conformance",
|
||||
"scenarios": [
|
||||
scenario("metadata-read", "manual-live-probe", "read a supplied Iceberg metadata location through DuckDB iceberg_scan"),
|
||||
scenario("read-table", "manual-live-probe", "read-path verification only"),
|
||||
scenario("write-table", "not-claimed", "DuckDB write/commit compatibility is not claimed"),
|
||||
scenario("catalog-attach", "generated-harness", "attach RustFS as a generic signed Iceberg REST catalog"),
|
||||
scenario("read-table", "manual-live-probe", "read an existing table through the attached catalog"),
|
||||
scenario("write-table", "manual-validation-required", "generated SQL does not promote DuckDB write compatibility"),
|
||||
],
|
||||
},
|
||||
{
|
||||
@@ -489,8 +490,66 @@ def duckdb_sql_probe(
|
||||
return "\n".join(statements) + "\n"
|
||||
|
||||
|
||||
def duckdb_command() -> str:
|
||||
return shell_join(["duckdb", "-c", ".read /tmp/rustfs-s3tables-duckdb-read.sql"])
|
||||
def duckdb_rest_catalog_sql(
|
||||
*,
|
||||
endpoint: str,
|
||||
warehouse: str,
|
||||
access_key: str,
|
||||
secret_key: str,
|
||||
region: str,
|
||||
catalog_name: str,
|
||||
namespace: str,
|
||||
table: str,
|
||||
rest_path: str,
|
||||
rest_signing_name: str,
|
||||
) -> str:
|
||||
parsed = re.match(r"^(https?)://(.+)$", normalized_endpoint(endpoint))
|
||||
if not parsed:
|
||||
raise ValueError("DuckDB REST catalog endpoint must include http:// or https://")
|
||||
scheme, endpoint_without_scheme = parsed.groups()
|
||||
rest_path = normalized_rest_path(rest_path)
|
||||
catalog_identifier = quote_double_identifier(catalog_name)
|
||||
table_identifier = ".".join(
|
||||
[catalog_identifier, quote_double_identifier(namespace), quote_double_identifier(table)]
|
||||
)
|
||||
statements = [
|
||||
"INSTALL httpfs;",
|
||||
"LOAD httpfs;",
|
||||
"INSTALL iceberg;",
|
||||
"LOAD iceberg;",
|
||||
"CREATE OR REPLACE SECRET rustfs_s3 (",
|
||||
" TYPE s3,",
|
||||
" PROVIDER config,",
|
||||
f" KEY_ID {sql_string(access_key)},",
|
||||
f" SECRET {sql_string(secret_key)},",
|
||||
f" REGION {sql_string(region)},",
|
||||
f" ENDPOINT {sql_string(endpoint_without_scheme)},",
|
||||
" URL_STYLE 'path',",
|
||||
f" USE_SSL {'true' if scheme == 'https' else 'false'},",
|
||||
f" SCOPE {sql_string(f's3://{warehouse}')}",
|
||||
");",
|
||||
f"ATTACH {sql_string(warehouse)} AS {catalog_identifier} (",
|
||||
" TYPE iceberg,",
|
||||
f" ENDPOINT {sql_string(f'{normalized_endpoint(endpoint)}{rest_path}')},",
|
||||
" AUTHORIZATION_TYPE 'sigv4',",
|
||||
" SECRET 'rustfs_s3',",
|
||||
f" SIGV4_REGION {sql_string(region)},",
|
||||
f" SIGV4_SERVICE {sql_string(rest_signing_name)},",
|
||||
" ACCESS_DELEGATION_MODE 'none',",
|
||||
" STAGE_CREATE_TABLES false,",
|
||||
" SKIP_CREATE_TABLE_METADATA_UPDATES true,",
|
||||
" DISABLE_MULTI_TABLE_COMMIT true,",
|
||||
" REMOVE_FILES_ON_DELETE false,",
|
||||
" PURGE_REQUESTED false,",
|
||||
" SUPPORT_NESTED_NAMESPACES false",
|
||||
");",
|
||||
f"SELECT COUNT(*) AS row_count FROM {table_identifier};",
|
||||
]
|
||||
return "\n".join(statements) + "\n"
|
||||
|
||||
|
||||
def duckdb_command(*, sql_file: str = "/tmp/rustfs-s3tables-duckdb-read.sql") -> str:
|
||||
return shell_join(["duckdb", "-c", f".read {sql_file}"])
|
||||
|
||||
|
||||
def snowflake_sql_template(*, endpoint: str, warehouse: str, rest_path: str, namespace: str, table: str) -> str:
|
||||
@@ -1425,6 +1484,18 @@ def live_conformance_harness(
|
||||
region=region,
|
||||
metadata_location=metadata_location,
|
||||
)
|
||||
duckdb_rest_sql = duckdb_rest_catalog_sql(
|
||||
endpoint=endpoint,
|
||||
warehouse=warehouse,
|
||||
access_key=access_key,
|
||||
secret_key=secret_key,
|
||||
region=region,
|
||||
catalog_name=catalog_name,
|
||||
namespace=namespace,
|
||||
table=table,
|
||||
rest_path=rest_path,
|
||||
rest_signing_name=rest_signing_name,
|
||||
)
|
||||
snowflake_sql = snowflake_sql_template(
|
||||
endpoint=endpoint,
|
||||
warehouse=warehouse,
|
||||
@@ -1553,13 +1624,24 @@ def live_conformance_harness(
|
||||
OrderedDict(
|
||||
[
|
||||
("name", "DuckDB Iceberg"),
|
||||
("status", "manual-live-read-probe"),
|
||||
("status", "manual-live-harness"),
|
||||
("version", duckdb_version),
|
||||
("metadata_location", metadata_location),
|
||||
("sql_file", "/tmp/rustfs-s3tables-duckdb-read.sql"),
|
||||
("sql", duckdb_sql),
|
||||
("command", duckdb_command()),
|
||||
("expected", "iceberg_scan returns row_count=2 when metadata_location points at the current Iceberg metadata JSON"),
|
||||
("rest_catalog_sql_file", "/tmp/rustfs-s3tables-duckdb-rest.sql"),
|
||||
("rest_catalog_sql", duckdb_rest_sql),
|
||||
(
|
||||
"rest_catalog_command",
|
||||
duckdb_command(sql_file="/tmp/rustfs-s3tables-duckdb-rest.sql"),
|
||||
),
|
||||
(
|
||||
"rest_catalog_expected",
|
||||
"generic Iceberg REST ATTACH returns row_count=2 for an existing RustFS table",
|
||||
),
|
||||
("rest_catalog_write_compatibility", "manual-live-validation-required"),
|
||||
("write_compatibility", "not-claimed"),
|
||||
]
|
||||
),
|
||||
@@ -1627,6 +1709,7 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
|
||||
parser.add_argument("--print-operations-guide", action="store_true")
|
||||
parser.add_argument("--print-spark-config", action="store_true")
|
||||
parser.add_argument("--print-spark-sql", action="store_true")
|
||||
parser.add_argument("--print-duckdb-rest-sql", action="store_true")
|
||||
return parser.parse_args(argv)
|
||||
|
||||
|
||||
@@ -1734,6 +1817,24 @@ def run(args: argparse.Namespace, output: StringIO | None = None) -> None:
|
||||
else:
|
||||
output.write(sql)
|
||||
printed = True
|
||||
if args.print_duckdb_rest_sql:
|
||||
sql = duckdb_rest_catalog_sql(
|
||||
endpoint=args.endpoint,
|
||||
warehouse=args.warehouse,
|
||||
access_key=args.access_key,
|
||||
secret_key=args.secret_key,
|
||||
region=args.region,
|
||||
catalog_name=args.catalog_name,
|
||||
namespace=args.namespace,
|
||||
table=args.table,
|
||||
rest_path=args.rest_path or "/iceberg",
|
||||
rest_signing_name=args.rest_signing_name or "s3",
|
||||
)
|
||||
if output is None:
|
||||
print(sql, end="")
|
||||
else:
|
||||
output.write(sql)
|
||||
printed = True
|
||||
if not printed:
|
||||
print_json({"engine_compatibility": engine_compatibility_matrix()}, output)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user