From 0887b96236191fdfec00dadb565f3f632c4952d1 Mon Sep 17 00:00:00 2001 From: t8y2 <1156263951@qq.com> Date: Fri, 5 Jun 2026 01:02:07 +0800 Subject: [PATCH] feat: add rqlite database support --- apps/desktop/public/icons/database/rqlite.png | Bin 0 -> 5029 bytes .../connection/ConnectionDialog.vue | 3 + .../src/components/icons/DatabaseIcon.vue | 1 + .../src/components/layout/ContentArea.vue | 1 + .../src/components/objects/ObjectBrowser.vue | 1 + .../desktop/src/lib/connectionPresentation.ts | 3 + .../desktop/src/lib/databaseCapabilitySets.ts | 7 + .../desktop/src/lib/databaseFeatureSupport.ts | 2 +- .../src/lib/databaseTableDataCapabilities.ts | 1 + apps/desktop/src/lib/objectRenameSql.ts | 2 +- apps/desktop/src/lib/sqlCompletion.ts | 2 + .../src/lib/tableStructureCapabilities.ts | 1 + .../src/lib/tableStructureEditorState.ts | 1 + apps/desktop/src/types/database.ts | 1 + .../assets/database-drivers.manifest.json | 10 + crates/dbx-core/src/connection.rs | 14 + crates/dbx-core/src/database_capabilities.rs | 1 + crates/dbx-core/src/db/mod.rs | 1 + crates/dbx-core/src/db/rqlite_driver.rs | 455 ++++++++++++++++++ crates/dbx-core/src/models/connection.rs | 29 ++ crates/dbx-core/src/query.rs | 15 +- crates/dbx-core/src/schema.rs | 14 + crates/dbx-core/src/table_structure_sql.rs | 30 +- .../connectionUrlPlaceholder.test.ts | 1 + .../app-tests/databaseCapabilities.test.ts | 10 + packages/app-tests/objectRenameSql.test.ts | 2 + .../tableStructureCapabilities.test.ts | 4 +- packages/cli/tests/bin.test.ts | 3 +- packages/cli/tests/cli.test.ts | 3 +- packages/mcp-server/src/index.ts | 4 +- packages/node-core/src/database.ts | 57 ++- packages/node-core/src/diagnostics.ts | 1 + .../node-core/tests/rqlite-direct.test.ts | 111 +++++ src-tauri/src/commands/connection.rs | 25 + 34 files changed, 804 insertions(+), 12 deletions(-) create mode 100644 apps/desktop/public/icons/database/rqlite.png create mode 100644 crates/dbx-core/src/db/rqlite_driver.rs create mode 100644 packages/node-core/tests/rqlite-direct.test.ts diff --git a/apps/desktop/public/icons/database/rqlite.png b/apps/desktop/public/icons/database/rqlite.png new file mode 100644 index 0000000000000000000000000000000000000000..a85d1215f84bc4f18320df24037cddac3f433ca9 GIT binary patch literal 5029 zcmai22T&B>lU|k~EZLQuL2{HFBnQbsf*?UbKt!Sf5+p4wQAJ#G20{D_{9wraObUx@M+p`px%szkdCCrY70a+<=yfoeBg3(Ha@*SrK*2 z-;}TH$y8k5GY)TXorJ9fBufIgFwMIL7+_+5J)u#1Y!v&_J%!8;vvNvTFC2tZB zc^er;RaqHTSw-=iH&t)mERcDB^xpu!{_b9n!v1%_whxyP5g_tk9)i4lJpzMVeFOdn zMovyuUhY4g;B1*pKp<$mk)F0q=v455i48ax z)#K;Ab-f`x2p|O|KohMaR$rp|E=D&Op3zZ|Ql|CxBvz~JD7`g=-@{-ii zB4|h{gJSz03yd5FU7eJE*u$^_xkWqulf_CKrAO!p%?r)w1yYUwTX+=UbF`nB7xn;D zOixy@`1)qpF3NgMI(l~ebpi> z0c>=Rn%0kkAEOw4prWmrkb9}Jmr=n-)bvaDr}pqE^(ePB>+*)h3|`E0QVqlVgm+cK zPp{ryw=Q(gsg}8_YZ}_v9Eo2_f4meO@k_e1LKLNs76^3{fsCaB zt_+SL17A_7}p?p|9aJ@1<=i2J0?_C@Y4HRMfSlf3m?DN6zQG8z&$NvO%k(x z`WeYavAK<9>D)JjQ3nI*+i~gPo>%dcRd=rrL@Hx+m(9&BYF~SX$6O2ygo$2mPbMxF zt#56rH=Idt9M{Izc(}@IXCU?#t!nLTB}4=fuVP;my$v>$wGhtP-cP2u!MUE!JH_B80T!lAeREakTE&0Ly2(+AJL<<%f3-^wuD+jw*>e?oby z8K{Abaar)fQhtoZ)$(Xu6{@Qza~UiaLcdTg{i?)vyF%*8+f}6E^|rK(6{T4$Mm#q{ ztig5w(Tydewm`K8VY77nJ4-lpj@-hOx;H?G3N64;0h|rIPqXsLHmqy1?VHc&eoFFr z8DeA00HNRx|BY$gZAVJi9j%WwzPrWP`8Ow<{My*_K1(?a25uwpP~5P&Ddodz!(l5Q ztn6&n_^})OC>9NMV$-?>S$l|J1fX*aM5iJwZb4F{v_jNi83+WTTesF8OCKg2L9L!189g^m z2Kh5wRAR>K<;_&r`brMUXdAN=8rxW5eUCjH-_$w2Ywy-Mw{Y1s)>JL-NdY#*jVbZj zx*h}1OkO69xTkcfill&u+fJvu@4bFrwm!Q3DM4qAT0PmsGb}aENTQibQpKe+g_6R zTv29bnonWqz>Xc8)dI1R(wQ3m4*4*B*-thz!{{+W4KXS)q%dqg2EABPHPK@7i3(L& z*y|(p=i=C%pGcn$Tj`TG_uaQ={9KkZ8AQC{gTS0T-da0nKk9Z|wxaAahq^eJ0m7y^;C}-p z?rV@If#vTehd03sd~X_vWhQhjsgU$A$@`i^tGR@!WQCI{Hj@0y8XMcEp+7`$><;$NSFI}gJ~|Io z!xv>BI$K~-T0Ngg(zt<&x6h3%k)TtBh!^4JwX)uX;e!Za2sz_s<&E^<%(uJ6s`(Mk z7yfsLEk-1mBb_Y4i!&3Fwf-_s_2D1^RB#e=95XWSN-xc@h;4M@R4H#@5z)L2b3GYRXTl z-uxsTT%(X~QjtV;QbQTq1fg&unG0s!28))mkCDVVY;3A@{J|%G-F<#kH1qp00ji_@ zmpSyHetlx{HMH4f#v0*CUt=v#Avo8P= zWz-w^-2VlpQ|pKX{hj;n3(x&I>0_biT;~(zq&)G-84Pry@;~(k^0M@cKmL<+{+O?@ zFx_gFx$_GX1nj=+E4Fd*D{E!q!|V)%(tz=0q}z=%>efdnaDn0Q zx$bWFbgHg@+w~|n-@}dB%JnyEEKTk*JuQRTnCl~gd(pF&6x~H`a*h^e?{D1OC;&QG z1Qht{$|~Wm0rz4$TZ9h}(nYxv#aCoW3wGHnEm*m9;7+|z12%`CO(d&{xhSOzGyqr67>4cbQp)vlY~Erp1BPO37i2K z;!o-)ug}b6?UCp`-0a)94IwAXmRs$z$9w8Yo!^x^-!D-=0HMh~L8GimXDW?pTzmi7 zzt&3?Z6F@;`3udGloL}o)dR2kod|T~At(6pxA?A)WMV!vJWh<=G@2=opJ-K8>z(vo z$yMW9e)rTiszd~Mv5CNW?YZ^kr_zsluPpa_qpN3Vx>$9~mW$hYE`Ly#%zQq`FO_v_ zht!o7cM>i3Sbt}_{Jat|8-lety0;EE$7ior;`2*CCxvvz-w-8WrC%$ot|j^_<3nB+ z!HEc`T)YaW24ff#Cb3K2E`I4@>Zy4=y*^_%eY96Kmxg!q1EyzDXKfPjDj0!mrNH+-U?-8ce{nWqiG4_LDFdyB5komb_AKSL9SrIEwG0Z8)NpRe zPn}WRLU(T-D=EDYjzBKYXxZ}FaBJTYAGDcf@2pM&loiBVJAUBa*g0|7#xSnZEp=q# z4dX(3C^!I0=_>LL8ea#47FTmoST&5_^4qK(I%8O>KD5dhsZ^T&lZH8XMV*P^F7Ha6 z?0oEXM+r$`E~&ch@1yGhvLEj{+wd>3=3z~m1Tz3eWrY{slaOt*91A}?AbBd`VSPB8JY+aW7I5pR5;(4cC_Fib=UAsLZ#Gc1>eWOHfMBCIymN(`PD== zHoLd}UQarCY%uNlmQG5BoglKu4*2~UWf%0diwzOpG8e3_CJxH!MGbloj*vbiMQK*Z z&bP&{!z^Rl@!~;o(l|u4f&i^wi?=9o=YPQypGtC)60|S`IM`25;~4{QG$3xN^P&`b z{Daz@m3k!kK;2^!=f~=;oP-KJWZm21;*%JF2YRXGmx;$RL$yG{W)O$3M);&5sMC#e zp%%@md~DV=(4`ajEe>+qISo_BpfSxG<7Y@i!LxD}16a5i;~T~qML!7O4GCh3D|UFY z<|Vt2n^u***I&AGhUZ2V0YflhID^k)%!jqr9O-6&$9}U_7Hdt$iqPD&x_t1rgT@@S zWM5LD23TbKYqFPe^#d(qOYadzlAKa8Mr3^4ei+S{qS}i>7Oj4M`^p;Fp$^tx9=#>Y z=|dTD@%L+8NOru=&o4`puQmre&bseFcw6*1vn0Vqj5^5rgkn=Eq)Q9HQ{THrP}O9) zTrAP4Rva!>`o5_)>oBcU#wgJ}^YUtBe-j^*QBN8ugq-t&+?MLST5Wz!7CR>%>SyjY zE(NK9)4xwZq35TGfdfp(n4TM@*JLgCzb)4*eYis`>Q4<)J%cOY zGC8lX5F*iv5k4{G5m!E|@i1xmm5ONos>10{&$~&rc8S|6)uuR0W-X&9KTWioxTL*0 z@7R!##An-TC&D{TrOwS1;q;2zq$?%?#E7&guj%rayQ6t5h7sRY1z_+7QNkppf~eV) zkDts-sqN3r6VewpWFoQIOxgqA!v9>#!+zFDM8D$CH9&@-f2a+S(Qrk8myqg@n}zyUcZ#mdn<%rAAiuukZ&I%S10CUW~`0 zztS(IaUwgfeth@SEPhu&Hh<5vlPTE!8O+^5^PJGZBvYgBy3>x4ObKUoKO#S1D~rZOW?B%xoPGDO<$2HSV_8P?G&QGaY`HVK zsWksJ#)p*@i6UUSy-Ea8I{oW&o%Lc_`GzfD3=mJmmP-KYvF$q$8yL4$sCqd=IzFo} z7mq4k#kctREa%03G8uyrn z;uoiIm;z9{ptjUGS|@CwldJ88Qd}7sUl#DOalEB>3sT-)z$&HOs0UcbW;e*KNqcM+eTq=2R&7H zDe`@1CBDU9Sy}Q;v4&|zNY7tju8~yV#|-=LmNkhf=h!B1h_UK4@2ZP)WtZGi;`H$o zH3f^G$LJV%Bwg;29h}^U@Tl>FRs?b~UsYgss8ad=`NltfxqmM?y(;e9cUe->@>=I< zcQO=n(s{#zX&%0Eg&{3HGbKO%R^tT#d8)$C%oHSFP>Jf>jC0}zj z{BlEoH-VnJCWAdU6VCWp(&FIw3qEdbb>+nTd%edq)XZ)1o_XFsdJ6K{QWq_PhmR9N uA9_goZ{PAvb-YizH1Pk(+Ks_0(jO~P<)eb+VB$|KkdeN*UX6}R?7smI?3U#K literal 0 HcmV?d00001 diff --git a/apps/desktop/src/components/connection/ConnectionDialog.vue b/apps/desktop/src/components/connection/ConnectionDialog.vue index 50ecbe6f5..3a952293d 100644 --- a/apps/desktop/src/components/connection/ConnectionDialog.vue +++ b/apps/desktop/src/components/connection/ConnectionDialog.vue @@ -229,6 +229,7 @@ const driverProfiles: Record< }, redis: { type: "redis", port: 6379, user: "", label: "Redis", icon: "redis" }, sqlite: { type: "sqlite", port: 0, user: "", label: "SQLite", icon: "sqlite" }, + rqlite: { type: "rqlite", port: 4001, user: "", label: "RQLite", icon: "rqlite" }, 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" }, @@ -551,6 +552,7 @@ const iconTypeMap: Record = { mysql: "mysql", postgres: "postgres", sqlite: "sqlite", + rqlite: "rqlite", access: "access", redis: "redis", mongodb: "mongodb", @@ -611,6 +613,7 @@ const dbOptions = [ { value: "mysql", label: "MySQL" }, { value: "postgres", label: "PostgreSQL" }, { value: "sqlite", label: "SQLite" }, + { value: "rqlite", label: "RQLite" }, { value: "access", label: "Microsoft Access" }, { value: "redis", label: "Redis" }, { value: "mongodb", label: "MongoDB" }, diff --git a/apps/desktop/src/components/icons/DatabaseIcon.vue b/apps/desktop/src/components/icons/DatabaseIcon.vue index a986f6113..236f0edcc 100644 --- a/apps/desktop/src/components/icons/DatabaseIcon.vue +++ b/apps/desktop/src/components/icons/DatabaseIcon.vue @@ -11,6 +11,7 @@ const assetIcons: Record = { postgres: "postgres", postgresql: "postgres", sqlite: "sqlite", + rqlite: "rqlite.png", redis: "redis", mongodb: "mongodb", clickhouse: "clickhouse", diff --git a/apps/desktop/src/components/layout/ContentArea.vue b/apps/desktop/src/components/layout/ContentArea.vue index 3c4c27132..7425c6261 100644 --- a/apps/desktop/src/components/layout/ContentArea.vue +++ b/apps/desktop/src/components/layout/ContentArea.vue @@ -158,6 +158,7 @@ const activeSqlFormatDialect = computed(() => { case "postgres": return "postgres"; case "sqlite": + case "rqlite": return "sqlite"; case "sqlserver": return "sqlserver"; diff --git a/apps/desktop/src/components/objects/ObjectBrowser.vue b/apps/desktop/src/components/objects/ObjectBrowser.vue index c5b6165d6..3f263eafe 100644 --- a/apps/desktop/src/components/objects/ObjectBrowser.vue +++ b/apps/desktop/src/components/objects/ObjectBrowser.vue @@ -192,6 +192,7 @@ const sourceFormatDialect = computed(() => { case "mysql": case "postgres": case "sqlite": + case "rqlite": case "sqlserver": return effectiveDatabaseType.value; case "gaussdb": diff --git a/apps/desktop/src/lib/connectionPresentation.ts b/apps/desktop/src/lib/connectionPresentation.ts index 541959e1a..0e75534d2 100644 --- a/apps/desktop/src/lib/connectionPresentation.ts +++ b/apps/desktop/src/lib/connectionPresentation.ts @@ -102,6 +102,9 @@ export function connectionUrlPlaceholder(dbType: DatabaseType): string { case "sqlite": return "sqlite:///absolute/path/to/database.db"; + case "rqlite": + return "http://user:password@host:4001"; + case "duckdb": return "duckdb:///absolute/path/to/database.duckdb"; diff --git a/apps/desktop/src/lib/databaseCapabilitySets.ts b/apps/desktop/src/lib/databaseCapabilitySets.ts index afc583d00..c3500588d 100644 --- a/apps/desktop/src/lib/databaseCapabilitySets.ts +++ b/apps/desktop/src/lib/databaseCapabilitySets.ts @@ -37,6 +37,7 @@ export const DIAGRAM_SUPPORTED_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "sqlserver", "oracle", "redshift", @@ -66,6 +67,7 @@ export const DATABASE_SEARCH_SUPPORTED_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "sqlserver", "oracle", "redshift", @@ -108,6 +110,7 @@ export const TABLE_IMPORT_SUPPORTED_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "duckdb", "clickhouse", "sqlserver", @@ -129,6 +132,7 @@ export const TABLE_STRUCTURE_SUPPORTED_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "duckdb", "clickhouse", "sqlserver", @@ -169,6 +173,7 @@ export const FIELD_LINEAGE_SUPPORTED_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "sqlserver", "oracle", "redshift", @@ -258,6 +263,7 @@ export const TRANSFER_SQL_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "sqlserver", "oracle", "clickhouse", @@ -272,6 +278,7 @@ export const DIAGRAM_SQL_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "sqlserver", "oracle", "redshift", diff --git a/apps/desktop/src/lib/databaseFeatureSupport.ts b/apps/desktop/src/lib/databaseFeatureSupport.ts index 53ca92232..e3a5fdf45 100644 --- a/apps/desktop/src/lib/databaseFeatureSupport.ts +++ b/apps/desktop/src/lib/databaseFeatureSupport.ts @@ -104,7 +104,7 @@ export function supportsObjectBrowserTreeNode(dbType: DatabaseType | undefined, } export function supportsTableTruncate(dbType?: DatabaseType): boolean { - return !!dbType && dbType !== "sqlite" && dbType !== "duckdb"; + return !!dbType && dbType !== "sqlite" && dbType !== "rqlite" && dbType !== "duckdb"; } export function usesPostgresLikeStructureCopy(dbType?: DatabaseType): boolean { diff --git a/apps/desktop/src/lib/databaseTableDataCapabilities.ts b/apps/desktop/src/lib/databaseTableDataCapabilities.ts index 42dd9e5ec..7aa80195f 100644 --- a/apps/desktop/src/lib/databaseTableDataCapabilities.ts +++ b/apps/desktop/src/lib/databaseTableDataCapabilities.ts @@ -49,6 +49,7 @@ const NAVICAT_STYLE_TABLE_DATA_TYPES = new Set([ "mysql", "postgres", "sqlite", + "rqlite", "duckdb", "sqlserver", "oracle", diff --git a/apps/desktop/src/lib/objectRenameSql.ts b/apps/desktop/src/lib/objectRenameSql.ts index 03906eb08..eb7205fc0 100644 --- a/apps/desktop/src/lib/objectRenameSql.ts +++ b/apps/desktop/src/lib/objectRenameSql.ts @@ -31,7 +31,7 @@ export function supportsObjectRename( if (objectType === "PROCEDURE" || objectType === "FUNCTION") { return false; } - if (databaseType === "sqlite") return objectType === "TABLE"; + if (databaseType === "sqlite" || databaseType === "rqlite") return objectType === "TABLE"; if (databaseType === "mysql" || databaseType === "goldendb") return objectType === "TABLE" || objectType === "VIEW"; if (postgresLikeRenameTypes.has(databaseType)) return objectType === "TABLE" || objectType === "VIEW"; if (oracleLikeRenameTypes.has(databaseType)) return objectType === "TABLE" || objectType === "VIEW"; diff --git a/apps/desktop/src/lib/sqlCompletion.ts b/apps/desktop/src/lib/sqlCompletion.ts index 4f814c2c4..415b84aa9 100644 --- a/apps/desktop/src/lib/sqlCompletion.ts +++ b/apps/desktop/src/lib/sqlCompletion.ts @@ -451,6 +451,7 @@ const DATABASE_SQL_KEYWORDS: Partial> = { mysql: MYSQL_SQL_KEYWORDS, postgres: POSTGRES_SQL_KEYWORDS, sqlite: SQLITE_SQL_KEYWORDS, + rqlite: SQLITE_SQL_KEYWORDS, sqlserver: SQLSERVER_SQL_KEYWORDS, }; @@ -848,6 +849,7 @@ const DATABASE_FUNCTION_SIGNATURES: Partial vastbase: postgresCapabilities, kingbase: postgresCapabilities, sqlite: sqliteCapabilities, + rqlite: sqliteCapabilities, duckdb: duckdbCapabilities, sqlserver: sqlserverCapabilities, oracle: oracleCapabilities, diff --git a/apps/desktop/src/lib/tableStructureEditorState.ts b/apps/desktop/src/lib/tableStructureEditorState.ts index dd4ee65d0..c5963b579 100644 --- a/apps/desktop/src/lib/tableStructureEditorState.ts +++ b/apps/desktop/src/lib/tableStructureEditorState.ts @@ -114,6 +114,7 @@ export const DATA_TYPE_OPTIONS: Record = { "oid", ], sqlite: ["integer", "real", "text", "blob", "numeric"], + rqlite: ["integer", "real", "text", "blob", "numeric"], sqlserver: [ "bit", "tinyint", diff --git a/apps/desktop/src/types/database.ts b/apps/desktop/src/types/database.ts index 9e6c1a545..c43b3ea04 100644 --- a/apps/desktop/src/types/database.ts +++ b/apps/desktop/src/types/database.ts @@ -2,6 +2,7 @@ export type DatabaseType = | "mysql" | "postgres" | "sqlite" + | "rqlite" | "redis" | "duckdb" | "clickhouse" diff --git a/crates/dbx-core/assets/database-drivers.manifest.json b/crates/dbx-core/assets/database-drivers.manifest.json index 10c1089af..521c625ba 100644 --- a/crates/dbx-core/assets/database-drivers.manifest.json +++ b/crates/dbx-core/assets/database-drivers.manifest.json @@ -30,6 +30,16 @@ "metadataConnectionScoped": false, "skipTcpProbe": true }, + { + "dbType": "rqlite", + "label": "RQLite", + "runtimeMode": "native", + "mcpMode": "direct", + "singleConnectionPool": true, + "metadataConnectionScoped": false, + "skipTcpProbe": false, + "defaultPort": 4001 + }, { "dbType": "redis", "label": "Redis", diff --git a/crates/dbx-core/src/connection.rs b/crates/dbx-core/src/connection.rs index 1b3e3339e..7c9ffd3c1 100644 --- a/crates/dbx-core/src/connection.rs +++ b/crates/dbx-core/src/connection.rs @@ -47,6 +47,7 @@ pub enum PoolKind { Mysql(db::mysql::MySqlPool, MysqlMode), Postgres(deadpool_postgres::Pool), Sqlite(db::sqlite::SqliteHandle), + Rqlite(db::rqlite_driver::RqliteClient), Redis(db::redis_driver::RedisConnection), DuckDb(Arc>), MongoDb(mongodb::Client), @@ -343,6 +344,18 @@ impl AppState { db::sqlite::connect_path_with_extensions(&expand_tilde(&db_config.host), extensions).await?, ) } + DatabaseType::Rqlite => { + let client = db::rqlite_driver::RqliteClient::new( + &url, + db_config.url_params.as_deref(), + &db_config.username, + &db_config.password, + db_config.ssl, + connect_timeout, + )?; + db::rqlite_driver::test_connection(&client, connect_timeout).await?; + PoolKind::Rqlite(client) + } DatabaseType::Redis => { let con = if db_config.uses_redis_cluster() { db::redis_driver::RedisConnection::Cluster(db::redis_driver::connect_cluster(&db_config).await?) @@ -919,6 +932,7 @@ pub async fn close_pool_kind(pool: PoolKind) { } PoolKind::Postgres(p) => p.close(), PoolKind::Sqlite(_) => {} + PoolKind::Rqlite(_) => {} PoolKind::Redis(_) => {} PoolKind::DuckDb(con) => { crate::db::duckdb_driver::close_connection(con); diff --git a/crates/dbx-core/src/database_capabilities.rs b/crates/dbx-core/src/database_capabilities.rs index 3066e5bd4..b243f2cd6 100644 --- a/crates/dbx-core/src/database_capabilities.rs +++ b/crates/dbx-core/src/database_capabilities.rs @@ -14,6 +14,7 @@ pub fn is_single_connection_pool(db_type: &DatabaseType) -> bool { db_type, DatabaseType::Sqlite | DatabaseType::DuckDb + | DatabaseType::Rqlite | DatabaseType::MongoDb | DatabaseType::Oracle | DatabaseType::Dameng diff --git a/crates/dbx-core/src/db/mod.rs b/crates/dbx-core/src/db/mod.rs index a8563b5e9..2084f1f35 100644 --- a/crates/dbx-core/src/db/mod.rs +++ b/crates/dbx-core/src/db/mod.rs @@ -9,6 +9,7 @@ pub mod ob_oracle; pub mod postgres; pub mod proxy_tunnel; pub mod redis_driver; +pub mod rqlite_driver; pub mod sqlite; pub mod sqlserver; pub mod ssh_tunnel; diff --git a/crates/dbx-core/src/db/rqlite_driver.rs b/crates/dbx-core/src/db/rqlite_driver.rs new file mode 100644 index 000000000..1cfeab004 --- /dev/null +++ b/crates/dbx-core/src/db/rqlite_driver.rs @@ -0,0 +1,455 @@ +use reqwest::Client as HttpClient; +use serde::Deserialize; +use std::time::{Duration, Instant}; + +use super::with_connection_timeout; +use crate::sql::starts_with_executable_sql_keyword; +use crate::types::{ + ColumnInfo, DatabaseInfo, ForeignKeyInfo, IndexInfo, ObjectSource, ObjectSourceKind, QueryResult, TableInfo, + TriggerInfo, +}; + +#[derive(Clone)] +pub struct RqliteClient { + http: HttpClient, + base_url: String, + query_params: String, + auth: Option<(String, String)>, +} + +impl RqliteClient { + pub fn new( + url: &str, + url_params: Option<&str>, + username: &str, + password: &str, + tls_enabled: bool, + timeout: Duration, + ) -> Result { + let mut builder = HttpClient::builder().connect_timeout(timeout); + if rqlite_accept_invalid_certs(tls_enabled, url_params) { + builder = builder.danger_accept_invalid_certs(true); + } + if rqlite_should_bypass_system_proxy(url) { + builder = builder.no_proxy(); + } + let http = builder.build().map_err(|e| format!("Failed to configure rqlite HTTP client: {e}"))?; + let auth = if username.trim().is_empty() { None } else { Some((username.to_string(), password.to_string())) }; + Ok(Self { + http, + base_url: url.trim_end_matches('/').split('?').next().unwrap_or(url).to_string(), + query_params: normalize_rqlite_url_params(url_params), + auth, + }) + } + + fn post_json(&self, path: &str, sql: &str) -> reqwest::RequestBuilder { + let req = self.http.post(self.endpoint(path)).json(&[sql]); + self.with_auth(req) + } + + fn with_auth(&self, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder { + if let Some((ref user, ref pass)) = self.auth { + req.basic_auth(user, Some(pass)) + } else { + req + } + } + + fn endpoint(&self, path: &str) -> String { + if self.query_params.is_empty() { + format!("{}{}", self.base_url, path) + } else { + format!("{}{}?{}", self.base_url, path, self.query_params) + } + } +} + +#[derive(Debug, Deserialize)] +struct RqliteResponse { + results: Vec, +} + +#[derive(Debug, Deserialize)] +struct RqliteResult { + #[serde(default)] + columns: Vec, + #[serde(default)] + values: Vec>, + #[serde(default)] + rows_affected: Option, + #[serde(default)] + error: Option, +} + +enum RqliteEndpoint { + Query, + Execute, +} + +impl RqliteEndpoint { + fn path(&self) -> &'static str { + match self { + Self::Query => "/db/query", + Self::Execute => "/db/execute", + } + } +} + +pub async fn test_connection(client: &RqliteClient, timeout: Duration) -> Result<(), String> { + with_connection_timeout("rqlite", timeout, async { query_one(client, "SELECT 1").await.map(|_| ()) }).await +} + +pub async fn list_databases(_client: &RqliteClient) -> Result, String> { + Ok(vec![DatabaseInfo { name: "main".to_string() }]) +} + +pub async fn list_tables(client: &RqliteClient, _schema: &str) -> Result, String> { + let result = query_one( + client, + "SELECT name, type FROM sqlite_master WHERE type IN ('table', 'view') AND name NOT LIKE 'sqlite_%' ORDER BY name", + ) + .await?; + Ok(result + .values + .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: &RqliteClient, _schema: &str, table: &str) -> Result, String> { + let result = query_one(client, &format!("PRAGMA table_info({})", sqlite_ident(table))).await?; + Ok(result + .values + .into_iter() + .map(|row| 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::().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::().ok()) + .unwrap_or(0) + > 0, + extra: None, + comment: None, + numeric_precision: None, + numeric_scale: None, + character_maximum_length: None, + }) + .collect()) +} + +pub async fn list_indexes(client: &RqliteClient, _schema: &str, table: &str) -> Result, String> { + let result = query_one(client, &format!("PRAGMA index_list({})", sqlite_ident(table))).await?; + let mut indexes = Vec::new(); + + for row in result.values { + 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::().ok()).unwrap_or(0) + != 0; + let origin = value_by_column(&result.columns, &row, "origin").unwrap_or_default(); + let column_result = query_one(client, &format!("PRAGMA index_info({})", sqlite_ident(&name))).await?; + let columns = column_result + .values + .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: &RqliteClient, + _schema: &str, + table: &str, +) -> Result, String> { + let result = query_one(client, &format!("PRAGMA foreign_key_list({})", sqlite_ident(table))).await?; + Ok(result + .values + .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(), + }) + .collect()) +} + +pub async fn list_triggers(client: &RqliteClient, _schema: &str, table: &str) -> Result, String> { + let result = query_one( + client, + &format!( + "SELECT name, sql FROM sqlite_master WHERE type = 'trigger' AND tbl_name = {} ORDER BY name", + sqlite_string(table) + ), + ) + .await?; + Ok(result + .values + .into_iter() + .map(|row| { + let sql_text = value_as_string(row.get(1)).unwrap_or_default().to_uppercase(); + let timing = if sql_text.contains("BEFORE") { + "BEFORE" + } else if sql_text.contains("AFTER") { + "AFTER" + } else { + "INSTEAD OF" + }; + let event = if sql_text.contains("INSERT") { + "INSERT" + } else if sql_text.contains("UPDATE") { + "UPDATE" + } else { + "DELETE" + }; + TriggerInfo { + name: value_as_string(row.first()).unwrap_or_default(), + event: event.to_string(), + timing: timing.to_string(), + } + }) + .collect()) +} + +pub async fn table_ddl(client: &RqliteClient, table: &str) -> Result { + first_string_cell( + query_one( + client, + &format!("SELECT sql FROM sqlite_master WHERE type='table' AND name={}", sqlite_string(table)), + ) + .await?, + ) +} + +pub async fn object_source( + client: &RqliteClient, + name: &str, + object_type: &ObjectSourceKind, +) -> Result { + let kind = match object_type { + ObjectSourceKind::View => "view", + _ => return Err("Object source is not supported for this rqlite object type".to_string()), + }; + let source = first_string_cell( + query_one( + 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 }) +} + +pub async fn execute_query(client: &RqliteClient, sql: &str) -> Result { + execute_query_with_max_rows(client, sql, None).await +} + +pub async fn execute_query_with_max_rows( + client: &RqliteClient, + sql: &str, + max_rows: Option, +) -> Result { + let start = Instant::now(); + if starts_with_executable_sql_keyword(sql, &["SELECT", "PRAGMA", "EXPLAIN", "WITH"]) { + let result = query_one(client, sql).await?; + Ok(query_result_from_rqlite_result(result, start.elapsed().as_millis(), max_rows)) + } else { + let result = execute_one(client, sql).await?; + let affected_rows = result.rows_affected.unwrap_or(0); + Ok(QueryResult { + columns: vec![], + rows: vec![], + affected_rows, + execution_time_ms: start.elapsed().as_millis(), + truncated: false, + session_id: None, + has_more: false, + }) + } +} + +async fn query_one(client: &RqliteClient, sql: &str) -> Result { + request_one(client, RqliteEndpoint::Query, sql).await +} + +async fn execute_one(client: &RqliteClient, sql: &str) -> Result { + request_one(client, RqliteEndpoint::Execute, sql).await +} + +async fn request_one(client: &RqliteClient, endpoint: RqliteEndpoint, sql: &str) -> Result { + let resp = + client.post_json(endpoint.path(), sql).send().await.map_err(|e| format!("rqlite request failed: {e}"))?; + let status = resp.status(); + let body = resp.text().await.map_err(|e| format!("rqlite response read failed: {e}"))?; + if !status.is_success() { + return Err(format!("rqlite error ({status}): {body}")); + } + let response: RqliteResponse = + serde_json::from_str(&body).map_err(|e| format!("rqlite parse error: {e}; body: {body}"))?; + let result = response.results.into_iter().next().ok_or_else(|| "rqlite returned no result".to_string())?; + if let Some(error) = result.error.as_ref().filter(|error| !error.is_empty()) { + return Err(format!("rqlite error: {error}")); + } + Ok(result) +} + +fn query_result_from_rqlite_result( + mut result: RqliteResult, + execution_time_ms: u128, + max_rows: Option, +) -> QueryResult { + let row_limit = max_rows.unwrap_or(crate::query::MAX_ROWS).max(1); + let truncated = result.values.len() > row_limit; + if truncated { + result.values.truncate(row_limit); + } + QueryResult { + columns: result.columns, + rows: result.values, + affected_rows: 0, + execution_time_ms, + truncated, + session_id: None, + has_more: false, + } +} + +fn first_string_cell(result: RqliteResult) -> Result { + result + .values + .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 { + 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 { + 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('\'', "''")) +} + +fn normalize_rqlite_url_params(params: Option<&str>) -> String { + params + .unwrap_or("") + .trim() + .trim_start_matches('?') + .split('&') + .filter(|part| !part.is_empty()) + .collect::>() + .join("&") +} + +fn rqlite_accept_invalid_certs(tls_enabled: bool, url_params: Option<&str>) -> bool { + tls_enabled + && url_params + .unwrap_or("") + .trim() + .trim_start_matches('?') + .split('&') + .filter_map(|pair| pair.split_once('=')) + .any(|(key, value)| { + matches!(key.trim().to_ascii_lowercase().as_str(), "insecure" | "tls_insecure" | "accept_invalid_certs") + && matches!(value.trim().to_ascii_lowercase().as_str(), "true" | "1" | "yes" | "on") + }) +} + +fn rqlite_should_bypass_system_proxy(base_url: &str) -> bool { + let Ok(parsed) = reqwest::Url::parse(base_url) else { + return false; + }; + let Some(host) = parsed.host_str() else { + return false; + }; + let host = host.trim_matches(['[', ']']); + host.eq_ignore_ascii_case("localhost") || host.parse::().is_ok_and(|ip| ip.is_loopback()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn converts_query_result_and_truncates_probe_rows() { + let result = RqliteResult { + columns: vec!["id".to_string(), "name".to_string()], + values: vec![ + vec![serde_json::json!(1), serde_json::json!("Ada")], + vec![serde_json::json!(2), serde_json::json!("Linus")], + ], + rows_affected: None, + error: None, + }; + + let result = query_result_from_rqlite_result(result, 8, Some(1)); + + assert_eq!(result.columns, vec!["id", "name"]); + assert_eq!(result.rows, vec![vec![serde_json::json!(1), serde_json::json!("Ada")]]); + assert_eq!(result.execution_time_ms, 8); + assert!(result.truncated); + } + + #[test] + fn maps_column_values_by_name() { + let columns = vec!["name".to_string(), "notnull".to_string(), "pk".to_string()]; + let row = vec![serde_json::json!("id"), serde_json::json!(1), serde_json::json!(1)]; + + assert_eq!(value_by_column(&columns, &row, "NAME").as_deref(), Some("id")); + assert_eq!(value_by_column(&columns, &row, "notnull").as_deref(), Some("1")); + } +} diff --git a/crates/dbx-core/src/models/connection.rs b/crates/dbx-core/src/models/connection.rs index 41257732d..1ebc23678 100644 --- a/crates/dbx-core/src/models/connection.rs +++ b/crates/dbx-core/src/models/connection.rs @@ -169,6 +169,7 @@ pub enum DatabaseType { Mysql, Postgres, Sqlite, + Rqlite, Redis, #[serde(rename = "duckdb")] DuckDb, @@ -301,6 +302,7 @@ impl ConnectionConfig { }, DatabaseType::Redshift => Some("dev"), DatabaseType::ClickHouse => Some("default"), + DatabaseType::Rqlite => Some("main"), DatabaseType::Gaussdb | DatabaseType::OpenGauss => Some("postgres"), DatabaseType::Kingbase | DatabaseType::Vastbase => Some("postgres"), DatabaseType::Highgo => Some("highgo"), @@ -386,6 +388,7 @@ impl ConnectionConfig { format!("postgres://{host}:{port}{db_part}{suffix}") } DatabaseType::ClickHouse => clickhouse_http_url(self, raw_host, port), + DatabaseType::Rqlite => rqlite_http_url(self, raw_host, port), DatabaseType::SqlServer => { format!("server=tcp:{host},{port};database={}", self.database.as_deref().unwrap_or("master")) } @@ -489,6 +492,7 @@ impl ConnectionConfig { format!("postgres://{}:{}@{host}:{port}{db_part}{suffix}", username, password) } DatabaseType::ClickHouse => clickhouse_http_url(self, raw_host, port), + DatabaseType::Rqlite => rqlite_http_url(self, raw_host, port), DatabaseType::SqlServer => format!( "server=tcp:{host},{port};user={};password={};database={}", self.username, @@ -825,6 +829,31 @@ fn clickhouse_http_url(config: &ConnectionConfig, host: &str, port: u16) -> Stri format!("{scheme}://{}:{port}", bracket_ipv6(trimmed)) } +fn rqlite_http_url(config: &ConnectionConfig, host: &str, port: u16) -> String { + let trimmed = host.trim(); + if let Some(rest) = trimmed.strip_prefix("https://") { + return format!("https://{}", trim_http_host_port(rest, port)); + } + if let Some(rest) = trimmed.strip_prefix("http://") { + let scheme = if config.ssl { "https" } else { "http" }; + return format!("{scheme}://{}", trim_http_host_port(rest, port)); + } + let scheme = if config.ssl { "https" } else { "http" }; + format!("{scheme}://{}:{port}", bracket_ipv6(trimmed)) +} + +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("]:") { + return format!("{authority}:{default_port}"); + } + if authority.rsplit_once(':').is_some() { + authority.to_string() + } else { + format!("{authority}:{default_port}") + } +} + fn trim_clickhouse_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("]:") { diff --git a/crates/dbx-core/src/query.rs b/crates/dbx-core/src/query.rs index 1c97da0f4..9d6fafda6 100644 --- a/crates/dbx-core/src/query.rs +++ b/crates/dbx-core/src/query.rs @@ -564,6 +564,17 @@ pub async fn do_execute( wait_for_query_opt(cancel_token, query_timeout, db::sqlite::execute_query_with_max_rows(&p, sql, max_rows)) .await } + PoolKind::Rqlite(client) => { + let client = client.clone(); + let max_rows = options.max_rows; + drop(connections); + wait_for_query_opt( + cancel_token, + query_timeout, + db::rqlite_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(); @@ -1107,7 +1118,9 @@ 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::ClickHouse(_) | PoolKind::SqlServer(_) | PoolKind::Agent(_) => TxPath::Explicit, + PoolKind::ClickHouse(_) | PoolKind::Rqlite(_) | PoolKind::SqlServer(_) | PoolKind::Agent(_) => { + TxPath::Explicit + } PoolKind::DuckDb(_) | PoolKind::Redis(_) | PoolKind::MongoDb(_) diff --git a/crates/dbx-core/src/schema.rs b/crates/dbx-core/src/schema.rs index 56773f320..bd85b02b1 100644 --- a/crates/dbx-core/src/schema.rs +++ b/crates/dbx-core/src/schema.rs @@ -300,6 +300,7 @@ async fn list_databases_once(state: &AppState, connection_id: &str) -> Result dispatch_mysql!(p, mode, db::mysql::list_databases, db::ob_oracle::list_databases), PoolKind::Postgres(p) => db::postgres::list_databases(p).await, PoolKind::Sqlite(p) => db::sqlite::list_databases(p).await, + PoolKind::Rqlite(client) => db::rqlite_driver::list_databases(client).await, PoolKind::DuckDb(con) => { let con = con.lock().map_err(|e| e.to_string())?; duckdb_list_databases_with_attached(&con, &duckdb_attached_names) @@ -428,6 +429,9 @@ async fn list_tables_once( PoolKind::Sqlite(p) => { db::sqlite::list_tables(p, schema).await.map(|tables| filter_table_infos(tables, filter, limit)) } + PoolKind::Rqlite(client) => { + db::rqlite_driver::list_tables(client, schema).await.map(|tables| filter_table_infos(tables, filter, limit)) + } PoolKind::MongoDb(client) => db::mongo_driver::list_collections(client, database) .await .map(|names| collection_names_to_tables(names, "COLLECTION")) @@ -825,6 +829,9 @@ pub async fn get_columns_core( } PoolKind::Postgres(p) => db::postgres::get_columns(p, schema, table).await.map(deduplicate_column_infos), PoolKind::Sqlite(p) => db::sqlite::get_columns(p, schema, table).await.map(deduplicate_column_infos), + PoolKind::Rqlite(client) => { + db::rqlite_driver::get_columns(client, schema, table).await.map(deduplicate_column_infos) + } _ => Ok(vec![]), } } @@ -896,6 +903,7 @@ pub async fn list_indexes_core( } PoolKind::Postgres(p) => db::postgres::list_indexes(p, schema, table).await, PoolKind::Sqlite(p) => db::sqlite::list_indexes(p, schema, table).await, + PoolKind::Rqlite(client) => db::rqlite_driver::list_indexes(client, schema, table).await, _ => Ok(vec![]), } } @@ -924,6 +932,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, _ => Ok(vec![]), } } @@ -952,6 +961,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, _ => Ok(vec![]), } } @@ -1019,6 +1029,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, _ => Err("DDL not supported for this database type".to_string()), } } @@ -1244,6 +1255,9 @@ pub async fn get_object_source_core( PoolKind::Sqlite(pool) => first_string_cell( db::sqlite::execute_query(pool, &sqlite_object_source_sql(name, &object_type)).await?, )?, + PoolKind::Rqlite(client) => { + return db::rqlite_driver::object_source(client, name, &object_type).await; + } PoolKind::ClickHouse(client) if matches!(object_type, db::ObjectSourceKind::View) => { let result = db::clickhouse_driver::execute_query( client, diff --git a/crates/dbx-core/src/table_structure_sql.rs b/crates/dbx-core/src/table_structure_sql.rs index 16ce82bfe..a6b6069ca 100644 --- a/crates/dbx-core/src/table_structure_sql.rs +++ b/crates/dbx-core/src/table_structure_sql.rs @@ -241,7 +241,7 @@ fn capabilities_for(database_type: Option) -> TableStructureCapabi comment: true, ..base }, - Some(DatabaseType::Sqlite) => TableStructureCapabilities { + Some(DatabaseType::Sqlite | DatabaseType::Rqlite) => TableStructureCapabilities { dialect: StructureDialect::Sqlite, add_column: true, drop_column: true, @@ -1966,6 +1966,34 @@ mod tests { ); } + #[test] + fn builds_rqlite_changes_with_sqlite_dialect() { + let mut email = column("email"); + email.data_type = "text".to_string(); + email.is_nullable = false; + let mut email_index = index("idx_users_email", &["email"]); + email_index.filter = "email IS NOT NULL".to_string(); + + let result = build_table_structure_change_sql(TableStructureSqlOptions { + database_type: Some(DatabaseType::Rqlite), + schema: None, + table_name: "users".to_string(), + columns: vec![email], + indexes: vec![email_index], + table_comment: None, + original_table_comment: None, + }); + + assert_eq!(result.warnings, Vec::::new()); + assert_eq!( + result.statements, + vec![ + "ALTER TABLE \"users\" ADD COLUMN \"email\" text NOT NULL;", + "CREATE INDEX \"idx_users_email\" ON \"users\" (\"email\") WHERE email IS NOT NULL;", + ] + ); + } + #[test] fn builds_mysql_column_reorder_statements() { let mut id = column("id"); diff --git a/packages/app-tests/connectionUrlPlaceholder.test.ts b/packages/app-tests/connectionUrlPlaceholder.test.ts index 970120ebe..13abcf121 100644 --- a/packages/app-tests/connectionUrlPlaceholder.test.ts +++ b/packages/app-tests/connectionUrlPlaceholder.test.ts @@ -11,6 +11,7 @@ const expected: Record = { redshift: "postgresql://user:password@host:port/database", redis: "redis://:password@host:port/0", sqlite: "sqlite:///absolute/path/to/database.db", + rqlite: "http://user:password@host:4001", duckdb: "duckdb:///absolute/path/to/database.duckdb", access: "jdbc:ucanaccess:///absolute/path/to/database.accdb", mongodb: "mongodb://user:password@host:port/database", diff --git a/packages/app-tests/databaseCapabilities.test.ts b/packages/app-tests/databaseCapabilities.test.ts index e766d1406..dc11f3949 100644 --- a/packages/app-tests/databaseCapabilities.test.ts +++ b/packages/app-tests/databaseCapabilities.test.ts @@ -175,6 +175,14 @@ test("uses Navicat-style table editing defaults for updateable SQL table engines requiresTransactionalTableForExistingRows: false, transaction: true, }); + assert.deepEqual(getDatabaseCapability("rqlite").tableData, { + insert: true, + updateRequiresPrimaryKey: false, + deleteRequiresPrimaryKey: false, + keylessRowPredicate: true, + requiresTransactionalTableForExistingRows: false, + transaction: true, + }); }); test("keeps conservative table editing defaults for unknown database types", () => { @@ -205,6 +213,7 @@ test("describes feature support through capability helpers", () => { assert.equal(supportsTableStructureEditing("opengauss"), true); assert.equal(supportsTableStructureEditing("redshift"), true); assert.equal(supportsTableStructureEditing("clickhouse"), true); + assert.equal(supportsTableStructureEditing("rqlite"), true); assert.equal(supportsTableStructureEditing("mongodb"), false); assert.equal(supportsDatabaseCreation("clickhouse"), true); assert.equal(supportsDatabaseCreation("sqlite"), false); @@ -220,6 +229,7 @@ test("describes feature support through capability helpers", () => { assert.equal(supportsObjectBrowser("mongodb"), false); assert.equal(supportsTableTruncate("mysql"), true); assert.equal(supportsTableTruncate("duckdb"), false); + assert.equal(supportsTableTruncate("rqlite"), false); }); test("object browser entry follows database tree shape", () => { diff --git a/packages/app-tests/objectRenameSql.test.ts b/packages/app-tests/objectRenameSql.test.ts index 35b431eb7..d92c7323f 100644 --- a/packages/app-tests/objectRenameSql.test.ts +++ b/packages/app-tests/objectRenameSql.test.ts @@ -10,6 +10,8 @@ test("recognizes object rename support for UI affordances", () => { assert.equal(supportsObjectRename("sqlserver", "PROCEDURE"), true); assert.equal(supportsObjectRename("sqlite", "TABLE"), true); assert.equal(supportsObjectRename("sqlite", "VIEW"), false); + assert.equal(supportsObjectRename("rqlite", "TABLE"), true); + assert.equal(supportsObjectRename("rqlite", "VIEW"), false); assert.equal(supportsObjectRename("oracle", "FUNCTION"), false); assert.equal(supportsObjectRename("dameng", "PROCEDURE"), false); assert.equal(supportsObjectRename("mysql", "PROCEDURE"), false); diff --git a/packages/app-tests/tableStructureCapabilities.test.ts b/packages/app-tests/tableStructureCapabilities.test.ts index ad04785f5..dbe174ab0 100644 --- a/packages/app-tests/tableStructureCapabilities.test.ts +++ b/packages/app-tests/tableStructureCapabilities.test.ts @@ -5,8 +5,8 @@ import { getTableStructureCapabilities, } from "../../apps/desktop/src/lib/tableStructureCapabilities.ts"; -test("sqlite and duckdb do not support table comments", () => { - for (const dbType of ["sqlite", "duckdb"] as const) { +test("sqlite-family and duckdb do not support table comments", () => { + for (const dbType of ["sqlite", "rqlite", "duckdb"] as const) { const caps = getTableStructureCapabilities(dbType); assert.equal(caps.comment, false, `${dbType} should not support comments`); assert.equal(caps.createTable, true, `${dbType} should still support creating tables`); diff --git a/packages/cli/tests/bin.test.ts b/packages/cli/tests/bin.test.ts index 6d5c9ff16..96e6f535c 100644 --- a/packages/cli/tests/bin.test.ts +++ b/packages/cli/tests/bin.test.ts @@ -30,7 +30,8 @@ test("prints capabilities when invoked through an npm-style symlink", async () = assert.equal(result.status, 0); assert.equal(result.stderr, ""); const payload = JSON.parse(result.stdout) as { directQueryTypes: string[]; bridgeRequiredTypes: string[] }; - assert.ok(payload.directQueryTypes.includes("postgres")); + assert.ok(payload.directQueryTypes.includes("postgres")); + assert.ok(payload.directQueryTypes.includes("rqlite")); assert.ok(payload.bridgeRequiredTypes.includes("oracle")); } finally { await rm(bin.dir, { recursive: true, force: true }); diff --git a/packages/cli/tests/cli.test.ts b/packages/cli/tests/cli.test.ts index 1779a4dab..43742b760 100644 --- a/packages/cli/tests/cli.test.ts +++ b/packages/cli/tests/cli.test.ts @@ -51,7 +51,7 @@ const diagnostics: DbxDiagnostics = { loadedConnectionCount: 2, bridgePortFile: "/tmp/dbx/mcp-bridge-port", bridgePortFileExists: false, - directQueryTypes: ["postgres", "mysql", "sqlite"], + directQueryTypes: ["postgres", "mysql", "sqlite", "rqlite"], bridgeRequiredTypes: ["oracle", "mongodb"], }; @@ -154,6 +154,7 @@ test("prints capabilities as json", async () => { const payload = JSON.parse(result.stdout) as { directQueryTypes: string[]; bridgeRequiredTypes: string[] }; assert.ok(payload.directQueryTypes.includes("postgres")); assert.ok(payload.directQueryTypes.includes("sqlite")); + assert.ok(payload.directQueryTypes.includes("rqlite")); assert.ok(payload.bridgeRequiredTypes.includes("oracle")); }); diff --git a/packages/mcp-server/src/index.ts b/packages/mcp-server/src/index.ts index b812634a0..945516f6b 100644 --- a/packages/mcp-server/src/index.ts +++ b/packages/mcp-server/src/index.ts @@ -46,7 +46,7 @@ function formatQueryToolResult(result: QueryResult, title?: string) { } export const DBX_CONNECTION_TYPE_DESCRIPTION = - "Database type: postgres, mysql, sqlite, redis, duckdb, clickhouse, sqlserver, mongodb, oracle, elasticsearch, doris, starrocks, redshift, dameng, kingbase, highgo, vastbase, goldendb, gaussdb, yashandb, databricks, saphana, teradata, vertica, firebird, exasol, opengauss, oceanbase-oracle, gbase, h2, snowflake, trino, hive, db2, informix, iris, neo4j, cassandra, bigquery, kylin, sundb, tdengine, xugu, jdbc, access"; + "Database type: postgres, mysql, sqlite, rqlite, redis, duckdb, clickhouse, sqlserver, mongodb, oracle, elasticsearch, doris, starrocks, redshift, dameng, kingbase, highgo, vastbase, goldendb, gaussdb, yashandb, databricks, saphana, teradata, vertica, firebird, exasol, opengauss, oceanbase-oracle, gbase, h2, snowflake, trino, hive, db2, informix, iris, neo4j, cassandra, bigquery, kylin, sundb, tdengine, xugu, jdbc, access"; export function createDbxMcpServer(backend: Backend, options: { isWebMode?: boolean } = {}): McpServer { const isWebMode = options.isWebMode ?? !!process.env.DBX_WEB_URL; @@ -181,7 +181,7 @@ export function createDbxMcpServer(backend: Backend, options: { isWebMode?: bool const existing = await backend.findConnection(name); if (existing) return text(`Connection "${name}" already exists.`); const FILE_BASED_TYPES = new Set(["sqlite", "duckdb", "access"]); - const DEFAULT_PORTS: Record = { tdengine: 6041, xugu: 5138 }; + const DEFAULT_PORTS: Record = { rqlite: 4001, tdengine: 6041, xugu: 5138 }; const resolvedPort = port ?? DEFAULT_PORTS[db_type] ?? (FILE_BASED_TYPES.has(db_type) ? 0 : undefined); if (resolvedPort === undefined) return text("Port is required for this database type."); const config = await backend.addConnection({ diff --git a/packages/node-core/src/database.ts b/packages/node-core/src/database.ts index abc2d4f6a..07db75ba6 100644 --- a/packages/node-core/src/database.ts +++ b/packages/node-core/src/database.ts @@ -42,6 +42,17 @@ interface PoolEntry { timer: ReturnType; } +interface RqliteResult { + columns?: string[]; + values?: unknown[][]; + rows_affected?: number; + error?: string; +} + +interface RqliteResponse { + results?: RqliteResult[]; +} + const pools = new Map(); const proxyTunnels = new Map(); @@ -365,6 +376,7 @@ async function mysqlQuery(config: ConnectionConfig, sql: string, params?: unknow async function query(config: ConnectionConfig, sql: string, params?: unknown[], options?: QueryOptions): Promise { if (config.db_type === "sqlite") return sqliteQuery(config, sql, options); + if (config.db_type === "rqlite") return rqliteQuery(config, sql, options); if (isMysqlType(config.db_type)) return mysqlQuery(config, sql, params, options); return pgQuery(config, sql, params, options); } @@ -398,6 +410,47 @@ function sqliteQuery(config: ConnectionConfig, sql: string, options?: QueryOptio } } +async function rqliteQuery(config: ConnectionConfig, sql: string, options?: QueryOptions): Promise { + const isReader = /^\s*(?:--[^\n]*\n|\s|\/\*[\s\S]*?\*\/)*(select|pragma|explain|with)\b/i.test(sql); + const endpoint = isReader ? "/db/query" : "/db/execute"; + const result = await rqliteRequest(config, endpoint, sql); + if (isReader) { + const columns = result.columns ?? []; + const rows = (result.values ?? []).slice(0, resolveMaxRows(options)).map((row) => { + const record: Record = {}; + columns.forEach((column, index) => { + record[column] = row[index]; + }); + return record; + }); + return { columns, rows, row_count: rows.length }; + } + return { columns: [], rows: [], row_count: result.rows_affected ?? 0 }; +} + +async function rqliteRequest(config: ConnectionConfig, endpoint: "/db/query" | "/db/execute", sql: string): Promise { + const { host, port } = await connectionEndpoint(config); + const scheme = config.ssl ? "https" : "http"; + const params = (config.url_params || "").trim().replace(/^\?/, ""); + const url = `${scheme}://${host}:${port}${endpoint}${params ? `?${params}` : ""}`; + const headers: Record = { "content-type": "application/json" }; + if (config.username) { + headers.authorization = `Basic ${Buffer.from(`${config.username}:${config.password || ""}`).toString("base64")}`; + } + const response = await fetch(url, { + method: "POST", + headers, + body: JSON.stringify([sql]), + }); + const text = await response.text(); + if (!response.ok) throw new Error(`rqlite error (${response.status}): ${text}`); + const payload = JSON.parse(text) as RqliteResponse; + const result = payload.results?.[0]; + if (!result) throw new Error("rqlite returned no result"); + if (result.error) throw new Error(`rqlite error: ${result.error}`); + return result; +} + export async function executeQuery(config: ConnectionConfig, sql: string, options?: QueryOptions): Promise { if (config.db_type === "mongodb") { const find = parseMongoFindCommand(sql); @@ -451,7 +504,7 @@ export async function listTables(config: ConnectionConfig, schema?: string): Pro }); return collections.map((name) => ({ name, type: "COLLECTION" })); } - if (config.db_type === "sqlite") { + if (config.db_type === "sqlite" || config.db_type === "rqlite") { const result = await query( config, `SELECT name, type FROM sqlite_master WHERE type IN ('table', 'view') AND name NOT LIKE 'sqlite_%' ORDER BY name`, @@ -484,7 +537,7 @@ export async function describeTable(config: ConnectionConfig, table: string, sch const result = await mongoFindDocuments(config, table, 0, 20, "{}"); return inferMongoColumns(result.documents); } - if (config.db_type === "sqlite") { + if (config.db_type === "sqlite" || config.db_type === "rqlite") { const result = await query(config, `PRAGMA table_info(${quoteSqliteIdentifier(table)})`); return result.rows.map((r) => ({ name: String(r.name || ""), diff --git a/packages/node-core/src/diagnostics.ts b/packages/node-core/src/diagnostics.ts index b60f2f898..eb6e7cd5c 100644 --- a/packages/node-core/src/diagnostics.ts +++ b/packages/node-core/src/diagnostics.ts @@ -9,6 +9,7 @@ export const DIRECT_QUERY_TYPES = [ "doris", "starrocks", "sqlite", + "rqlite", "gaussdb", "opengauss", ] as const; diff --git a/packages/node-core/tests/rqlite-direct.test.ts b/packages/node-core/tests/rqlite-direct.test.ts new file mode 100644 index 000000000..543b25d45 --- /dev/null +++ b/packages/node-core/tests/rqlite-direct.test.ts @@ -0,0 +1,111 @@ +import assert from "node:assert/strict"; +import { createServer, type IncomingMessage, type ServerResponse } from "node:http"; +import test from "node:test"; +import type { ConnectionConfig } from "../src/connections.js"; +import { describeTable, executeQuery, listTables } from "../src/database.js"; + +function rqliteConfig(port: number): ConnectionConfig { + return { + id: "rqlite-test", + name: "local-rqlite", + db_type: "rqlite", + host: "127.0.0.1", + port, + username: "dbx", + password: "secret", + database: "main", + ssh_enabled: false, + ssl: false, + }; +} + +async function withRqliteServer(handler: (req: IncomingMessage, res: ServerResponse, body: string) => void) { + const server = createServer((req, res) => { + let body = ""; + req.setEncoding("utf8"); + req.on("data", (chunk) => { + body += chunk; + }); + req.on("end", () => handler(req, res, body)); + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const address = server.address(); + assert.ok(address && typeof address === "object"); + return { + port: address.port, + close: () => new Promise((resolve, reject) => server.close((err) => (err ? reject(err) : resolve()))), + }; +} + +function json(res: ServerResponse, body: unknown) { + res.writeHead(200, { "content-type": "application/json" }); + res.end(JSON.stringify(body)); +} + +test("executes rqlite query through the HTTP API", async () => { + const server = await withRqliteServer((req, res, body) => { + assert.equal(req.url, "/db/query"); + assert.equal(req.headers.authorization, "Basic ZGJ4OnNlY3JldA=="); + assert.deepEqual(JSON.parse(body), ["select id, name from users"]); + json(res, { results: [{ columns: ["id", "name"], values: [[1, "Ada"], [2, "Linus"]] }] }); + }); + + try { + const result = await executeQuery(rqliteConfig(server.port), "select id, name from users", { maxRows: 1 }); + assert.deepEqual(result, { columns: ["id", "name"], rows: [{ id: 1, name: "Ada" }], row_count: 1 }); + } finally { + await server.close(); + } +}); + +test("lists rqlite tables and describes columns", async () => { + const server = await withRqliteServer((_req, res, body) => { + const sql = JSON.parse(body)[0] as string; + if (sql.includes("sqlite_master")) { + json(res, { results: [{ columns: ["name", "type"], values: [["users", "table"], ["active_users", "view"]] }] }); + return; + } + if (sql.includes("PRAGMA table_info")) { + json(res, { + results: [ + { + columns: ["cid", "name", "type", "notnull", "dflt_value", "pk"], + values: [ + [0, "id", "INTEGER", 1, null, 1], + [1, "name", "TEXT", 0, "'unknown'", 0], + ], + }, + ], + }); + return; + } + json(res, { results: [{ columns: [], values: [] }] }); + }); + + try { + assert.deepEqual(await listTables(rqliteConfig(server.port)), [ + { name: "users", type: "table" }, + { name: "active_users", type: "view" }, + ]); + assert.deepEqual(await describeTable(rqliteConfig(server.port), "users"), [ + { + name: "id", + data_type: "INTEGER", + is_nullable: false, + column_default: null, + is_primary_key: true, + comment: null, + }, + { + name: "name", + data_type: "TEXT", + is_nullable: true, + column_default: "'unknown'", + is_primary_key: false, + comment: null, + }, + ]); + } finally { + await server.close(); + } +}); diff --git a/src-tauri/src/commands/connection.rs b/src-tauri/src/commands/connection.rs index 5be1349ed..7d80ba777 100644 --- a/src-tauri/src/commands/connection.rs +++ b/src-tauri/src/commands/connection.rs @@ -395,6 +395,19 @@ pub async fn test_connection(state: State<'_, Arc>, config: Connection .await .map(|_| "Connection successful".to_string()) } + DatabaseType::Rqlite => { + let client = db::rqlite_driver::RqliteClient::new( + &url, + config.url_params.as_deref(), + &config.username, + &config.password, + config.ssl, + connect_timeout, + )?; + db::rqlite_driver::test_connection(&client, connect_timeout) + .await + .map(|_| "Connection successful".to_string()) + } db_type if database_capabilities::is_agent_type(&db_type) => { test_agent_connection(state.inner(), &config, &host, port).await } @@ -549,6 +562,18 @@ pub async fn connect_db(state: State<'_, Arc>, config: ConnectionConfi db::elasticsearch_driver::test_connection(&mut client, connect_timeout).await?; PoolKind::Elasticsearch(client) } + DatabaseType::Rqlite => { + let client = db::rqlite_driver::RqliteClient::new( + &url, + db_config.url_params.as_deref(), + &db_config.username, + &db_config.password, + db_config.ssl, + connect_timeout, + )?; + db::rqlite_driver::test_connection(&client, connect_timeout).await?; + PoolKind::Rqlite(client) + } db_type if database_capabilities::is_agent_type(&db_type) => { connect_agent_pool(state.inner(), &db_config, &host, port).await? }