mirror of
https://github.com/t8y2/dbx.git
synced 2026-10-02 02:34:42 +08:00
feat(plugins): hand SOCKS5 transport proxy routes to connections
This commit is contained in:
@@ -403,6 +403,8 @@ export interface PluginConnectionProviderContribution {
|
||||
filesystem_provider?: string;
|
||||
capabilities?: PluginConnectionCapability[];
|
||||
actions?: PluginConnectionActionContribution[];
|
||||
/** Multi-endpoint providers (Kafka advertised.listeners) receive a SOCKS5 runtime.proxy route over transport layers. */
|
||||
proxy_route?: boolean;
|
||||
}
|
||||
|
||||
export interface PluginWorkbenchContribution {
|
||||
|
||||
@@ -77,19 +77,19 @@ async fn run() -> Result<(), String> {
|
||||
}))
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
let test = host.test_connection(&config, "localhost", 22).await?;
|
||||
let test = host.test_connection(&config, "localhost", 22, None).await?;
|
||||
if !test.message.contains("localhost:22") {
|
||||
return Err(format!("Unexpected connection-test result: {}", test.message));
|
||||
}
|
||||
|
||||
let action = host.invoke_connection_action(&config, "suggest-greeting", "localhost", 22).await?;
|
||||
let action = host.invoke_connection_action(&config, "suggest-greeting", "localhost", 22, None).await?;
|
||||
if action.message.as_deref() != Some("Greeting updated by the plugin action.")
|
||||
|| action.field_values.get("greeting").and_then(Value::as_str) != Some("Hello from plugin action")
|
||||
{
|
||||
return Err(format!("Unexpected connection-action result: {action:?}"));
|
||||
}
|
||||
|
||||
let handle = host.connect_connection(&config, "localhost", 22).await?;
|
||||
let handle = host.connect_connection(&config, "localhost", 22, None).await?;
|
||||
if !handle.is_running() {
|
||||
return Err("Plugin connection handle is not running".to_string());
|
||||
}
|
||||
|
||||
@@ -39,7 +39,7 @@ use crate::nacos::config::{NACOS_CONSOLE_SESSION_PASSWORD, NACOS_PRIMARY_SESSION
|
||||
use crate::path_utils::expand_tilde;
|
||||
use crate::plugins::{
|
||||
PluginConnectionActionResult, PluginConnectionHandle, PluginDriverSession, PluginHost, PluginRegistry,
|
||||
PluginRuntimeEnv,
|
||||
PluginRuntimeEnv, PluginRuntimeProxy,
|
||||
};
|
||||
use crate::query_cancel::RunningQueries;
|
||||
use crate::session_credentials::SessionCredentialStore;
|
||||
@@ -2334,7 +2334,9 @@ impl AppState {
|
||||
|
||||
validate_h2_file_connection(&db_config)?;
|
||||
self.ensure_current_connection_attempt(connection_id, connection_attempt).await?;
|
||||
let (host, port) = self.connection_host_port(connection_id, &db_config).await?;
|
||||
let endpoint = self.connection_endpoint(connection_id, &db_config).await?;
|
||||
let (host, port) = (endpoint.host, endpoint.port);
|
||||
let runtime_proxy = endpoint.proxy;
|
||||
if let Err(err) = self.ensure_current_connection_attempt(connection_id, connection_attempt).await {
|
||||
self.reset_connection_transport_for_config(connection_id, &db_config).await;
|
||||
return Err(err);
|
||||
@@ -2927,9 +2929,9 @@ impl AppState {
|
||||
}
|
||||
self.external_driver_pool("jdbc", &jdbc_config).await?
|
||||
}
|
||||
DatabaseType::Plugin => {
|
||||
PoolKind::PluginConnection(self.plugin_host.connect_connection(&db_config, &host, port).await?)
|
||||
}
|
||||
DatabaseType::Plugin => PoolKind::PluginConnection(
|
||||
self.plugin_host.connect_connection(&db_config, &host, port, runtime_proxy).await?,
|
||||
),
|
||||
#[cfg(feature = "mq-admin")]
|
||||
DatabaseType::MessageQueue => {
|
||||
// MQ admin connections don't hold a data query pool. We just test
|
||||
@@ -3161,9 +3163,29 @@ impl AppState {
|
||||
connection_id: &str,
|
||||
config: &ConnectionConfig,
|
||||
) -> Result<(String, u16), String> {
|
||||
let endpoint = self.connection_endpoint(connection_id, config).await?;
|
||||
Ok((endpoint.host, endpoint.port))
|
||||
}
|
||||
|
||||
/// Resolves the runtime dial endpoint for a plugin connection, including
|
||||
/// the host-managed SOCKS5 route when the provider declares
|
||||
/// `proxy_route` and transport layers are configured.
|
||||
pub async fn plugin_connection_endpoint(
|
||||
&self,
|
||||
connection_id: &str,
|
||||
config: &ConnectionConfig,
|
||||
) -> Result<ConnectionEndpoint, String> {
|
||||
self.connection_endpoint(connection_id, config).await
|
||||
}
|
||||
|
||||
async fn connection_endpoint(
|
||||
&self,
|
||||
connection_id: &str,
|
||||
config: &ConnectionConfig,
|
||||
) -> Result<ConnectionEndpoint, String> {
|
||||
let transport_layers = self.resolved_transport_layers(config).await?;
|
||||
if transport_layers.is_empty() || db::sqlite_worker::sqlite_ssh_worker_requested(config) {
|
||||
return Ok((config.host.clone(), config.port));
|
||||
return Ok(ConnectionEndpoint::direct(config.host.clone(), config.port));
|
||||
}
|
||||
if config.uses_oracle_tns() {
|
||||
// A TNS descriptor may contain several failover addresses, so rewriting it
|
||||
@@ -3180,10 +3202,31 @@ impl AppState {
|
||||
== crate::mq::types::MqSystemKind::RocketMq
|
||||
{
|
||||
self.rocketmq_socks_proxy_for_transport_layers(connection_id, &transport_layers).await?;
|
||||
return Ok((config.host.clone(), config.port));
|
||||
return Ok(ConnectionEndpoint::direct(config.host.clone(), config.port));
|
||||
}
|
||||
|
||||
// Multi-endpoint plugin providers (Kafka bootstrap + advertised
|
||||
// listeners) route every endpoint through a host-managed SOCKS5
|
||||
// dialer instead of a static tunnel, which can only reach a single
|
||||
// broker. The payload keeps the logical endpoint so the plugin can
|
||||
// still resolve its own seed list and metadata names.
|
||||
if config.db_type == DatabaseType::Plugin && self.plugin_host.wants_proxy_route(config).await {
|
||||
if let Some(proxy) = self.socks5_route_for_transport_layers(connection_id, &transport_layers).await? {
|
||||
return Ok(ConnectionEndpoint { host: config.host.clone(), port: config.port, proxy: Some(proxy) });
|
||||
}
|
||||
}
|
||||
|
||||
let (remote_host, remote_port) = connection_remote_endpoint(config);
|
||||
// Plugin providers commonly declare no host/port binding (Kafka keeps
|
||||
// its endpoints in provider fields instead), so a static tunnel would
|
||||
// silently forward to an empty target and every downstream dial would
|
||||
// time out with no actionable hint. Fail here instead.
|
||||
if config.db_type == DatabaseType::Plugin && remote_host.is_empty() {
|
||||
return Err(
|
||||
"Transport layers for this plugin connection need a remote host and port. The connection provider must declare host/port fields or support proxy_route (SOCKS5 routing); otherwise remove the SSH/proxy/HTTP tunnel layer."
|
||||
.to_string(),
|
||||
);
|
||||
}
|
||||
let local_port = db::transport_layer_tunnel::start_transport_layers(
|
||||
connection_id,
|
||||
&transport_layers,
|
||||
@@ -3195,7 +3238,67 @@ impl AppState {
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(("127.0.0.1".to_string(), local_port))
|
||||
Ok(ConnectionEndpoint { host: "127.0.0.1".to_string(), port: local_port, proxy: None })
|
||||
}
|
||||
|
||||
/// Builds the host-managed SOCKS5 route from the transport chain for
|
||||
/// plugin providers declaring `proxy_route` (mirrors the
|
||||
/// RocketMQ proxy path). `None` = fall back to the static tunnel path.
|
||||
async fn socks5_route_for_transport_layers(
|
||||
&self,
|
||||
connection_id: &str,
|
||||
transport_layers: &[TransportLayerConfig],
|
||||
) -> Result<Option<PluginRuntimeProxy>, String> {
|
||||
use crate::models::connection::ProxyType;
|
||||
|
||||
let Some(final_layer) = transport_layers.last() else {
|
||||
return Ok(None);
|
||||
};
|
||||
match final_layer {
|
||||
TransportLayerConfig::Ssh(_) => {
|
||||
// The final SSH hop exposes a dynamic SOCKS5 endpoint so every
|
||||
// advertised broker is reachable through one tunnel.
|
||||
let local_port = db::transport_layer_tunnel::start_transport_layers_with_final_ssh_socks5(
|
||||
connection_id,
|
||||
transport_layers,
|
||||
&self.tunnels,
|
||||
&self.proxy_tunnels,
|
||||
&self.http_tunnels,
|
||||
)
|
||||
.await?;
|
||||
Ok(Some(PluginRuntimeProxy::socks5("127.0.0.1".to_string(), local_port, String::new(), String::new())))
|
||||
}
|
||||
TransportLayerConfig::Proxy(proxy) if proxy.proxy_type == ProxyType::Socks5 => {
|
||||
if transport_layers.len() == 1 {
|
||||
Ok(Some(PluginRuntimeProxy::socks5(
|
||||
proxy.host.clone(),
|
||||
proxy.port,
|
||||
proxy.username.clone(),
|
||||
proxy.password.clone(),
|
||||
)))
|
||||
} else {
|
||||
let local_port = db::transport_layer_tunnel::start_transport_layers(
|
||||
connection_id,
|
||||
&transport_layers[..transport_layers.len() - 1],
|
||||
&proxy.host,
|
||||
proxy.port,
|
||||
&self.tunnels,
|
||||
&self.proxy_tunnels,
|
||||
&self.http_tunnels,
|
||||
)
|
||||
.await?;
|
||||
Ok(Some(PluginRuntimeProxy::socks5(
|
||||
"127.0.0.1".to_string(),
|
||||
local_port,
|
||||
proxy.username.clone(),
|
||||
proxy.password.clone(),
|
||||
)))
|
||||
}
|
||||
}
|
||||
// HTTP-tunnel chains cannot serve arbitrary endpoints; fall back
|
||||
// to the static tunnel path (guarded below for empty endpoints).
|
||||
TransportLayerConfig::Proxy(_) | TransportLayerConfig::HttpTunnel(_) => Ok(None),
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn invoke_plugin_connection_action(
|
||||
@@ -3210,8 +3313,12 @@ impl AppState {
|
||||
let transport_id = format!("{}:plugin-action:{action_id}", config.id);
|
||||
let has_transport_layers = config.has_effective_transport_layers();
|
||||
let connection_id = if has_transport_layers { transport_id.as_str() } else { config.id.as_str() };
|
||||
let result = match self.connection_host_port(connection_id, &config).await {
|
||||
Ok((host, port)) => self.plugin_host.invoke_connection_action(&config, action_id, &host, port).await,
|
||||
let result = match self.plugin_connection_endpoint(connection_id, &config).await {
|
||||
Ok(endpoint) => {
|
||||
self.plugin_host
|
||||
.invoke_connection_action(&config, action_id, &endpoint.host, endpoint.port, endpoint.proxy)
|
||||
.await
|
||||
}
|
||||
Err(error) => Err(error),
|
||||
};
|
||||
if has_transport_layers {
|
||||
@@ -5560,6 +5667,24 @@ fn is_agent_validate_connection_unsupported(err: &str) -> bool {
|
||||
lower.contains("validate_connection") && (lower.contains("unknown method") || lower.contains("method not found"))
|
||||
}
|
||||
|
||||
/// Runtime dial endpoint handed to a plugin lifecycle call: the logical
|
||||
/// `host:port` plus an optional host-managed SOCKS5 route for providers
|
||||
/// declaring `proxy_route`. When `proxy` is set the plugin is
|
||||
/// expected to dial every endpoint (its seed list and metadata names) through
|
||||
/// the route, keeping the logical endpoint only for metadata discovery.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ConnectionEndpoint {
|
||||
pub host: String,
|
||||
pub port: u16,
|
||||
pub proxy: Option<PluginRuntimeProxy>,
|
||||
}
|
||||
|
||||
impl ConnectionEndpoint {
|
||||
fn direct(host: String, port: u16) -> Self {
|
||||
Self { host, port, proxy: None }
|
||||
}
|
||||
}
|
||||
|
||||
fn connection_remote_endpoint(config: &ConnectionConfig) -> (String, u16) {
|
||||
if config.db_type == DatabaseType::MongoDb {
|
||||
config
|
||||
|
||||
@@ -90,7 +90,10 @@ pub async fn start_transport_layers_with_final_ssh_local_port(
|
||||
.await
|
||||
}
|
||||
|
||||
#[cfg(feature = "mq-admin")]
|
||||
/// Starts the layer chain with the final SSH hop exposing a dynamic SOCKS5
|
||||
/// endpoint (ssh -D) instead of a static forward, so callers can route many
|
||||
/// remote endpoints through one tunnel (RocketMQ name servers/brokers, plugin
|
||||
/// providers with `proxy_route`, ...).
|
||||
pub async fn start_transport_layers_with_final_ssh_socks5(
|
||||
connection_id: &str,
|
||||
layers: &[TransportLayerConfig],
|
||||
|
||||
@@ -23,7 +23,9 @@ pub use filesystem::{
|
||||
PLUGIN_FILESYSTEM_CREATE_DIRECTORY_METHOD, PLUGIN_FILESYSTEM_DELETE_METHOD, PLUGIN_FILESYSTEM_LIST_METHOD,
|
||||
PLUGIN_FILESYSTEM_READ_METHOD, PLUGIN_FILESYSTEM_RENAME_METHOD, PLUGIN_FILESYSTEM_WRITE_METHOD,
|
||||
};
|
||||
pub use host::{ActivePluginSession, PluginConnectionActionResult, PluginConnectionHandle, PluginHost};
|
||||
pub use host::{
|
||||
ActivePluginSession, PluginConnectionActionResult, PluginConnectionHandle, PluginHost, PluginRuntimeProxy,
|
||||
};
|
||||
pub use installer::{
|
||||
PluginInstallPolicy, PluginInstallResponse, PluginInstallResult, PluginPackageInstaller, PluginRollbackResponse,
|
||||
PluginRollbackResult, PluginSignatureStatus, PluginTrustStore, PluginTrustedKey, DBXP_EXTENSION,
|
||||
|
||||
@@ -4,7 +4,7 @@ use std::time::Duration;
|
||||
|
||||
use crate::models::connection::{ConnectionConfig, ConnectionTestResult, DatabaseType};
|
||||
use serde::de::DeserializeOwned;
|
||||
use serde::Serialize;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::sync::{broadcast, Mutex, RwLock};
|
||||
|
||||
use super::{
|
||||
@@ -191,13 +191,26 @@ impl PluginHost {
|
||||
/// to an installed connection provider.
|
||||
pub fn connection_params_standalone(&self, config: &ConnectionConfig) -> Result<serde_json::Value, String> {
|
||||
let (_, provider) = self.resolve_connection_provider(config)?;
|
||||
plugin_connection_params(config, &provider, &config.host, config.port)
|
||||
plugin_connection_params(config, &provider, &config.host, config.port, None)
|
||||
}
|
||||
|
||||
/// Reports whether the connection's provider declared
|
||||
/// `proxy_route` (multi-endpoint targets that want a SOCKS5
|
||||
/// runtime route over transport layers instead of a static tunnel).
|
||||
/// Unresolvable configs report `false` so the caller falls back to the
|
||||
/// static-tunnel path (and its empty-endpoint guard).
|
||||
pub async fn wants_proxy_route(&self, config: &ConnectionConfig) -> bool {
|
||||
match self.resolve_connection_provider(config) {
|
||||
Ok((_, provider)) => provider.proxy_route,
|
||||
Err(_) => false,
|
||||
}
|
||||
}
|
||||
pub async fn test_connection(
|
||||
&self,
|
||||
config: &ConnectionConfig,
|
||||
runtime_host: &str,
|
||||
runtime_port: u16,
|
||||
runtime_proxy: Option<PluginRuntimeProxy>,
|
||||
) -> Result<ConnectionTestResult, String> {
|
||||
let _activity = self
|
||||
.inner
|
||||
@@ -214,7 +227,7 @@ impl PluginHost {
|
||||
let result: serde_json::Value = session
|
||||
.invoke_with_timeout(
|
||||
PLUGIN_CONNECTION_TEST_METHOD,
|
||||
plugin_connection_params(config, &provider, runtime_host, runtime_port)?,
|
||||
plugin_connection_params(config, &provider, runtime_host, runtime_port, runtime_proxy.as_ref())?,
|
||||
None,
|
||||
Some(plugin_connect_deadline(config, &provider)),
|
||||
)
|
||||
@@ -227,6 +240,7 @@ impl PluginHost {
|
||||
config: &ConnectionConfig,
|
||||
runtime_host: &str,
|
||||
runtime_port: u16,
|
||||
runtime_proxy: Option<PluginRuntimeProxy>,
|
||||
) -> Result<PluginConnectionHandle, String> {
|
||||
let activity = self
|
||||
.inner
|
||||
@@ -235,7 +249,7 @@ impl PluginHost {
|
||||
.begin_connection(config.plugin_id.as_deref().unwrap_or_default(), &config.name)?;
|
||||
let (_, provider) = self.resolve_connection_provider(config)?;
|
||||
validate_plugin_connection_values(config, &provider)?;
|
||||
let params = plugin_connection_params(config, &provider, runtime_host, runtime_port)?;
|
||||
let params = plugin_connection_params(config, &provider, runtime_host, runtime_port, runtime_proxy.as_ref())?;
|
||||
let needs_session = provider.has_capability(PluginConnectionCapability::Connect)
|
||||
|| provider.has_capability(PluginConnectionCapability::Disconnect);
|
||||
let session = if needs_session {
|
||||
@@ -274,6 +288,7 @@ impl PluginHost {
|
||||
action_id: &str,
|
||||
runtime_host: &str,
|
||||
runtime_port: u16,
|
||||
runtime_proxy: Option<PluginRuntimeProxy>,
|
||||
) -> Result<PluginConnectionActionResult, String> {
|
||||
let _activity = self
|
||||
.inner
|
||||
@@ -284,7 +299,8 @@ impl PluginHost {
|
||||
let action = plugin_invoke_connection_action(&provider, action_id)?;
|
||||
validate_plugin_connection_values_for_action(config, &provider, action.requires_valid_form)?;
|
||||
let session = self.activate(config.plugin_id.as_deref().unwrap_or_default()).await?;
|
||||
let mut params = plugin_connection_params(config, &provider, runtime_host, runtime_port)?;
|
||||
let mut params =
|
||||
plugin_connection_params(config, &provider, runtime_host, runtime_port, runtime_proxy.as_ref())?;
|
||||
params
|
||||
.as_object_mut()
|
||||
.ok_or("Plugin connection action params must be an object")?
|
||||
@@ -433,21 +449,50 @@ fn plugin_connection_params(
|
||||
provider: &PluginConnectionProviderContribution,
|
||||
runtime_host: &str,
|
||||
runtime_port: u16,
|
||||
runtime_proxy: Option<&PluginRuntimeProxy>,
|
||||
) -> Result<serde_json::Value, String> {
|
||||
let mut runtime = serde_json::json!({
|
||||
"host": runtime_host,
|
||||
"port": runtime_port,
|
||||
});
|
||||
if let Some(proxy) = runtime_proxy {
|
||||
runtime["proxy"] = serde_json::to_value(proxy).map_err(|error| error.to_string())?;
|
||||
}
|
||||
Ok(serde_json::json!({
|
||||
"provider": {
|
||||
"id": provider.id,
|
||||
"databaseType": provider.database_type,
|
||||
},
|
||||
"connection": serde_json::to_value(config).map_err(|error| error.to_string())?,
|
||||
"runtime": {
|
||||
"host": runtime_host,
|
||||
"port": runtime_port,
|
||||
},
|
||||
"runtime": runtime,
|
||||
"operationId": uuid::Uuid::new_v4().to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
/// A host-managed SOCKS5 route handed to a plugin through
|
||||
/// `runtime.proxy`. Providers declaring `proxy_route` (multi-endpoint
|
||||
/// targets such as Kafka) dial every advertised broker through this route
|
||||
/// instead of a static tunnel, which can only reach a single endpoint.
|
||||
/// Credentials ride the same encrypted lifecycle channel as connection
|
||||
/// secrets and must never be logged by the plugin.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct PluginRuntimeProxy {
|
||||
#[serde(rename = "type")]
|
||||
pub proxy_type: String,
|
||||
pub host: String,
|
||||
pub port: u16,
|
||||
#[serde(default, skip_serializing_if = "String::is_empty")]
|
||||
pub username: String,
|
||||
#[serde(default, skip_serializing_if = "String::is_empty")]
|
||||
pub password: String,
|
||||
}
|
||||
|
||||
impl PluginRuntimeProxy {
|
||||
pub fn socks5(host: String, port: u16, username: String, password: String) -> Self {
|
||||
Self { proxy_type: "socks5".to_string(), host, port, username, password }
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_plugin_connection_values(
|
||||
config: &ConnectionConfig,
|
||||
provider: &PluginConnectionProviderContribution,
|
||||
@@ -775,9 +820,9 @@ fn ensure_permission(plugin: &super::InstalledPlugin, required_permission: Optio
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{
|
||||
plugin_connect_deadline, plugin_connection_action_result, plugin_field_is_visible,
|
||||
plugin_connect_deadline, plugin_connection_action_result, plugin_connection_params, plugin_field_is_visible,
|
||||
plugin_invoke_connection_action, validate_plugin_connection_values,
|
||||
validate_plugin_connection_values_for_action,
|
||||
validate_plugin_connection_values_for_action, PluginRuntimeProxy,
|
||||
};
|
||||
use crate::models::connection::ConnectionConfig;
|
||||
use crate::plugins::PluginConnectionProviderContribution;
|
||||
@@ -824,9 +869,9 @@ mod tests {
|
||||
let lifecycle = registry.lifecycle();
|
||||
let host = super::PluginHost::new(registry);
|
||||
let update = lifecycle.begin_update("sample.ui").unwrap();
|
||||
assert!(host.connect_connection(&config, "localhost", 0).await.is_err());
|
||||
assert!(host.connect_connection(&config, "localhost", 0, None).await.is_err());
|
||||
drop(update);
|
||||
let connection = host.connect_connection(&config, "localhost", 0).await.unwrap();
|
||||
let connection = host.connect_connection(&config, "localhost", 0, None).await.unwrap();
|
||||
assert!(lifecycle.begin_update("sample.ui").unwrap_err().contains("Saved UI connection"));
|
||||
connection.disconnect().await.unwrap();
|
||||
drop(connection);
|
||||
@@ -885,6 +930,41 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connection_params_embed_optional_socks5_proxy_route() {
|
||||
let config: ConnectionConfig = serde_json::from_value(serde_json::json!({
|
||||
"id": "plugin-connection",
|
||||
"name": "Plugin connection",
|
||||
"db_type": "plugin",
|
||||
"host": "",
|
||||
"port": 0,
|
||||
"username": "",
|
||||
"password": "",
|
||||
"database": null,
|
||||
"plugin_id": "sample",
|
||||
"plugin_connection_provider": "sample.connection",
|
||||
"plugin_connection_type": "sample"
|
||||
}))
|
||||
.unwrap();
|
||||
let provider: PluginConnectionProviderContribution = serde_json::from_value(serde_json::json!({
|
||||
"id": "sample.connection",
|
||||
"label": "Sample",
|
||||
"database_type": "sample",
|
||||
"fields": []
|
||||
}))
|
||||
.unwrap();
|
||||
|
||||
let without = plugin_connection_params(&config, &provider, "", 0, None).unwrap();
|
||||
assert!(without["runtime"].get("proxy").is_none());
|
||||
|
||||
let proxy = PluginRuntimeProxy::socks5("127.0.0.1".to_string(), 1080, "user".to_string(), "pass".to_string());
|
||||
let with = plugin_connection_params(&config, &provider, "k1", 9092, Some(&proxy)).unwrap();
|
||||
assert_eq!(with["runtime"]["proxy"]["type"], "socks5");
|
||||
assert_eq!(with["runtime"]["proxy"]["host"], "127.0.0.1");
|
||||
assert_eq!(with["runtime"]["proxy"]["port"], 1080);
|
||||
assert_eq!(with["runtime"]["proxy"]["username"], "user");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rejects_zero_port_for_port_bound_connection_field() {
|
||||
let config: ConnectionConfig = serde_json::from_value(serde_json::json!({
|
||||
|
||||
@@ -560,6 +560,14 @@ pub struct PluginConnectionProviderContribution {
|
||||
pub capabilities: Vec<PluginConnectionCapability>,
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub actions: Vec<PluginConnectionActionContribution>,
|
||||
/// Providers whose targets have multiple reachable endpoints (Kafka
|
||||
/// bootstrap + advertised listeners) declare this flag so the host hands
|
||||
/// them a SOCKS5 `runtime.proxy` route instead of a static tunnel, which
|
||||
/// can only reach a single endpoint. Without the flag, transport layers
|
||||
/// keep today's static-tunnel behavior (fine for single-endpoint
|
||||
/// providers such as SSH or LDAP, which declare binding host/port fields).
|
||||
#[serde(default, skip_serializing_if = "is_false")]
|
||||
pub proxy_route: bool,
|
||||
}
|
||||
|
||||
impl PluginConnectionProviderContribution {
|
||||
@@ -1470,6 +1478,26 @@ mod tests {
|
||||
PluginManifest,
|
||||
};
|
||||
|
||||
#[test]
|
||||
fn connection_provider_proxy_route_defaults_false_and_parses() {
|
||||
let provider: PluginConnectionProviderContribution = serde_json::from_value(serde_json::json!({
|
||||
"id": "sample.connection",
|
||||
"database_type": "sample",
|
||||
"fields": []
|
||||
}))
|
||||
.unwrap();
|
||||
assert!(!provider.proxy_route);
|
||||
|
||||
let provider: PluginConnectionProviderContribution = serde_json::from_value(serde_json::json!({
|
||||
"id": "sample.connection",
|
||||
"database_type": "sample",
|
||||
"fields": [],
|
||||
"proxy_route": true
|
||||
}))
|
||||
.unwrap();
|
||||
assert!(provider.proxy_route);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parses_only_strict_https_network_permissions() {
|
||||
assert_eq!(
|
||||
|
||||
@@ -273,8 +273,8 @@ async fn run_temporary_connection_test(
|
||||
|
||||
if config.db_type == DatabaseType::Plugin {
|
||||
let result = async {
|
||||
let (host, port) = app.connection_host_port(&temp_id, &config).await?;
|
||||
app.plugin_host.test_connection(&config, &host, port).await
|
||||
let endpoint = app.plugin_connection_endpoint(&temp_id, &config).await?;
|
||||
app.plugin_host.test_connection(&config, &endpoint.host, endpoint.port, endpoint.proxy).await
|
||||
}
|
||||
.await;
|
||||
app.reset_connection_transport_for_config(&temp_id, &config).await;
|
||||
@@ -428,7 +428,7 @@ pub async fn connect_db(
|
||||
app.configs.write().await.insert(connection_id.clone(), runtime_config);
|
||||
|
||||
if config.db_type == dbx_core::models::connection::DatabaseType::Plugin {
|
||||
let (host, port) = match app.connection_host_port(&connection_id, &config).await {
|
||||
let endpoint = match app.plugin_connection_endpoint(&connection_id, &config).await {
|
||||
Ok(endpoint) => endpoint,
|
||||
Err(error) => {
|
||||
app.reset_connection_transport_for_config(&connection_id, &config).await;
|
||||
@@ -441,14 +441,15 @@ pub async fn connect_db(
|
||||
rollback_session_credential_writes(app, &session_credential_writes);
|
||||
return Err(AppError::from(error));
|
||||
}
|
||||
let handle = match app.plugin_host.connect_connection(&config, &host, port).await {
|
||||
Ok(handle) => handle,
|
||||
Err(error) => {
|
||||
app.reset_connection_transport_for_config(&connection_id, &config).await;
|
||||
rollback_session_credential_writes(app, &session_credential_writes);
|
||||
return Err(AppError::from(error));
|
||||
}
|
||||
};
|
||||
let handle =
|
||||
match app.plugin_host.connect_connection(&config, &endpoint.host, endpoint.port, endpoint.proxy).await {
|
||||
Ok(handle) => handle,
|
||||
Err(error) => {
|
||||
app.reset_connection_transport_for_config(&connection_id, &config).await;
|
||||
rollback_session_credential_writes(app, &session_credential_writes);
|
||||
return Err(AppError::from(error));
|
||||
}
|
||||
};
|
||||
let pool = PoolKind::PluginConnection(handle);
|
||||
if let Err(error) =
|
||||
app.insert_connection_pool_for_attempt(&connection_id, attempt, connection_id.clone(), pool, &config).await
|
||||
|
||||
@@ -199,6 +199,8 @@ Manifest 中的路径相对于包根目录,不能以 `/` 开头、不能包含
|
||||
|
||||
连接生命周期使用固定的 `connection/test`、`connection/connect`、`connection/disconnect`,自定义表单动作使用 `connection/action`,参数包含 `provider`、`connection`、`runtime`;动作额外包含 `action: { id }`。后端使用 `connection.id` 保存会话,并使用 `runtime.host`/`runtime.port` 连接经过 DBX 隧道或代理转换后的端点。仅在后端生命周期请求中接收补齐 Secret 的连接;前端通过 `connectionId` 使用已经建立的会话。
|
||||
|
||||
多端点协议(Kafka `advertised.listeners`、集群发现等)在贡献点上声明 `proxy_route`。配置了传输层后,DBX 会在 payload 中下发 SOCKS5 路由而非静态转发:`runtime.host`/`runtime.port` 保留逻辑端点,插件通过 `runtime.proxy`(`{ "type": "socks5", "host": "...", "port": 1080, "username": "...", "password": "..." }`)拨号每个广播端点——SSH 末层暴露其动态 SOCKS5 端点,SOCKS5 代理层直接使用;代理凭据与连接 Secret 走同一加密通道,插件不得记录日志。未声明该旗标时传输层保持静态转发行为,要求标准 `host`/`port` 字段给出单一远端端点;对会隧道到空端点的插件连接,DBX 直接给出可行动的报错而非静默超时。
|
||||
|
||||
动作的 `when` 可为 `always`、`create`、`edit`;`variant` 可为 `default`、`outline`、`secondary`、`destructive`、`ghost`;`timeout_ms` 为 1–120000。`requires_valid_form` 决定是否要求完整表单,`close_on_success` 决定成功后是否关闭。返回 `{ "success": true, "message": "...", "fieldValues": { "port": 443 } }` 可以更新已声明的字段;返回 `success: false` 表示失败。
|
||||
|
||||
### 其他 Contributions
|
||||
|
||||
@@ -199,6 +199,8 @@ The result is `{ "action": "submit", "value": "123456" }`, `{ "action": "cancel"
|
||||
|
||||
The fixed lifecycle methods are `connection/test`, `connection/connect`, and `connection/disconnect`; custom form actions use `connection/action`. Parameters include `provider`, `connection`, and `runtime`, plus `action: { id }` for actions. Keep sessions keyed by `connection.id` and connect to `runtime.host`/`runtime.port`, which include DBX tunnel/proxy resolution. Hydrated secrets are provided only to backend lifecycle requests; the frontend uses `connectionId` to access an established session.
|
||||
|
||||
Multi-endpoint protocols (Kafka `advertised.listeners`, cluster discovery) declare `proxy_route` on the contribution. With transport layers configured, DBX then delivers a SOCKS5 route in the payload instead of a static forward, keeping `runtime.host`/`runtime.port` at the logical endpoint while the plugin dials every advertised endpoint through `runtime.proxy` (`{ "type": "socks5", "host": "...", "port": 1080, "username": "...", "password": "..." }`; an SSH final hop exposes its dynamic SOCKS5 endpoint, a SOCKS5 proxy layer is used directly, and proxy credentials ride the same encrypted channel as connection secrets and must never be logged). Without the flag, transport layers keep the static-tunnel behavior, which requires a single remote endpoint from the standard `host`/`port` fields; DBX rejects plugin connections that would tunnel to an empty endpoint instead of timing out silently.
|
||||
|
||||
Action `when` accepts `always`, `create`, or `edit`; `variant` accepts `default`, `outline`, `secondary`, `destructive`, or `ghost`; `timeout_ms` is 1–120000. `requires_valid_form` requires a complete form and `close_on_success` closes the dialog on success. Return `{ "success": true, "message": "...", "fieldValues": { "port": 443 } }` to update declared fields; `success: false` is a failure.
|
||||
|
||||
### Other Contributions
|
||||
|
||||
@@ -288,6 +288,37 @@ Lifecycle methods receive:
|
||||
|
||||
`runtime.host` and `runtime.port` are the final endpoint after DBX transport layers. A protocol plugin must connect to this endpoint instead of rebuilding DBX tunnels itself.
|
||||
|
||||
##### Transport proxy route for multi-endpoint targets
|
||||
|
||||
A static tunnel forwards exactly one remote endpoint. Protocols whose server advertises additional endpoints a client must dial (Kafka `advertised.listeners`, cluster discovery, etc.) cannot be served that way: the bootstrap endpoint connects, but every advertised broker is unreachable. Such providers declare `proxy_route` on the connection-provider contribution:
|
||||
|
||||
```json
|
||||
{
|
||||
"type": "connection-provider",
|
||||
"id": "vendor.kafka.connection",
|
||||
"database_type": "kafka",
|
||||
"proxy_route": true
|
||||
}
|
||||
```
|
||||
|
||||
When transport layers are configured, DBX then delivers a SOCKS5 route instead of a static forward:
|
||||
|
||||
```json
|
||||
{
|
||||
"provider": { "...": "..." },
|
||||
"connection": { "...": "..." },
|
||||
"runtime": {
|
||||
"host": "",
|
||||
"port": 0,
|
||||
"proxy": { "type": "socks5", "host": "127.0.0.1", "port": 49153, "username": "", "password": "" }
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
- With SSH as the final transport layer the route is the hop's dynamic SOCKS5 endpoint (`ssh -D`); with a SOCKS5 proxy layer the route is that proxy, tunneled through any preceding layers. `username`/`password` are omitted when empty.
|
||||
- `runtime.host`/`runtime.port` stay at the connection's logical endpoint, which the plugin should keep using as its seed/metadata source while dialing every endpoint through the SOCKS5 route. Credentials ride the same encrypted lifecycle channel as connection secrets and must never be logged by the plugin.
|
||||
- Without the flag, transport layers keep the static-tunnel behavior, which requires the connection to resolve a single remote endpoint (providers should declare `host`/`port` bindings, as the SSH and LDAP plugins do); DBX rejects plugin connections that would tunnel to an empty endpoint instead of timing out silently.
|
||||
|
||||
#### Connection dialog actions
|
||||
|
||||
Connection providers may add ordered custom actions before DBX-owned lifecycle buttons:
|
||||
|
||||
@@ -284,6 +284,7 @@
|
||||
"uniqueItems": true,
|
||||
"items": { "enum": ["test", "connect", "disconnect"] }
|
||||
},
|
||||
"proxy_route": { "type": "boolean", "default": false },
|
||||
"actions": { "type": "array", "items": { "$ref": "#/$defs/connectionAction" } }
|
||||
}
|
||||
},
|
||||
|
||||
@@ -1168,8 +1168,13 @@ async fn test_connection_with_info_inner(
|
||||
let tunnel_id = format!("{}:test", config.id);
|
||||
let has_transport_layers = config.has_effective_transport_layers();
|
||||
let connection_id = if has_transport_layers { tunnel_id.as_str() } else { config.id.as_str() };
|
||||
let (host, port) = state.connection_host_port(connection_id, &config).await?;
|
||||
let probe_result = probe_connection_endpoint(&config, &host, port).await;
|
||||
let endpoint = state.plugin_connection_endpoint(connection_id, &config).await?;
|
||||
let (host, port) = (endpoint.host.clone(), endpoint.port);
|
||||
let runtime_proxy = endpoint.proxy;
|
||||
// SOCKS-routed plugin tests dial through the route themselves; probing the
|
||||
// logical endpoint (possibly empty host/port) would mislead.
|
||||
let probe_result =
|
||||
if runtime_proxy.is_some() { Ok(()) } else { probe_connection_endpoint(&config, &host, port).await };
|
||||
let url = connection_url_for_endpoint(&config, &host, port);
|
||||
let target = redacted_connection_url_for_endpoint(&config, &host, port);
|
||||
let connect_timeout = std::time::Duration::from_secs(config.effective_connect_timeout_secs());
|
||||
@@ -1177,7 +1182,7 @@ async fn test_connection_with_info_inner(
|
||||
let gaussdb_m_jdbc_config = gaussdb_m_jdbc_command_config(&config, &host, port);
|
||||
log::info!("[test_connection] db_type={:?} target={}", config.db_type, target);
|
||||
if config.db_type == DatabaseType::Plugin {
|
||||
let result = state.plugin_host.test_connection(&config, &host, port).await;
|
||||
let result = state.plugin_host.test_connection(&config, &host, port, runtime_proxy).await;
|
||||
if has_transport_layers {
|
||||
state.reset_connection_transport_for_config(&tunnel_id, &config).await;
|
||||
}
|
||||
@@ -1722,7 +1727,9 @@ pub async fn connect_db(
|
||||
drop_nacos_adapters_for_connection_ids(state.inner(), std::slice::from_ref(&id)).await;
|
||||
state.reset_connection_transport_for_config(&id, &db_config).await;
|
||||
|
||||
let (host, port) = state.connection_host_port(&id, &db_config).await?;
|
||||
let endpoint = state.plugin_connection_endpoint(&id, &db_config).await?;
|
||||
let (host, port) = (endpoint.host, endpoint.port);
|
||||
let runtime_proxy = endpoint.proxy;
|
||||
if let Err(err) = state.ensure_current_connection_attempt(&id, Some(attempt)).await {
|
||||
state.reset_connection_transport_for_config(&id, &db_config).await;
|
||||
return Err(err);
|
||||
@@ -2082,9 +2089,9 @@ pub async fn connect_db(
|
||||
DatabaseType::Mqtt => {
|
||||
return Err("MQTT support is not compiled in this build. Rebuild with the 'mq-admin' feature.".to_string());
|
||||
}
|
||||
DatabaseType::Plugin => {
|
||||
PoolKind::PluginConnection(state.plugin_host.connect_connection(&db_config, &host, port).await?)
|
||||
}
|
||||
DatabaseType::Plugin => PoolKind::PluginConnection(
|
||||
state.plugin_host.connect_connection(&db_config, &host, port, runtime_proxy).await?,
|
||||
),
|
||||
db_type => return Err(format!("Unsupported database type: {db_type:?}")),
|
||||
};
|
||||
|
||||
|
||||
Reference in New Issue
Block a user