feat: add InfluxDB support (#980)
This commit is contained in:
parent
c13efc65f8
commit
0bd07ccf46
|
|
@ -0,0 +1,10 @@
|
|||
<?xml version="1.0" encoding="UTF-8" standalone="no"?>
|
||||
<svg width="154" height="156" viewBox="0 0 154 156" fill="none" xmlns="http://www.w3.org/2000/svg">
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="m 111.76244,105.55719 33.75203,-7.710196 c 0.51964,-0.104857 1.01231,-0.316521 1.4455,-0.621423 0.43411,-0.304901 0.7985,-0.696439 1.0718,-1.149608 0.27329,-0.453077 0.44899,-0.958023 0.51591,-1.482398 0.066,-0.524375 0.0223,-1.056837 -0.12921,-1.563456 l -14.3536,-62.024578 c -0.25936,-1.023558 -0.91564,-1.902939 -1.82477,-2.445254 -0.90912,-0.542316 -1.99673,-0.703319 -3.02484,-0.447685 l -33.738648,7.709824 c -1.021327,0.261025 -1.897268,0.914332 -2.43633,1.817045 -0.538969,0.902713 -0.697276,1.981302 -0.440062,2.99984 l 14.31261,62.025039 c 0.25936,1.02347 0.91471,1.90285 1.82383,2.44479 0.90913,0.54195 1.99767,0.70369 3.02578,0.44806 z" fill="#d30971" />
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="m 100.13156,143.8846 40.92187,-37.74642 c 1.5431,-1.55054 1.15361,-2.50522 -0.97234,-1.74389 l -28.12437,6.36667 c -1.1629,0.28631 -2.2347,0.85986 -3.11687,1.6686 -0.88124,0.8078 -1.54403,1.82383 -1.92702,2.95512 l -8.532589,27.35097 c -0.58331,2.11757 0.195211,2.69856 1.751319,1.14895 z" fill="#d30971" />
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="m 25.47096,131.56491 61.279242,18.86671 c 1.046984,0.19428 2.129198,0.0111 3.053291,-0.51685 0.924,-0.52707 1.629642,-1.36462 1.990783,-2.36205 l 10.268954,-32.74995 c 0.14873,-0.51034 0.19428,-1.04578 0.13479,-1.57378 -0.0586,-0.52893 -0.22217,-1.04019 -0.4806,-1.50591 -0.25842,-0.46572 -0.60608,-0.87566 -1.02346,-1.20752 -0.41738,-0.33093 -0.897043,-0.57727 -1.41017,-0.72322 L 38.004733,91.105883 c -1.037315,-0.295606 -2.149927,-0.170392 -3.095028,0.34822 -0.945009,0.518704 -1.645632,1.388696 -1.949046,2.420155 L 22.761334,126.55449 c -0.302949,1.02253 -0.188147,2.12222 0.31931,3.0611 0.507549,0.93794 1.366665,1.63884 2.390316,1.94932 z" fill="#d30971" />
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="m 5.6834977,64.52991 12.4226043,54.1901 c 0.389028,2.11758 1.389533,2.11758 1.931477,0 l 8.531849,-27.35115 c 0.290772,-1.165411 0.30267,-2.3826 0.03477,-3.553403 C 28.33629,86.644654 27.796019,85.552772 27.02698,84.628028 L 7.4343376,63.560918 C 6.0725724,61.969203 5.0581988,62.412053 5.6834977,64.52991 Z" fill="#d30971" />
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="M 53.595516,6.602436 6.6564671,49.885345 c -0.7530573,0.714101 -1.1951843,1.693783 -1.2315519,2.728775 -0.036368,1.035084 0.3359197,2.043118 1.0370099,2.807882 L 29.931445,80.710788 c 0.723304,0.764299 1.72009,1.213378 2.773859,1.249632 1.05377,0.03635 2.079373,-0.343107 2.853897,-1.055815 l 46.92517,-43.352152 c 0.759558,-0.710384 1.206313,-1.691181 1.242752,-2.728404 0.03644,-1.037222 -0.340411,-2.046744 -1.048285,-2.808254 L 59.250973,6.7408501 C 58.890669,6.3587779 58.457672,6.0518109 57.977266,5.8377701 57.496767,5.6237302 56.978435,5.5068801 56.452294,5.4940029 55.926153,5.4811258 55.402708,5.5724784 54.912263,5.7627578 54.42191,5.9530368 53.974319,6.2384524 53.595516,6.602436 Z" fill="#d30971" />
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="m 98.964945,104.58764 c 2.139885,0.58192 3.487775,-0.56704 2.917945,-2.76828 L 88.30678,43.269262 C 87.723098,41.151498 85.972332,40.57014 84.429884,42.106546 L 40.214339,83.022555 c -1.556299,1.536406 -1.167177,3.266629 0.958767,3.847986 z" fill="#d30971" />
|
||||
<path fill-rule="evenodd" clip-rule="evenodd" d="M 122.0445,22.922017 68.519141,6.6026591 C 66.393197,6.0213059 66.004075,6.7964388 67.740991,8.5266523 L 87.333698,29.524435 c 0.878916,0.818493 1.941981,1.415188 3.099956,1.740262 1.158067,0.325166 2.377394,0.369135 3.556006,0.128374 l 28.12456,-6.353291 c 2.07017,-0.581357 2.07017,-1.550256 -0.0697,-2.117763 z" fill="#d30971" />
|
||||
</svg>
|
||||
|
After Width: | Height: | Size: 3.7 KiB |
|
|
@ -416,6 +416,7 @@ const driverProfiles: Record<
|
|||
iotdb: { type: "iotdb", port: 6667, user: "root", label: "Apache IoTDB", icon: "iotdb" },
|
||||
etcd: { type: "etcd", port: 2379, user: "", label: "etcd", icon: "etcd" },
|
||||
iris: { type: "iris", port: 1972, user: "_SYSTEM", label: "IRIS", icon: "iris" },
|
||||
influxdb: { type: "influxdb", port: 8086, user: "", label: "InfluxDB", icon: "InfluxDB" },
|
||||
custom_mysql: {
|
||||
type: "mysql",
|
||||
port: 3306,
|
||||
|
|
@ -715,6 +716,7 @@ const iconTypeMap: Record<string, string> = {
|
|||
bigquery: "bigquery",
|
||||
kylin: "kylin",
|
||||
sundb: "sundb",
|
||||
influxdb: "influxdb",
|
||||
jdbc: "jdbc",
|
||||
custom_mysql: "mysql",
|
||||
custom_postgres: "postgres",
|
||||
|
|
@ -777,6 +779,7 @@ const dbOptions = [
|
|||
{ value: "xugu", label: "虚谷 XuguDB" },
|
||||
{ value: "iotdb", label: "Apache IoTDB" },
|
||||
{ value: "etcd", label: "etcd" },
|
||||
{ value: "influxdb", label: "InfluxDB" },
|
||||
{ value: "iris", label: "IRIS" },
|
||||
{ value: "jdbc", label: "JDBC" },
|
||||
{ value: "custom_mysql", label: "Custom (MySQL)" },
|
||||
|
|
@ -824,7 +827,7 @@ const sqliteExtensionPaths = computed({
|
|||
form.value.url_params = setSqliteExtensionPaths(form.value.url_params, value);
|
||||
},
|
||||
});
|
||||
const tlsCapableDatabaseTypes = new Set<DatabaseType>(["mysql", "postgres", "redshift", "gaussdb", "kwdb", "opengauss", "redis", "etcd", "clickhouse", "elasticsearch"]);
|
||||
const tlsCapableDatabaseTypes = new Set<DatabaseType>(["mysql", "postgres", "redshift", "gaussdb", "kwdb", "opengauss", "redis", "etcd", "clickhouse", "elasticsearch", "influxdb"]);
|
||||
const supportsTlsToggle = computed(() => tlsCapableDatabaseTypes.has(form.value.db_type));
|
||||
const supportsCaCertificatePath = computed(() => form.value.db_type === "clickhouse");
|
||||
const bareMysqlProfiles = new Set(["doris", "starrocks", "selectdb", "oceanbase"]);
|
||||
|
|
|
|||
|
|
@ -73,6 +73,7 @@ const assetIcons: Record<string, string> = {
|
|||
iotdb: "iotdb",
|
||||
etcd: "etcd",
|
||||
iris: "iris.png",
|
||||
influxdb: "influxdb",
|
||||
};
|
||||
|
||||
const letterIcons: Record<string, { letter: string; color: string }> = {};
|
||||
|
|
|
|||
|
|
@ -184,7 +184,11 @@ function getIconInfo(node: TreeNode): { icon: any; colorClass: string } | null {
|
|||
case "view":
|
||||
return { icon: Eye, colorClass: "text-purple-500" };
|
||||
case "column":
|
||||
return { icon: Columns3, colorClass: "text-muted-foreground" };
|
||||
if ((node.meta as ColumnInfo).is_primary_key) {
|
||||
return { icon: Columns3, colorClass: "text-orange-400" };
|
||||
} else {
|
||||
return { icon: Columns3, colorClass: "text-muted-foreground" };
|
||||
}
|
||||
case "group-columns":
|
||||
return { icon: ListTree, colorClass: "text-green-400" };
|
||||
case "group-indexes":
|
||||
|
|
|
|||
|
|
@ -145,6 +145,9 @@ export function connectionUrlPlaceholder(dbType: DatabaseType): string {
|
|||
case "iris":
|
||||
return "iris://user:password@host:port/namespace";
|
||||
|
||||
case "influxdb":
|
||||
return "influxdb://user:password@host:port/database";
|
||||
|
||||
case "jdbc":
|
||||
return "jdbc:mysql://host:3306/database";
|
||||
|
||||
|
|
|
|||
|
|
@ -145,7 +145,7 @@ export const TABLE_STRUCTURE_SUPPORTED_TYPES = new Set<DatabaseType>([
|
|||
"access",
|
||||
]);
|
||||
|
||||
export const CREATE_DATABASE_SUPPORTED_TYPES = new Set<DatabaseType>(["mysql", "postgres", "sqlserver", "clickhouse", "oracle", "gaussdb", "kwdb", "opengauss", "oceanbase-oracle", "doris", "starrocks", "redshift"]);
|
||||
export const CREATE_DATABASE_SUPPORTED_TYPES = new Set<DatabaseType>(["mysql", "postgres", "sqlserver", "clickhouse", "oracle", "gaussdb", "kwdb", "opengauss", "oceanbase-oracle", "doris", "starrocks", "redshift", "influxdb"]);
|
||||
|
||||
export const FIELD_LINEAGE_SUPPORTED_TYPES = new Set<DatabaseType>(["mysql", "postgres", "sqlite", "rqlite", "turso", "sqlserver", "oracle", "redshift", "dameng", "gaussdb", "kwdb", "opengauss", "oceanbase-oracle"]);
|
||||
|
||||
|
|
|
|||
|
|
@ -96,7 +96,7 @@ export function supportsObjectBrowserTreeNode(dbType: DatabaseType | undefined,
|
|||
}
|
||||
|
||||
export function supportsTableTruncate(dbType?: DatabaseType): boolean {
|
||||
return !!dbType && dbType !== "sqlite" && dbType !== "rqlite" && dbType !== "turso" && dbType !== "duckdb";
|
||||
return !!dbType && dbType !== "sqlite" && dbType !== "rqlite" && dbType !== "turso" && dbType !== "duckdb" && dbType !== "influxdb";
|
||||
}
|
||||
|
||||
export function usesPostgresLikeStructureCopy(dbType?: DatabaseType): boolean {
|
||||
|
|
|
|||
|
|
@ -137,6 +137,18 @@ const DATABASE_CAPABILITY_OVERRIDES: Partial<Record<DatabaseType, Partial<Databa
|
|||
transaction: false,
|
||||
},
|
||||
},
|
||||
influxdb: {
|
||||
tableData: {
|
||||
insert: false,
|
||||
updateRequiresPrimaryKey: false,
|
||||
deleteRequiresPrimaryKey: true,
|
||||
keylessRowPredicate: false,
|
||||
requiresTransactionalTableForExistingRows: false,
|
||||
existingRowsReadonly: true,
|
||||
transaction: false,
|
||||
readonly: true,
|
||||
},
|
||||
},
|
||||
};
|
||||
|
||||
function defaultTableDataCapability(dbType?: DatabaseType): TableDataCapability {
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ const profileMap: Record<string, ConnectionProfile> = {
|
|||
gaussdb: { dbType: "gaussdb", profile: "gaussdb", label: "GaussDB", port: 5432, user: "gaussdb" },
|
||||
kwdb: { dbType: "kwdb", profile: "kwdb", label: "KWDB", port: 26257, user: "root" },
|
||||
opengauss: { dbType: "gaussdb", profile: "opengauss", label: "openGauss", port: 5432, user: "gaussdb" },
|
||||
influxdb: { dbType: "influxdb", profile: "influxdb", label: "InfluxDB", port: 8086, user: "" },
|
||||
};
|
||||
|
||||
function normalizeKey(value: unknown) {
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
import type { DatabaseType } from "@/types/database";
|
||||
|
||||
export type TableStructureDialect = "mysql" | "postgres" | "sqlite" | "duckdb" | "sqlserver" | "oracle" | "h2" | "clickhouse" | "unsupported";
|
||||
export type TableStructureDialect = "mysql" | "postgres" | "sqlite" | "duckdb" | "sqlserver" | "oracle" | "h2" | "clickhouse" | "influxdb" | "unsupported";
|
||||
|
||||
export interface TableStructureCapabilities {
|
||||
dialect: TableStructureDialect;
|
||||
|
|
@ -199,6 +199,20 @@ const accessCapabilities = capabilities({
|
|||
createIndex: true,
|
||||
});
|
||||
|
||||
const influxdbCapabilities = capabilities({
|
||||
dialect: "influxdb",
|
||||
createTable: false,
|
||||
addColumn: false,
|
||||
dropColumn: false,
|
||||
renameColumn: false,
|
||||
alterExistingColumn: false,
|
||||
alterType: false,
|
||||
alterNullability: false,
|
||||
alterDefault: false,
|
||||
reorderColumn: false,
|
||||
comment: false,
|
||||
});
|
||||
|
||||
const capabilityByType: Partial<Record<DatabaseType, TableStructureCapabilities>> = {
|
||||
mysql: mysqlCapabilities,
|
||||
doris: mysqlCapabilities,
|
||||
|
|
@ -225,6 +239,7 @@ const capabilityByType: Partial<Record<DatabaseType, TableStructureCapabilities>
|
|||
h2: h2Capabilities,
|
||||
access: accessCapabilities,
|
||||
clickhouse: clickhouseCapabilities,
|
||||
influxdb: influxdbCapabilities,
|
||||
};
|
||||
|
||||
export function getTableStructureCapabilities(dbType?: DatabaseType): TableStructureCapabilities {
|
||||
|
|
|
|||
|
|
@ -257,6 +257,7 @@ export const useConnectionStore = defineStore("connection", () => {
|
|||
bigquery: "BigQuery",
|
||||
kylin: "Kylin",
|
||||
sundb: "SunDB",
|
||||
influxdb: "InfluxDB",
|
||||
};
|
||||
|
||||
const profile = config.driver_profile || config.db_type;
|
||||
|
|
@ -1201,7 +1202,8 @@ export const useConnectionStore = defineStore("connection", () => {
|
|||
},
|
||||
];
|
||||
|
||||
if (node.type === "table") {
|
||||
const config = getConfig(connectionId);
|
||||
if (node.type === "table" && config?.db_type !== "influxdb") {
|
||||
children.push(
|
||||
{
|
||||
id: `${parentId}:__indexes`,
|
||||
|
|
|
|||
|
|
@ -49,6 +49,7 @@ export type DatabaseType =
|
|||
| "iotdb"
|
||||
| "etcd"
|
||||
| "iris"
|
||||
| "influxdb"
|
||||
| "jdbc";
|
||||
|
||||
export interface SqlSnippet {
|
||||
|
|
|
|||
|
|
@ -526,6 +526,17 @@
|
|||
"skipTcpProbe": true,
|
||||
"defaultPort": 1972
|
||||
},
|
||||
{
|
||||
"dbType": "influxdb",
|
||||
"label": "InfluxDB",
|
||||
"runtimeMode": "agent",
|
||||
"mcpMode": "bridge",
|
||||
"agentKey": "influxdb",
|
||||
"singleConnectionPool": false,
|
||||
"metadataConnectionScoped": false,
|
||||
"skipTcpProbe": true,
|
||||
"defaultPort": 8086
|
||||
},
|
||||
{
|
||||
"dbType": "jdbc",
|
||||
"label": "JDBC",
|
||||
|
|
|
|||
|
|
@ -48,6 +48,7 @@ pub enum PoolKind {
|
|||
ClickHouse(db::clickhouse_driver::ChClient),
|
||||
SqlServer(Arc<tokio::sync::Mutex<db::sqlserver::SqlServerClient>>),
|
||||
Elasticsearch(db::elasticsearch_driver::EsClient),
|
||||
InfluxDb(db::influxdb_driver::InfluxdbClient),
|
||||
Agent(Arc<tokio::sync::Mutex<db::agent_driver::AgentDriverClient>>),
|
||||
ExternalTabular(Arc<external::ExternalPool>),
|
||||
ExternalDriver { driver_id: String, config: Arc<ConnectionConfig>, session: Arc<PluginDriverSession> },
|
||||
|
|
@ -482,6 +483,19 @@ impl AppState {
|
|||
db::elasticsearch_driver::test_connection(&mut client, connect_timeout).await?;
|
||||
PoolKind::Elasticsearch(client)
|
||||
}
|
||||
DatabaseType::InfluxDb => {
|
||||
let username = if db_config.username.is_empty() { None } else { Some(db_config.username.clone()) };
|
||||
let password = if db_config.password.is_empty() { None } else { Some(db_config.password.clone()) };
|
||||
let client = db::influxdb_driver::InfluxdbClient::new_with_ca_cert(
|
||||
&url,
|
||||
username,
|
||||
password,
|
||||
Some(&db_config.ca_cert_path),
|
||||
connect_timeout,
|
||||
)?;
|
||||
db::influxdb_driver::test_connection(&client, connect_timeout).await?;
|
||||
PoolKind::InfluxDb(client)
|
||||
}
|
||||
DatabaseType::Dameng
|
||||
| DatabaseType::Kingbase
|
||||
| DatabaseType::Highgo
|
||||
|
|
@ -960,6 +974,7 @@ pub async fn close_pool_kind(pool: PoolKind) {
|
|||
PoolKind::ClickHouse(_) => {}
|
||||
PoolKind::SqlServer(_) => {}
|
||||
PoolKind::Elasticsearch(_) => {}
|
||||
PoolKind::InfluxDb(_) => {}
|
||||
PoolKind::Agent(client) => {
|
||||
let mut client = client.lock().await;
|
||||
let _ = client.disconnect().await;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,313 @@
|
|||
use percent_encoding::{utf8_percent_encode, NON_ALPHANUMERIC};
|
||||
use reqwest::{Certificate, Client as HttpClient};
|
||||
use serde::Deserialize;
|
||||
use std::fs;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use super::with_connection_timeout;
|
||||
use crate::sql::starts_with_executable_sql_keyword;
|
||||
use crate::types::{ColumnInfo, DatabaseInfo, QueryResult, TableInfo};
|
||||
|
||||
pub struct InfluxdbClient {
|
||||
http: HttpClient,
|
||||
base_url: String,
|
||||
username: Option<String>,
|
||||
password: Option<String>,
|
||||
}
|
||||
|
||||
impl InfluxdbClient {
|
||||
pub fn new(url: &str, username: Option<String>, password: Option<String>, timeout: Duration) -> Self {
|
||||
let http = HttpClient::builder().connect_timeout(timeout).build().unwrap_or_else(|_| HttpClient::new());
|
||||
Self { http, base_url: url.trim_end_matches('/').to_string(), username, password }
|
||||
}
|
||||
|
||||
pub fn new_with_ca_cert(
|
||||
url: &str,
|
||||
username: Option<String>,
|
||||
password: Option<String>,
|
||||
ca_cert_path: Option<&str>,
|
||||
timeout: Duration,
|
||||
) -> Result<Self, String> {
|
||||
let mut builder = HttpClient::builder().connect_timeout(timeout);
|
||||
if let Some(path) = ca_cert_path.map(str::trim).filter(|path| !path.is_empty()) {
|
||||
let path = expand_cert_path(path);
|
||||
let cert_bytes =
|
||||
fs::read(&path).map_err(|e| format!("Failed to read InfluxDB CA certificate at {path}: {e}"))?;
|
||||
let cert = Certificate::from_pem(&cert_bytes)
|
||||
.or_else(|_| Certificate::from_der(&cert_bytes))
|
||||
.map_err(|e| format!("Failed to parse InfluxDB CA certificate at {path}: {e}"))?;
|
||||
builder = builder.add_root_certificate(cert);
|
||||
}
|
||||
let http = builder.build().map_err(|e| format!("Failed to configure InfluxDB HTTP client: {e}"))?;
|
||||
Ok(Self { http, base_url: url.trim_end_matches('/').to_string(), username, password })
|
||||
}
|
||||
}
|
||||
|
||||
fn expand_cert_path(path: &str) -> String {
|
||||
let home = || std::env::var(if cfg!(windows) { "USERPROFILE" } else { "HOME" }).ok();
|
||||
if path == "~" || path.starts_with("~/") || path.starts_with("~\\") {
|
||||
if let Some(home) = home() {
|
||||
return format!("{}{}", home, &path[1..]);
|
||||
}
|
||||
}
|
||||
if let Some(rest) = path.strip_prefix("$HOME") {
|
||||
if let Some(home) = home() {
|
||||
return format!("{home}{rest}");
|
||||
}
|
||||
}
|
||||
if let Some(rest) = path.strip_prefix("${HOME}") {
|
||||
if let Some(home) = home() {
|
||||
return format!("{home}{rest}");
|
||||
}
|
||||
}
|
||||
if let Some(rest) = path.strip_prefix("%USERPROFILE%") {
|
||||
if let Ok(home) = std::env::var("USERPROFILE") {
|
||||
return format!("{home}{rest}");
|
||||
}
|
||||
}
|
||||
path.to_string()
|
||||
}
|
||||
|
||||
impl Clone for InfluxdbClient {
|
||||
fn clone(&self) -> Self {
|
||||
Self {
|
||||
http: self.http.clone(),
|
||||
base_url: self.base_url.clone(),
|
||||
username: self.username.clone(),
|
||||
password: self.password.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Default)]
|
||||
struct InfluxErrorResult {
|
||||
#[serde(default)]
|
||||
error: String,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct InfluxJsonResult {
|
||||
results: Vec<InfluxQueryResult>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
#[allow(dead_code)]
|
||||
struct InfluxQueryResult {
|
||||
statement_id: usize,
|
||||
#[serde(default)]
|
||||
#[allow(dead_code)]
|
||||
series: Vec<InfluxSeries>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
#[allow(dead_code)]
|
||||
struct InfluxSeries {
|
||||
#[serde(default)]
|
||||
#[allow(dead_code)]
|
||||
name: String,
|
||||
columns: Vec<String>,
|
||||
values: Vec<Vec<serde_json::Value>>,
|
||||
}
|
||||
|
||||
fn build_query_url(base_url: &str, database: Option<&str>, sql: &str) -> String {
|
||||
let mut url = format!("{}/query", base_url);
|
||||
let mut has_param = false;
|
||||
if let Some(db) = database {
|
||||
url.push_str(&format!("?db={db}"));
|
||||
has_param = true;
|
||||
}
|
||||
if has_param {
|
||||
url.push('&');
|
||||
} else {
|
||||
url.push('?');
|
||||
}
|
||||
let encoded_sql = utf8_percent_encode(sql, NON_ALPHANUMERIC);
|
||||
url.push_str(&format!("q={encoded_sql}"));
|
||||
url
|
||||
}
|
||||
|
||||
fn build_request(client: &InfluxdbClient, req: reqwest::RequestBuilder) -> reqwest::RequestBuilder {
|
||||
match (&client.username, &client.password) {
|
||||
(Some(u), Some(p)) if !u.is_empty() => req.basic_auth(u, Some(p)),
|
||||
(Some(u), None) if !u.is_empty() => req.basic_auth(u, None::<&str>),
|
||||
_ => req,
|
||||
}
|
||||
}
|
||||
|
||||
async fn influx_query(client: &InfluxdbClient, sql: &str, database: Option<&str>) -> Result<InfluxJsonResult, String> {
|
||||
let url = build_query_url(&client.base_url, database, sql);
|
||||
log::info!("[influxdb] query url={url} username={:?} password={}", client.username, client.password.is_some());
|
||||
|
||||
let req = if starts_with_executable_sql_keyword(sql, &["SELECT", "SHOW"]) {
|
||||
build_request(client, client.http.get(&url))
|
||||
} else {
|
||||
build_request(client, client.http.post(&url))
|
||||
};
|
||||
|
||||
let resp = req.send().await.map_err(|e| format!("InfluxDB request failed: {e}"))?;
|
||||
log::info!("[influxdb] response status={}", resp.status());
|
||||
if !resp.status().is_success() {
|
||||
let error_json = resp.json::<InfluxErrorResult>().await.unwrap_or_default();
|
||||
let msg = error_json.error;
|
||||
log::error!("[influxdb] error: {msg}");
|
||||
return Err(format!("InfluxDB error: {msg}"));
|
||||
}
|
||||
resp.json::<InfluxJsonResult>().await.map_err(|e| format!("InfluxDB parse error: {e}"))
|
||||
}
|
||||
|
||||
pub async fn test_connection(client: &InfluxdbClient, timeout: Duration) -> Result<(), String> {
|
||||
let url = format!("{}/query?q=SHOW DATABASES", client.base_url);
|
||||
let req = build_request(client, client.http.get(&url));
|
||||
let resp = with_connection_timeout("InfluxDB", timeout, async {
|
||||
req.send().await.map_err(|e| format!("InfluxDB connection failed: {e}"))
|
||||
})
|
||||
.await?;
|
||||
if !resp.status().is_success() {
|
||||
let body = resp.text().await.unwrap_or_default();
|
||||
return Err(format!("InfluxDB error: {body}"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub async fn list_databases(client: &InfluxdbClient) -> Result<Vec<DatabaseInfo>, String> {
|
||||
let result = influx_query(client, "SHOW DATABASES", None).await?;
|
||||
Ok(result
|
||||
.results
|
||||
.iter()
|
||||
.flat_map(|r| &r.series)
|
||||
.flat_map(|s| &s.values)
|
||||
.map(|row| DatabaseInfo { name: row[0].as_str().unwrap_or("").to_string() })
|
||||
.collect())
|
||||
}
|
||||
|
||||
pub async fn list_tables(client: &InfluxdbClient, database: &str) -> Result<Vec<TableInfo>, String> {
|
||||
let result = influx_query(client, "SHOW MEASUREMENTS", Some(database)).await?;
|
||||
let empty = vec![];
|
||||
let series = result.results.first().map(|r| &r.series).unwrap_or(&empty);
|
||||
if series.is_empty() {
|
||||
return Ok(vec![]);
|
||||
}
|
||||
let first_series = &series[0];
|
||||
Ok(first_series
|
||||
.values
|
||||
.iter()
|
||||
.map(|row| TableInfo {
|
||||
name: row[0].as_str().unwrap_or("").to_string(),
|
||||
table_type: "TABLE".to_string(),
|
||||
comment: None,
|
||||
parent_schema: None,
|
||||
parent_name: None,
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
pub async fn get_columns(client: &InfluxdbClient, database: &str, table: &str) -> Result<Vec<ColumnInfo>, String> {
|
||||
let empty = vec![];
|
||||
|
||||
let tag_sql = format!("SHOW TAG KEYS FROM \"{}\"", table);
|
||||
let tag_result = influx_query(client, &tag_sql, Some(database)).await?;
|
||||
let tag_series = tag_result.results.first().map(|r| &r.series).unwrap_or(&empty);
|
||||
|
||||
let field_sql = format!("SHOW FIELD KEYS FROM \"{}\"", table);
|
||||
let field_result = influx_query(client, &field_sql, Some(database)).await?;
|
||||
let field_series = field_result.results.first().map(|r| &r.series).unwrap_or(&empty);
|
||||
|
||||
let time_col = ColumnInfo {
|
||||
name: "time".to_string(),
|
||||
data_type: "timestamp".to_string(),
|
||||
is_nullable: false,
|
||||
column_default: None,
|
||||
is_primary_key: true,
|
||||
extra: None,
|
||||
comment: None,
|
||||
numeric_precision: None,
|
||||
numeric_scale: None,
|
||||
character_maximum_length: None,
|
||||
};
|
||||
|
||||
let cols: Vec<ColumnInfo> = std::iter::once(time_col)
|
||||
.chain(tag_series.first().into_iter().flat_map(|s| s.values.iter()).map(|row| ColumnInfo {
|
||||
name: row[0].as_str().unwrap_or("").to_string(),
|
||||
data_type: "string".to_string(),
|
||||
is_nullable: true,
|
||||
column_default: None,
|
||||
is_primary_key: true,
|
||||
extra: None,
|
||||
comment: None,
|
||||
numeric_precision: None,
|
||||
numeric_scale: None,
|
||||
character_maximum_length: None,
|
||||
}))
|
||||
.chain(field_series.first().into_iter().flat_map(|s| s.values.iter()).map(|row| {
|
||||
let data_type = row.get(1).and_then(|v| v.as_str()).unwrap_or("unknown").to_string();
|
||||
ColumnInfo {
|
||||
name: row[0].as_str().unwrap_or("").to_string(),
|
||||
data_type,
|
||||
is_nullable: true,
|
||||
column_default: None,
|
||||
is_primary_key: false,
|
||||
extra: None,
|
||||
comment: None,
|
||||
numeric_precision: None,
|
||||
numeric_scale: None,
|
||||
character_maximum_length: None,
|
||||
}
|
||||
}))
|
||||
.collect();
|
||||
|
||||
Ok(cols)
|
||||
}
|
||||
|
||||
pub async fn execute_query(client: &InfluxdbClient, database: &str, sql: &str) -> Result<QueryResult, String> {
|
||||
let start = Instant::now();
|
||||
let url = build_query_url(&client.base_url, Some(database), sql);
|
||||
let req = if starts_with_executable_sql_keyword(sql, &["SELECT", "SHOW"]) {
|
||||
build_request(client, client.http.get(&url))
|
||||
} else {
|
||||
build_request(client, client.http.post(&url))
|
||||
};
|
||||
let resp = req.send().await.map_err(|e| format!("InfluxDB request failed: {e}"))?;
|
||||
if !resp.status().is_success() {
|
||||
let error_json = resp.json::<InfluxErrorResult>().await.unwrap_or_default();
|
||||
let msg = error_json.error;
|
||||
return Err(format!("InfluxDB error: {msg}"));
|
||||
}
|
||||
let json = resp.json::<InfluxJsonResult>().await.map_err(|e| format!("InfluxDB parse error: {e}"))?;
|
||||
let series = json.results.iter().flat_map(|r| &r.series).next();
|
||||
match series {
|
||||
Some(s) => Ok(QueryResult {
|
||||
columns: s.columns.clone(),
|
||||
column_types: vec![],
|
||||
column_sortables: s.columns.iter().map(|_| false).collect(),
|
||||
rows: s.values.clone(),
|
||||
affected_rows: s.values.len() as u64,
|
||||
execution_time_ms: start.elapsed().as_millis(),
|
||||
truncated: false,
|
||||
session_id: None,
|
||||
has_more: false,
|
||||
}),
|
||||
None => Ok(QueryResult {
|
||||
columns: vec![],
|
||||
column_types: vec![],
|
||||
column_sortables: vec![],
|
||||
rows: vec![],
|
||||
affected_rows: 0,
|
||||
execution_time_ms: start.elapsed().as_millis(),
|
||||
truncated: false,
|
||||
session_id: None,
|
||||
has_more: false,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn query_url() {
|
||||
let url = build_query_url("http://localhost:8086", Some("sample"), "SHOW DATABASES");
|
||||
|
||||
assert_eq!(url, "http://localhost:8086/query?db=sample&q=SHOW%20DATABASES");
|
||||
}
|
||||
}
|
||||
|
|
@ -4,6 +4,7 @@ pub mod duckdb_driver;
|
|||
pub mod elasticsearch_driver;
|
||||
pub mod elasticsearch_sql;
|
||||
pub mod file_validator;
|
||||
pub mod influxdb_driver;
|
||||
pub mod mongo_driver;
|
||||
pub mod mysql;
|
||||
pub mod ob_oracle;
|
||||
|
|
|
|||
|
|
@ -165,6 +165,8 @@ pub fn build_drop_table_sql(options: TableAdminSqlOptions) -> String {
|
|||
let table = qualified_name(options.database_type, options.schema.as_deref(), &options.table_name);
|
||||
if matches!(options.database_type, Some(DatabaseType::Iotdb)) {
|
||||
return format!("DELETE TIMESERIES {};", iotdb_timeseries_pattern(&table));
|
||||
} else if matches!(options.database_type, Some(DatabaseType::InfluxDb)) {
|
||||
return format!("DROP MEASUREMENT {};", table);
|
||||
}
|
||||
format!("DROP TABLE {table};")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -274,6 +274,8 @@ pub enum DatabaseType {
|
|||
Iris,
|
||||
#[serde(rename = "turso")]
|
||||
Turso,
|
||||
#[serde(rename = "influxdb")]
|
||||
InfluxDb,
|
||||
Jdbc,
|
||||
}
|
||||
|
||||
|
|
@ -728,6 +730,10 @@ impl ConnectionConfig {
|
|||
format!("etcd://{host}:{port}")
|
||||
}
|
||||
DatabaseType::Iris => format!("iris://{host}:{port}{db_part}"),
|
||||
DatabaseType::InfluxDb => {
|
||||
let scheme = if self.ssl { "https" } else { "http" };
|
||||
format!("{scheme}://{host}:{port}")
|
||||
}
|
||||
DatabaseType::Jdbc => "jdbc:<redacted>".to_string(),
|
||||
}
|
||||
}
|
||||
|
|
@ -918,6 +924,10 @@ impl ConnectionConfig {
|
|||
DatabaseType::Iris => {
|
||||
format!("iris://{}:{}@{host}:{port}{db_part}", username, password)
|
||||
}
|
||||
DatabaseType::InfluxDb => {
|
||||
let scheme = if self.ssl { "https" } else { "http" };
|
||||
format!("{scheme}://{host}:{port}")
|
||||
}
|
||||
DatabaseType::Jdbc => {
|
||||
self.connection_string.as_deref().filter(|value| !value.is_empty()).unwrap_or("jdbc:").to_string()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -749,6 +749,15 @@ pub async fn do_execute(
|
|||
}
|
||||
PoolKind::Redis(_) => Err("Use Redis-specific commands".to_string()),
|
||||
PoolKind::MongoDb(_) => Err("Use MongoDB-specific commands".to_string()),
|
||||
PoolKind::InfluxDb(client) => {
|
||||
let client = client.clone();
|
||||
let database = pool_key.split(':').nth(1).unwrap_or("default").to_string();
|
||||
let max_rows = options.max_rows;
|
||||
drop(connections);
|
||||
wait_for_query_opt(cancel_token, query_timeout, db::influxdb_driver::execute_query(&client, &database, sql))
|
||||
.await
|
||||
.map(|result| truncate_result_with_max_rows(result, max_rows))
|
||||
}
|
||||
PoolKind::Agent(client) => {
|
||||
let client = client.clone();
|
||||
let sql = sql.to_string();
|
||||
|
|
@ -1297,6 +1306,7 @@ pub async fn execute_statements_in_transaction(
|
|||
| PoolKind::Redis(_)
|
||||
| PoolKind::MongoDb(_)
|
||||
| PoolKind::Elasticsearch(_)
|
||||
| PoolKind::InfluxDb(_)
|
||||
| PoolKind::ExternalTabular(_)
|
||||
| PoolKind::ExternalDriver { .. } => TxPath::None,
|
||||
})
|
||||
|
|
|
|||
|
|
@ -275,6 +275,10 @@ async fn list_databases_once(state: &AppState, connection_id: &str) -> Result<Ve
|
|||
drop(connections);
|
||||
return db::clickhouse_driver::list_databases(&client).await;
|
||||
}
|
||||
if let Some(client) = extract_pool!(&connections, connection_id, InfluxDb) {
|
||||
drop(connections);
|
||||
return db::influxdb_driver::list_databases(&client).await;
|
||||
}
|
||||
try_sqlserver!(connections, connection_id, list_databases);
|
||||
if let Some(client) = extract_pool!(&connections, connection_id, Agent) {
|
||||
let is_mongo =
|
||||
|
|
@ -448,6 +452,10 @@ async fn list_tables_once(
|
|||
drop(connections);
|
||||
return db::clickhouse_driver::list_tables(&client, clickhouse_metadata_database(database, schema)).await;
|
||||
}
|
||||
if let Some(client) = extract_pool!(&connections, &pool_key, InfluxDb) {
|
||||
drop(connections);
|
||||
return db::influxdb_driver::list_tables(&client, database).await;
|
||||
}
|
||||
try_sqlserver!(connections, &pool_key, list_tables, schema, filter, limit);
|
||||
if let Some(client) = extract_pool!(&connections, &pool_key, Agent) {
|
||||
let fallback_config = db_config.clone();
|
||||
|
|
@ -1236,6 +1244,10 @@ pub async fn get_columns_core(
|
|||
.await
|
||||
.map(deduplicate_column_infos);
|
||||
}
|
||||
if let Some(client) = extract_pool!(&connections, &pool_key, InfluxDb) {
|
||||
drop(connections);
|
||||
return db::influxdb_driver::get_columns(&client, database, table).await.map(deduplicate_column_infos);
|
||||
}
|
||||
try_sqlserver!(connections, &pool_key, get_columns, schema, table);
|
||||
if let Some(client) = extract_pool!(&connections, &pool_key, Agent) {
|
||||
let fallback_config = db_config.clone();
|
||||
|
|
|
|||
|
|
@ -58,7 +58,10 @@ pub fn build_table_data_select_sql(options: TableDataSelectSqlOptions) -> String
|
|||
let row_id_alias =
|
||||
if options.include_row_id && database_type == Some(DatabaseType::Oracle) { Some("t") } else { None };
|
||||
let default_order_alias = if database_type == Some(DatabaseType::Jdbc) { Some("dbx_t") } else { row_id_alias };
|
||||
let default_order_by = if !options.primary_keys.is_empty() {
|
||||
let default_order_by = if database_type == Some(DatabaseType::InfluxDb) {
|
||||
// InfluxQL only allows sorting of timestamp column
|
||||
Some("time DESC".to_string())
|
||||
} else if !options.primary_keys.is_empty() {
|
||||
Some(
|
||||
options
|
||||
.primary_keys
|
||||
|
|
|
|||
|
|
@ -1822,6 +1822,13 @@ pub async fn get_columns_for_transfer(
|
|||
let mut client = client.lock().await;
|
||||
return db::sqlserver::get_columns(&mut client, &schema, &table).await;
|
||||
}
|
||||
if let Some(PoolKind::InfluxDb(client)) = connections.get(pool_key) {
|
||||
let client = client.clone();
|
||||
let database = database.to_string();
|
||||
let table = table.to_string();
|
||||
drop(connections);
|
||||
return db::influxdb_driver::get_columns(&client, &database, &table).await;
|
||||
}
|
||||
if let Some(PoolKind::Agent(client)) = connections.get(pool_key) {
|
||||
let client = client.clone();
|
||||
let database = database.to_string();
|
||||
|
|
|
|||
|
|
@ -426,6 +426,20 @@ pub async fn test_connection(state: State<'_, Arc<AppState>>, config: Connection
|
|||
.await
|
||||
.map(|_| "Connection successful".to_string())
|
||||
}
|
||||
DatabaseType::InfluxDb => {
|
||||
let username = if config.username.is_empty() { None } else { Some(config.username.clone()) };
|
||||
let password = if config.password.is_empty() { None } else { Some(config.password.clone()) };
|
||||
let client = db::influxdb_driver::InfluxdbClient::new_with_ca_cert(
|
||||
&url,
|
||||
username,
|
||||
password,
|
||||
Some(&config.ca_cert_path),
|
||||
connect_timeout,
|
||||
)?;
|
||||
db::influxdb_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
|
||||
}
|
||||
|
|
@ -617,6 +631,19 @@ pub async fn connect_db(state: State<'_, Arc<AppState>>, config: ConnectionConfi
|
|||
db::turso_driver::test_connection(&client, connect_timeout).await?;
|
||||
PoolKind::Turso(client)
|
||||
}
|
||||
DatabaseType::InfluxDb => {
|
||||
let username = if db_config.username.is_empty() { None } else { Some(db_config.username.clone()) };
|
||||
let password = if db_config.password.is_empty() { None } else { Some(db_config.password.clone()) };
|
||||
let client = db::influxdb_driver::InfluxdbClient::new_with_ca_cert(
|
||||
&url,
|
||||
username,
|
||||
password,
|
||||
Some(&db_config.ca_cert_path),
|
||||
connect_timeout,
|
||||
)?;
|
||||
db::influxdb_driver::test_connection(&client, connect_timeout).await?;
|
||||
PoolKind::InfluxDb(client)
|
||||
}
|
||||
db_type if database_capabilities::is_agent_type(&db_type) => {
|
||||
connect_agent_pool(state.inner(), &db_config, &host, port).await?
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue