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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 23 additions & 4 deletions src/driver/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,10 +14,13 @@ pub struct ConnectParams {
pub connection_string: Option<String>,
pub host: Option<String>,
pub port: Option<u16>,
#[serde(alias = "username")]
pub user: Option<String>,
pub database: Option<String>,
#[serde(alias = "safe_mode", alias = "safeMode")]
pub readonly: Option<bool>,
#[serde(alias = "sslMode")]
pub ssl_mode: Option<String>,
}

#[derive(Debug)]
Expand All @@ -44,8 +47,17 @@ impl ClickHouseClient {
} else {
let parsed = Url::parse(&cs)?;
let host = parsed.host_str().unwrap_or("localhost");
let port = parsed.port().unwrap_or(8123);
let scheme = parsed.scheme();
let port_str = match parsed.port() {
Some(p) => format!(":{}", p),
None => {
if scheme == "http" {
":8123".to_string()
} else {
String::new()
}
}
};
let user = if !parsed.username().is_empty() {
parsed.username().to_string()
} else {
Expand All @@ -58,17 +70,24 @@ impl ClickHouseClient {
params.database.unwrap_or_else(|| "default".to_string())
};
let readonly = params.readonly.unwrap_or(false);
let base = format!("{}://{}:{}", scheme, host, port);
let base = format!("{}://{}{}", scheme, host, port_str);
(base, user, database, readonly)
}
} else {
let host = params.host.unwrap_or_else(|| "localhost".to_string());
let port = params.port.unwrap_or(8123);
let scheme = match params.ssl_mode.as_deref() {
Some("prefer") | Some("require") => "https",
_ => "http",
};
let mut port = params.port.unwrap_or(8123);
if scheme == "https" && port == 8123 {
port = 8443;
}
let user = params.user.unwrap_or_else(|| "default".to_string());
let database = params.database.unwrap_or_else(|| "default".to_string());
let readonly = params.readonly.unwrap_or(false);
(
format!("http://{}:{}", host, port),
format!("{}://{}:{}", scheme, host, port),
user,
database,
readonly,
Expand Down
79 changes: 37 additions & 42 deletions src/rpc/handlers/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -328,21 +328,19 @@ pub async fn handle_get_server_stats(params: Option<Value>) -> Result<Value, Dri
.unwrap_or_else(|_| r#"[0]"#.to_string());

let mut version_str = "ClickHouse".to_string();
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&version_text, 0) {
if let Some(row) = parsed.rows.first() {
if let Some(v) = row.first().and_then(|x| x.as_str()) {
version_str = format!("ClickHouse {}", v);
}
}
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&version_text, 0)
&& let Some(row) = parsed.rows.first()
&& let Some(v) = row.first().and_then(|x| x.as_str())
{
version_str = format!("ClickHouse {}", v);
}

let mut uptime_sec = 0;
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&uptime_text, 0) {
if let Some(row) = parsed.rows.first() {
if let Some(v) = row.first().and_then(|x| x.as_u64()) {
uptime_sec = v;
}
}
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&uptime_text, 0)
&& let Some(row) = parsed.rows.first()
&& let Some(v) = row.first().and_then(|x| x.as_u64())
{
uptime_sec = v;
}

let db_sizes_text = client
Expand Down Expand Up @@ -396,14 +394,13 @@ pub async fn handle_get_object_metadata(params: Option<Value>) -> Result<Value,
.ok_or_else(|| DriverError::ConnectionNotFound(p.connection_id))?;

let parts: Vec<&str> = p.node_id.split('.').collect();
let (db_name, tbl_name) =
if parts.len() >= 3 && (parts[0] == "table" || parts[0] == "view") {
(parts[1], parts[2])
} else if parts.len() >= 2 {
(parts[0], parts[1])
} else {
("default", p.node_id.as_str())
};
let (db_name, tbl_name) = if parts.len() >= 3 && (parts[0] == "table" || parts[0] == "view") {
(parts[1], parts[2])
} else if parts.len() >= 2 {
(parts[0], parts[1])
} else {
("default", p.node_id.as_str())
};

if client.base_url.starts_with("mock://") || client.base_url.starts_with("test://") {
return Ok(json!({
Expand All @@ -425,35 +422,33 @@ pub async fn handle_get_object_metadata(params: Option<Value>) -> Result<Value,
db_name, tbl_name
);
let mut ddl_str = String::new();
if let Ok(text) = client.post_sql(&ddl_sql, |_| {}).await {
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0) {
if let Some(row) = parsed.rows.first() {
if let Some(v) = row.first().and_then(|x| x.as_str()) {
ddl_str = v.to_string();
}
}
}
if let Ok(text) = client.post_sql(&ddl_sql, |_| {}).await
&& let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0)
&& let Some(row) = parsed.rows.first()
&& let Some(v) = row.first().and_then(|x| x.as_str())
{
ddl_str = v.to_string();
}

let cols_sql = format!(
"SELECT name, type, comment FROM system.columns WHERE database = '{}' AND table = '{}' ORDER BY position FORMAT JSONCompactEachRowWithNamesAndTypes",
db_name, tbl_name
);
let mut columns = Vec::new();
if let Ok(text) = client.post_sql(&cols_sql, |_| {}).await {
if let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0) {
for row in parsed.rows {
let name = row.first().and_then(|x| x.as_str()).unwrap_or("unknown");
let col_type = row.get(1).and_then(|x| x.as_str()).unwrap_or("String");
let comment = row.get(2).and_then(|x| x.as_str()).unwrap_or("");
let is_nullable = col_type.starts_with("Nullable(");
columns.push(json!({
"name": name,
"dataType": col_type,
"isNullable": is_nullable,
"comment": comment
}));
}
if let Ok(text) = client.post_sql(&cols_sql, |_| {}).await
&& let Ok(parsed) = crate::mapper::row_compact::parse_compact_output(&text, 0)
{
for row in parsed.rows {
let name = row.first().and_then(|x| x.as_str()).unwrap_or("unknown");
let col_type = row.get(1).and_then(|x| x.as_str()).unwrap_or("String");
let comment = row.get(2).and_then(|x| x.as_str()).unwrap_or("");
let is_nullable = col_type.starts_with("Nullable(");
columns.push(json!({
"name": name,
"dataType": col_type,
"isNullable": is_nullable,
"comment": comment
}));
}
}

Expand Down
4 changes: 3 additions & 1 deletion src/rpc/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@ pub async fn dispatch(method: &str, params: Option<Value>) -> Result<Value, Driv
"sdui.contextActions" => schema::handle_context_actions(params).await,
"db.getCapabilities" => schema::handle_get_capabilities(params).await,
"db.getServerStats" => schema::handle_get_server_stats(params).await,
"db.getObjectMetadata" | "db.getObjectDDL" => schema::handle_get_object_metadata(params).await,
"db.getObjectMetadata" | "db.getObjectDDL" => {
schema::handle_get_object_metadata(params).await
}
_ => Err(DriverError::Rpc {
code: -32601,
message: format!("Method not found: {}", method),
Expand Down
Loading