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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 20 additions & 18 deletions duckdb/src/catalog/duckdb_catalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 <format>
Expand All @@ -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;
Expand All @@ -48,26 +45,32 @@ void DuckDBCatalog::init() {
"table_schema = '{}' order by table_name;",
catalogName, defaultSchemaName);
auto result = connector.executeQuery(query);
std::unique_ptr<duckdb::DataChunk> 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<std::string> tableNames;
while (true) {
std::unique_ptr<duckdb::DataChunk> 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<std::string>());
}
}
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<common::string_t>(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) {
Expand All @@ -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<common::string_t>(i).getAsString();
for (auto& tableName : tableNames) {
auto lowerName = tableName;
common::StringUtils::toLower(lowerName);
if (lowerName.rfind("rel_", 0) == 0) {
Expand Down
1 change: 0 additions & 1 deletion duckdb/src/include/catalog/duckdb_catalog.h
Original file line number Diff line number Diff line change
Expand Up @@ -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_;
Expand Down
95 changes: 83 additions & 12 deletions duckdb/src/include/storage/attached_duckdb_database.h
Original file line number Diff line number Diff line change
@@ -1,11 +1,56 @@
#pragma once
#include <string>
#include <vector>

#include "connector/duckdb_connector.h"
#include "main/attached_database.h"
#include <format>

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.
//
// 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<std::string> splitQualifiedTableName(const std::string& tableName) {
std::vector<std::string> 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;
}

inline 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,
Expand All @@ -22,20 +67,46 @@ class AttachedDuckDBDatabase : public main::AttachedDatabase {
}

std::vector<std::string> getTableColumnNames(const std::string& tableName) const override {
std::string query = std::format("SELECT column_name FROM information_schema.columns WHERE "
"table_name = '{}' ORDER BY ordinal_position",
tableName);

auto result = connector->executeQuery(query);
if (!result || result->RowCount() == 0) {
return {};
// 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();
// 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<std::string> filterCandidates;
if (parts.size() == 3) {
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) {
filterCandidates.push_back(
std::format(" AND table_schema = '{}'", escapeSingleQuotes(parts[0])));
}

std::vector<std::string> columnNames;
for (auto i = 0u; i < result->RowCount(); i++) {
columnNames.push_back(result->GetValue(0, i).GetValue<std::string>());
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<std::string> columnNames;
for (auto i = 0u; i < result->RowCount(); i++) {
columnNames.push_back(result->GetValue(0, i).GetValue<std::string>());
}
return columnNames;
}
}
return columnNames;
return {};
}

protected:
Expand Down
28 changes: 28 additions & 0 deletions duckdb/test/test_files/duckdb_rel.test
Original file line number Diff line number Diff line change
Expand Up @@ -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
1 change: 1 addition & 0 deletions iceberg/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
28 changes: 28 additions & 0 deletions iceberg/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,34 @@ 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='<bearer-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. 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
`SKIP_UNSUPPORTED_TABLE = true`.

### Options

| Option | Description |
Expand Down
46 changes: 36 additions & 10 deletions iceberg/src/connector/iceberg_connector.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -6,14 +6,38 @@
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 '<warehouse>' AS <alias> (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;
// 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{
"Cannot attach an Iceberg catalog without a warehouse. Give the warehouse "
"identifier as the ATTACH path (ATTACH '<warehouse>' AS <alias> (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)};
}
} 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."};
}
// Creates an in-memory duckdb instance, then install iceberg and httpfs.
instance = std::make_unique<duckdb::DuckDB>(nullptr);
Expand All @@ -24,17 +48,19 @@ void IcebergConnector::connect(const std::string& /*dbPath*/, const std::string&
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.
// 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
Expand Down
3 changes: 2 additions & 1 deletion iceberg/src/include/options/iceberg_options.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
24 changes: 24 additions & 0 deletions iceberg/src/include/storage/iceberg_storage.h
Original file line number Diff line number Diff line change
@@ -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
2 changes: 2 additions & 0 deletions iceberg/src/main/iceberg_extension.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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<IcebergStorageExtension>(db));
ExtensionUtils::addTableFunc<IcebergScanFunction>(db);
ExtensionUtils::addTableFunc<IcebergMetadataFunction>(db);
ExtensionUtils::addTableFunc<IcebergSnapshotsFunction>(db);
Expand Down
Loading
Loading