From f5e1ca53ff0149ae103f8c24003c99a022e6369c Mon Sep 17 00:00:00 2001 From: Arun Sharma Date: Fri, 18 Sep 2026 17:44:23 -0700 Subject: [PATCH 1/2] Iceberg/Unity Catalog ATTACH support and SQL pushdown generalization Iceberg and Unity Catalog are SQL engines, so besides LOAD FROM they now support ATTACH as graph-queryable databases, reusing the shared DuckDB-backed catalog machinery: - iceberg: new IcebergStorageExtension (DBTYPE ICEBERG). ATTACH '' AS attaches the REST catalog inside the embedded DuckDB instance under the attached alias; connection options come from the iceberg_* settings. IcebergConnector::connect() honors the attach path/alias instead of ignoring them, with server-independent validation (missing warehouse/endpoint, bad SCHEMA/SKIP_UNSUPPORTED_TABLE). - unity_catalog: fix DuckDBCatalog using the attach path instead of the attached alias as its catalog name ( broke ATTACH ... AS whenever alias != path); LOAD-first/INSTALL-fallback for the uc_catalog DuckDB extension so pre-installed copies work offline without a repository origin conflict; drop duplicated delta install. - shared DuckDB catalog: init() collects tables across all result chunks; AttachedDuckDBDatabase::getTableColumnNames() accepts qualified catalog[.schema].table references and scopes information_schema accordingly (unqualified fallback preserved). Tests: new iceberg_attach.test (runnable ATTACH validation + SKIPped live REST-catalog flow); unity_catalog.test gains a SKIPped aliased-attach and pushdown case (both need live servers, per repo convention). --- duckdb/src/catalog/duckdb_catalog.cpp | 38 +++++----- duckdb/src/include/catalog/duckdb_catalog.h | 1 - .../storage/attached_duckdb_database.h | 64 +++++++++++++++- iceberg/CMakeLists.txt | 1 + iceberg/README.md | 24 ++++++ iceberg/src/connector/iceberg_connector.cpp | 50 ++++++++++--- iceberg/src/include/options/iceberg_options.h | 3 +- iceberg/src/include/storage/iceberg_storage.h | 24 ++++++ iceberg/src/main/iceberg_extension.cpp | 2 + iceberg/src/options/iceberg_options.cpp | 5 +- iceberg/src/storage/CMakeLists.txt | 12 +++ iceberg/src/storage/iceberg_storage.cpp | 46 ++++++++++++ iceberg/test/test_files/iceberg_attach.test | 75 +++++++++++++++++++ .../src/connector/unity_catalog_connector.cpp | 23 ++++-- .../src/storage/unity_catalog_storage.cpp | 11 ++- .../test/test_files/unity_catalog.test | 28 +++++++ 16 files changed, 365 insertions(+), 42 deletions(-) create mode 100644 iceberg/src/include/storage/iceberg_storage.h create mode 100644 iceberg/src/storage/CMakeLists.txt create mode 100644 iceberg/src/storage/iceberg_storage.cpp create mode 100644 iceberg/test/test_files/iceberg_attach.test diff --git a/duckdb/src/catalog/duckdb_catalog.cpp b/duckdb/src/catalog/duckdb_catalog.cpp index 4f13a4fb..66496160 100644 --- a/duckdb/src/catalog/duckdb_catalog.cpp +++ b/duckdb/src/catalog/duckdb_catalog.cpp @@ -14,7 +14,6 @@ #include "common/string_utils.h" #include "connector/duckdb_type_converter.h" #include "function/duckdb_scan.h" -#include "storage/buffer_manager/memory_manager.h" #include "storage/duckdb_storage.h" #include "storage/storage_manager.h" #include @@ -26,9 +25,7 @@ DuckDBCatalog::DuckDBCatalog(std::string dbPath, std::string catalogName, std::string defaultSchemaName, main::ClientContext* context, const DuckDBConnector& connector, const binder::AttachOption& attachOption, std::string attachedDbName) : CatalogExtension{}, dbPath{std::move(dbPath)}, catalogName{std::move(catalogName)}, - defaultSchemaName{std::move(defaultSchemaName)}, - dbName{std::move(attachedDbName)}, - tableNamesVector{common::LogicalType::STRING(), storage::MemoryManager::Get(*context)}, + defaultSchemaName{std::move(defaultSchemaName)}, dbName{std::move(attachedDbName)}, connector{connector}, context_{context} { skipUnsupportedTable = DuckDBStorageExtension::SKIP_UNSUPPORTED_TABLE_DEFAULT_VAL; auto& options = attachOption.options; @@ -48,26 +45,32 @@ void DuckDBCatalog::init() { "table_schema = '{}' order by table_name;", catalogName, defaultSchemaName); auto result = connector.executeQuery(query); - std::unique_ptr resultChunk; - try { - resultChunk = result->Fetch(); - } catch (std::exception& e) { - throw common::BinderException(e.what()); + // Collect every table name across all result chunks: catalogs with more + // tables than fit a single DuckDB data chunk must not lose tables. + std::vector tableNames; + while (true) { + std::unique_ptr resultChunk; + try { + resultChunk = result->Fetch(); + } catch (std::exception& e) { + throw common::BinderException(e.what()); + } + if (resultChunk == nullptr || resultChunk->size() == 0) { + break; + } + for (auto i = 0u; i < resultChunk->size(); i++) { + tableNames.push_back(resultChunk->GetValue(0, i).GetValue()); + } } - if (resultChunk == nullptr || resultChunk->size() == 0) { + if (tableNames.empty()) { return; } - duckdb_conversion_func_t conversionFunc; - DuckDBResultConverter::getDuckDBVectorConversionFunc(common::PhysicalTypeID::STRING, - conversionFunc); - conversionFunc(resultChunk->data[0], tableNamesVector, resultChunk->size()); // Two-pass initialization: node tables must be registered before rel tables // so that rel tables can resolve their src/dst node table IDs. The table // enumeration order is alphabetical, which can put rel_* tables before the // node tables they reference. // First pass: register node tables (everything that is not a rel table). - for (auto i = 0u; i < resultChunk->size(); i++) { - auto tableName = tableNamesVector.getValue(i).getAsString(); + for (auto& tableName : tableNames) { auto lowerName = tableName; common::StringUtils::toLower(lowerName); if (lowerName.rfind("rel_", 0) == 0 || lowerName.rfind("csr_rel_", 0) == 0) { @@ -76,8 +79,7 @@ void DuckDBCatalog::init() { createForeignTable(tableName); } // Second pass: register rel tables. - for (auto i = 0u; i < resultChunk->size(); i++) { - auto tableName = tableNamesVector.getValue(i).getAsString(); + for (auto& tableName : tableNames) { auto lowerName = tableName; common::StringUtils::toLower(lowerName); if (lowerName.rfind("rel_", 0) == 0) { diff --git a/duckdb/src/include/catalog/duckdb_catalog.h b/duckdb/src/include/catalog/duckdb_catalog.h index a94dfe63..6ec7e5d1 100644 --- a/duckdb/src/include/catalog/duckdb_catalog.h +++ b/duckdb/src/include/catalog/duckdb_catalog.h @@ -61,7 +61,6 @@ class DuckDBCatalog : public extension::CatalogExtension { // entries must store this name -- not the schema-qualified catalog name. std::string defaultSchemaName; std::string dbName; - common::ValueVector tableNamesVector; bool skipUnsupportedTable; const DuckDBConnector& connector; main::ClientContext* context_; diff --git a/duckdb/src/include/storage/attached_duckdb_database.h b/duckdb/src/include/storage/attached_duckdb_database.h index d8916c5b..fbfe3509 100644 --- a/duckdb/src/include/storage/attached_duckdb_database.h +++ b/duckdb/src/include/storage/attached_duckdb_database.h @@ -1,4 +1,7 @@ #pragma once +#include +#include + #include "connector/duckdb_connector.h" #include "main/attached_database.h" #include @@ -6,6 +9,42 @@ namespace lbug { namespace duckdb_extension { +// Splits a possibly qualified SQL table reference into its parts, honouring +// double-quoted identifiers: `"catalog".schema.table` -> [catalog, schema, +// table]. Surrounding quotes are stripped from each part. +static std::vector splitQualifiedTableName(const std::string& tableName) { + std::vector parts; + std::string current; + bool inQuotes = false; + for (auto c : tableName) { + if (c == '"') { + inQuotes = !inQuotes; + continue; + } + if (c == '.' && !inQuotes) { + parts.push_back(current); + current.clear(); + continue; + } + current.push_back(c); + } + parts.push_back(current); + return parts; +} + +static std::string escapeSingleQuotes(const std::string& value) { + std::string result; + result.reserve(value.size()); + for (auto c : value) { + if (c == '\'') { + result += "''"; + } else { + result += c; + } + } + return result; +} + class AttachedDuckDBDatabase : public main::AttachedDatabase { public: AttachedDuckDBDatabase(std::string dbName, std::string dbType, @@ -22,11 +61,32 @@ class AttachedDuckDBDatabase : public main::AttachedDatabase { } std::vector getTableColumnNames(const std::string& tableName) const override { + // Accepts both bare table names and qualified `catalog[.schema].table` + // references (as produced by the SQL push-down optimizer). Scoping the + // information_schema lookup by catalog/schema keeps same-named tables + // in different schemas or catalogs of one attached database apart. + auto parts = splitQualifiedTableName(tableName); + auto unqualified = parts.back(); + std::string filters; + if (parts.size() == 3) { + filters = std::format(" AND table_catalog = '{}' AND table_schema = '{}'", + escapeSingleQuotes(parts[0]), escapeSingleQuotes(parts[1])); + } else if (parts.size() == 2) { + filters = std::format(" AND table_schema = '{}'", escapeSingleQuotes(parts[0])); + } std::string query = std::format("SELECT column_name FROM information_schema.columns WHERE " - "table_name = '{}' ORDER BY ordinal_position", - tableName); + "table_name = '{}'{} ORDER BY ordinal_position", + escapeSingleQuotes(unqualified), filters); auto result = connector->executeQuery(query); + if ((!result || result->RowCount() == 0) && !filters.empty()) { + // Fall back to the unqualified lookup: some engines report + // catalog/schema names differently than the attached alias. + query = std::format("SELECT column_name FROM information_schema.columns WHERE " + "table_name = '{}' ORDER BY ordinal_position", + escapeSingleQuotes(unqualified)); + result = connector->executeQuery(query); + } if (!result || result->RowCount() == 0) { return {}; } diff --git a/iceberg/CMakeLists.txt b/iceberg/CMakeLists.txt index bfc60036..cc5acb48 100644 --- a/iceberg/CMakeLists.txt +++ b/iceberg/CMakeLists.txt @@ -18,6 +18,7 @@ include_directories( add_subdirectory(src/function) add_subdirectory(src/connector) add_subdirectory(src/options) +add_subdirectory(src/storage) add_subdirectory(src/installer) add_subdirectory(src/main) diff --git a/iceberg/README.md b/iceberg/README.md index 0eae0fd5..85dd63d3 100644 --- a/iceberg/README.md +++ b/iceberg/README.md @@ -43,6 +43,30 @@ alias `iceberg_catalog`. Table references can use either `iceberg_catalog.namespace.table` (3-part) or `namespace.table` (2-part, which is prefixed with the alias automatically). Any other 3-part name is rejected. +## Attaching a REST catalog as a database + +An Iceberg REST catalog is a SQL engine, so besides `LOAD FROM` it can be +attached as a database and queried with graph patterns. Attached tables behave +like any other foreign tables: `MATCH` scans push filters, projections, +limits and ordering down to the catalog, and multi-hop patterns over +relationship tables are rewritten into a single SQL join (see the foreign join +push-down optimizer). + +```cypher +CALL iceberg_endpoint='https://rest-catalog.example.com'; +CALL iceberg_token=''; +ATTACH 'warehouse' AS ice (DBTYPE ICEBERG); +LOAD FROM ice.events RETURN count(*); +MATCH (e:ice.events) WHERE e.ts > timestamp('2026-01-01 00:00:00') RETURN count(*); +``` + +The ATTACH path is the warehouse identifier; connection and authentication +options come from the `iceberg_*` options above. Tables are enumerated from +the `default` namespace unless overridden with the `SCHEMA` attach option +(e.g. `ATTACH 'warehouse' AS ice (DBTYPE ICEBERG, SCHEMA = 'analytics')`), and +tables with unsupported column types can be skipped with +`SKIP_UNSUPPORTED_TABLE = true`. + ### Options | Option | Description | diff --git a/iceberg/src/connector/iceberg_connector.cpp b/iceberg/src/connector/iceberg_connector.cpp index eeaa3a6c..4b868dd2 100644 --- a/iceberg/src/connector/iceberg_connector.cpp +++ b/iceberg/src/connector/iceberg_connector.cpp @@ -6,14 +6,47 @@ namespace lbug { namespace iceberg_extension { -void IcebergConnector::connect(const std::string& /*dbPath*/, const std::string& /*catalogName*/, +void IcebergConnector::connect(const std::string& dbPath, const std::string& catalogName, const std::string& /*schemaName*/, main::ClientContext* context) { auto config = IcebergOptions::getRestCatalogConfig(context); - if (!config.restCatalogConfigured() && config.hasAnyOption()) { - throw common::RuntimeException{std::format( - "Iceberg REST catalog options were set but '{}' is empty. Set '{}' to the warehouse " - "identifier to enable Iceberg REST catalog access.", - IcebergWarehouse::NAME, IcebergWarehouse::NAME)}; + // ATTACH '' AS (DBTYPE ICEBERG) passes the warehouse and + // the attached alias through; the LOAD FROM table-function path calls + // connect() with empty arguments and relies on the global options instead. + const bool attaching = !dbPath.empty() || !catalogName.empty(); + if (!dbPath.empty()) { + config.warehouse = dbPath; + } + const auto catalogAlias = + catalogName.empty() ? IcebergSecretManager::CATALOG_ALIAS : catalogName; + if (config.warehouse.empty()) { + if (attaching) { + throw common::RuntimeException{ + "Cannot attach an Iceberg catalog without a warehouse. Give the warehouse " + "identifier as the ATTACH path (ATTACH '' AS (DBTYPE " + "ICEBERG)) or set the 'iceberg_warehouse' option."}; + } + if (config.hasAnyOption()) { + throw common::RuntimeException{std::format( + "Iceberg REST catalog options were set but '{}' is empty. Set '{}' to the " + "warehouse " + "identifier to enable Iceberg REST catalog access.", + IcebergWarehouse::NAME, IcebergWarehouse::NAME)}; + } + // Creates an in-memory duckdb instance, then install iceberg and httpfs. + instance = std::make_unique(nullptr); + connection = std::make_unique(*instance); + // Install the Desired Extension on DuckDB + executeQuery("install iceberg;"); + executeQuery("load iceberg;"); + executeQuery("install httpfs;"); + executeQuery("load httpfs;"); + initRemoteFSSecrets(context); + return; + } + if (attaching && config.endpoint.empty()) { + throw common::RuntimeException{ + "Cannot attach an Iceberg REST catalog without an endpoint. Set the " + "'iceberg_endpoint' option to the REST catalog URL before attaching."}; } // Creates an in-memory duckdb instance, then install iceberg and httpfs. instance = std::make_unique(nullptr); @@ -28,13 +61,10 @@ void IcebergConnector::connect(const std::string& /*dbPath*/, const std::string& // DuckDB instance, so iceberg tables can be referenced by fully qualified // name (e.g. iceberg_catalog.default.events) instead of a filesystem path. // This avoids the version-hint.text dependency of filesystem-based catalogs. - if (!config.restCatalogConfigured()) { - return; - } if (config.hasAuth()) { executeQuery(IcebergSecretManager::getSecret(config)); } - executeQuery(IcebergSecretManager::getAttachQuery(config)); + executeQuery(IcebergSecretManager::getAttachQuery(config, catalogAlias)); } } // namespace iceberg_extension diff --git a/iceberg/src/include/options/iceberg_options.h b/iceberg/src/include/options/iceberg_options.h index 3f85cda0..3bd077ce 100644 --- a/iceberg/src/include/options/iceberg_options.h +++ b/iceberg/src/include/options/iceberg_options.h @@ -106,7 +106,8 @@ struct IcebergSecretManager { static constexpr const char* SECRET_NAME = "iceberg_rest_secret"; static std::string getSecret(const IcebergRestCatalogConfig& config); - static std::string getAttachQuery(const IcebergRestCatalogConfig& config); + static std::string getAttachQuery(const IcebergRestCatalogConfig& config, + const std::string& catalogAlias); }; } // namespace iceberg_extension diff --git a/iceberg/src/include/storage/iceberg_storage.h b/iceberg/src/include/storage/iceberg_storage.h new file mode 100644 index 00000000..bc82fe17 --- /dev/null +++ b/iceberg/src/include/storage/iceberg_storage.h @@ -0,0 +1,24 @@ +#pragma once + +#include "storage/storage_extension.h" + +namespace lbug { +namespace main { +class Database; +} // namespace main + +namespace iceberg_extension { + +class IcebergStorageExtension final : public storage::StorageExtension { +public: + static constexpr const char* DB_TYPE = "ICEBERG"; + + static constexpr const char* DEFAULT_SCHEMA_NAME = "default"; + + explicit IcebergStorageExtension(main::Database& database); + + bool canHandleDB(std::string dbType) const override; +}; + +} // namespace iceberg_extension +} // namespace lbug diff --git a/iceberg/src/main/iceberg_extension.cpp b/iceberg/src/main/iceberg_extension.cpp index f2cc1c4c..4c3399a4 100644 --- a/iceberg/src/main/iceberg_extension.cpp +++ b/iceberg/src/main/iceberg_extension.cpp @@ -5,6 +5,7 @@ #include "main/database.h" #include "main/duckdb_extension.h" #include "options/iceberg_options.h" +#include "storage/iceberg_storage.h" namespace lbug { namespace iceberg_extension { @@ -13,6 +14,7 @@ using namespace lbug::extension; void IcebergExtension::load(main::ClientContext* context) { auto& db = *context->getDatabase(); + db.registerStorageExtension(EXTENSION_NAME, std::make_unique(db)); ExtensionUtils::addTableFunc(db); ExtensionUtils::addTableFunc(db); ExtensionUtils::addTableFunc(db); diff --git a/iceberg/src/options/iceberg_options.cpp b/iceberg/src/options/iceberg_options.cpp index d04c71a4..3f6ee724 100644 --- a/iceberg/src/options/iceberg_options.cpp +++ b/iceberg/src/options/iceberg_options.cpp @@ -96,7 +96,8 @@ std::string IcebergSecretManager::getSecret(const IcebergRestCatalogConfig& conf return std::format("CREATE SECRET {} ({} TYPE ICEBERG);", SECRET_NAME, options); } -std::string IcebergSecretManager::getAttachQuery(const IcebergRestCatalogConfig& config) { +std::string IcebergSecretManager::getAttachQuery(const IcebergRestCatalogConfig& config, + const std::string& catalogAlias) { std::string options = "TYPE ICEBERG"; if (config.hasAuth()) { options += std::format(", SECRET {}", SECRET_NAME); @@ -112,7 +113,7 @@ std::string IcebergSecretManager::getAttachQuery(const IcebergRestCatalogConfig& std::format(", AUTHORIZATION_TYPE '{}'", escapeSingleQuotes(config.authorizationType)); } return std::format("ATTACH '{}' AS {} ({});", escapeSingleQuotes(config.warehouse), - CATALOG_ALIAS, options); + catalogAlias, options); } } // namespace iceberg_extension diff --git a/iceberg/src/storage/CMakeLists.txt b/iceberg/src/storage/CMakeLists.txt new file mode 100644 index 00000000..25e1ecb4 --- /dev/null +++ b/iceberg/src/storage/CMakeLists.txt @@ -0,0 +1,12 @@ +add_library(lbug_iceberg_storage + OBJECT + iceberg_storage.cpp + ${PROJECT_SOURCE_DIR}/extension/duckdb/src/catalog/duckdb_catalog.cpp + ${PROJECT_SOURCE_DIR}/extension/duckdb/src/function/duckdb_scan.cpp + ${PROJECT_SOURCE_DIR}/extension/duckdb/src/function/clear_cache.cpp + ${PROJECT_SOURCE_DIR}/extension/duckdb/src/catalog/duckdb_table_catalog_entry.cpp +) + +set(ICEBERG_EXTENSION_OBJECT_FILES + ${ICEBERG_EXTENSION_OBJECT_FILES} $ + PARENT_SCOPE) diff --git a/iceberg/src/storage/iceberg_storage.cpp b/iceberg/src/storage/iceberg_storage.cpp new file mode 100644 index 00000000..eafbcade --- /dev/null +++ b/iceberg/src/storage/iceberg_storage.cpp @@ -0,0 +1,46 @@ +#include "storage/iceberg_storage.h" + +#include "catalog/duckdb_catalog.h" +#include "common/string_utils.h" +#include "connector/iceberg_connector.h" +#include "extension/extension.h" +#include "function/clear_cache.h" +#include "storage/attached_duckdb_database.h" + +namespace lbug { +namespace iceberg_extension { + +std::unique_ptr attachIceberg(std::string dbName, std::string dbPath, + main::ClientContext* clientContext, const binder::AttachOption& attachOption) { + if (dbName == "") { + dbName = dbPath; + } + auto schemaName = duckdb_extension::DuckDBCatalog::bindSchemaName(attachOption, + IcebergStorageExtension::DEFAULT_SCHEMA_NAME); + auto connector = std::make_unique(); + // The DuckDB-side catalog is attached AS dbName, so DuckDBCatalog must use + // dbName (not dbPath) as its catalog name: information_schema lookups and + // generated SQL reference the attached alias. The catalog is constructed + // before connecting so that ATTACH option validation happens before any + // network I/O. + auto catalog = std::make_unique(dbPath, dbName, schemaName, + clientContext, *connector, attachOption, dbName); + connector->connect(dbPath, dbName, schemaName, clientContext); + catalog->init(); + return std::make_unique(dbName, + IcebergStorageExtension::DB_TYPE, std::move(catalog), std::move(connector)); +} + +IcebergStorageExtension::IcebergStorageExtension(main::Database& database) + : StorageExtension{attachIceberg} { + extension::ExtensionUtils::addStandaloneTableFunc( + database); +} + +bool IcebergStorageExtension::canHandleDB(std::string dbType_) const { + common::StringUtils::toUpper(dbType_); + return dbType_ == DB_TYPE; +} + +} // namespace iceberg_extension +} // namespace lbug diff --git a/iceberg/test/test_files/iceberg_attach.test b/iceberg/test/test_files/iceberg_attach.test new file mode 100644 index 00000000..80d6ccef --- /dev/null +++ b/iceberg/test/test_files/iceberg_attach.test @@ -0,0 +1,75 @@ +-DATASET CSV empty + +-- + +# ATTACH validation that does not require a running Iceberg REST catalog. + +-CASE AttachValidation +-STATEMENT load extension "${LBUG_ROOT_DIRECTORY}/extension/iceberg/build/libiceberg.lbug_extension"; +---- ok +-LOG AttachWithoutEndpoint +-STATEMENT ATTACH 'my_warehouse' AS ice (DBTYPE ICEBERG); +---- error +Runtime exception: Cannot attach an Iceberg REST catalog without an endpoint. Set the 'iceberg_endpoint' option to the REST catalog URL before attaching. +-LOG AttachWithoutWarehouse +-STATEMENT CALL iceberg_endpoint='http://127.0.0.1:8181'; +---- ok +-STATEMENT ATTACH '' AS ice (DBTYPE ICEBERG); +---- error +Runtime exception: Cannot attach an Iceberg catalog without a warehouse. Give the warehouse identifier as the ATTACH path (ATTACH '' AS (DBTYPE ICEBERG)) or set the 'iceberg_warehouse' option. +-LOG AttachWithInvalidOptions +-STATEMENT ATTACH 'my_warehouse' AS ice (DBTYPE ICEBERG, SKIP_UNSUPPORTED_TABLE = 'yes'); +---- error +Runtime exception: Invalid option value for SKIP_UNSUPPORTED_TABLE +-STATEMENT ATTACH 'my_warehouse' AS ice (DBTYPE ICEBERG, SCHEMA = 5); +---- error +Runtime exception: Invalid option value for SCHEMA +-STATEMENT CALL iceberg_endpoint=''; +---- ok +-STATEMENT LOAD FROM '${LBUG_ROOT_DIRECTORY}/extension/iceberg/test/iceberg_tables/person_table' (file_format='iceberg', allow_moved_paths=true) RETURN count(*); +---- 1 +5 + +# Requires a running Iceberg REST catalog (e.g. Lakekeeper) at +# http://127.0.0.1:8181 that serves the lineitem_iceberg test table as +# default.lineitem_iceberg. Tuple counts assume the catalog's current snapshot +# is the same as the local table's. +-CASE AttachRestCatalog +-SKIP +-STATEMENT load extension "${LBUG_ROOT_DIRECTORY}/extension/iceberg/build/libiceberg.lbug_extension"; +---- ok +-STATEMENT CALL iceberg_warehouse='warehouse'; +---- ok +-STATEMENT CALL iceberg_endpoint='http://127.0.0.1:8181'; +---- ok +-STATEMENT CALL iceberg_token='test-token'; +---- ok +-LOG AttachIcebergCatalog +-STATEMENT ATTACH 'warehouse' AS ice (DBTYPE ICEBERG); +---- 1 +Attached database successfully. +-STATEMENT LOAD FROM ice.lineitem_iceberg RETURN count(*); +---- 1 +60175 +-STATEMENT MATCH (l:ice.lineitem_iceberg) RETURN count(*); +---- 1 +60175 +-STATEMENT CALL TABLE_INFO('ice.lineitem_iceberg') RETURN *; +---- ok +-STATEMENT CALL SHOW_TABLES() RETURN *; +---- ok +-LOG JoinPushdownOverAttachedCatalog +-STATEMENT CREATE REL TABLE lineitem_rel (FROM ice.lineitem_iceberg TO ice.lineitem_iceberg) WITH (storage = 'ice.lineitem_iceberg'); +---- 1 +Table lineitem_rel has been created. +-STATEMENT EXPLAIN MATCH (a:ice.lineitem_iceberg)-[b:lineitem_rel]->(c:ice.lineitem_iceberg) RETURN count(*); +---- ok +-STATEMENT DETACH ice; +---- 1 +Detached database successfully. +-STATEMENT CALL iceberg_warehouse=''; +---- ok +-STATEMENT CALL iceberg_endpoint=''; +---- ok +-STATEMENT CALL iceberg_token=''; +---- ok diff --git a/unity_catalog/src/connector/unity_catalog_connector.cpp b/unity_catalog/src/connector/unity_catalog_connector.cpp index bf1aeb56..d6111b18 100644 --- a/unity_catalog/src/connector/unity_catalog_connector.cpp +++ b/unity_catalog/src/connector/unity_catalog_connector.cpp @@ -6,16 +6,29 @@ namespace lbug { namespace unity_catalog_extension { +// Ensures a DuckDB extension is loaded. LOAD is attempted first so that an +// already-installed extension is reused as-is (no network, no repository +// origin check). INSTALL runs only as a fallback when the extension is +// missing, e.g. on a fresh machine. +static void ensureExtensionLoaded(const duckdb_extension::DuckDBConnector& connector, + const std::string& extensionName) { + try { + connector.executeQuery(std::format("load {};", extensionName)); + return; + } catch (const common::Exception&) { + // Not installed yet: fall through to INSTALL. + } + connector.executeQuery(std::format("install {};", extensionName)); + connector.executeQuery(std::format("load {};", extensionName)); +} + void UnityCatalogConnector::connect(const std::string& dbPath, const std::string& catalogName, const std::string& /*schemaName*/, main::ClientContext* context) { // Creates an in-memory duckdb instance, then install httpfs and attach postgres. instance = std::make_unique(nullptr); connection = std::make_unique(*instance); - executeQuery("install uc_catalog from core_nightly;"); - executeQuery("load uc_catalog;"); - executeQuery("install delta;"); - executeQuery("load delta;"); - executeQuery("install delta;"); + ensureExtensionLoaded(*this, "uc_catalog"); + ensureExtensionLoaded(*this, "delta"); executeQuery(DuckDBUnityCatalogSecretManager::getSecret(context)); executeQuery( std::format("attach '{}' as {} (TYPE UC_CATALOG, read_only);", dbPath, catalogName)); diff --git a/unity_catalog/src/storage/unity_catalog_storage.cpp b/unity_catalog/src/storage/unity_catalog_storage.cpp index 8c8431e5..a8dd237d 100644 --- a/unity_catalog/src/storage/unity_catalog_storage.cpp +++ b/unity_catalog/src/storage/unity_catalog_storage.cpp @@ -16,11 +16,16 @@ std::unique_ptr attachUnityCatalog(std::string dbName, s dbName = dbPath; } auto connector = std::make_unique(); - connector->connect(dbPath, dbName, UnityCatalogStorageExtension::DEFAULT_SCHEMA_NAME, - clientContext); - auto catalog = std::make_unique(dbPath, dbPath, + // The DuckDB-side catalog is attached AS dbName, so DuckDBCatalog must use + // dbName (not dbPath) as its catalog name: information_schema lookups and + // generated SQL reference the attached alias. The catalog is constructed + // before connecting so that ATTACH option validation happens before any + // network I/O. + auto catalog = std::make_unique(dbPath, dbName, UnityCatalogStorageExtension::DEFAULT_SCHEMA_NAME, clientContext, *connector, attachOption, dbName); + connector->connect(dbPath, dbName, UnityCatalogStorageExtension::DEFAULT_SCHEMA_NAME, + clientContext); catalog->init(); return std::make_unique(dbName, UnityCatalogStorageExtension::DB_TYPE, std::move(catalog), std::move(connector)); diff --git a/unity_catalog/test/test_files/unity_catalog.test b/unity_catalog/test/test_files/unity_catalog.test index 8088d7c2..989cca1f 100644 --- a/unity_catalog/test/test_files/unity_catalog.test +++ b/unity_catalog/test/test_files/unity_catalog.test @@ -21,3 +21,31 @@ Dan|False|9|20|19|7.700000|9.200000 # Duckdb has an issue when scanning this table, enable after duckdb fixes it. #-STATEMENT LOAD FROM university.person RETURN *; #---- 1 + +# Requires a running Unity Catalog server at http://127.0.0.1:8080 serving the +# university catalog (see test/setup/create.bash). Exercises ATTACH under an +# alias different from the catalog name, MATCH scans, and SQL pushdown. +-CASE AttachWithAliasAndPushdown +-SKIP +-STATEMENT load extension "${LBUG_ROOT_DIRECTORY}/extension/unity_catalog/build/libunity_catalog.lbug_extension"; +---- ok +-LOG AttachUnderAlias +-STATEMENT ATTACH 'university' AS u (DBTYPE UC_CATALOG); +---- ok +-STATEMENT LOAD FROM u.grades RETURN count(*); +---- 1 +4 +-STATEMENT MATCH (g:u.grades) RETURN count(*); +---- 1 +4 +-STATEMENT CALL TABLE_INFO('u.grades') RETURN *; +---- ok +-LOG JoinPushdownOverAttachedCatalog +-STATEMENT CREATE REL TABLE grade_rel (FROM u.grades TO u.grades) WITH (storage = 'u.grades'); +---- 1 +Table grade_rel has been created. +-STATEMENT EXPLAIN MATCH (a:u.grades)-[b:grade_rel]->(c:u.grades) RETURN count(*); +---- ok +-STATEMENT DETACH u; +---- 1 +Detached database successfully. From 1da29c58032e1d5aaa8c51cbac1942812496caa4 Mon Sep 17 00:00:00 2001 From: Arun Sharma Date: Sat, 19 Sep 2026 12:00:15 -0700 Subject: [PATCH 2/2] Address review findings: tiered column lookup, connector dedup, hermetic pushdown test - attached_duckdb_database.h: narrow the getTableColumnNames fallback to most-to-least specific filters (catalog+schema, schema-only, then unqualified last resort) with comments documenting the 2-part/3-part convention shared with the push-down optimizer; mark header helpers inline instead of static. - iceberg_connector.cpp: validate configuration before creating the embedded DuckDB instance instead of duplicating the setup block in both branches; behavior unchanged. - iceberg storage CMakeLists: drop duckdb_scan.cpp, already provided by the delta_connector static lib linked into the extension. - README: document that an empty ATTACH path falls back to the iceberg_warehouse option. - duckdb_rel.test: new hermetic DuckDBAttachAliasJoinPushdown case covering attach-under-alias, catalog init, column resolution and join pushdown with property projections (no external server needed). --- .../storage/attached_duckdb_database.h | 65 +++++++++++-------- duckdb/test/test_files/duckdb_rel.test | 28 ++++++++ iceberg/README.md | 6 +- iceberg/src/connector/iceberg_connector.cpp | 20 +++--- iceberg/src/storage/CMakeLists.txt | 1 - 5 files changed, 79 insertions(+), 41 deletions(-) diff --git a/duckdb/src/include/storage/attached_duckdb_database.h b/duckdb/src/include/storage/attached_duckdb_database.h index fbfe3509..2ffba3e6 100644 --- a/duckdb/src/include/storage/attached_duckdb_database.h +++ b/duckdb/src/include/storage/attached_duckdb_database.h @@ -12,7 +12,13 @@ namespace duckdb_extension { // Splits a possibly qualified SQL table reference into its parts, honouring // double-quoted identifiers: `"catalog".schema.table` -> [catalog, schema, // table]. Surrounding quotes are stripped from each part. -static std::vector splitQualifiedTableName(const std::string& tableName) { +// +// Callers treat a 2-part reference as `schema.table` and a 3-part reference as +// `catalog.schema.table`, matching the convention used by the SQL push-down +// optimizer (foreign table scans are described as `"catalog".schema.table`). +// References with more than 3 parts have no defined meaning and are looked up +// by their unqualified table name. +inline std::vector splitQualifiedTableName(const std::string& tableName) { std::vector parts; std::string current; bool inQuotes = false; @@ -32,7 +38,7 @@ static std::vector splitQualifiedTableName(const std::string& table return parts; } -static std::string escapeSingleQuotes(const std::string& value) { +inline std::string escapeSingleQuotes(const std::string& value) { std::string result; result.reserve(value.size()); for (auto c : value) { @@ -67,35 +73,40 @@ class AttachedDuckDBDatabase : public main::AttachedDatabase { // in different schemas or catalogs of one attached database apart. auto parts = splitQualifiedTableName(tableName); auto unqualified = parts.back(); - std::string filters; + // Candidate filters from most to least specific. Some engines report + // the catalog name differently than the attached alias, so a fully + // qualified lookup may miss while a schema-scoped one hits. The + // unqualified lookup is the last resort: it keeps attached databases + // from older extension builds (which only match bare names) working, + // but it can match a same-named table in another schema when the + // qualification is wrong, so it must stay last. + std::vector filterCandidates; if (parts.size() == 3) { - filters = std::format(" AND table_catalog = '{}' AND table_schema = '{}'", - escapeSingleQuotes(parts[0]), escapeSingleQuotes(parts[1])); + filterCandidates.push_back(std::format(" AND table_catalog = '{}' AND table_schema " + "= '{}'", + escapeSingleQuotes(parts[0]), escapeSingleQuotes(parts[1]))); + filterCandidates.push_back( + std::format(" AND table_schema = '{}'", escapeSingleQuotes(parts[1]))); } else if (parts.size() == 2) { - filters = std::format(" AND table_schema = '{}'", escapeSingleQuotes(parts[0])); + filterCandidates.push_back( + std::format(" AND table_schema = '{}'", escapeSingleQuotes(parts[0]))); } - std::string query = std::format("SELECT column_name FROM information_schema.columns WHERE " - "table_name = '{}'{} ORDER BY ordinal_position", - escapeSingleQuotes(unqualified), filters); - - auto result = connector->executeQuery(query); - if ((!result || result->RowCount() == 0) && !filters.empty()) { - // Fall back to the unqualified lookup: some engines report - // catalog/schema names differently than the attached alias. - query = std::format("SELECT column_name FROM information_schema.columns WHERE " - "table_name = '{}' ORDER BY ordinal_position", - escapeSingleQuotes(unqualified)); - result = connector->executeQuery(query); - } - if (!result || result->RowCount() == 0) { - return {}; - } - - std::vector columnNames; - for (auto i = 0u; i < result->RowCount(); i++) { - columnNames.push_back(result->GetValue(0, i).GetValue()); + filterCandidates.emplace_back(""); + for (auto& filters : filterCandidates) { + std::string query = + std::format("SELECT column_name FROM information_schema.columns WHERE table_name " + "= '{}'{} ORDER BY ordinal_position", + escapeSingleQuotes(unqualified), filters); + auto result = connector->executeQuery(query); + if (result && result->RowCount() != 0) { + std::vector columnNames; + for (auto i = 0u; i < result->RowCount(); i++) { + columnNames.push_back(result->GetValue(0, i).GetValue()); + } + return columnNames; + } } - return columnNames; + return {}; } protected: diff --git a/duckdb/test/test_files/duckdb_rel.test b/duckdb/test/test_files/duckdb_rel.test index f874ef44..a0bc911d 100644 --- a/duckdb/test/test_files/duckdb_rel.test +++ b/duckdb/test/test_files/duckdb_rel.test @@ -62,3 +62,31 @@ Attached database successfully. 2 -STATEMENT DETACH g; ---- ok + +# Exercises the shared attached-database machinery (catalog init under an alias, +# column-name resolution and the foreign join push-down optimizer) without any +# external server: the results are only correct if the pushed-down SQL join +# resolves the rel endpoint columns and the node ID columns. +-CASE DuckDBAttachAliasJoinPushdown +-SKIP_FSM_LEAK_CHECK +-LOAD_DYNAMIC_EXTENSION duckdb +-STATEMENT ATTACH '${LBUG_ROOT_DIRECTORY}/extension/duckdb/test/test_files/users_sessions.db' as adel (dbtype duckdb); +---- 1 +Attached database successfully. +-STATEMENT CREATE REL TABLE owns_rel (FROM adel.users TO adel.sessions) WITH (storage = 'adel.rel_user_owns_session'); +---- 1 +Table owns_rel has been created. +-STATEMENT EXPLAIN MATCH (u:adel.users)-[b:owns_rel]->(s:adel.sessions) RETURN count(*); +---- ok +-STATEMENT MATCH (u:adel.users)-[b:owns_rel]->(s:adel.sessions) RETURN count(*); +---- 1 +5 +-STATEMENT MATCH (u:adel.users)-[b:owns_rel]->(s:adel.sessions) RETURN u.name, s.device ORDER BY u.name, s.device; +---- 5 +Alice|laptop +Alice|mobile +Bob|desktop +Carol|mobile +Carol|tablet +-STATEMENT DETACH adel; +---- ok diff --git a/iceberg/README.md b/iceberg/README.md index 85dd63d3..1adcbbc5 100644 --- a/iceberg/README.md +++ b/iceberg/README.md @@ -61,7 +61,11 @@ MATCH (e:ice.events) WHERE e.ts > timestamp('2026-01-01 00:00:00') RETURN count( ``` The ATTACH path is the warehouse identifier; connection and authentication -options come from the `iceberg_*` options above. Tables are enumerated from +options come from the `iceberg_*` options above. If the path is empty, the +`iceberg_warehouse` option is used as the warehouse instead, so +`ATTACH '' AS ice (DBTYPE ICEBERG)` attaches the globally configured warehouse. +Attaching without a warehouse from either source fails with an error. +Tables are enumerated from the `default` namespace unless overridden with the `SCHEMA` attach option (e.g. `ATTACH 'warehouse' AS ice (DBTYPE ICEBERG, SCHEMA = 'analytics')`), and tables with unsupported column types can be skipped with diff --git a/iceberg/src/connector/iceberg_connector.cpp b/iceberg/src/connector/iceberg_connector.cpp index 4b868dd2..5f0a730d 100644 --- a/iceberg/src/connector/iceberg_connector.cpp +++ b/iceberg/src/connector/iceberg_connector.cpp @@ -18,6 +18,8 @@ void IcebergConnector::connect(const std::string& dbPath, const std::string& cat } const auto catalogAlias = catalogName.empty() ? IcebergSecretManager::CATALOG_ALIAS : catalogName; + // Validate the configuration before creating the embedded DuckDB instance, + // so misconfiguration fails without any setup side effects. if (config.warehouse.empty()) { if (attaching) { throw common::RuntimeException{ @@ -32,18 +34,7 @@ void IcebergConnector::connect(const std::string& dbPath, const std::string& cat "identifier to enable Iceberg REST catalog access.", IcebergWarehouse::NAME, IcebergWarehouse::NAME)}; } - // Creates an in-memory duckdb instance, then install iceberg and httpfs. - instance = std::make_unique(nullptr); - connection = std::make_unique(*instance); - // Install the Desired Extension on DuckDB - executeQuery("install iceberg;"); - executeQuery("load iceberg;"); - executeQuery("install httpfs;"); - executeQuery("load httpfs;"); - initRemoteFSSecrets(context); - return; - } - if (attaching && config.endpoint.empty()) { + } else if (attaching && config.endpoint.empty()) { throw common::RuntimeException{ "Cannot attach an Iceberg REST catalog without an endpoint. Set the " "'iceberg_endpoint' option to the REST catalog URL before attaching."}; @@ -57,6 +48,11 @@ void IcebergConnector::connect(const std::string& dbPath, const std::string& cat executeQuery("install httpfs;"); executeQuery("load httpfs;"); initRemoteFSSecrets(context); + if (config.warehouse.empty()) { + // No REST catalog is configured: the table functions scan iceberg + // files directly from the filesystem instead. + return; + } // If the Iceberg REST catalog is configured, attach it inside the embedded // DuckDB instance, so iceberg tables can be referenced by fully qualified // name (e.g. iceberg_catalog.default.events) instead of a filesystem path. diff --git a/iceberg/src/storage/CMakeLists.txt b/iceberg/src/storage/CMakeLists.txt index 25e1ecb4..e3d24bd7 100644 --- a/iceberg/src/storage/CMakeLists.txt +++ b/iceberg/src/storage/CMakeLists.txt @@ -2,7 +2,6 @@ add_library(lbug_iceberg_storage OBJECT iceberg_storage.cpp ${PROJECT_SOURCE_DIR}/extension/duckdb/src/catalog/duckdb_catalog.cpp - ${PROJECT_SOURCE_DIR}/extension/duckdb/src/function/duckdb_scan.cpp ${PROJECT_SOURCE_DIR}/extension/duckdb/src/function/clear_cache.cpp ${PROJECT_SOURCE_DIR}/extension/duckdb/src/catalog/duckdb_table_catalog_entry.cpp )