feat(d1): add Cloudflare D1 database support

This commit is contained in:
Anubis 2026-07-12 22:03:57 +08:00 committed by GitHub
parent 977e7f78ca
commit 2a92746dc0
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
61 changed files with 1732 additions and 69 deletions

View File

@ -91,7 +91,7 @@
### 60+ Databases, One Tool
MySQL, PostgreSQL, SQLite, Redis, MongoDB, DuckDB, ClickHouse, SQL Server, Oracle, Elasticsearch, Qdrant, Milvus, Weaviate, MariaDB, TiDB, OceanBase, openGauss, GaussDB, KWDB, KingBase, Vastbase, GoldenDB, Doris, SelectDB, StarRocks, Manticore Search, Redshift, DM, TDengine, XuguDB, CockroachDB, Access, HighGo, and more. Agent/JDBC-oriented profiles extend DBX to H2, Snowflake, Trino, PrestoSQL, Hive, DB2, Informix, Neo4j, Cassandra, BigQuery, Kylin, SunDB, and custom JDBC connections. New native and agent-driven drivers also cover Databricks, SAP HANA, Teradata, Vertica, Firebird, Exasol, YashanDB, GBase 8a/8s, Databend, RQLite, Turso, InfluxDB, QuestDB, IoTDB, etcd, ZooKeeper, Nacos, IRIS, and more. Message queue admin is also available for Pulsar, Kafka, and RocketMQ. All in a single ~20 MB app. No bundled Chromium.
MySQL, PostgreSQL, SQLite, Cloudflare D1, Redis, MongoDB, DuckDB, ClickHouse, SQL Server, Oracle, Elasticsearch, Qdrant, Milvus, Weaviate, MariaDB, TiDB, OceanBase, openGauss, GaussDB, KWDB, KingBase, Vastbase, GoldenDB, Doris, SelectDB, StarRocks, Manticore Search, Redshift, DM, TDengine, XuguDB, CockroachDB, Access, HighGo, and more. Agent/JDBC-oriented profiles extend DBX to H2, Snowflake, Trino, PrestoSQL, Hive, DB2, Informix, Neo4j, Cassandra, BigQuery, Kylin, SunDB, and custom JDBC connections. New native and agent-driven drivers also cover Databricks, SAP HANA, Teradata, Vertica, Firebird, Exasol, YashanDB, GBase 8a/8s, Databend, RQLite, Turso, InfluxDB, QuestDB, IoTDB, etcd, ZooKeeper, Nacos, IRIS, and more. Message queue admin is also available for Pulsar, Kafka, and RocketMQ. All in a single ~20 MB app. No bundled Chromium.
### Query Editor
@ -404,7 +404,7 @@ DBX is 20 MB with no runtime dependencies (no Java, no Python). It includes AI a
<details>
<summary><strong>What databases are supported?</strong></summary>
MySQL, PostgreSQL, SQLite, Redis, MongoDB, DuckDB, ClickHouse, SQL Server, Oracle, Elasticsearch, Qdrant, Milvus, Weaviate, MariaDB, TiDB, OceanBase, openGauss, GaussDB, KWDB, KingBase, Vastbase, GoldenDB, Doris, SelectDB, StarRocks, Manticore Search, Redshift, DM, TDengine, XuguDB, CockroachDB, Access, HighGo, and more. Agent/JDBC-oriented profiles extend support to H2, Snowflake, Trino, PrestoSQL, Hive, DB2, Informix, Neo4j, Cassandra, BigQuery, Kylin, SunDB, Databricks, SAP HANA, Teradata, Vertica, Firebird, Exasol, YashanDB, GBase 8a/8s, Databend, RQLite, Turso, InfluxDB, QuestDB, IoTDB, etcd, ZooKeeper, Nacos, IRIS, and custom JDBC connections. Message queue admin (Pulsar, Kafka, RocketMQ) is also supported.
MySQL, PostgreSQL, SQLite, Cloudflare D1, Redis, MongoDB, DuckDB, ClickHouse, SQL Server, Oracle, Elasticsearch, Qdrant, Milvus, Weaviate, MariaDB, TiDB, OceanBase, openGauss, GaussDB, KWDB, KingBase, Vastbase, GoldenDB, Doris, SelectDB, StarRocks, Manticore Search, Redshift, DM, TDengine, XuguDB, CockroachDB, Access, HighGo, and more. Agent/JDBC-oriented profiles extend support to H2, Snowflake, Trino, PrestoSQL, Hive, DB2, Informix, Neo4j, Cassandra, BigQuery, Kylin, SunDB, Databricks, SAP HANA, Teradata, Vertica, Firebird, Exasol, YashanDB, GBase 8a/8s, Databend, RQLite, Turso, InfluxDB, QuestDB, IoTDB, etcd, ZooKeeper, Nacos, IRIS, and custom JDBC connections. Message queue admin (Pulsar, Kafka, RocketMQ) is also supported.
</details>
<details>

View File

@ -91,7 +91,7 @@
### 60+ 种数据库,一个工具搞定
MySQL、PostgreSQL、SQLite、Redis、MongoDB、DuckDB、ClickHouse、SQL Server、Oracle、Elasticsearch、MariaDB、TiDB、OceanBase、openGauss、GaussDB、KWDB、KingBase、Vastbase、GoldenDB、Doris、SelectDB、StarRocks、Manticore Search、Redshift、DM、TDengine、虚谷 XuguDB、CockroachDB、Access、HighGo 等数据库都能直接连接。Agent/JDBC 方向的配置还可扩展到 H2、Snowflake、Trino、Hive、DB2、Informix、Neo4j、Cassandra、BigQuery、Kylin、SunDB 和自定义 JDBC。新增的原生与 Agent 驱动还覆盖了 Databricks、SAP HANA、Teradata、Vertica、Firebird、Exasol、崖山 YashanDB、GBase、Databend、RQLite、Turso、InfluxDB、QuestDB、IoTDB、etcd、IRIS 等。全部装进约 20 MB 的应用里,不内嵌 Chromium。
MySQL、PostgreSQL、SQLite、Cloudflare D1、Redis、MongoDB、DuckDB、ClickHouse、SQL Server、Oracle、Elasticsearch、MariaDB、TiDB、OceanBase、openGauss、GaussDB、KWDB、KingBase、Vastbase、GoldenDB、Doris、SelectDB、StarRocks、Manticore Search、Redshift、DM、TDengine、虚谷 XuguDB、CockroachDB、Access、HighGo 等数据库都能直接连接。Agent/JDBC 方向的配置还可扩展到 H2、Snowflake、Trino、Hive、DB2、Informix、Neo4j、Cassandra、BigQuery、Kylin、SunDB 和自定义 JDBC。新增的原生与 Agent 驱动还覆盖了 Databricks、SAP HANA、Teradata、Vertica、Firebird、Exasol、崖山 YashanDB、GBase、Databend、RQLite、Turso、InfluxDB、QuestDB、IoTDB、etcd、IRIS 等。全部装进约 20 MB 的应用里,不内嵌 Chromium。
### 查询编辑器
@ -398,7 +398,7 @@ DBX 仅 20 MB无需运行时依赖无需 Java、无需 Python。AI 和
<details>
<summary><strong>支持哪些数据库?</strong></summary>
MySQL、PostgreSQL、SQLite、Redis、MongoDB、DuckDB、ClickHouse、SQL Server、Oracle、Elasticsearch、Qdrant、Milvus、Weaviate、MariaDB、TiDB、OceanBase、openGauss、GaussDB、KWDB、KingBase、Vastbase、GoldenDB、Doris、SelectDB、StarRocks、Manticore Search、Redshift、DM、TDengine、虚谷 XuguDB、CockroachDB、Access、HighGo 等。JDBC 方向配置可扩展到 H2、Snowflake、Trino、PrestoSQL、Hive、DB2、Informix、Neo4j、Cassandra、BigQuery、Kylin、SunDB、Databricks、SAP HANA、Teradata、Vertica、Firebird、Exasol、崖山 YashanDB、GBase 8a/8s、Databend、RQLite、Turso、InfluxDB、QuestDB、IoTDB、etcd、ZooKeeper、Nacos、IRIS 及自定义 JDBC 连接并支持消息队列管理Pulsar、Kafka、RocketMQ
MySQL、PostgreSQL、SQLite、Cloudflare D1、Redis、MongoDB、DuckDB、ClickHouse、SQL Server、Oracle、Elasticsearch、Qdrant、Milvus、Weaviate、MariaDB、TiDB、OceanBase、openGauss、GaussDB、KWDB、KingBase、Vastbase、GoldenDB、Doris、SelectDB、StarRocks、Manticore Search、Redshift、DM、TDengine、虚谷 XuguDB、CockroachDB、Access、HighGo 等。JDBC 方向配置可扩展到 H2、Snowflake、Trino、PrestoSQL、Hive、DB2、Informix、Neo4j、Cassandra、BigQuery、Kylin、SunDB、Databricks、SAP HANA、Teradata、Vertica、Firebird、Exasol、崖山 YashanDB、GBase 8a/8s、Databend、RQLite、Turso、InfluxDB、QuestDB、IoTDB、etcd、ZooKeeper、Nacos、IRIS 及自定义 JDBC 连接并支持消息队列管理Pulsar、Kafka、RocketMQ
</details>
<details>

View File

@ -0,0 +1,4 @@
<svg xmlns="http://www.w3.org/2000/svg" viewBox="54 3 50 23">
<path fill="#f48120" d="M88.1 24c.3-1 .2-2-.3-2.6-.5-.6-1.2-1-2.1-1.1l-17.4-.2c-.1 0-.2-.1-.3-.1-.1-.1-.1-.2 0-.3.1-.2.2-.3.4-.3l17.5-.2c2.1-.1 4.3-1.8 5.1-3.8l1-2.6c0-.1.1-.2 0-.3-1.1-5.1-5.7-8.9-11.1-8.9-5 0-9.3 3.2-10.8 7.7-1-.7-2.2-1.1-3.6-1-2.4.2-4.3 2.2-4.6 4.6-.1.6 0 1.2.1 1.8-3.9.1-7.1 3.3-7.1 7.3 0 .4 0 .7.1 1.1 0 .2.2.3.3.3h32.1c.2 0 .4-.1.4-.3l.3-1.1z" />
<path fill="#faad3f" d="M93.6 12.8h-.5c-.1 0-.2.1-.3.2l-.7 2.4c-.3 1-.2 2 .3 2.6.5.6 1.2 1 2.1 1.1l3.7.2c.1 0 .2.1.3.1.1.1.1.2 0 .3-.1.2-.2.3-.4.3l-3.8.2c-2.1.1-4.3 1.8-5.1 3.8l-.2.9c-.1.1 0 .3.2.3h13.2c.2 0 .3-.1.3-.3.2-.8.4-1.7.4-2.6 0-5.2-4.3-9.5-9.5-9.5" />
</svg>

After

Width:  |  Height:  |  Size: 704 B

View File

@ -0,0 +1,33 @@
<script setup lang="ts">
import { useI18n } from "vue-i18n";
import { Input } from "@/components/ui/input";
import { Label } from "@/components/ui/label";
import PasswordInput from "@/components/ui/PasswordInput.vue";
const accountId = defineModel<string>("accountId", { required: true });
const databaseId = defineModel<string | undefined>("databaseId");
const apiToken = defineModel<string>("apiToken", { required: true });
const { t } = useI18n();
</script>
<template>
<div class="grid grid-cols-4 items-center gap-4">
<Label class="justify-self-start text-left">{{ t("connection.d1AccountId") }}</Label>
<Input v-model="accountId" class="col-span-3" placeholder="023e105f4ecef8ad9ca31a8372d0c353" />
</div>
<div class="grid grid-cols-4 items-center gap-4">
<Label class="justify-self-start text-left">{{ t("connection.d1DatabaseId") }}</Label>
<Input v-model="databaseId" class="col-span-3" placeholder="xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx" />
</div>
<div class="grid grid-cols-4 items-center gap-4">
<Label class="justify-self-start text-left">{{ t("connection.d1ApiToken") }}</Label>
<PasswordInput v-model="apiToken" class="col-span-3" placeholder="Cloudflare API Token" />
</div>
<div class="grid grid-cols-4 items-start gap-4">
<span />
<p class="col-span-3 text-xs text-muted-foreground">{{ t("connection.d1TokenHint") }}</p>
</div>
</template>

View File

@ -50,9 +50,11 @@ import { buildDraftVisibleDatabasesConnectionId, connectionCanChooseVisibleDatab
import { canSaveVisibleDatabaseSelection, connectionUsesVisibleSchemaFilter, filterDatabaseNamesForVisiblePicker, isSystemDatabaseName, normalizeVisibleDatabaseSelection, buildDraftVisibleSchemasConnectionId, normalizeVisibleSchemaSelection } from "@/lib/database/visibleDatabases";
import { isSchemaAware, isSingleDatabase } from "@/lib/database/databaseFeatureSupport";
import VisibleSchemasDialog from "@/components/sidebar/VisibleSchemasDialog.vue";
import CloudflareD1ConnectionFields from "@/components/connection/CloudflareD1ConnectionFields.vue";
import { oceanbaseModeConnectionPatch, oceanbaseSubModeFromConfig } from "@/lib/database/oceanbaseConnectionMode";
import { translateBackendError } from "@/i18n/backend-errors";
import { applyHiveKerberosSubmitConfig, hiveKerberosFormConfig, type HiveKerberosAuthMode } from "@/lib/database/hiveKerberosOptions";
import { hasCloudflareD1Credentials, isCloudflareD1Connection, normalizeCloudflareD1Connection } from "@/lib/connection/cloudflareD1";
type DbOption = { value: string; label: string };
type DbCategory = { key: string; title: string; options: DbOption[] };
@ -547,6 +549,7 @@ const driverProfiles: Record<
sqlite: { type: "sqlite", port: 0, user: "", label: "SQLite", icon: "sqlite" },
rqlite: { type: "rqlite", port: 4001, user: "", label: "RQLite", icon: "rqlite" },
turso: { type: "turso", port: 443, user: "", label: "Turso", icon: "turso" },
"cloudflare-d1": { type: "cloudflare-d1", port: 443, user: "", label: "Cloudflare D1", icon: "cloudflare-d1" },
duckdb: { type: "duckdb", port: 0, user: "", label: "DuckDB", icon: "duckdb" },
access: { type: "access", port: 0, user: "", label: "Microsoft Access", icon: "access" },
mongodb: { type: "mongodb", port: 27017, user: "", label: "MongoDB", icon: "mongodb" },
@ -1635,6 +1638,7 @@ const iconTypeMap: Record<string, string> = {
sqlite: "sqlite",
rqlite: "rqlite",
turso: "turso",
"cloudflare-d1": "cloudflare-d1",
access: "access",
redis: "redis",
mongodb: "mongodb",
@ -1724,6 +1728,7 @@ const dbOptions: DbOption[] = [
{ value: "dm", label: "DM (Dameng)" },
{ value: "opengauss", label: "openGauss" },
{ value: "turso", label: "Turso" },
{ value: "cloudflare-d1", label: "Cloudflare D1" },
{ value: "duckdb", label: "DuckDB" },
{ value: "rqlite", label: "RQLite" },
{ value: "access", label: "Microsoft Access" },
@ -1915,7 +1920,7 @@ const zookeeperConnectString = computed({
form.value.connection_string = normalizeZooKeeperConnectString(value);
},
});
const canUseTransportLayers = computed(() => form.value.db_type !== "sqlite" && form.value.db_type !== "access" && !isH2FileMode.value);
const canUseTransportLayers = computed(() => form.value.db_type !== "sqlite" && form.value.db_type !== "access" && !isCloudflareD1Connection(form.value) && !isH2FileMode.value);
const shouldShowAgentDriverInstallHint = computed(() => showAgentDriverInstallHint(form.value.db_type, agentDrivers.value, form.value.driver_profile));
const h2DriverMissing = computed(() => form.value.db_type === "h2" && isH2FileMode.value && agentDrivers.value.find((d) => d.db_type === "h2")?.installed !== true);
const canChooseVisibleDatabases = computed(() => connectionCanChooseVisibleDatabases(form.value));
@ -2040,6 +2045,7 @@ const hasRequiredConnectionTarget = computed(() => {
}
if (form.value.db_type === "zookeeper") return !!(form.value.host || form.value.connection_string || connectionUrlInput.value.trim());
if (form.value.db_type === "nacos") return !!nacosServerAddr.value.trim();
if (isCloudflareD1Connection(form.value)) return hasCloudflareD1Credentials(form.value);
if (isH2FileMode.value) return !!(form.value.host.trim() || h2FilePathFromJdbcUrl(form.value.connection_string));
return !!(form.value.host || (mongoUseUrl.value && form.value.connection_string) || (form.value.db_type === "jdbc" && form.value.connection_string) || connectionUrlInput.value.trim());
});
@ -2290,6 +2296,12 @@ function connectionConfigForSubmit(id: string): ConnectionConfig {
throw new Error(t("connection.kingbaseDatabaseRequired"));
}
}
if (isCloudflareD1Connection(config)) {
normalizeCloudflareD1Connection(config);
if (!hasCloudflareD1Credentials(config)) {
throw new Error(t("connection.d1FieldsRequired"));
}
}
config.transport_layers = (config.transport_layers || []).map(normalizeTransportLayer);
config.transport_layers = config.transport_layers.map((layer) => {
if (layer.type !== "ssh") return layer;
@ -4744,6 +4756,10 @@ function openExternalUrl(url: string) {
</div>
</template>
<template v-else-if="form.db_type === 'cloudflare-d1'">
<CloudflareD1ConnectionFields v-model:account-id="form.host" v-model:database-id="form.database" v-model:api-token="form.password" />
</template>
<!-- MySQL / PostgreSQL: host, port, user, password, database -->
<template v-else>
<div class="grid grid-cols-4 items-center gap-4">

View File

@ -14,6 +14,7 @@ const assetIcons: Record<string, string> = {
sqlite: "sqlite",
rqlite: "rqlite.png",
turso: "turso.png",
cloudflare_d1: "cloudflare-d1",
redis: "redis",
mongodb: "mongodb",
mongodb_legacy: "mongodb",

View File

@ -474,6 +474,8 @@ function connectionTooltipScheme(config: Pick<ConnectionConfig, "db_type" | "ssl
case "turso":
case "mq":
return config.ssl ? "https" : "http";
case "cloudflare-d1":
return "https";
case "dameng":
return "dm";
default:
@ -4825,7 +4827,7 @@ function treeItemMenuItems(): ContextMenuItem[] {
items.push({ label: t("contextMenu.newQuery"), action: newQuery, icon: TerminalSquare });
const sqlHistoryMenu = savedSqlHistorySubmenu();
if (sqlHistoryMenu) items.push(sqlHistoryMenu);
if (node.type === "database") {
if (node.type === "database" && currentDatabaseType() !== "cloudflare-d1") {
if (!isNodeDefaultDatabase.value) {
items.push({ label: t("contextMenu.setDefaultDatabase"), action: setNodeAsDefaultDatabase, icon: Database });
} else {

View File

@ -198,6 +198,11 @@ export default {
driverName: "Driver Name",
driverNamePlaceholder: "Vendor or environment name",
urlParams: "URL Params",
d1AccountId: "Account ID",
d1DatabaseId: "Database ID",
d1ApiToken: "API Token",
d1TokenHint: "The token needs D1 Read permission; writes, DDL, and imports also require D1 Write.",
d1FieldsRequired: "Cloudflare Account ID, D1 Database ID, and API Token are required.",
hiveAuthMode: "Auth",
hiveAuthNone: "None",
hivePrincipal: "Principal",

View File

@ -200,6 +200,11 @@ export default withEnglishFallback({
driverName: "Nombre del driver",
driverNamePlaceholder: "Nombre del proveedor o entorno",
urlParams: "Parámetros de URL",
d1AccountId: "ID de cuenta",
d1DatabaseId: "ID de base de datos",
d1ApiToken: "Token de API",
d1TokenHint: "El token necesita al menos el permiso D1 Read; las escrituras, el DDL y las importaciones también requieren D1 Write.",
d1FieldsRequired: "Se requieren el ID de cuenta de Cloudflare, el ID de la base de datos D1 y el token de API.",
gbaseServer: "GBASEDBTSERVER",
informixServer: "INFORMIXSERVER",
sslEnable: "Habilitar conexión cifrada",

View File

@ -199,6 +199,11 @@ export default withEnglishFallback({
driverName: "Nome Driver",
driverNamePlaceholder: "Nome del fornitore o dell'ambiente",
urlParams: "Parametri URL",
d1AccountId: "ID account",
d1DatabaseId: "ID database",
d1ApiToken: "Token API",
d1TokenHint: "Il token richiede almeno l'autorizzazione D1 Read; per scritture, DDL e importazioni è necessaria anche D1 Write.",
d1FieldsRequired: "Sono obbligatori l'ID account Cloudflare, l'ID database D1 e il token API.",
gbaseServer: "GBASEDBTSERVER",
informixServer: "INFORMIXSERVER",
sslEnable: "Abilita connessione crittografata",

View File

@ -199,6 +199,11 @@ export default withEnglishFallback({
driverName: "ドライバー名",
driverNamePlaceholder: "ベンダー名または環境名",
urlParams: "URLパラメータ",
d1AccountId: "アカウントID",
d1DatabaseId: "データベースID",
d1ApiToken: "APIトークン",
d1TokenHint: "トークンには少なくともD1 Read権限が必要です。書き込み、DDL、インポートにはD1 Write権限も必要です。",
d1FieldsRequired: "CloudflareアカウントID、D1データベースID、APIトークンは必須です。",
gbaseServer: "GBASEDBTSERVER",
informixServer: "INFORMIXSERVER",
sslEnable: "暗号化接続を有効にする",

View File

@ -200,6 +200,11 @@ export default withEnglishFallback({
driverName: "Nome do Driver",
driverNamePlaceholder: "Nome do fornecedor ou ambiente",
urlParams: "Parâmetros da URL",
d1AccountId: "ID da conta",
d1DatabaseId: "ID do banco de dados",
d1ApiToken: "Token da API",
d1TokenHint: "O token precisa de pelo menos a permissão D1 Read; gravações, DDL e importações também exigem D1 Write.",
d1FieldsRequired: "O ID da conta Cloudflare, o ID do banco de dados D1 e o token da API são obrigatórios.",
gbaseServer: "GBASEDBTSERVER",
informixServer: "INFORMIXSERVER",
sslEnable: "Habilitar conexão criptografada",

View File

@ -200,6 +200,11 @@ export default withEnglishFallback({
driverName: "驱动名称",
driverNamePlaceholder: "厂商或环境名称",
urlParams: "URL 参数",
d1AccountId: "Account ID",
d1DatabaseId: "Database ID",
d1ApiToken: "API Token",
d1TokenHint: "Token 至少需要 D1 Read 权限执行写入、DDL 或导入时还需要 D1 Write 权限。",
d1FieldsRequired: "必须填写 Cloudflare Account ID、D1 Database ID 和 API Token。",
hiveAuthMode: "认证",
hiveAuthNone: "无",
hivePrincipal: "Principal",

View File

@ -200,6 +200,11 @@ export default withEnglishFallback({
driverName: "驅動程式名稱",
driverNamePlaceholder: "廠商或環境名稱",
urlParams: "URL 參數",
d1AccountId: "帳戶 ID",
d1DatabaseId: "資料庫 ID",
d1ApiToken: "API Token",
d1TokenHint: "Token 至少需要 D1 Read 權限執行寫入、DDL 或匯入時還需要 D1 Write 權限。",
d1FieldsRequired: "必須填寫 Cloudflare Account ID、D1 Database ID 和 API Token。",
gbaseServer: "GBASEDBTSERVER",
informixServer: "INFORMIXSERVER",
sslEnable: "啟用加密連線",

View File

@ -16,6 +16,7 @@ describe("supportsTransaction", () => {
expect(supportsTransaction("duckdb")).toBe(false);
expect(supportsTransaction("qdrant")).toBe(false);
expect(supportsTransaction("turso")).toBe(false);
expect(supportsTransaction("cloudflare-d1")).toBe(false);
expect(supportsTransaction("sqlite")).toBe(false);
expect(supportsTransaction("clickhouse")).toBe(false);
expect(supportsTransaction("sqlserver")).toBe(false);

View File

@ -9,7 +9,7 @@ describe("sqlFormatter", () => {
});
it("maps SQLite-compatible database types to the sqlite formatter dialect", () => {
for (const dbType of ["sqlite", "rqlite", "turso"]) {
for (const dbType of ["sqlite", "rqlite", "turso", "cloudflare-d1"]) {
expect(sqlFormatDialectForDbType(dbType)).toBe("sqlite");
}
});

View File

@ -0,0 +1,23 @@
import type { ConnectionConfig } from "@/types/database";
type CloudflareD1Config = Pick<ConnectionConfig, "db_type" | "host" | "database" | "password">;
type MutableCloudflareD1Config = CloudflareD1Config & Pick<ConnectionConfig, "port" | "username" | "ssl" | "url_params" | "transport_layers">;
export function isCloudflareD1Connection(config: Pick<ConnectionConfig, "db_type">): boolean {
return config.db_type === "cloudflare-d1";
}
export function hasCloudflareD1Credentials(config: CloudflareD1Config): boolean {
return !!config.host.trim() && !!config.database?.trim() && !!config.password.trim();
}
export function normalizeCloudflareD1Connection(config: MutableCloudflareD1Config): void {
config.host = config.host.trim();
config.database = config.database?.trim() || undefined;
config.password = config.password.trim();
config.port = 443;
config.username = "";
config.ssl = true;
config.url_params = "";
config.transport_layers = [];
}

View File

@ -17,6 +17,7 @@ export function connectionDriverLabel(connection?: Pick<ConnectionConfig, "db_ty
export function connectionEndpointLabel(connection?: ConnectionPresentationConfig): string {
if (!connection) return "";
if (connection.db_type === "cloudflare-d1") return [connection.host, connection.database].filter(Boolean).join("/");
if (LOCAL_DATABASE_TYPES.has(connection.db_type) || (connection.db_type === "h2" && connection.port === 0)) {
return connection.host || connection.database || "local";
}
@ -45,6 +46,7 @@ function redactConnectionHost(host: string): string {
export function connectionRedactedEndpointLabel(connection?: ConnectionPresentationConfig): string {
if (!connection) return "";
if (connection.db_type === "cloudflare-d1") return `${REDACTED_HOST_SEGMENT}/${REDACTED_HOST_SEGMENT}`;
if (LOCAL_DATABASE_TYPES.has(connection.db_type) || (connection.db_type === "h2" && connection.port === 0)) {
return connectionEndpointLabel(connection);
}
@ -111,6 +113,9 @@ export function connectionUrlPlaceholder(dbType: DatabaseType): string {
case "turso":
return "https://[your-db]-[org].turso.io";
case "cloudflare-d1":
return "https://api.cloudflare.com/client/v4/accounts/{account_id}/d1/database/{database_id}";
case "duckdb":
return "duckdb:///absolute/path/to/database.duckdb";

View File

@ -3,7 +3,7 @@ import { filterDatabaseNamesForVisiblePicker, normalizeVisibleDatabaseSelection
const DRAFT_VISIBLE_DATABASES_PREFIX = "__visible_draft_";
const UNSUPPORTED_VISIBLE_DATABASE_TYPES = new Set<DatabaseType>(["elasticsearch", "qdrant", "milvus", "weaviate", "chromadb", "etcd", "zookeeper"]);
const UNSUPPORTED_VISIBLE_DATABASE_TYPES = new Set<DatabaseType>(["cloudflare-d1", "elasticsearch", "qdrant", "milvus", "weaviate", "chromadb", "etcd", "zookeeper"]);
type VisibleDatabaseConnectionFields = Pick<
ConnectionConfig,

View File

@ -75,4 +75,4 @@ export const DATABASE_OBJECT_TREE_TYPES = new Set<DatabaseType>(["jdbc"]);
export const PG_LIKE_STRUCTURE_TYPES = new Set<DatabaseType>(["postgres", "redshift", "gaussdb", "kwdb", "opengauss", "questdb"]);
export const DIAGRAM_SQL_TYPES = new Set<DatabaseType>(["mysql", "postgres", "sqlite", "rqlite", "turso", "sqlserver", "oracle", "redshift", "dameng", "gaussdb", "kwdb", "opengauss", "questdb", "oceanbase-oracle"]);
export const DIAGRAM_SQL_TYPES = new Set<DatabaseType>(["mysql", "postgres", "sqlite", "rqlite", "turso", "cloudflare-d1", "sqlserver", "oracle", "redshift", "dameng", "gaussdb", "kwdb", "opengauss", "questdb", "oceanbase-oracle"]);

View File

@ -118,7 +118,7 @@ export function supportsObjectBrowserTreeNode(dbType: DatabaseType | undefined,
}
export function supportsTableTruncate(dbType?: DatabaseType): boolean {
return !!dbType && dbType !== "sqlite" && dbType !== "rqlite" && dbType !== "turso" && dbType !== "duckdb" && dbType !== "influxdb" && dbType !== "manticoresearch";
return !!dbType && dbType !== "sqlite" && dbType !== "rqlite" && dbType !== "turso" && dbType !== "cloudflare-d1" && dbType !== "duckdb" && dbType !== "influxdb" && dbType !== "manticoresearch";
}
export function usesPostgresLikeStructureCopy(dbType?: DatabaseType): boolean {

View File

@ -20,6 +20,7 @@ export const DATABASE_NAMESPACE_CREATION_MATRIX = {
sqlite: { deferred: "file-backed; create a new connection/file instead" },
rqlite: { deferred: "single SQLite-compatible database per node" },
turso: { deferred: "remote libSQL database lifecycle is provider-managed" },
"cloudflare-d1": { deferred: "Cloudflare D1 database lifecycle is provider-managed" },
redis: { deferred: "numbered logical databases are server-configured" },
duckdb: { connection: "attach" },
clickhouse: { connection: "database" },

View File

@ -35,6 +35,7 @@ const DATABASE_TYPE_OBJECTS = new Map<DatabaseType, SidebarObjectKind[]>([
["sqlite", TABLE_VIEW_OBJECTS],
["rqlite", TABLE_VIEW_OBJECTS],
["turso", TABLE_VIEW_OBJECTS],
["cloudflare-d1", TABLE_VIEW_OBJECTS],
["duckdb", TABLE_VIEW_OBJECTS],
["clickhouse", TABLE_VIEW_OBJECTS],
["doris", TABLE_VIEW_OBJECTS],

View File

@ -21,6 +21,7 @@ export const DATABASE_PROPERTY_EDITING_MATRIX = {
sqlite: { deferred: "file-backed database properties are not edited in-place" },
rqlite: { deferred: "single SQLite-compatible database per node" },
turso: { deferred: "remote libSQL database lifecycle is provider-managed" },
"cloudflare-d1": { deferred: "Cloudflare D1 database lifecycle is provider-managed" },
redis: { deferred: "numbered logical databases are server-configured" },
duckdb: { deferred: "attached database file properties need a dedicated DuckDB workflow" },
clickhouse: { deferred: "database property editing not verified for first pass" },

View File

@ -52,6 +52,7 @@ const NAVICAT_STYLE_TABLE_DATA_TYPES = new Set<DatabaseType>([
"sqlite",
"rqlite",
"turso",
"cloudflare-d1",
"duckdb",
"sqlserver",
"oracle",

View File

@ -4,7 +4,8 @@ import { usesTreeSchemaMode } from "@/lib/database/databaseCapabilities";
export const TREE_SCHEMA_DEFAULT_DATABASE_SELECT_VALUE = "__dbx_tree_schema_default_database__";
export const EMPTY_DATABASE_SELECT_VALUE = "__dbx_empty_database__";
export function resolveDefaultDatabase(connection: Pick<ConnectionConfig, "database">, options: string[]): string {
export function resolveDefaultDatabase(connection: Pick<ConnectionConfig, "database"> & Partial<Pick<ConnectionConfig, "db_type">>, options: string[]): string {
if (connection.db_type === "cloudflare-d1") return "main";
return connection.database || options[0] || "";
}
@ -29,6 +30,7 @@ export function formatDatabaseLabel(connection: Pick<ConnectionConfig, "db_type"
return database || labels.noDatabase;
}
export function isDefaultDatabase(connection: Pick<ConnectionConfig, "database"> | undefined, database: string): boolean {
export function isDefaultDatabase(connection: (Pick<ConnectionConfig, "database"> & Partial<Pick<ConnectionConfig, "db_type">>) | undefined, database: string): boolean {
if (connection?.db_type === "cloudflare-d1") return database === "main";
return !!connection?.database && !!database && connection.database === database;
}

View File

@ -0,0 +1 @@
export const CLOUDFLARE_D1_COMMON_FUNCTION_NAMES = new Set(["ABS", "AVG", "CAST", "CEIL", "COALESCE", "CONCAT", "COUNT", "FLOOR", "LENGTH", "LOWER", "LTRIM", "MAX", "MIN", "MOD", "NULLIF", "REPLACE", "ROUND", "RTRIM", "SIGN", "SQRT", "SUBSTR", "SUBSTRING", "SUM", "TRIM", "UPPER"]);

View File

@ -155,6 +155,7 @@ export function sqlSemanticDialectFor(options: { databaseType?: DatabaseType; di
case "sqlite":
case "rqlite":
case "turso":
case "cloudflare-d1":
return SQL_SEMANTIC_DIALECTS.sqlite;
case "duckdb":
return SQL_SEMANTIC_DIALECTS.duckdb;

View File

@ -1,6 +1,7 @@
import { Cassandra, MariaSQL, MSSQL, MySQL, PLSQL, PostgreSQL, SQLite, StandardSQL } from "@codemirror/lang-sql";
import type { DatabaseType, SqlSnippet } from "@/types/database";
import { buildMongoCompletionItemsFromContext, type MongoCompletionItem } from "@/lib/mongo/mongoCompletion";
import { CLOUDFLARE_D1_COMMON_FUNCTION_NAMES } from "@/lib/sql/cloudflareD1";
const SQL_KEYWORDS = [
"SELECT",
@ -470,6 +471,7 @@ const DATABASE_SQL_KEYWORDS: Partial<Record<DatabaseType, string[]>> = {
sqlite: SQLITE_SQL_KEYWORDS,
rqlite: SQLITE_SQL_KEYWORDS,
turso: SQLITE_SQL_KEYWORDS,
"cloudflare-d1": SQLITE_SQL_KEYWORDS,
sqlserver: SQLSERVER_SQL_KEYWORDS,
manticoresearch: MANTICORESEARCH_SQL_KEYWORDS,
};
@ -896,6 +898,8 @@ const SQLITE_FUNCTION_SIGNATURES = new Map<string, string[]>([
["NOW", []],
]);
const CLOUDFLARE_D1_FUNCTION_SIGNATURES = new Map(Array.from(SQLITE_FUNCTION_SIGNATURES.entries()).filter(([name]) => name !== "NOW"));
const SQLSERVER_FUNCTION_SIGNATURES = new Map<string, string[]>([
["TRY_CAST", ["expression AS type"]],
["TRY_CONVERT", ["type", "expression"]],
@ -958,6 +962,7 @@ const DATABASE_FUNCTION_SIGNATURES: Partial<Record<DatabaseType, Map<string, str
sqlite: SQLITE_FUNCTION_SIGNATURES,
rqlite: SQLITE_FUNCTION_SIGNATURES,
turso: SQLITE_FUNCTION_SIGNATURES,
"cloudflare-d1": CLOUDFLARE_D1_FUNCTION_SIGNATURES,
sqlserver: SQLSERVER_FUNCTION_SIGNATURES,
manticoresearch: MANTICORESEARCH_FUNCTION_SIGNATURES,
};
@ -3727,7 +3732,8 @@ function buildSnippetItems(prefix: string, snippets: SqlSnippet[], keywordCase?:
}
function activeFunctionSignatures(databaseType?: DatabaseType): Map<string, string[]> {
const signatures = databaseType ? new Map(Array.from(SQL_FUNCTION_SIGNATURES.entries()).filter(([name]) => COMMON_SQL_FUNCTION_NAMES.has(name))) : new Map(SQL_FUNCTION_SIGNATURES);
const commonFunctionNames = databaseType === "cloudflare-d1" ? CLOUDFLARE_D1_COMMON_FUNCTION_NAMES : COMMON_SQL_FUNCTION_NAMES;
const signatures = databaseType ? new Map(Array.from(SQL_FUNCTION_SIGNATURES.entries()).filter(([name]) => commonFunctionNames.has(name))) : new Map(SQL_FUNCTION_SIGNATURES);
const databaseSignatures = databaseType ? DATABASE_FUNCTION_SIGNATURES[databaseType] : undefined;
if (databaseSignatures) {
for (const [name, parameters] of databaseSignatures) signatures.set(name, parameters);

View File

@ -30,6 +30,7 @@ export function sqlFormatDialectForDbType(dbType: string | null | undefined): Sq
case "sqlite":
case "rqlite":
case "turso":
case "cloudflare-d1":
return "sqlite";
case "sqlserver":
return "sqlserver";

View File

@ -21,7 +21,7 @@ export function supportsObjectRename(databaseType: DatabaseType | undefined, obj
if (objectType === "PROCEDURE" || objectType === "FUNCTION") {
return false;
}
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso") return objectType === "TABLE";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso" || databaseType === "cloudflare-d1") return objectType === "TABLE";
if (databaseType === "mysql" || databaseType === "goldendb") return objectType === "TABLE" || objectType === "VIEW";
if (postgresLikeRenameTypes.has(databaseType)) return objectType === "TABLE" || objectType === "VIEW" || objectType === "MATERIALIZED_VIEW";
if (oracleLikeRenameTypes.has(databaseType)) return objectType === "TABLE" || objectType === "VIEW" || objectType === "MATERIALIZED_VIEW";

View File

@ -151,29 +151,29 @@ export function importDataTypeForDatabase(inferredType: ImportInferredType, data
case "boolean":
if (["mysql", "doris", "starrocks", "goldendb", "sundb", "databend"].includes(databaseType || "")) return "TINYINT(1)";
if (databaseType === "sqlserver") return "BIT";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso") return "INTEGER";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso" || databaseType === "cloudflare-d1") return "INTEGER";
if (databaseType === "oracle" || databaseType === "oceanbase-oracle" || databaseType === "dameng") return "NUMBER(1)";
if (databaseType === "clickhouse") return "UInt8";
return "BOOLEAN";
case "integer":
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso") return "INTEGER";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso" || databaseType === "cloudflare-d1") return "INTEGER";
if (databaseType === "oracle" || databaseType === "oceanbase-oracle" || databaseType === "dameng") return "NUMBER(19)";
if (databaseType === "clickhouse") return "Int64";
return "BIGINT";
case "decimal":
if (["postgres", "gaussdb", "opengauss", "redshift", "kingbase", "highgo", "kwdb", "vastbase"].includes(databaseType || "")) return "DOUBLE PRECISION";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso") return "REAL";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso" || databaseType === "cloudflare-d1") return "REAL";
if (databaseType === "oracle" || databaseType === "oceanbase-oracle" || databaseType === "dameng") return "BINARY_DOUBLE";
if (databaseType === "clickhouse") return "Float64";
return "DOUBLE";
case "date":
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso") return "TEXT";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso" || databaseType === "cloudflare-d1") return "TEXT";
if (databaseType === "clickhouse") return "Date";
return "DATE";
case "timestamp":
if (["mysql", "doris", "starrocks", "goldendb", "sundb", "databend"].includes(databaseType || "")) return "DATETIME";
if (databaseType === "sqlserver") return "DATETIME2";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso") return "TEXT";
if (databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "turso" || databaseType === "cloudflare-d1") return "TEXT";
if (databaseType === "clickhouse") return "DateTime64";
return "TIMESTAMP";
case "json":

View File

@ -129,7 +129,7 @@ function columnPlaceholderValue(column: ColumnInfo, databaseType?: DatabaseType)
const colName = column.name;
if (typeLooksNumeric(dataType)) return "0";
if (typeLooksBoolean(dataType)) return databaseType === "mysql" || databaseType === "sqlite" || databaseType === "rqlite" ? "1" : "TRUE";
if (typeLooksBoolean(dataType)) return databaseType === "mysql" || databaseType === "sqlite" || databaseType === "rqlite" || databaseType === "cloudflare-d1" ? "1" : "TRUE";
if (typeLooksDateOnly(dataType)) return "'2024-01-01'";
if (typeLooksTimeOnly(dataType)) return "'12:00:00'";
if (typeLooksTimestamp(dataType)) return databaseType === "tdengine" ? "NOW" : "CURRENT_TIMESTAMP";

View File

@ -1692,6 +1692,7 @@ export const useConnectionStore = defineStore("connection", () => {
async function setDefaultDatabase(connectionId: string, database: string) {
const config = getConfig(connectionId);
if (config?.db_type === "cloudflare-d1") return;
if (!config || config.database === database) return;
await updateConnection({
...config,
@ -1701,6 +1702,7 @@ export const useConnectionStore = defineStore("connection", () => {
async function clearDefaultDatabase(connectionId: string) {
const config = getConfig(connectionId);
if (config?.db_type === "cloudflare-d1") return;
if (!config || !config.database) return;
await updateConnection({
...config,
@ -1709,7 +1711,9 @@ export const useConnectionStore = defineStore("connection", () => {
}
function isDefaultDatabase(connectionId: string, database: string): boolean {
return getConfig(connectionId)?.database === database && database !== "";
const config = getConfig(connectionId);
if (config?.db_type === "cloudflare-d1") return database === "main";
return config?.database === database && database !== "";
}
async function setVisibleDatabases(connectionId: string, databaseNames: string[]) {

View File

@ -4,6 +4,7 @@ export type DatabaseType =
| "sqlite"
| "rqlite"
| "turso"
| "cloudflare-d1"
| "redis"
| "duckdb"
| "clickhouse"

View File

@ -116,6 +116,35 @@
"driverManagement": false
}
},
{
"dbType": "cloudflare-d1",
"label": "Cloudflare D1",
"runtimeMode": "native",
"mcpMode": "bridge",
"singleConnectionPool": true,
"metadataConnectionScoped": false,
"skipTcpProbe": true,
"defaultPort": 443,
"supportLevel": "operate",
"capabilities": {
"queryExecution": true,
"metadataBrowse": true,
"objectBrowser": true,
"objectSource": true,
"schemaSearch": true,
"diagram": true,
"tableDataEdit": true,
"tableStructureEdit": false,
"tableImport": true,
"dataTransfer": true,
"sqlFileExecution": true,
"databaseCreate": false,
"fieldLineage": true,
"sqlExplain": false,
"userAdmin": false,
"driverManagement": false
}
},
{
"dbType": "redis",
"label": "Redis",

View File

@ -86,6 +86,7 @@ pub enum PoolKind {
Sqlite(db::sqlite::SqliteHandle),
Rqlite(db::rqlite_driver::RqliteClient),
Turso(db::turso_driver::TursoClient),
CloudflareD1(db::cloudflare_d1_driver::CloudflareD1Client),
Redis(db::redis_driver::RedisConnection),
DuckDb(DuckDbHandle),
DuckDbWorker(DuckDbWorkerHandle),
@ -244,7 +245,11 @@ pub fn database_connection_config(config: &ConnectionConfig, database: Option<&s
if let Some(db) = database {
if !matches!(
db_config.db_type,
DatabaseType::Oracle | DatabaseType::Dameng | DatabaseType::MongoDb | DatabaseType::OceanbaseOracle
DatabaseType::Oracle
| DatabaseType::Dameng
| DatabaseType::MongoDb
| DatabaseType::OceanbaseOracle
| DatabaseType::CloudflareD1
) {
db_config.database = Some(db.to_string());
}
@ -1083,6 +1088,9 @@ impl AppState {
db::turso_driver::test_connection(&client, connect_timeout).await?;
PoolKind::Turso(client)
}
DatabaseType::CloudflareD1 => {
PoolKind::CloudflareD1(db::cloudflare_d1_driver::connect(&db_config, connect_timeout).await?)
}
DatabaseType::Redis => {
let con = if db_config.uses_redis_cluster() {
db::redis_driver::RedisConnection::Cluster(
@ -1848,6 +1856,18 @@ impl AppState {
}
}
}
PoolKind::CloudflareD1(client) => {
let client = client.clone();
drop(connections);
let timeout = crate::db::connection_timeout();
match db::cloudflare_d1_driver::test_connection(&client, timeout).await {
Ok(()) => false,
Err(err) => {
log::warn!("Cloudflare D1 connection pool '{pool_key}' is stale: {err}");
true
}
}
}
PoolKind::Agent(client) => {
let client = client.clone();
drop(connections);
@ -2360,6 +2380,15 @@ impl AppState {
false
}
},
PoolKind::CloudflareD1(client) => {
match db::cloudflare_d1_driver::test_connection(client, timeout).await {
Ok(()) => true,
Err(e) => {
log::warn!("Cloudflare D1 connection pool '{key}' is unhealthy: {e}");
false
}
}
}
PoolKind::Agent(client) => {
let mut agent = client.lock().await;
match agent.test_connection(serde_json::json!({})).await {
@ -2769,12 +2798,13 @@ fn session_scoped_pool_key_for(
client_session_id: Option<&str>,
) -> String {
let shares_base_pool = config.is_some_and(|config| {
config.db_type == DatabaseType::DuckDb
matches!(config.db_type, DatabaseType::DuckDb | DatabaseType::CloudflareD1)
|| (config.db_type == DatabaseType::Sqlite && db::sqlite::is_memory_database_path(&config.host))
});
if shares_base_pool {
// In-memory SQLite databases only exist inside one connection. A session-scoped
// handle would silently point query/data tabs at a different empty database.
// DuckDB and D1 already use connection-scoped handles. In-memory SQLite databases
// only exist inside one connection, so a session-scoped handle would point tabs at
// a different empty database.
return base_pool_key;
}
session_scoped_pool_key(base_pool_key, client_session_id)
@ -2787,6 +2817,7 @@ fn clone_pool_kind(pool: &PoolKind) -> PoolKind {
PoolKind::Sqlite(p) => PoolKind::Sqlite(p.clone()),
PoolKind::Rqlite(client) => PoolKind::Rqlite(client.clone()),
PoolKind::Turso(client) => PoolKind::Turso(client.clone()),
PoolKind::CloudflareD1(client) => PoolKind::CloudflareD1(client.clone()),
#[cfg(feature = "duckdb-bundled")]
PoolKind::DuckDb(con) => PoolKind::DuckDb(con.clone()),
#[cfg(feature = "duckdb-bundled")]
@ -2821,6 +2852,7 @@ pub async fn close_pool_kind(pool: PoolKind) {
PoolKind::Sqlite(_) => {}
PoolKind::Rqlite(_) => {}
PoolKind::Turso(_) => {}
PoolKind::CloudflareD1(_) => {}
PoolKind::Redis(conn) => {
drop(conn);
}
@ -3939,6 +3971,20 @@ mod tests {
assert_eq!(scoped.database.as_deref(), Some("admin"));
}
#[test]
fn cloudflare_d1_query_namespace_does_not_replace_database_id() {
let mut config = mysql_config(Some("database-uuid"));
config.db_type = DatabaseType::CloudflareD1;
let scoped = database_connection_config(&config, Some("main"));
assert_eq!(scoped.database.as_deref(), Some("database-uuid"));
assert_eq!(
super::base_pool_key_for(Some(DatabaseType::CloudflareD1), "d1-conn", Some("main"), false),
"d1-conn"
);
}
#[test]
fn oracle_database_connection_ignores_requested_database() {
let mut config = mysql_config(Some("ORCL"));
@ -4001,7 +4047,6 @@ mod tests {
super::session_scoped_pool_key_for(Some(&duckdb), "duckdb-conn".to_string(), Some("tab-1")),
"duckdb-conn"
);
let mut sqlite_memory = mysql_config(None);
sqlite_memory.db_type = DatabaseType::Sqlite;
sqlite_memory.host = " :MeMoRy: ".to_string();
@ -4015,6 +4060,13 @@ mod tests {
super::session_scoped_pool_key_for(Some(&sqlite_memory), "sqlite-file".to_string(), Some("tab-1")),
"sqlite-file:session:tab-1"
);
let mut cloudflare_d1 = mysql_config(None);
cloudflare_d1.db_type = DatabaseType::CloudflareD1;
assert_eq!(
super::session_scoped_pool_key_for(Some(&cloudflare_d1), "d1-conn".to_string(), Some("tab-1")),
"d1-conn"
);
}
#[test]

View File

@ -2088,6 +2088,7 @@ fn uses_keyless_row_predicate(database_type: Option<DatabaseType>) -> bool {
| DatabaseType::Sqlite
| DatabaseType::Rqlite
| DatabaseType::Turso
| DatabaseType::CloudflareD1
| DatabaseType::DuckDb
| DatabaseType::SqlServer
| DatabaseType::Oracle

View File

@ -16,6 +16,7 @@ pub fn is_single_connection_pool(db_type: &DatabaseType) -> bool {
| DatabaseType::DuckDb
| DatabaseType::Rqlite
| DatabaseType::Turso
| DatabaseType::CloudflareD1
| DatabaseType::MongoDb
| DatabaseType::Oracle
| DatabaseType::Dameng
@ -43,6 +44,7 @@ pub fn skips_tcp_probe(db_type: &DatabaseType) -> bool {
DatabaseType::Sqlite
| DatabaseType::DuckDb
| DatabaseType::Turso
| DatabaseType::CloudflareD1
| DatabaseType::Jdbc
| DatabaseType::MessageQueue
) || is_agent_type(db_type)
@ -75,6 +77,13 @@ mod tests {
assert!(!is_local_file_db_type(&DatabaseType::Redis));
assert!(!is_local_file_db_type(&DatabaseType::MongoDb));
assert!(!is_local_file_db_type(&DatabaseType::Turso));
assert!(!is_local_file_db_type(&DatabaseType::CloudflareD1));
assert!(!is_local_file_db_type(&DatabaseType::Rqlite));
}
#[test]
fn cloudflare_d1_uses_a_single_http_pool_without_tcp_probe() {
assert!(is_single_connection_pool(&DatabaseType::CloudflareD1));
assert!(skips_tcp_probe(&DatabaseType::CloudflareD1));
}
}

View File

@ -0,0 +1,94 @@
use crate::models::connection::DatabaseType;
use crate::table_import::{mapping_indexes_for_columns, ImportSqlBatch, TableImportColumnMapping};
use crate::transfer::generate_insert_typed;
use super::sql_limits::build_sql_batches;
pub(crate) fn build_import_insert_batches(
rows: &[Vec<serde_json::Value>],
source_columns: &[String],
mappings: &[TableImportColumnMapping],
target_column_types: &[(String, String)],
table: &str,
schema: &str,
max_rows: usize,
) -> Result<Vec<ImportSqlBatch>, String> {
let mapped = mapping_indexes_for_columns(source_columns, mappings)?;
let columns = mapped.iter().map(|(_, target)| target.clone()).collect::<Vec<_>>();
let column_types = columns
.iter()
.map(|column| {
target_column_types
.iter()
.find(|(name, _)| name.eq_ignore_ascii_case(column))
.map(|(_, data_type)| data_type.clone())
})
.collect::<Vec<_>>();
let rows = rows
.iter()
.map(|row| {
mapped
.iter()
.map(|(source_index, _)| row.get(*source_index).cloned().unwrap_or(serde_json::Value::Null))
.collect::<Vec<_>>()
})
.collect::<Vec<_>>();
build_sql_batches(rows.len(), max_rows, "import row", |range| {
generate_insert_typed(&columns, &column_types, &rows[range], table, schema, &DatabaseType::CloudflareD1)
})
.map(|batches| {
batches.into_iter().map(|batch| ImportSqlBatch { sql: batch.sql, row_count: batch.item_count }).collect()
})
}
pub(crate) fn build_streaming_import_insert_batch(
rows: &[Vec<serde_json::Value>],
source_columns: &[String],
mappings: &[TableImportColumnMapping],
target_column_types: &[(String, String)],
table: &str,
schema: &str,
max_rows: usize,
) -> Result<Option<ImportSqlBatch>, String> {
let batches =
build_import_insert_batches(rows, source_columns, mappings, target_column_types, table, schema, max_rows)?;
if batches.is_empty() {
return Ok(None);
}
let row_count = batches.iter().map(|batch| batch.row_count).sum();
let sql = batches.into_iter().map(|batch| batch.sql).collect::<Vec<_>>().join(";\n");
Ok(Some(ImportSqlBatch { sql, row_count }))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn streaming_import_joins_individually_size_limited_statements() {
let rows = vec![vec![serde_json::json!("a".repeat(60_000))], vec![serde_json::json!("b".repeat(60_000))]];
let columns = vec!["value".to_string()];
let mappings = vec![TableImportColumnMapping {
source_column: "value".to_string(),
target_column: "value".to_string(),
target_data_type: None,
}];
let batch = crate::table_import::build_import_insert_batch_from_rows(
&rows,
&columns,
&mappings,
&[],
"events",
"main",
&DatabaseType::CloudflareD1,
)
.unwrap()
.unwrap();
assert_eq!(batch.row_count, 2);
assert_eq!(batch.sql.split(";\n").count(), 2);
assert!(batch.sql.split(";\n").all(|statement| statement.len() <= super::super::MAX_SQL_STATEMENT_BYTES));
}
}

View File

@ -0,0 +1,628 @@
use reqwest::Client as HttpClient;
use serde::Deserialize;
use std::time::{Duration, Instant};
use super::{http_client_builder, json_value_for_js, with_connection_timeout};
use crate::models::connection::ConnectionConfig;
use crate::types::{
ColumnInfo, DatabaseInfo, ForeignKeyInfo, IndexInfo, ObjectSource, ObjectSourceKind, QueryResult, TableInfo,
TriggerInfo,
};
mod import;
mod sql_guard;
mod sql_lexer;
mod sql_limits;
#[cfg(test)]
mod transfer_tests;
mod trigger;
pub(crate) use import::{build_import_insert_batches, build_streaming_import_insert_batch};
pub use sql_limits::MAX_SQL_STATEMENT_BYTES;
const CLOUDFLARE_API_BASE_URL: &str = "https://api.cloudflare.com/client/v4";
#[derive(Clone)]
pub struct CloudflareD1Client {
http: HttpClient,
endpoint: String,
api_token: String,
}
impl CloudflareD1Client {
pub fn new(account_id: &str, database_id: &str, api_token: &str, timeout: Duration) -> Result<Self, String> {
let account_id = required_identifier(account_id, "Cloudflare Account ID")?;
let database_id = required_identifier(database_id, "Cloudflare D1 Database ID")?;
let endpoint = format!("{CLOUDFLARE_API_BASE_URL}/accounts/{account_id}/d1/database/{database_id}/raw");
Self::with_endpoint(endpoint, api_token, timeout)
}
fn with_endpoint(endpoint: String, api_token: &str, timeout: Duration) -> Result<Self, String> {
let api_token = required_field(api_token, "Cloudflare API Token")?;
let http = http_client_builder(timeout)
.build()
.map_err(|error| format!("Failed to configure Cloudflare D1 HTTP client: {error}"))?;
Ok(Self { http, endpoint, api_token: api_token.to_string() })
}
fn post_raw(&self, sql: &str) -> reqwest::RequestBuilder {
self.http
.post(&self.endpoint)
.bearer_auth(&self.api_token)
.header(reqwest::header::CONTENT_TYPE, "application/json")
.json(&serde_json::json!({ "sql": sql }))
}
}
#[derive(Debug, Deserialize)]
struct D1ApiResponse {
#[serde(default)]
success: Option<bool>,
#[serde(default)]
errors: Vec<D1ApiError>,
#[serde(default)]
result: Vec<D1StatementResult>,
}
#[derive(Debug, Deserialize)]
struct D1ApiError {
#[serde(default)]
code: Option<u64>,
message: String,
}
#[derive(Debug, Default, Deserialize)]
struct D1StatementResult {
#[serde(default)]
success: Option<bool>,
#[serde(default)]
results: D1RawResult,
#[serde(default)]
meta: D1Meta,
}
#[derive(Debug, Default, Deserialize)]
struct D1RawResult {
#[serde(default)]
columns: Vec<String>,
#[serde(default)]
rows: Vec<Vec<serde_json::Value>>,
}
#[derive(Debug, Default, Deserialize)]
struct D1Meta {
#[serde(default)]
changes: Option<u64>,
}
pub async fn test_connection(client: &CloudflareD1Client, timeout: Duration) -> Result<(), String> {
with_connection_timeout("Cloudflare D1", timeout, async { send_raw(client, "SELECT 1").await.map(|_| ()) }).await
}
pub async fn connect(config: &ConnectionConfig, timeout: Duration) -> Result<CloudflareD1Client, String> {
let client = CloudflareD1Client::new(
&config.host,
config.database.as_deref().unwrap_or_default(),
&config.password,
timeout,
)?;
test_connection(&client, timeout).await?;
Ok(client)
}
pub async fn list_databases(_client: &CloudflareD1Client) -> Result<Vec<DatabaseInfo>, String> {
Ok(vec![DatabaseInfo { name: "main".to_string() }])
}
pub async fn list_tables(client: &CloudflareD1Client, _schema: &str) -> Result<Vec<TableInfo>, String> {
let result = query_inner(
client,
"SELECT name, type FROM sqlite_master WHERE type IN ('table', 'view') AND name NOT LIKE 'sqlite\\_%' ESCAPE '\\' AND name <> '_cf_KV' ORDER BY name",
)
.await?;
Ok(result
.rows
.into_iter()
.map(|row| {
let table_type = value_as_string(row.get(1)).unwrap_or_else(|| "table".to_string());
TableInfo {
name: value_as_string(row.first()).unwrap_or_default(),
table_type: if table_type.eq_ignore_ascii_case("view") { "VIEW" } else { "BASE TABLE" }.to_string(),
comment: None,
parent_schema: None,
parent_name: None,
}
})
.collect())
}
pub async fn get_columns(client: &CloudflareD1Client, _schema: &str, table: &str) -> Result<Vec<ColumnInfo>, String> {
// table_info omits generated columns; table_xinfo exposes them and marks
// virtual/stored generated columns with hidden values 2 and 3.
let result = query_inner(client, &format!("PRAGMA table_xinfo({})", sqlite_ident(table))).await?;
Ok(result
.rows
.into_iter()
.map(|row| {
let hidden = value_by_column(&result.columns, &row, "hidden")
.and_then(|value| value.parse::<i64>().ok())
.unwrap_or(0);
ColumnInfo {
name: value_by_column(&result.columns, &row, "name").unwrap_or_default(),
data_type: value_by_column(&result.columns, &row, "type").unwrap_or_default(),
is_nullable: value_by_column(&result.columns, &row, "notnull")
.and_then(|value| value.parse::<i64>().ok())
.unwrap_or(0)
== 0,
column_default: value_by_column(&result.columns, &row, "dflt_value"),
is_primary_key: value_by_column(&result.columns, &row, "pk")
.and_then(|value| value.parse::<i64>().ok())
.unwrap_or(0)
> 0,
extra: match hidden {
2 => Some("generated always as virtual".to_string()),
3 => Some("generated always as stored".to_string()),
_ => None,
},
comment: None,
numeric_precision: None,
numeric_scale: None,
character_maximum_length: None,
enum_values: None,
..Default::default()
}
})
.collect())
}
pub async fn list_indexes(client: &CloudflareD1Client, _schema: &str, table: &str) -> Result<Vec<IndexInfo>, String> {
let result = query_inner(client, &format!("PRAGMA index_list({})", sqlite_ident(table))).await?;
let mut indexes = Vec::new();
for row in result.rows {
let name = value_by_column(&result.columns, &row, "name").unwrap_or_default();
if name.is_empty() {
continue;
}
let is_unique =
value_by_column(&result.columns, &row, "unique").and_then(|value| value.parse::<i64>().ok()).unwrap_or(0)
!= 0;
let origin = value_by_column(&result.columns, &row, "origin").unwrap_or_default();
let column_result = query_inner(client, &format!("PRAGMA index_info({})", sqlite_ident(&name))).await?;
let columns =
column_result.rows.iter().filter_map(|row| value_by_column(&column_result.columns, row, "name")).collect();
indexes.push(IndexInfo {
name,
columns,
is_unique,
is_primary: origin == "pk",
filter: None,
index_type: None,
included_columns: None,
comment: None,
});
}
Ok(indexes)
}
pub async fn list_foreign_keys(
client: &CloudflareD1Client,
_schema: &str,
table: &str,
) -> Result<Vec<ForeignKeyInfo>, String> {
let result = query_inner(client, &format!("PRAGMA foreign_key_list({})", sqlite_ident(table))).await?;
Ok(result
.rows
.into_iter()
.map(|row| ForeignKeyInfo {
name: format!("fk_{}", value_by_column(&result.columns, &row, "id").unwrap_or_else(|| "0".to_string())),
column: value_by_column(&result.columns, &row, "from").unwrap_or_default(),
ref_schema: None,
ref_table: value_by_column(&result.columns, &row, "table").unwrap_or_default(),
ref_column: value_by_column(&result.columns, &row, "to").unwrap_or_default(),
on_update: value_by_column(&result.columns, &row, "on_update"),
on_delete: value_by_column(&result.columns, &row, "on_delete"),
})
.collect())
}
pub async fn list_triggers(
client: &CloudflareD1Client,
_schema: &str,
table: &str,
) -> Result<Vec<TriggerInfo>, String> {
let result = query_inner(
client,
&format!(
"SELECT name, sql FROM sqlite_master WHERE type = 'trigger' AND tbl_name = {} ORDER BY name",
sqlite_string(table)
),
)
.await?;
Ok(result
.rows
.into_iter()
.map(|row| {
let statement = value_as_string(row.get(1));
let metadata = trigger::metadata_from_sql(statement.as_deref().unwrap_or_default());
TriggerInfo {
name: value_as_string(row.first()).unwrap_or_default(),
event: metadata.event.to_string(),
timing: metadata.timing.to_string(),
statement,
}
})
.collect())
}
pub async fn table_ddl(client: &CloudflareD1Client, table: &str) -> Result<String, String> {
first_string_cell(
query_inner(
client,
&format!("SELECT sql FROM sqlite_master WHERE type='table' AND name={}", sqlite_string(table)),
)
.await?,
)
}
pub async fn object_source(
client: &CloudflareD1Client,
name: &str,
object_type: &ObjectSourceKind,
) -> Result<ObjectSource, String> {
let kind = match object_type {
ObjectSourceKind::View => "view",
_ => return Err("Object source is not supported for this Cloudflare D1 object type".to_string()),
};
let source = first_string_cell(
query_inner(
client,
&format!(
"SELECT sql FROM sqlite_master WHERE type={} AND name={}",
sqlite_string(kind),
sqlite_string(name)
),
)
.await?,
)?;
Ok(ObjectSource { name: name.to_string(), object_type: object_type.clone(), schema: None, source, editable: None })
}
pub async fn execute_query_with_max_rows(
client: &CloudflareD1Client,
sql: &str,
max_rows: Option<usize>,
) -> Result<QueryResult, String> {
let start = Instant::now();
let sql = sql.trim();
if sql.is_empty() {
return Ok(empty_query_result(start.elapsed().as_millis()));
}
sql_guard::validate_sql(sql)?;
sql_limits::validate_statement_sizes(sql)?;
let mut statements = send_raw(client, sql).await?;
let affected_rows: u64 = statements.iter().filter_map(|statement| statement.meta.changes).sum();
let display = statements.pop().unwrap_or_default();
Ok(query_result(display.results, affected_rows, start.elapsed().as_millis(), max_rows))
}
async fn query_inner(client: &CloudflareD1Client, sql: &str) -> Result<D1RawResult, String> {
send_raw(client, sql)
.await?
.into_iter()
.next()
.map(|statement| statement.results)
.ok_or_else(|| "Cloudflare D1 returned no result".to_string())
}
async fn send_raw(client: &CloudflareD1Client, sql: &str) -> Result<Vec<D1StatementResult>, String> {
let response =
client.post_raw(sql).send().await.map_err(|error| format!("Cloudflare D1 request failed: {error}"))?;
let status = response.status();
let body = response.text().await.map_err(|error| format!("Cloudflare D1 response read failed: {error}"))?;
let parsed: D1ApiResponse = serde_json::from_str(&body)
.map_err(|error| format!("Cloudflare D1 response parse failed: {error}; body: {}", response_excerpt(&body)))?;
if !status.is_success() || parsed.success == Some(false) {
return Err(format_d1_error(status.as_u16(), &parsed.errors, &body));
}
if parsed.result.iter().any(|result| result.success == Some(false)) {
return Err(format_d1_error(status.as_u16(), &parsed.errors, &body));
}
if parsed.result.is_empty() {
return Err("Cloudflare D1 returned no statement results".to_string());
}
Ok(parsed.result)
}
fn query_result(
mut result: D1RawResult,
affected_rows: u64,
execution_time_ms: u128,
max_rows: Option<usize>,
) -> QueryResult {
result.rows = result.rows.into_iter().map(|row| row.into_iter().map(json_value_for_js).collect()).collect();
let row_limit = max_rows.unwrap_or(crate::query::MAX_ROWS).max(1);
let truncated = result.rows.len() > row_limit;
if truncated {
result.rows.truncate(row_limit);
}
QueryResult {
columns: result.columns,
column_types: Vec::new(),
column_sortables: vec![],
rows: result.rows,
affected_rows,
execution_time_ms,
truncated,
session_id: None,
has_more: false,
}
}
fn empty_query_result(execution_time_ms: u128) -> QueryResult {
query_result(D1RawResult::default(), 0, execution_time_ms, None)
}
fn format_d1_error(status: u16, errors: &[D1ApiError], body: &str) -> String {
let details = errors
.iter()
.map(|error| match error.code {
Some(code) => format!("{} ({code})", error.message),
None => error.message.clone(),
})
.collect::<Vec<_>>()
.join("; ");
if details.is_empty() {
format!("Cloudflare D1 API error ({status}): {}", response_excerpt(body))
} else {
format!("Cloudflare D1 API error ({status}): {details}")
}
}
fn response_excerpt(body: &str) -> String {
const MAX_CHARS: usize = 1_000;
let mut excerpt = body.chars().take(MAX_CHARS).collect::<String>();
if body.chars().count() > MAX_CHARS {
excerpt.push('…');
}
excerpt
}
fn required_field<'a>(value: &'a str, label: &str) -> Result<&'a str, String> {
let value = value.trim();
if value.is_empty() {
Err(format!("{label} is required"))
} else {
Ok(value)
}
}
fn required_identifier<'a>(value: &'a str, label: &str) -> Result<&'a str, String> {
let value = required_field(value, label)?;
if value.chars().all(|character| character.is_ascii_alphanumeric() || character == '-' || character == '_') {
Ok(value)
} else {
Err(format!("{label} contains invalid characters"))
}
}
fn first_string_cell(result: D1RawResult) -> Result<String, String> {
result
.rows
.first()
.and_then(|row| row.first())
.and_then(|value| value_as_string(Some(value)))
.ok_or_else(|| "Object not found".to_string())
}
fn value_by_column(columns: &[String], row: &[serde_json::Value], name: &str) -> Option<String> {
columns
.iter()
.position(|column| column.eq_ignore_ascii_case(name))
.and_then(|index| row.get(index))
.and_then(|value| value_as_string(Some(value)))
}
fn value_as_string(value: Option<&serde_json::Value>) -> Option<String> {
match value? {
serde_json::Value::Null => None,
serde_json::Value::String(value) => Some(value.clone()),
serde_json::Value::Number(value) => Some(value.to_string()),
serde_json::Value::Bool(value) => Some(value.to_string()),
other => Some(other.to_string()),
}
}
fn sqlite_ident(value: &str) -> String {
format!("\"{}\"", value.replace('"', "\"\""))
}
fn sqlite_string(value: &str) -> String {
format!("'{}'", value.replace('\'', "''"))
}
#[cfg(test)]
mod tests {
use super::*;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
#[test]
fn validates_required_connection_fields() {
let timeout = Duration::from_secs(1);
assert!(CloudflareD1Client::new("", "database", "token", timeout).err().unwrap().contains("Account ID"));
assert!(CloudflareD1Client::new("account", "", "token", timeout).err().unwrap().contains("Database ID"));
assert!(CloudflareD1Client::new("account", "database", "", timeout).err().unwrap().contains("API Token"));
assert!(CloudflareD1Client::new("account/path", "database", "token", timeout)
.err()
.unwrap()
.contains("invalid characters"));
}
#[test]
fn parses_raw_api_response() {
let response: D1ApiResponse = serde_json::from_value(serde_json::json!({
"success": true,
"errors": [],
"result": [{
"success": true,
"results": { "columns": ["id", "name"], "rows": [[1, "Ada"]] },
"meta": { "changes": 0 }
}]
}))
.unwrap();
assert_eq!(response.success, Some(true));
assert_eq!(response.result[0].results.columns, ["id", "name"]);
assert_eq!(response.result[0].results.rows[0], [serde_json::json!(1), serde_json::json!("Ada")]);
}
#[test]
fn converts_and_truncates_query_result() {
let raw = D1RawResult {
columns: vec!["id".to_string()],
rows: vec![vec![serde_json::json!(1)], vec![serde_json::json!(2)]],
};
let result = query_result(raw, 0, 5, Some(1));
assert_eq!(result.rows, vec![vec![serde_json::json!(1)]]);
assert!(result.truncated);
}
#[test]
fn limits_error_response_excerpts() {
let excerpt = response_excerpt(&"x".repeat(1_100));
assert_eq!(excerpt.chars().count(), 1_001);
assert!(excerpt.ends_with('…'));
}
#[tokio::test]
async fn sends_bearer_authenticated_raw_query_and_maps_response() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let request = read_http_request(&mut socket).await;
assert!(request.starts_with("POST /raw HTTP/1.1"));
assert!(request.to_ascii_lowercase().contains("authorization: bearer test-token"));
assert!(request.contains("\"sql\":\"SELECT 1\""));
let body = r#"{"success":true,"errors":[],"result":[{"success":true,"results":{"columns":["1"],"rows":[[1]]},"meta":{"changes":0}}]}"#;
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
socket.write_all(response.as_bytes()).await.unwrap();
});
let client =
CloudflareD1Client::with_endpoint(format!("http://{address}/raw"), "test-token", Duration::from_secs(2))
.unwrap();
let result = execute_query_with_max_rows(&client, "SELECT 1", None).await.unwrap();
assert_eq!(result.columns, ["1"]);
assert_eq!(result.rows, vec![vec![serde_json::json!(1)]]);
server.await.unwrap();
}
#[tokio::test]
async fn lists_tables_and_views_from_sqlite_master() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let request = read_http_request(&mut socket).await;
assert!(request.contains("sqlite_master"));
assert!(request.contains("name <> '_cf_KV'"));
assert!(!request.contains("d1_%"));
let body = r#"{"success":true,"errors":[],"result":[{"success":true,"results":{"columns":["name","type"],"rows":[["users","table"],["active_users","view"]]},"meta":{"changes":0}}]}"#;
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
socket.write_all(response.as_bytes()).await.unwrap();
});
let client =
CloudflareD1Client::with_endpoint(format!("http://{address}/raw"), "test-token", Duration::from_secs(2))
.unwrap();
let tables = list_tables(&client, "main").await.unwrap();
assert_eq!(tables.len(), 2);
assert_eq!(tables[0].name, "users");
assert_eq!(tables[0].table_type, "BASE TABLE");
assert_eq!(tables[1].name, "active_users");
assert_eq!(tables[1].table_type, "VIEW");
server.await.unwrap();
}
#[tokio::test]
async fn lists_generated_columns_from_table_xinfo() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.unwrap();
let request = read_http_request(&mut socket).await;
assert!(request.contains("PRAGMA table_xinfo"));
let body = r#"{"success":true,"errors":[],"result":[{"success":true,"results":{"columns":["cid","name","type","notnull","dflt_value","pk","hidden"],"rows":[[0,"price","INTEGER",1,null,0,0],[1,"tax","INTEGER",0,null,0,2],[2,"total","INTEGER",0,null,0,3]]},"meta":{"changes":0}}]}"#;
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
socket.write_all(response.as_bytes()).await.unwrap();
});
let client =
CloudflareD1Client::with_endpoint(format!("http://{address}/raw"), "test-token", Duration::from_secs(2))
.unwrap();
let columns = get_columns(&client, "main", "orders").await.unwrap();
assert_eq!(columns.len(), 3);
assert_eq!(columns[0].extra, None);
assert_eq!(columns[1].extra.as_deref(), Some("generated always as virtual"));
assert_eq!(columns[2].extra.as_deref(), Some("generated always as stored"));
server.await.unwrap();
}
#[test]
fn rejects_oversized_statements_but_allows_large_multi_statement_batches() {
let oversized = format!("SELECT '{}'", "x".repeat(MAX_SQL_STATEMENT_BYTES));
let error = sql_limits::validate_statement_sizes(&oversized).unwrap_err();
assert!(error.contains("statement 1"));
assert!(error.contains("100000 bytes"));
let first = format!("SELECT '{}'", "a".repeat(60_000));
let second = format!("SELECT '{}'", "b".repeat(60_000));
assert!(sql_limits::validate_statement_sizes(&format!("{first}; {second}")).is_ok());
let trigger = format!(
"CREATE TRIGGER audit AFTER INSERT ON users BEGIN INSERT INTO logs VALUES ('{}'); INSERT INTO logs VALUES ('{}'); END",
"a".repeat(60_000),
"b".repeat(60_000)
);
assert!(sql_limits::validate_statement_sizes(&trigger).is_err());
}
async fn read_http_request(socket: &mut tokio::net::TcpStream) -> String {
let mut bytes = Vec::new();
let mut buffer = [0_u8; 2048];
loop {
let read = socket.read(&mut buffer).await.unwrap();
if read == 0 {
break;
}
bytes.extend_from_slice(&buffer[..read]);
let Some(header_end) = bytes.windows(4).position(|window| window == b"\r\n\r\n") else {
continue;
};
let headers = String::from_utf8_lossy(&bytes[..header_end]);
let content_length = headers
.lines()
.find_map(|line| line.split_once(':').filter(|(name, _)| name.eq_ignore_ascii_case("content-length")))
.and_then(|(_, value)| value.trim().parse::<usize>().ok())
.unwrap_or(0);
if bytes.len() >= header_end + 4 + content_length {
break;
}
}
String::from_utf8(bytes).unwrap()
}
}

View File

@ -0,0 +1,168 @@
use super::sql_lexer::{complete_sql_statements, lex_sql, SqlLexeme};
const SUPPORTED_PRAGMAS: &[&str] = &[
"case_sensitive_like",
"defer_foreign_keys",
"foreign_key_check",
"foreign_key_list",
"foreign_keys",
"ignore_check_constraints",
"index_info",
"index_list",
"index_xinfo",
"legacy_alter_table",
"optimize",
"quick_check",
"recursive_triggers",
"reverse_unordered_selects",
"table_info",
"table_list",
"table_xinfo",
];
pub(super) fn validate_sql(sql: &str) -> Result<(), String> {
for statement in complete_sql_statements(sql) {
validate_statement(statement)?;
}
Ok(())
}
fn validate_statement(statement: &str) -> Result<(), String> {
let lexemes = lex_sql(statement);
let words = lexemes
.iter()
.filter_map(|lexeme| match lexeme {
SqlLexeme::Word(word) => Some(*word),
SqlLexeme::Number(_) | SqlLexeme::Symbol(_) | SqlLexeme::Semicolon(_) => None,
})
.collect::<Vec<_>>();
let Some(first) = words.first() else {
return Ok(());
};
if ["BEGIN", "COMMIT", "END", "ROLLBACK", "SAVEPOINT", "RELEASE", "START"]
.iter()
.any(|keyword| first.eq_ignore_ascii_case(keyword))
{
return Err(
"Cloudflare D1 does not support explicit transaction or savepoint statements. Submit the SQL statements together without BEGIN, COMMIT, ROLLBACK, SAVEPOINT, or RELEASE."
.to_string(),
);
}
if first.eq_ignore_ascii_case("ATTACH") || first.eq_ignore_ascii_case("DETACH") {
return Err(
"Cloudflare D1 does not support ATTACH or DETACH. Each connection targets one provider-managed D1 database."
.to_string(),
);
}
if first.eq_ignore_ascii_case("CREATE") {
validate_create_statement(&words)?;
}
if first.eq_ignore_ascii_case("PRAGMA") {
validate_pragma(&words, &lexemes)?;
}
Ok(())
}
fn validate_create_statement(words: &[&str]) -> Result<(), String> {
if words.get(1).is_some_and(|word| word.eq_ignore_ascii_case("TEMP") || word.eq_ignore_ascii_case("TEMPORARY")) {
return Err(
"Cloudflare D1 does not support temporary tables, indexes, views, or triggers. Create a persistent object instead."
.to_string(),
);
}
let creates_virtual_table = words
.get(1..3)
.is_some_and(|prefix| prefix[0].eq_ignore_ascii_case("VIRTUAL") && prefix[1].eq_ignore_ascii_case("TABLE"));
if !creates_virtual_table {
return Ok(());
}
let module =
words.iter().position(|word| word.eq_ignore_ascii_case("USING")).and_then(|index| words.get(index + 1));
if module.is_some_and(|module| module.eq_ignore_ascii_case("fts5") || module.eq_ignore_ascii_case("fts5vocab")) {
return Ok(());
}
Err("Cloudflare D1 only supports FTS5 and fts5vocab virtual table modules.".to_string())
}
fn validate_pragma(words: &[&str], lexemes: &[SqlLexeme<'_>]) -> Result<(), String> {
let mut pragma = words.get(1).copied();
if pragma.is_some_and(|word| word.eq_ignore_ascii_case("main") || word.eq_ignore_ascii_case("temp")) {
pragma = words.get(2).copied();
}
let Some(pragma) = pragma else {
return Err("Cloudflare D1 requires a supported PRAGMA name.".to_string());
};
if !SUPPORTED_PRAGMAS.iter().any(|supported| pragma.eq_ignore_ascii_case(supported)) {
return Err(format!(
"Cloudflare D1 does not support PRAGMA {pragma}. Use one of the documented D1-compatible PRAGMA statements."
));
}
if pragma.eq_ignore_ascii_case("optimize") && contains_negative_one(lexemes) {
return Err("Cloudflare D1 does not support PRAGMA optimize(-1); use PRAGMA optimize instead.".to_string());
}
Ok(())
}
fn contains_negative_one(lexemes: &[SqlLexeme<'_>]) -> bool {
lexemes
.windows(2)
.any(|pair| matches!(pair[0], SqlLexeme::Symbol(b'-')) && matches!(pair[1], SqlLexeme::Number("1")))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn accepts_supported_d1_sql_features() {
assert!(validate_sql(
"CREATE TRIGGER audit AFTER INSERT ON users BEGIN INSERT INTO logs VALUES (NEW.id); END;"
)
.is_ok());
assert!(
validate_sql("CREATE VIRTUAL TABLE docs USING fts5(title, body); SELECT json_extract('{}', '$');").is_ok()
);
assert!(validate_sql("PRAGMA table_info(users); PRAGMA optimize;").is_ok());
}
#[test]
fn rejects_explicit_transaction_control() {
for sql in [
"BEGIN; INSERT INTO users VALUES (1); COMMIT;",
"SAVEPOINT before_update",
"ROLLBACK TO before_update",
"RELEASE before_update",
] {
let error = validate_sql(sql).unwrap_err();
assert!(error.contains("does not support explicit transaction"));
}
}
#[test]
fn rejects_attached_temporary_and_unknown_virtual_databases() {
assert!(validate_sql("ATTACH DATABASE 'other.db' AS other").unwrap_err().contains("ATTACH"));
assert!(validate_sql("CREATE TEMP TABLE scratch(id INTEGER)").unwrap_err().contains("temporary"));
assert!(validate_sql("CREATE VIRTUAL TABLE places USING rtree(id, min_x, max_x)")
.unwrap_err()
.contains("only supports FTS5"));
}
#[test]
fn rejects_unsupported_pragma_forms() {
assert!(validate_sql("PRAGMA journal_mode=WAL").unwrap_err().contains("PRAGMA journal_mode"));
assert!(validate_sql("PRAGMA data_version").unwrap_err().contains("PRAGMA data_version"));
assert!(validate_sql("PRAGMA optimize(-1)").unwrap_err().contains("optimize(-1)"));
assert!(validate_sql("PRAGMA optimize /* -1 is only a comment */").is_ok());
}
}

View File

@ -0,0 +1,180 @@
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum SqlLexeme<'a> {
Word(&'a str),
Number(&'a str),
Symbol(u8),
Semicolon(usize),
}
pub(super) fn lex_sql(sql: &str) -> Vec<SqlLexeme<'_>> {
let bytes = sql.as_bytes();
let mut lexemes = Vec::new();
let mut index = 0;
while index < bytes.len() {
match bytes[index] {
b'-' if bytes.get(index + 1) == Some(&b'-') => {
index += 2;
while index < bytes.len() && bytes[index] != b'\n' {
index += 1;
}
}
b'/' if bytes.get(index + 1) == Some(&b'*') => {
index += 2;
while index + 1 < bytes.len() && !(bytes[index] == b'*' && bytes[index + 1] == b'/') {
index += 1;
}
index = (index + 2).min(bytes.len());
}
quote @ (b'\'' | b'"' | b'`') => {
index += 1;
while index < bytes.len() {
if bytes[index] != quote {
index += 1;
continue;
}
if bytes.get(index + 1) == Some(&quote) {
index += 2;
} else {
index += 1;
break;
}
}
}
b'[' => {
index += 1;
while index < bytes.len() {
if bytes[index] != b']' {
index += 1;
continue;
}
if bytes.get(index + 1) == Some(&b']') {
index += 2;
} else {
index += 1;
break;
}
}
}
b';' => {
lexemes.push(SqlLexeme::Semicolon(index));
index += 1;
}
byte if byte.is_ascii_digit() => {
let start = index;
index += 1;
while index < bytes.len() && bytes[index].is_ascii_digit() {
index += 1;
}
lexemes.push(SqlLexeme::Number(&sql[start..index]));
}
byte if byte.is_ascii_alphabetic() || byte == b'_' => {
let start = index;
index += 1;
while index < bytes.len()
&& (bytes[index].is_ascii_alphanumeric() || matches!(bytes[index], b'_' | b'$'))
{
index += 1;
}
lexemes.push(SqlLexeme::Word(&sql[start..index]));
}
byte if matches!(byte, b'(' | b')' | b'=' | b'.' | b',' | b'-' | b'+') => {
lexemes.push(SqlLexeme::Symbol(byte));
index += 1;
}
_ => index += 1,
}
}
lexemes
}
pub(super) fn complete_sql_statements(sql: &str) -> Vec<&str> {
let mut statements = Vec::new();
let mut statement_start = 0;
let mut prefix = Vec::with_capacity(3);
let mut is_trigger = false;
let mut trigger_body_started = false;
let mut trigger_depth = 0_usize;
for lexeme in lex_sql(sql) {
match lexeme {
SqlLexeme::Word(word) => {
if prefix.len() < 3 {
prefix.push(word);
is_trigger = is_create_trigger_prefix(&prefix);
}
if !is_trigger {
continue;
}
if !trigger_body_started && word.eq_ignore_ascii_case("BEGIN") {
trigger_body_started = true;
trigger_depth = 1;
} else if trigger_body_started
&& (word.eq_ignore_ascii_case("BEGIN") || word.eq_ignore_ascii_case("CASE"))
{
trigger_depth += 1;
} else if trigger_body_started && word.eq_ignore_ascii_case("END") {
trigger_depth = trigger_depth.saturating_sub(1);
}
}
SqlLexeme::Semicolon(position) => {
if is_trigger && trigger_body_started && trigger_depth > 0 {
continue;
}
push_statement(sql, statement_start, position, &mut statements);
statement_start = position + 1;
prefix.clear();
is_trigger = false;
trigger_body_started = false;
trigger_depth = 0;
}
SqlLexeme::Number(_) | SqlLexeme::Symbol(_) => {}
}
}
push_statement(sql, statement_start, sql.len(), &mut statements);
statements
}
fn is_create_trigger_prefix(words: &[&str]) -> bool {
words.first().is_some_and(|word| word.eq_ignore_ascii_case("CREATE"))
&& (words.get(1).is_some_and(|word| word.eq_ignore_ascii_case("TRIGGER"))
|| (words
.get(1)
.is_some_and(|word| word.eq_ignore_ascii_case("TEMP") || word.eq_ignore_ascii_case("TEMPORARY"))
&& words.get(2).is_some_and(|word| word.eq_ignore_ascii_case("TRIGGER"))))
}
fn push_statement<'a>(sql: &'a str, start: usize, end: usize, statements: &mut Vec<&'a str>) {
let statement = sql[start..end].trim();
if !statement.is_empty() {
statements.push(statement);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn keeps_trigger_body_semicolons_in_one_statement() {
let sql = "CREATE TRIGGER audit AFTER UPDATE ON users BEGIN INSERT INTO logs VALUES ('a;b'); UPDATE stats SET value = CASE WHEN value > 0 THEN value ELSE 0 END; END; SELECT 1;";
let statements = complete_sql_statements(sql);
assert_eq!(statements.len(), 2);
assert!(statements[0].starts_with("CREATE TRIGGER"));
assert!(statements[0].ends_with("END"));
assert_eq!(statements[1], "SELECT 1");
}
#[test]
fn ignores_trigger_words_in_comments_and_strings() {
let sql = "SELECT 'CREATE TRIGGER fake; BEGIN END'; -- CREATE TRIGGER ignored\nSELECT 2;";
assert_eq!(
complete_sql_statements(sql),
["SELECT 'CREATE TRIGGER fake; BEGIN END'", "-- CREATE TRIGGER ignored\nSELECT 2"]
);
}
}

View File

@ -0,0 +1,115 @@
use std::ops::Range;
pub const MAX_SQL_STATEMENT_BYTES: usize = 100_000;
#[derive(Debug, PartialEq, Eq)]
pub(super) struct SqlBatch {
pub sql: String,
pub item_count: usize,
}
pub(super) fn build_sql_batches<F>(
item_count: usize,
max_items: usize,
item_label: &str,
render: F,
) -> Result<Vec<SqlBatch>, String>
where
F: Fn(Range<usize>) -> String,
{
let mut batches = Vec::new();
let mut start = 0;
let max_items = max_items.max(1);
while start < item_count {
let mut end = start + 1;
let mut accepted = render(start..end);
if accepted.len() > MAX_SQL_STATEMENT_BYTES {
return Err(format!(
"Cloudflare D1 {item_label} {} generates a {}-byte SQL statement; the maximum is {MAX_SQL_STATEMENT_BYTES} bytes. Reduce the value size because a single item cannot be split safely.",
start + 1,
accepted.len()
));
}
while end < item_count && end - start < max_items {
let candidate = render(start..end + 1);
if candidate.len() > MAX_SQL_STATEMENT_BYTES {
break;
}
accepted = candidate;
end += 1;
}
if !accepted.trim().is_empty() {
batches.push(SqlBatch { sql: accepted, item_count: end - start });
}
start = end;
}
Ok(batches)
}
pub(super) fn validate_statement_sizes(sql: &str) -> Result<(), String> {
let statements = super::sql_lexer::complete_sql_statements(sql);
for (index, statement) in statements.iter().enumerate() {
if statement.len() > MAX_SQL_STATEMENT_BYTES {
return Err(statement_size_error(index + 1, statement.len()));
}
}
Ok(())
}
fn statement_size_error(index: usize, byte_len: usize) -> String {
format!(
"Cloudflare D1 SQL statement {index} is {byte_len} bytes; the maximum is {MAX_SQL_STATEMENT_BYTES} bytes. Split the statement or reduce the batch/value size."
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn splits_utf8_sql_batches_by_bytes() {
let rows = ["".repeat(20_000), "".repeat(20_000)];
let batches = build_sql_batches(rows.len(), 100, "import row", |range| {
format!("INSERT INTO events VALUES ('{}')", rows[range].join("'), ('"))
})
.unwrap();
assert_eq!(batches.len(), 2);
assert!(batches.iter().all(|batch| batch.sql.len() <= MAX_SQL_STATEMENT_BYTES));
}
#[test]
fn rejects_a_single_oversized_item() {
let error = build_sql_batches(1, 100, "import row", |_| {
format!("INSERT INTO events VALUES ('{}')", "x".repeat(MAX_SQL_STATEMENT_BYTES))
})
.unwrap_err();
assert!(error.contains("import row 1"));
assert!(error.contains("maximum is 100000 bytes"));
}
#[test]
fn validates_trigger_and_other_batch_statements_individually() {
let trigger = "CREATE TRIGGER audit AFTER UPDATE ON users BEGIN INSERT INTO logs VALUES (NEW.id); END";
let first = format!("SELECT '{}'", "a".repeat(60_000));
let second = format!("SELECT '{}'", "b".repeat(60_000));
assert!(validate_statement_sizes(&format!("{trigger}; {first}; {second}")).is_ok());
}
#[test]
fn rejects_an_oversized_complete_trigger_statement() {
let trigger = format!(
"CREATE TRIGGER audit AFTER UPDATE ON users BEGIN INSERT INTO logs VALUES ('{}'); INSERT INTO logs VALUES ('{}'); END",
"a".repeat(60_000),
"b".repeat(60_000)
);
assert!(validate_statement_sizes(&trigger).is_err());
}
}

View File

@ -0,0 +1,19 @@
use serde_json::json;
use crate::models::connection::DatabaseType;
use crate::transfer::generate_upsert_typed;
#[test]
fn uses_sqlite_upsert_syntax() {
let sql = generate_upsert_typed(
&[String::from("id"), String::from("name")],
&[Some(String::from("integer")), Some(String::from("text"))],
&[vec![json!(1), json!("Ada")]],
"users",
"main",
&DatabaseType::CloudflareD1,
&[String::from("id")],
);
assert!(sql.contains("ON CONFLICT (\"id\") DO UPDATE SET"));
}

View File

@ -0,0 +1,71 @@
use super::sql_lexer::{lex_sql, SqlLexeme};
pub(super) struct TriggerMetadata {
pub timing: &'static str,
pub event: &'static str,
}
pub(super) fn metadata_from_sql(sql: &str) -> TriggerMetadata {
let words = lex_sql(sql)
.into_iter()
.filter_map(|lexeme| match lexeme {
SqlLexeme::Word(word) => Some(word),
SqlLexeme::Number(_) | SqlLexeme::Symbol(_) | SqlLexeme::Semicolon(_) => None,
})
.collect::<Vec<_>>();
let declaration_start =
words.iter().position(|word| word.eq_ignore_ascii_case("TRIGGER")).map_or(0, |index| index + 1);
let declaration_end = words[declaration_start..]
.iter()
.position(|word| word.eq_ignore_ascii_case("ON") || word.eq_ignore_ascii_case("BEGIN"))
.map_or(words.len(), |index| declaration_start + index);
let declaration = &words[declaration_start..declaration_end];
let timing = if declaration
.windows(2)
.any(|words| words[0].eq_ignore_ascii_case("INSTEAD") && words[1].eq_ignore_ascii_case("OF"))
{
"INSTEAD OF"
} else if declaration.iter().any(|word| word.eq_ignore_ascii_case("AFTER")) {
"AFTER"
} else {
// SQLite defaults to BEFORE when no timing keyword is present.
"BEFORE"
};
let event = declaration
.iter()
.find_map(|word| ["DELETE", "INSERT", "UPDATE"].into_iter().find(|event| word.eq_ignore_ascii_case(event)))
.unwrap_or("UNKNOWN");
TriggerMetadata { timing, event }
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn reads_event_from_declaration_not_trigger_body() {
let metadata = metadata_from_sql(
"CREATE TRIGGER audit AFTER UPDATE ON users BEGIN INSERT INTO audit_log VALUES (NEW.id); END",
);
assert_eq!(metadata.timing, "AFTER");
assert_eq!(metadata.event, "UPDATE");
}
#[test]
fn supports_default_and_instead_of_timing() {
let default_timing = metadata_from_sql(
"CREATE TRIGGER update_customer UPDATE OF name ON customers BEGIN DELETE FROM cache; END",
);
let instead_of = metadata_from_sql(
"CREATE TRIGGER [insert] INSTEAD OF INSERT ON customer_view BEGIN UPDATE customers SET name = NEW.name; END",
);
assert_eq!(default_timing.timing, "BEFORE");
assert_eq!(default_timing.event, "UPDATE");
assert_eq!(instead_of.timing, "INSTEAD OF");
assert_eq!(instead_of.event, "INSERT");
}
}

View File

@ -1,5 +1,7 @@
pub mod agent_driver;
pub mod clickhouse_driver;
pub mod cloudflare_d1;
pub use cloudflare_d1 as cloudflare_d1_driver;
pub mod duckdb_driver;
#[cfg(feature = "duckdb-bundled")]
pub mod duckdb_worker_process;

View File

@ -439,6 +439,8 @@ pub enum DatabaseType {
Iris,
#[serde(rename = "turso")]
Turso,
#[serde(rename = "cloudflare-d1")]
CloudflareD1,
#[serde(rename = "influxdb")]
InfluxDb,
#[serde(rename = "questdb")]
@ -761,7 +763,7 @@ impl ConnectionConfig {
},
DatabaseType::Redshift => Some("dev"),
DatabaseType::ClickHouse => Some("default"),
DatabaseType::Rqlite | DatabaseType::Turso => Some("main"),
DatabaseType::Rqlite | DatabaseType::Turso | DatabaseType::CloudflareD1 => Some("main"),
DatabaseType::Gaussdb | DatabaseType::OpenGauss => Some("postgres"),
DatabaseType::Kwdb => Some("defaultdb"),
DatabaseType::Vastbase => Some("postgres"),
@ -872,6 +874,7 @@ impl ConnectionConfig {
DatabaseType::ClickHouse => clickhouse_http_url(self, raw_host, port),
DatabaseType::Rqlite => rqlite_http_url(self, raw_host, port),
DatabaseType::Turso => turso_http_url(self, raw_host, port),
DatabaseType::CloudflareD1 => cloudflare_d1_api_url(self),
DatabaseType::SqlServer => {
format!("server=tcp:{host},{port};database={}", self.database.as_deref().unwrap_or("master"))
}
@ -1009,6 +1012,7 @@ impl ConnectionConfig {
DatabaseType::ClickHouse => clickhouse_http_url(self, raw_host, port),
DatabaseType::Rqlite => rqlite_http_url(self, raw_host, port),
DatabaseType::Turso => turso_http_url(self, raw_host, port),
DatabaseType::CloudflareD1 => cloudflare_d1_api_url(self),
DatabaseType::SqlServer => format!(
"server=tcp:{host},{port};user={};password={};database={}",
self.username,
@ -1695,6 +1699,12 @@ fn turso_http_url(config: &ConnectionConfig, host: &str, port: u16) -> String {
format!("{scheme}://{}:{port}", bracket_ipv6(trimmed))
}
fn cloudflare_d1_api_url(config: &ConnectionConfig) -> String {
let account_id = config.host.trim();
let database_id = config.database.as_deref().unwrap_or_default().trim();
format!("https://api.cloudflare.com/client/v4/accounts/{account_id}/d1/database/{database_id}")
}
fn trim_http_host_port(value: &str, default_port: u16) -> String {
let authority = value.trim_end_matches('/').split('/').next().unwrap_or(value).split('?').next().unwrap_or(value);
if authority.starts_with('[') && !authority.contains("]:") {
@ -1935,6 +1945,22 @@ mod tests {
assert_eq!(serde_json::from_str::<DatabaseType>("\"zookeeper\"").unwrap(), DatabaseType::ZooKeeper);
}
#[test]
fn cloudflare_d1_database_type_and_api_url_are_stable() {
assert_eq!(serde_json::to_string(&DatabaseType::CloudflareD1).unwrap(), "\"cloudflare-d1\"");
assert_eq!(serde_json::from_str::<DatabaseType>("\"cloudflare-d1\"").unwrap(), DatabaseType::CloudflareD1);
let mut config = mysql_config("", "secret-token", Some("database-id"));
config.db_type = DatabaseType::CloudflareD1;
config.host = "account-id".to_string();
config.port = 443;
assert_eq!(
config.connection_url(),
"https://api.cloudflare.com/client/v4/accounts/account-id/d1/database/database-id"
);
assert!(!config.connection_url().contains("secret-token"));
}
#[test]
fn connection_config_defaults_missing_agent_java_options_to_empty() {
let config: ConnectionConfig = serde_json::from_value(serde_json::json!({

View File

@ -862,6 +862,7 @@ fn should_discard_pool_after_query_timeout(db_type: Option<DatabaseType>) -> boo
| DatabaseType::SqlServer
| DatabaseType::Rqlite
| DatabaseType::Turso
| DatabaseType::CloudflareD1
| DatabaseType::Elasticsearch
| DatabaseType::Qdrant
| DatabaseType::Milvus
@ -1285,6 +1286,17 @@ pub async fn do_execute(
)
.await
}
PoolKind::CloudflareD1(client) => {
let client = client.clone();
let max_rows = options.max_rows;
drop(connections);
wait_for_query_opt(
cancel_token,
query_timeout,
db::cloudflare_d1_driver::execute_query_with_max_rows(&client, sql, max_rows),
)
.await
}
PoolKind::ClickHouse(client) => {
let client = client.clone();
let database = pool_key.split(':').nth(1).unwrap_or("default").to_string();
@ -1821,13 +1833,16 @@ pub async fn execute_multi_core_with_options_for_client(
.map(|results| results.into_iter().map(Into::into).collect());
}
let is_turso = {
let is_http_sqlite = {
let configs = state.configs.read().await;
configs.get(connection_id).is_some_and(|c| c.db_type == DatabaseType::Turso)
configs
.get(connection_id)
.is_some_and(|c| matches!(c.db_type, DatabaseType::Turso | DatabaseType::CloudflareD1))
};
// Turso sends all statements in one HTTP pipeline for transactional integrity.
if is_turso {
// HTTP SQLite providers send all statements in one request so the provider
// can preserve batch ordering and atomicity.
if is_http_sqlite {
let result =
execute_sql_statement_with_options(state, connection_id, database, sql, schema, cancel_token, options)
.await?;
@ -2313,6 +2328,7 @@ pub async fn execute_statements_in_transaction(
PoolKind::Postgres(pg) => TxPath::Pg(pg.clone()),
PoolKind::Mysql(mp, _mode) => TxPath::Mysql(mp.clone(), false),
PoolKind::Sqlite(sq) => TxPath::Sqlite(sq.clone()),
PoolKind::CloudflareD1(client) => TxPath::CloudflareD1(client.clone()),
PoolKind::ClickHouse(_)
| PoolKind::Rqlite(_)
| PoolKind::Turso(_)
@ -2351,6 +2367,15 @@ pub async fn execute_statements_in_transaction(
exec_tx_mysql_inner(state, &pool_key, pool, statements, start, operation_budget.clone()).await
}
Some(TxPath::Sqlite(pool)) => exec_tx_sqlite_inner(pool, statements, start).await,
Some(TxPath::CloudflareD1(client)) => {
let sql = statements.join(";\n");
wait_for_query_opt(
None,
operation_budget.query_timeout,
db::cloudflare_d1_driver::execute_query_with_max_rows(&client, &sql, None),
)
.await
}
Some(TxPath::Explicit) => {
let mysql_dialect = connection_mysql_query_dialect(state, connection_id).await;
exec_tx_explicit_inner(state, &pool_key, mysql_dialect, Some(database), statements, schema, start).await
@ -2376,6 +2401,7 @@ enum TxPath {
Pg(deadpool_postgres::Pool),
Mysql(mysql_async::Pool, bool),
Sqlite(db::sqlite::SqliteHandle),
CloudflareD1(db::cloudflare_d1_driver::CloudflareD1Client),
Explicit,
None,
}

View File

@ -808,6 +808,7 @@ async fn list_databases_once(state: &AppState, connection_id: &str) -> Result<Ve
drop(connections);
client.list_databases().await
}
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::list_databases(client).await,
_ => Ok(vec![]),
}
}
@ -1908,6 +1909,9 @@ async fn list_tables_once(
.await
.map(|infos| collection_names_to_tables(infos.into_iter().map(|i| i.name).collect(), "COLLECTION"))
.map(|tables| filter_table_infos(tables, filter, limit, offset, object_types)),
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::list_tables(client, schema)
.await
.map(|tables| filter_table_infos(tables, filter, limit, offset, object_types)),
_ => Ok(vec![]),
}
}
@ -4197,6 +4201,9 @@ pub async fn get_columns_core(
PoolKind::Rqlite(client) => {
db::rqlite_driver::get_columns(client, schema, table).await.map(deduplicate_column_infos)
}
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::get_columns(client, schema, table)
.await
.map(deduplicate_column_infos),
_ => Ok(vec![]),
}
})
@ -4293,6 +4300,7 @@ pub async fn list_indexes_core(
PoolKind::Sqlite(p) => db::sqlite::list_indexes(p, schema, table).await,
PoolKind::Rqlite(client) => db::rqlite_driver::list_indexes(client, schema, table).await,
PoolKind::MongoDb(client) => db::mongo_driver::list_indexes(client, database, table).await,
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::list_indexes(client, schema, table).await,
_ => Ok(vec![]),
}
})
@ -4339,6 +4347,7 @@ pub async fn list_foreign_keys_core(
PoolKind::Postgres(p) => db::postgres::list_foreign_keys(p, schema, table).await,
PoolKind::Sqlite(p) => db::sqlite::list_foreign_keys(p, schema, table).await,
PoolKind::Rqlite(client) => db::rqlite_driver::list_foreign_keys(client, schema, table).await,
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::list_foreign_keys(client, schema, table).await,
_ => Ok(vec![]),
}
})
@ -4383,6 +4392,7 @@ pub async fn list_triggers_core(
PoolKind::Postgres(p) => db::postgres::list_triggers(p, schema, table).await,
PoolKind::Sqlite(p) => db::sqlite::list_triggers(p, schema, table).await,
PoolKind::Rqlite(client) => db::rqlite_driver::list_triggers(client, schema, table).await,
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::list_triggers(client, schema, table).await,
_ => Ok(vec![]),
}
})
@ -4655,6 +4665,7 @@ pub async fn get_table_ddl_core(
PoolKind::Postgres(p) => pg_ddl(p, schema, table).await,
PoolKind::Sqlite(p) => sqlite_ddl(p, table).await,
PoolKind::Rqlite(client) => db::rqlite_driver::table_ddl(client, table).await,
PoolKind::CloudflareD1(client) => db::cloudflare_d1_driver::table_ddl(client, table).await,
_ => Err("DDL not supported for this database type".to_string()),
}
}
@ -5074,6 +5085,9 @@ async fn get_object_source_once(
.await?;
first_string_cell(result)?
}
PoolKind::CloudflareD1(client) => {
return db::cloudflare_d1_driver::object_source(client, name, &object_type).await;
}
_ => return Err("Object source is not supported for this database type".to_string()),
}
}

View File

@ -744,6 +744,17 @@ pub fn build_import_insert_batch_from_rows(
if rows.is_empty() {
return Ok(None);
}
if *db_type == DatabaseType::CloudflareD1 {
return crate::db::cloudflare_d1::build_streaming_import_insert_batch(
rows,
columns,
mappings,
target_column_types,
table,
schema,
rows.len(),
);
}
let mapped = mapping_indexes_for_columns(columns, mappings)?;
let target_columns = mapped.iter().map(|(_, target)| target.clone()).collect::<Vec<_>>();
let column_types = target_columns
@ -781,6 +792,17 @@ pub fn build_import_insert_batches(
db_type: &DatabaseType,
batch_size: usize,
) -> Result<Vec<ImportSqlBatch>, String> {
if *db_type == DatabaseType::CloudflareD1 {
return crate::db::cloudflare_d1::build_import_insert_batches(
&data.rows,
&data.columns,
mappings,
target_column_types,
table,
schema,
batch_size.clamp(1, 100),
);
}
let mapped = mapping_indexes(data, mappings)?;
let columns = mapped.iter().map(|(_, target)| target.clone()).collect::<Vec<_>>();
let column_types = columns
@ -819,7 +841,7 @@ pub fn build_import_insert_batches(
pub fn truncate_sql(table: &str, schema: &str, db_type: &DatabaseType) -> String {
let full_table = qualified_table(table, schema, db_type);
match db_type {
DatabaseType::Sqlite => format!("DELETE FROM {full_table}"),
DatabaseType::Sqlite | DatabaseType::CloudflareD1 => format!("DELETE FROM {full_table}"),
_ => format!("TRUNCATE TABLE {full_table}"),
}
}
@ -937,7 +959,7 @@ fn text_data_type(db_type: &DatabaseType) -> &'static str {
fn integer_data_type(db_type: &DatabaseType) -> &'static str {
match db_type {
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso => "INTEGER",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso | DatabaseType::CloudflareD1 => "INTEGER",
DatabaseType::Oracle | DatabaseType::OceanbaseOracle | DatabaseType::Dameng => "NUMBER(19)",
DatabaseType::ClickHouse => "Int64",
_ => "BIGINT",
@ -954,7 +976,7 @@ fn decimal_data_type(db_type: &DatabaseType) -> &'static str {
| DatabaseType::Highgo
| DatabaseType::Kwdb
| DatabaseType::Vastbase => "DOUBLE PRECISION",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso => "REAL",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso | DatabaseType::CloudflareD1 => "REAL",
DatabaseType::Oracle | DatabaseType::OceanbaseOracle | DatabaseType::Dameng => "BINARY_DOUBLE",
DatabaseType::ClickHouse => "Float64",
_ => "DOUBLE",
@ -970,7 +992,7 @@ fn boolean_data_type(db_type: &DatabaseType) -> &'static str {
| DatabaseType::Sundb
| DatabaseType::Databend => "TINYINT(1)",
DatabaseType::SqlServer => "BIT",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso => "INTEGER",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso | DatabaseType::CloudflareD1 => "INTEGER",
DatabaseType::Oracle | DatabaseType::OceanbaseOracle | DatabaseType::Dameng => "NUMBER(1)",
DatabaseType::ClickHouse => "UInt8",
_ => "BOOLEAN",
@ -979,7 +1001,7 @@ fn boolean_data_type(db_type: &DatabaseType) -> &'static str {
fn date_data_type(db_type: &DatabaseType) -> &'static str {
match db_type {
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso => "TEXT",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso | DatabaseType::CloudflareD1 => "TEXT",
DatabaseType::ClickHouse => "Date",
_ => "DATE",
}
@ -994,7 +1016,7 @@ fn timestamp_data_type(db_type: &DatabaseType) -> &'static str {
| DatabaseType::Sundb
| DatabaseType::Databend => "DATETIME",
DatabaseType::SqlServer => "DATETIME2",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso => "TEXT",
DatabaseType::Sqlite | DatabaseType::Rqlite | DatabaseType::Turso | DatabaseType::CloudflareD1 => "TEXT",
DatabaseType::ClickHouse => "DateTime64",
_ => "TIMESTAMP",
}
@ -1382,6 +1404,7 @@ where
};
let effective_batch_size = match db_type {
DatabaseType::Oracle | DatabaseType::OceanbaseOracle => 1,
DatabaseType::CloudflareD1 => batch_size.clamp(1, 100),
_ => batch_size.max(1),
};
let mut rows_imported = 0;

View File

@ -1020,6 +1020,7 @@ pub fn escape_value_typed(val: &serde_json::Value, db_type: &DatabaseType, colum
serde_json::Value::Bool(b) => match db_type {
DatabaseType::Mysql
| DatabaseType::Sqlite
| DatabaseType::CloudflareD1
| DatabaseType::DuckDb
| DatabaseType::Doris
| DatabaseType::StarRocks => {
@ -1881,7 +1882,10 @@ pub fn generate_upsert_typed(
}
match db_type {
db_type if is_postgres_transfer_dialect(db_type) => {
db_type
if is_postgres_transfer_dialect(db_type)
|| matches!(db_type, DatabaseType::Sqlite | DatabaseType::CloudflareD1 | DatabaseType::DuckDb) =>
{
let pk_list = pk_columns.iter().map(|c| quote_identifier(c, db_type)).collect::<Vec<_>>().join(", ");
let mut sql = format!("INSERT INTO {full_table} ({col_list}) VALUES\n{}", value_rows.join(",\n"));
if non_pk_columns.is_empty() {
@ -2069,6 +2073,10 @@ fn generate_transfer_write_sql_batches(
}
let max_rows = max_transfer_write_rows(db_type, mode);
let max_sql_bytes = match db_type {
DatabaseType::CloudflareD1 => crate::db::cloudflare_d1::MAX_SQL_STATEMENT_BYTES,
_ => MAX_TRANSFER_WRITE_SQL_BYTES,
};
let mut statements = Vec::new();
let mut start = 0;
@ -2096,7 +2104,7 @@ fn generate_transfer_write_sql_batches(
db_type,
pk_columns,
);
if candidate.len() > MAX_TRANSFER_WRITE_SQL_BYTES && !accepted.is_empty() {
if candidate.len() > max_sql_bytes && !accepted.is_empty() {
break;
}
accepted = candidate;
@ -3896,7 +3904,9 @@ where
if request.mode == TransferMode::Overwrite {
let full_table = qualified_table(&target_table, &request.target_schema, target_db_type);
let truncate_sql = match target_db_type {
DatabaseType::Sqlite | DatabaseType::DuckDb => format!("DELETE FROM {full_table}"),
DatabaseType::Sqlite | DatabaseType::CloudflareD1 | DatabaseType::DuckDb => {
format!("DELETE FROM {full_table}")
}
_ => format!("TRUNCATE TABLE {full_table}"),
};
execute_on_pool(state, target_pool_key, &truncate_sql)
@ -4174,7 +4184,9 @@ where
if request.mode == TransferMode::Overwrite {
let full_table = qualified_table(&target_table, &request.target_schema, target_db_type);
let truncate_sql = match target_db_type {
DatabaseType::Sqlite | DatabaseType::DuckDb => format!("DELETE FROM {full_table}"),
DatabaseType::Sqlite | DatabaseType::CloudflareD1 | DatabaseType::DuckDb => {
format!("DELETE FROM {full_table}")
}
_ => format!("TRUNCATE TABLE {full_table}"),
};
execute_on_pool(state, target_pool_key, &truncate_sql).await.map_err(|e| format!("Failed to truncate: {e}"))?;

View File

@ -16,7 +16,7 @@ description: 了解 DBX 可以连接哪些数据库,以及每类高级功能
| MySQL 兼容类型 | MariaDB、TiDB、OceanBase、Doris、SelectDB、StarRocks、Manticore Search、GoldenDB、TDSQL、PolarDB、GreatSQL、自定义 MySQL | 复用 MySQL 风格连接能力 |
| PostgreSQL 兼容类型 | openGauss、KingBase、HighGo、Vastbase、CockroachDB、自定义 PostgreSQL | 复用 PostgreSQL 风格连接能力 |
| 文件型数据库 | SQLite、DuckDB、Microsoft Access、RQLite、Turso | 选择本地数据库文件或 HTTP 端点,不填写主机和端口 |
| 时序与边缘数据库 | InfluxDB、IoTDB、TDengine、Databend、QuestDB | 面向时序、IoT 或边缘场景优化 |
| 时序与边缘数据库 | Cloudflare D1、InfluxDB、IoTDB、TDengine、Databend、QuestDB | 托管 SQLite、时序、IoT 或边缘场景 |
| Agent/JDBC 扩展类型 | H2、Snowflake、Trino、PrestoSQL、Hive、DB2、Informix、Neo4j、Cassandra、BigQuery、Kylin、SunDB、XuguDB、Databricks、IRIS、JDBC | 功能覆盖取决于对应驱动路径 |
## 默认端口
@ -49,6 +49,7 @@ description: 了解 DBX 可以连接哪些数据库,以及每类高级功能
| Exasol | 8563 |
| SAP HANA | 39015 |
| Databricks | 443 |
| Cloudflare D1 REST API | 443 |
| InfluxDB | 8086 |
| IoTDB | 6667 |
| etcd | 2379 |
@ -81,14 +82,14 @@ DBX 只会在具备足够元数据和 SQL 生成能力的数据库上启用高
| 功能 | 支持类型 |
| ------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Schema 感知树 | PostgreSQL、SQL Server、Oracle、Redshift、DM、GaussDB、KWDB、KingBase、HighGo、Vastbase、JDBC、H2、Snowflake、Trino、DB2、TDengine、XuguDB、Teradata、Vertica、Exasol、Firebird、SAP HANA、YashanDB、GBase |
| ER 图 | MySQL、PostgreSQL、SQLite、SQL Server、Oracle、Redshift、DM、GaussDB、KWDB、KingBase、HighGo、Vastbase、GoldenDB、Access、H2、DB2、Teradata、Vertica、Firebird、Exasol、GBase、YashanDB |
| ER 图 | MySQL、PostgreSQL、SQLite、Cloudflare D1、SQL Server、Oracle、Redshift、DM、GaussDB、KWDB、KingBase、HighGo、Vastbase、GoldenDB、Access、H2、DB2、Teradata、Vertica、Firebird、Exasol、GBase、YashanDB |
| 数据库搜索 | MySQL、PostgreSQL、SQLite、SQL Server、Oracle、Redshift、DuckDB、ClickHouse、DM、GaussDB、KWDB、KingBase、HighGo、Vastbase、GoldenDB、Access、H2、Snowflake、Trino、Hive、DB2、Informix、Neo4j、Cassandra、BigQuery、Kylin、SunDB、TDengine、XuguDB、Databricks、Teradata、Vertica |
| 表数据导入 | MySQL、PostgreSQL、SQLite、DuckDB、ClickHouse、SQL Server、Oracle、Doris、StarRocks、Redshift、DM、GaussDB、KWDB、KingBase、HighGo、Vastbase、GoldenDB、Access、Databend |
| 表数据导入 | MySQL、PostgreSQL、SQLite、Cloudflare D1、DuckDB、ClickHouse、SQL Server、Oracle、Doris、StarRocks、Redshift、DM、GaussDB、KWDB、KingBase、HighGo、Vastbase、GoldenDB、Access、Databend |
| 几何字段地图预览 | PostgreSQL、HighGo、KingBase、Vastbase、openGauss、GaussDBPostGIS `geometry` / `geography` 列在结果集中以 WKT 显示,并提供网格工具栏的"地图预览"按钮) |
| 表结构编辑器 | MySQL、PostgreSQL、SQLite、SQL Server、KWDB |
| 创建数据库 | MySQL、PostgreSQL、SQL Server、ClickHouse、Oracle、DM、GaussDB、KWDB、Doris、StarRocks、Redshift、Teradata、Vertica |
| 字段血缘 | MySQL、PostgreSQL、SQLite、SQL Server、Oracle、Redshift、DM、GaussDB、KWDB |
| 数据传输 | MySQL、PostgreSQL、SQLite、SQL Server、Oracle、ClickHouse、DuckDB、DM、GaussDB、KWDB |
| 字段血缘 | MySQL、PostgreSQL、SQLite、Cloudflare D1、SQL Server、Oracle、Redshift、DM、GaussDB、KWDB |
| 数据传输 | MySQL、PostgreSQL、SQLite、Cloudflare D1、SQL Server、Oracle、ClickHouse、DuckDB、DM、GaussDB、KWDB |
| 不支持 SQL 文件执行 | Redis、MongoDB、Elasticsearch、etcd、Nacos |
<Callout type="warn">如果某个数据库可以连接,但没有出现在某个功能行中,该功能可能会被隐藏或受限。这是为了避免 DBX 生成无法可靠审查的 SQL。</Callout>
@ -99,6 +100,10 @@ DBX 只会在具备足够元数据和 SQL 生成能力的数据库上启用高
DBX 也可以解析常见连接 URL例如 MySQL、PostgreSQL、Redis、MongoDB、ClickHouse、SQL Server、Oracle、Elasticsearch、DM、GaussDB、KWDB、openGauss、TDengine、XuguDB、Access、Teradata、Vertica、Firebird、Exasol、GBase、YashanDB、SAP HANA、InfluxDB、QuestDB、IoTDB、etcd、Nacos、RQLite 和 Databricks。
### Cloudflare D1
新建 **Cloudflare D1** 连接时填写 Account ID、D1 Database IDUUID和 API Token。浏览元数据和查询至少需要 `D1 Read` 权限插入、更新、DDL、导入或传输还需要 `D1 Write` 权限。DBX 会直接调用 Cloudflare HTTPS API因此该类型不提供 SSH 或代理传输层。
### Nacos 控制台
Nacos 连接会打开专用管理控制台,用于管理命名空间、配置、服务、实例、配置历史和 Raw API 请求。创建连接时请填写 Nacos 控制台/admin API 地址。Nacos 3 Docker 部署通常使用 `8085` 暴露控制台;旧版本也可能与服务端口共用 `8848`。

View File

@ -16,7 +16,7 @@ The connection dialog exposes ready-to-use profiles with default ports and drive
| MySQL-compatible profiles | MariaDB, TiDB, OceanBase, Doris, SelectDB, StarRocks, Manticore Search, GoldenDB, TDSQL, PolarDB, GreatSQL, custom MySQL | Reuse MySQL-style connection handling where the engine speaks a compatible protocol |
| PostgreSQL-compatible profiles | openGauss, KingBase, HighGo, Vastbase, CockroachDB, custom PostgreSQL | Reuse PostgreSQL-style connection handling where the engine speaks a compatible protocol |
| File-based engines | SQLite, DuckDB, Microsoft Access, RQLite, Turso | Choose a local database file or HTTP endpoint instead of host and port |
| Time-series and edge | InfluxDB, IoTDB, TDengine, Databend, QuestDB | Optimized for time-series, IoT, or edge workloads |
| Time-series and edge | Cloudflare D1, InfluxDB, IoTDB, TDengine, Databend, QuestDB | Managed SQLite, time-series, IoT, or edge workloads |
| Agent/JDBC-oriented engines | H2, Snowflake, Trino, PrestoSQL, Hive, DB2, Informix, Neo4j, Cassandra, BigQuery, Kylin, SunDB, XuguDB, Databricks, IRIS, JDBC | Feature coverage depends on the driver path used by that engine |
## Default Ports
@ -49,6 +49,7 @@ The connection dialog exposes ready-to-use profiles with default ports and drive
| Exasol | 8563 |
| SAP HANA | 39015 |
| Databricks | 443 |
| Cloudflare D1 REST API | 443 |
| InfluxDB | 8086 |
| IoTDB | 6667 |
| etcd | 2379 |
@ -81,14 +82,14 @@ DBX intentionally enables advanced workflows only where the app has enough metad
| Feature | Supported Types |
| ------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Schema-aware tree | PostgreSQL, SQL Server, Oracle, Redshift, DM, GaussDB, KWDB, KingBase, HighGo, Vastbase, JDBC, H2, Snowflake, Trino, DB2, TDengine, XuguDB, Teradata, Vertica, Exasol, Firebird, SAP HANA, YashanDB, GBase |
| ER diagram | MySQL, PostgreSQL, SQLite, SQL Server, Oracle, Redshift, DM, GaussDB, KWDB, KingBase, HighGo, Vastbase, GoldenDB, Access, H2, DB2, Teradata, Vertica, Firebird, Exasol, GBase, YashanDB |
| ER diagram | MySQL, PostgreSQL, SQLite, Cloudflare D1, SQL Server, Oracle, Redshift, DM, GaussDB, KWDB, KingBase, HighGo, Vastbase, GoldenDB, Access, H2, DB2, Teradata, Vertica, Firebird, Exasol, GBase, YashanDB |
| Database search | MySQL, PostgreSQL, SQLite, SQL Server, Oracle, Redshift, DuckDB, ClickHouse, DM, GaussDB, KWDB, KingBase, HighGo, Vastbase, GoldenDB, Access, H2, Snowflake, Trino, Hive, DB2, Informix, Neo4j, Cassandra, BigQuery, Kylin, SunDB, TDengine, XuguDB, Databricks, Teradata, Vertica |
| Table import | MySQL, PostgreSQL, SQLite, DuckDB, ClickHouse, SQL Server, Oracle, Doris, StarRocks, Redshift, DM, GaussDB, KWDB, KingBase, HighGo, Vastbase, GoldenDB, Access, Databend |
| Table import | MySQL, PostgreSQL, SQLite, Cloudflare D1, DuckDB, ClickHouse, SQL Server, Oracle, Doris, StarRocks, Redshift, DM, GaussDB, KWDB, KingBase, HighGo, Vastbase, GoldenDB, Access, Databend |
| Geometry map preview | PostgreSQL, HighGo, KingBase, Vastbase, openGauss, GaussDB (PostGIS `geometry` / `geography` columns are returned as WKT in the result grid and unlock the grid-toolbar map preview button) |
| Table structure editor | MySQL, PostgreSQL, SQLite, SQL Server, KWDB |
| Create database | MySQL, PostgreSQL, SQL Server, ClickHouse, Oracle, DM, GaussDB, KWDB, Doris, StarRocks, Redshift, Teradata, Vertica |
| Field lineage | MySQL, PostgreSQL, SQLite, SQL Server, Oracle, Redshift, DM, GaussDB, KWDB |
| Data transfer | MySQL, PostgreSQL, SQLite, SQL Server, Oracle, ClickHouse, DuckDB, DM, GaussDB, KWDB |
| Field lineage | MySQL, PostgreSQL, SQLite, Cloudflare D1, SQL Server, Oracle, Redshift, DM, GaussDB, KWDB |
| Data transfer | MySQL, PostgreSQL, SQLite, Cloudflare D1, SQL Server, Oracle, ClickHouse, DuckDB, DM, GaussDB, KWDB |
| SQL file execution unavailable | Redis, MongoDB, Elasticsearch, etcd, Nacos |
<Callout type="warn">If a database can connect but does not appear in a feature row, that feature may be hidden or limited for that engine. This prevents DBX from generating SQL it cannot review reliably.</Callout>
@ -99,6 +100,10 @@ Most network databases support host, port, username, password, default database,
DBX can also parse common connection URLs for engines such as MySQL, PostgreSQL, Redis, MongoDB, ClickHouse, SQL Server, Oracle, Elasticsearch, DM, GaussDB, KWDB, openGauss, TDengine, XuguDB, Access, Teradata, Vertica, Firebird, Exasol, GBase, YashanDB, SAP HANA, InfluxDB, QuestDB, IoTDB, etcd, Nacos, RQLite, and Databricks.
### Cloudflare D1
Create a **Cloudflare D1** connection with the Account ID, D1 Database ID (UUID), and an API Token. The token needs `D1 Read` permission for browsing and queries; add `D1 Write` for inserts, updates, DDL, imports, or transfers. DBX uses Cloudflare's HTTPS API directly, so SSH and proxy transport layers are not available for this profile.
### Nacos Console
Nacos connections open a dedicated admin console for namespaces, configs, services, instances, config history, and raw API requests. Use the Nacos console/admin API address when creating the connection. Nacos 3 Docker deployments usually expose that console on `8085`; older deployments may share `8848` with the service port.

View File

@ -77,6 +77,21 @@ test("redacts network endpoint labels for quick connection cards", () => {
);
});
test("presents Cloudflare D1 account and database identifiers without treating them as a host and port", () => {
const d1Connection: ConnectionConfig = {
...baseConnection,
db_type: "cloudflare-d1",
driver_profile: undefined,
driver_label: "Cloudflare D1",
host: "account-id",
port: 443,
database: "database-id",
};
assert.equal(connectionEndpointLabel(d1Connection), "account-id/database-id");
assert.equal(connectionRedactedEndpointLabel(d1Connection), "***/***");
});
test("redacts host-like quick connection names", () => {
assert.equal(
connectionRedactedNameLabel({

View File

@ -82,16 +82,17 @@ test("ZooKeeper connections do not offer visible database selection", () => {
assert.equal(connectionCanChooseVisibleDatabases(config({ db_type: "zookeeper" })), false);
});
test("Cloudflare D1 does not offer a visible database filter for its fixed main namespace", () => {
assert.equal(connectionCanChooseVisibleDatabases(config({ db_type: "cloudflare-d1" })), false);
});
test("OceanBase Oracle uses schema filtering for visible object selection", () => {
assert.equal(connectionUsesVisibleSchemaFilter(config({ db_type: "oceanbase-oracle" })), true);
assert.equal(connectionUsesVisibleSchemaFilter(config({ db_type: "mysql", driver_profile: "oceanbase" })), false);
});
test("Dameng default SYSDBA user remains selectable", () => {
assert.deepEqual(
filterDatabaseNamesForConnection(["SYS", "SYSDBA", "SYSAUDITOR"], config({ db_type: "dameng" })),
["SYSDBA"],
);
assert.deepEqual(filterDatabaseNamesForConnection(["SYS", "SYSDBA", "SYSAUDITOR"], config({ db_type: "dameng" })), ["SYSDBA"]);
});
test("visible database selection is stale when connection target changes", () => {

View File

@ -14,6 +14,12 @@ test("没有默认数据库且无候选项时返回空字符串", () => {
assert.equal(resolveDefaultDatabase({ database: undefined }, []), "");
});
test("Cloudflare D1 使用 SQLite main 命名空间而不是连接凭据中的 Database ID", () => {
assert.equal(resolveDefaultDatabase({ db_type: "cloudflare-d1", database: "d1-database-uuid" }, []), "main");
assert.equal(isDefaultDatabase({ db_type: "cloudflare-d1", database: "d1-database-uuid" }, "main"), true);
assert.equal(isDefaultDatabase({ db_type: "cloudflare-d1", database: "d1-database-uuid" }, "d1-database-uuid"), false);
});
test("判断当前数据库是否为默认数据库", () => {
assert.equal(isDefaultDatabase({ database: "analytics" }, "analytics"), true);
assert.equal(isDefaultDatabase({ database: "analytics" }, "app"), false);

View File

@ -762,6 +762,28 @@ test("suggests SQL Server IIF and CHOOSE scalar functions", () => {
);
});
test("only suggests Cloudflare D1 supported common functions", () => {
const supported = buildSqlCompletionItems("SELECT SQ", "SELECT SQ".length, {
tables,
columnsByTable,
databaseType: "cloudflare-d1",
});
const unsupportedNow = buildSqlCompletionItems("SELECT NO", "SELECT NO".length, {
tables,
columnsByTable,
databaseType: "cloudflare-d1",
});
const unsupportedPower = buildSqlCompletionItems("SELECT PO", "SELECT PO".length, {
tables,
columnsByTable,
databaseType: "cloudflare-d1",
});
assert.ok(supported.some((item) => item.label === "SQRT"));
assert.ok(!unsupportedNow.some((item) => item.label === "NOW"));
assert.ok(!unsupportedPower.some((item) => item.label === "POWER"));
});
test("suggests SQL Server IDENTITY_INSERT after SET", () => {
const sql = "set iden";
const items = buildSqlCompletionItems(sql, sql.length, {

View File

@ -91,7 +91,7 @@ function formatRedisCommandToolResult(result: RedisCommandResult) {
}
export const DBX_CONNECTION_TYPE_DESCRIPTION =
"Database type: postgres, mysql, sqlite, rqlite, redis, duckdb, clickhouse, sqlserver, mongodb, oracle, elasticsearch, etcd, doris, starrocks, manticoresearch, milvus, qdrant, weaviate, chromadb, redshift, dameng, kingbase, highgo, vastbase, goldendb, databend, gaussdb, kwdb, yashandb, databricks, saphana, teradata, vertica, firebird, exasol, opengauss, oceanbase-oracle, questdb, gbase, h2, snowflake, trino, prestosql, hive, spark, db2, informix, influxdb, iris, neo4j, cassandra, bigquery, kylin, sundb, oscar, tdengine, iotdb, xugu, zookeeper, jdbc, access, mq";
"Database type: postgres, mysql, sqlite, rqlite, cloudflare-d1, redis, duckdb, clickhouse, sqlserver, mongodb, oracle, elasticsearch, etcd, doris, starrocks, manticoresearch, milvus, qdrant, weaviate, chromadb, redshift, dameng, kingbase, highgo, vastbase, goldendb, databend, gaussdb, kwdb, yashandb, databricks, saphana, teradata, vertica, firebird, exasol, opengauss, oceanbase-oracle, questdb, gbase, h2, snowflake, trino, prestosql, hive, spark, db2, informix, influxdb, iris, neo4j, cassandra, bigquery, kylin, sundb, oscar, tdengine, iotdb, xugu, zookeeper, jdbc, access, mq";
const FILE_CAPABLE_CONNECTION_TYPES = new Set(["sqlite", "duckdb", "access", "h2"]);
interface McpScope {
@ -127,12 +127,7 @@ async function loadScopedConnections(backend: Backend, scope: McpScope): Promise
return connections.filter((config) => connectionMatchesScope(config, scope));
}
async function resolveConnection(
backend: Backend,
scope: McpScope,
requestedId?: string,
requestedName?: string,
): Promise<{ config?: ConnectionConfig; error?: ReturnType<typeof toolError> }> {
async function resolveConnection(backend: Backend, scope: McpScope, requestedId?: string, requestedName?: string): Promise<{ config?: ConnectionConfig; error?: ReturnType<typeof toolError> }> {
// connection_id takes priority over connection_name when both are provided.
if (requestedId?.trim()) {
const connections = await backend.loadConnections();
@ -153,10 +148,7 @@ async function resolveConnection(
if (matching.length > 1) {
const lines = matching.map((c) => `- ${c.id}: ${c.db_type} @ ${c.host}:${c.port}`);
return {
error: toolError(
"AMBIGUOUS_CONNECTION",
`Multiple connections found with name "${requestedName}". Please specify connection_id:\n${lines.join("\n")}`,
),
error: toolError("AMBIGUOUS_CONNECTION", `Multiple connections found with name "${requestedName}". Please specify connection_id:\n${lines.join("\n")}`),
};
}
return { config: matching[0] };
@ -340,11 +332,11 @@ export function createDbxMcpServer(backend: Backend, options: { isWebMode?: bool
{
name: z.string().describe("Connection name"),
db_type: z.string().describe(DBX_CONNECTION_TYPE_DESCRIPTION),
host: z.string().describe("Database host"),
host: z.string().describe("Database host; for cloudflare-d1, use the Cloudflare Account ID"),
port: z.number().optional().describe("Database port (TDengine defaults to 6041, IoTDB defaults to 6667, XuguDB defaults to 5138)"),
username: z.string().default("").describe("Username"),
password: z.string().default("").describe("Password"),
database: z.string().optional().describe("Default database name"),
password: z.string().default("").describe("Password; for cloudflare-d1, use the API Token"),
database: z.string().optional().describe("Default database name; for cloudflare-d1, use the D1 Database ID"),
ssl: z.boolean().default(false).describe("Enable SSL"),
driver_profile: z.string().optional().describe("Driver profile (e.g. 'gbase8a', 'gbase8s')"),
},
@ -354,6 +346,7 @@ export function createDbxMcpServer(backend: Backend, options: { isWebMode?: bool
const DEFAULT_PORTS: Record<string, number> = {
kwdb: 26257,
rqlite: 4001,
"cloudflare-d1": 443,
tdengine: 6041,
oscar: 2003,
iotdb: 6667,

View File

@ -13,6 +13,7 @@ export function isDirectQueryType(dbType: string): dbType is DirectQueryType {
}
export const BRIDGE_REQUIRED_TYPES = [
"cloudflare-d1",
"redis",
"mongodb",
"duckdb",

View File

@ -823,6 +823,9 @@ pub async fn test_connection(state: State<'_, Arc<AppState>>, config: Connection
.await
.map(|_| "Connection successful".to_string())
}
DatabaseType::CloudflareD1 => db::cloudflare_d1_driver::connect(&config, connect_timeout)
.await
.map(|_| "Connection successful".to_string()),
DatabaseType::InfluxDb => {
let client = db::influxdb_driver::InfluxdbClient::new_for_config(&url, &config, connect_timeout)?;
db::influxdb_driver::test_connection(&client, connect_timeout)
@ -1110,6 +1113,9 @@ pub async fn connect_db(
db::turso_driver::test_connection(&client, connect_timeout).await?;
PoolKind::Turso(client)
}
DatabaseType::CloudflareD1 => {
PoolKind::CloudflareD1(db::cloudflare_d1_driver::connect(&db_config, connect_timeout).await?)
}
DatabaseType::InfluxDb => {
let client = db::influxdb_driver::InfluxdbClient::new_for_config(&url, &db_config, connect_timeout)?;
db::influxdb_driver::test_connection(&client, connect_timeout).await?;