diff --git a/agents/drivers/rabbitmq/helpers.go b/agents/drivers/rabbitmq/helpers.go index fb2f8e26a..c01784f61 100644 --- a/agents/drivers/rabbitmq/helpers.go +++ b/agents/drivers/rabbitmq/helpers.go @@ -244,6 +244,19 @@ func integerProperty(properties jsonObject, key string) (int, bool) { return *value, true } +func endpointOverride(config jsonObject, key string) (*address, error) { + override := objectOrNil(config, key) + if override == nil { + return nil, nil + } + host := strings.TrimSpace(stringOrEmpty(override, "host")) + port := integerOrNull(override, "port") + if host == "" || port == nil || *port < 1 || *port > 65535 { + return nil, fmt.Errorf("%s must include a host and a port between 1 and 65535", key) + } + return &address{Host: host, Port: *port}, nil +} + func boolProperty(config jsonObject, key string) bool { return boolOrDefault(objectOrNil(config, "properties"), key, false) } diff --git a/agents/drivers/rabbitmq/helpers_test.go b/agents/drivers/rabbitmq/helpers_test.go index 98b4d2f2e..4b77bddd2 100644 --- a/agents/drivers/rabbitmq/helpers_test.go +++ b/agents/drivers/rabbitmq/helpers_test.go @@ -3,8 +3,10 @@ package main import ( "encoding/base64" "encoding/json" + "net" "strings" "testing" + "time" ) func mustObject(t *testing.T, source string) jsonObject { @@ -134,6 +136,53 @@ func TestTLSAndManagementConfiguration(t *testing.T) { if managementPort(mustObject(t, `{"properties":{"management_port":55672}}`), false) != 55672 { t.Fatal("management port override ignored") } + override, err := endpointOverride(mustObject(t, `{"connect_override":{"host":"127.0.0.1","port":45672}}`), "connect_override") + if err != nil || override == nil || override.Host != "127.0.0.1" || override.Port != 45672 { + t.Fatalf("unexpected endpoint override %#v, %v", override, err) + } + if _, err := endpointOverride(mustObject(t, `{"connect_override":{"host":"","port":0}}`), "connect_override"); err == nil { + t.Fatal("invalid endpoint override was accepted") + } +} + +func TestDialAddressUsesConnectOverride(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer listener.Close() + accepted := make(chan string, 1) + go func() { + connection, acceptError := listener.Accept() + if acceptError != nil { + accepted <- "" + return + } + defer connection.Close() + header := make([]byte, 8) + read, _ := connection.Read(header) + accepted <- string(header[:read]) + }() + + localPort := listener.Addr().(*net.TCPAddr).Port + config := jsonObject{ + "connect_override": jsonObject{"host": "127.0.0.1", "port": localPort}, + "properties": jsonObject{ + "connection_timeout_ms": 1000, + "handshake_timeout_ms": 1000, + }, + } + if _, err := dialAddress(config, address{Host: "rabbit.invalid", Port: 5672}); err == nil { + t.Fatal("fake AMQP endpoint unexpectedly completed the handshake") + } + select { + case header := <-accepted: + if !strings.HasPrefix(header, "AMQP") { + t.Fatalf("tunnel endpoint did not receive the AMQP protocol header: %q", header) + } + case <-time.After(2 * time.Second): + t.Fatal("AMQP connection did not reach the tunnel endpoint") + } } func TestCredentialAndAuthHelpers(t *testing.T) { diff --git a/agents/drivers/rabbitmq/main.go b/agents/drivers/rabbitmq/main.go index bab7c4a3f..4911a8e0a 100644 --- a/agents/drivers/rabbitmq/main.go +++ b/agents/drivers/rabbitmq/main.go @@ -356,6 +356,10 @@ func openConnection(config jsonObject) (*amqp.Connection, error) { func dialAddress(config jsonObject, endpoint address) (*amqp.Connection, error) { properties := objectOrNil(config, "properties") + connectOverride, err := endpointOverride(config, "connect_override") + if err != nil { + return nil, err + } connectionTimeout := durationMilliseconds(config, "request_timeout_ms", defaultRequestTimeout) if configured, ok := integerProperty(properties, "connection_timeout_ms"); ok { connectionTimeout = time.Duration(configured) * time.Millisecond @@ -385,6 +389,9 @@ func dialAddress(config jsonObject, endpoint address) (*amqp.Connection, error) ChannelMax: defaultChannelMax, Heartbeat: time.Duration(heartbeat) * time.Second, Dial: func(network, target string) (net.Conn, error) { + if connectOverride != nil { + target = net.JoinHostPort(connectOverride.Host, strconv.Itoa(connectOverride.Port)) + } dialer := net.Dialer{Timeout: connectionTimeout} connection, err := dialer.Dial(network, target) if err != nil { diff --git a/agents/drivers/rabbitmq/management.go b/agents/drivers/rabbitmq/management.go index d78056afc..16eda1d38 100644 --- a/agents/drivers/rabbitmq/management.go +++ b/agents/drivers/rabbitmq/management.go @@ -19,6 +19,7 @@ import ( const ( defaultManagementPort = 15672 defaultManagementTLSPort = 15671 + defaultAMQPTLSPort = 5671 managementPageSize = 100 managementConnectTimeout = 10 * time.Second managementRequestTimeout = 20 * time.Second @@ -90,8 +91,23 @@ func managementRequestOnce(baseURL string, connection jsonObject, method, path s if body != nil { request.Header.Set("Content-Type", "application/json") } + connectOverride, err := endpointOverride(connection, "management_connect_override") + if err != nil { + return nil, err + } + dialer := &net.Dialer{Timeout: managementConnectTimeout} + dialContext := dialer.DialContext + if connectOverride != nil { + dialContext = func(ctx context.Context, network, _ string) (net.Conn, error) { + return dialer.DialContext( + ctx, + network, + net.JoinHostPort(connectOverride.Host, strconv.Itoa(connectOverride.Port)), + ) + } + } transport := &http.Transport{ - DialContext: (&net.Dialer{Timeout: managementConnectTimeout}).DialContext, + DialContext: dialContext, TLSHandshakeTimeout: managementConnectTimeout, ResponseHeaderTimeout: managementConnectTimeout, TLSClientConfig: &tls.Config{InsecureSkipVerify: tlsSkipVerify(connection)}, @@ -161,11 +177,22 @@ func managementBaseURLs(connection jsonObject) ([]string, error) { return []string{normalizeManagementURL(*explicit)}, nil } tlsEnabled := managementTLS(connection) - port := managementPort(connection, tlsEnabled) addresses, err := resolveAddresses(connection) if err != nil { return nil, err } + port, configured := configuredManagementPort(connection, tlsEnabled) + if !configured { + for _, endpoint := range addresses { + isDefaultAMQPPort := endpoint.Port == defaultAMQPPort || (tlsEnabled && endpoint.Port == defaultAMQPTLSPort) + if !isDefaultAMQPPort { + return nil, fmt.Errorf( + "RabbitMQ Management API URL is required when AMQP uses non-default port %d because the Management listener port is configured independently", + endpoint.Port, + ) + } + } + } baseURLs := make([]string, 0, len(addresses)) for _, endpoint := range addresses { baseURLs = append(baseURLs, managementBaseURL(endpoint.Host, port, tlsEnabled)) @@ -190,13 +217,18 @@ func managementTLS(connection jsonObject) bool { } func managementPort(connection jsonObject, tlsEnabled bool) int { + port, _ := configuredManagementPort(connection, tlsEnabled) + return port +} + +func configuredManagementPort(connection jsonObject, tlsEnabled bool) (int, bool) { if configured, ok := integerProperty(objectOrNil(connection, "properties"), "management_port"); ok { - return configured + return configured, true } if tlsEnabled { - return defaultManagementTLSPort + return defaultManagementTLSPort, false } - return defaultManagementPort + return defaultManagementPort, false } func credentialOrGuest(connection jsonObject, key string) string { diff --git a/agents/drivers/rabbitmq/management_test.go b/agents/drivers/rabbitmq/management_test.go index eba67c121..c412bc146 100644 --- a/agents/drivers/rabbitmq/management_test.go +++ b/agents/drivers/rabbitmq/management_test.go @@ -21,10 +21,18 @@ func TestManagementBaseURLs(t *testing.T) { if err != nil || withoutAddresses[0] != "http://mgmt:15672" { t.Fatalf("unexpected URL %#v, %v", withoutAddresses, err) } - derived, err := managementBaseURLs(mustObject(t, `{"addresses":"mq1:5672,mq2:5673"}`)) + derived, err := managementBaseURLs(mustObject(t, `{"addresses":"mq1:5672,mq2:5672"}`)) if err != nil || len(derived) != 2 || derived[0] != "http://mq1:15672" || derived[1] != "http://mq2:15672" { t.Fatalf("unexpected derived URLs %#v, %v", derived, err) } + _, err = managementBaseURLs(mustObject(t, `{"addresses":"mq1:5673"}`)) + if err == nil || !strings.Contains(err.Error(), "Management API URL is required") || !strings.Contains(err.Error(), "5673") { + t.Fatalf("unexpected custom AMQP port error %v", err) + } + customPort, err := managementBaseURLs(mustObject(t, `{"addresses":"mq1:5673","properties":{"management_port":15673}}`)) + if err != nil || len(customPort) != 1 || customPort[0] != "http://mq1:15673" { + t.Fatalf("unexpected custom management URL %#v, %v", customPort, err) + } tlsDerived, err := managementBaseURLs(mustObject(t, `{"addresses":"mq1","tls":{}}`)) if err != nil || tlsDerived[0] != "https://mq1:15671" { t.Fatalf("unexpected TLS URLs %#v, %v", tlsDerived, err) @@ -35,6 +43,40 @@ func TestManagementBaseURLs(t *testing.T) { } } +func TestManagementRequestUsesConnectOverride(t *testing.T) { + var observedHost string + server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) { + observedHost = request.Host + writer.Header().Set("Content-Type", "application/json") + _, _ = writer.Write([]byte(`{"items":[]}`)) + })) + defer server.Close() + + localAddress := strings.TrimPrefix(server.URL, "http://") + localHost, localPortText, err := net.SplitHostPort(localAddress) + if err != nil { + t.Fatal(err) + } + localPort, err := strconv.Atoi(localPortText) + if err != nil { + t.Fatal(err) + } + managementURL := "http://rabbit.internal:" + localPortText + connection := jsonObject{ + "management_url": managementURL, + "management_connect_override": jsonObject{ + "host": localHost, + "port": localPort, + }, + } + if _, err := managementGet(connection, "/api/queues"); err != nil { + t.Fatal(err) + } + if observedHost != "rabbit.internal:"+localPortText { + t.Fatalf("management Host header changed to %q", observedHost) + } +} + func TestManagementErrorMessages(t *testing.T) { for _, status := range []int{401, 403} { message := managementErrorMessage(status, http.MethodGet, "/api/queues") diff --git a/apps/desktop/src/i18n/locales/en.ts b/apps/desktop/src/i18n/locales/en.ts index 9e414b563..f6cc9628b 100644 --- a/apps/desktop/src/i18n/locales/en.ts +++ b/apps/desktop/src/i18n/locales/en.ts @@ -441,7 +441,7 @@ export default { mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "Leave empty to derive from AMQP addresses; reverse-proxy path prefixes like https://proxy/rmq are supported", + mqRabbitmqAdminUrlHint: "Leave empty only for the default AMQP port; custom AMQP ports require the independently configured Management API URL. Reverse-proxy paths like https://proxy/rmq are supported", mqRabbitmqUsernamePlaceholder: "Defaults to guest", mqRabbitmqPasswordPlaceholder: "Defaults to guest", mqSecurity: "Security", diff --git a/apps/desktop/src/i18n/locales/es.ts b/apps/desktop/src/i18n/locales/es.ts index a2a700f0a..c8fe020ad 100644 --- a/apps/desktop/src/i18n/locales/es.ts +++ b/apps/desktop/src/i18n/locales/es.ts @@ -603,7 +603,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "Déjalo vacío para derivarlo de las direcciones AMQP; admite prefijos de ruta de proxy inverso como https://proxy/rmq", + mqRabbitmqAdminUrlHint: "Déjalo vacío solo con el puerto AMQP predeterminado; los puertos AMQP personalizados requieren la URL de la API de administración configurada por separado. Admite rutas de proxy inverso como https://proxy/rmq", mqRabbitmqUsernamePlaceholder: "Por defecto guest", mqRabbitmqPasswordPlaceholder: "Por defecto guest", mqSecurity: "Security", diff --git a/apps/desktop/src/i18n/locales/it.ts b/apps/desktop/src/i18n/locales/it.ts index 30d22e317..c7eac48df 100644 --- a/apps/desktop/src/i18n/locales/it.ts +++ b/apps/desktop/src/i18n/locales/it.ts @@ -601,7 +601,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "Lascia vuoto per derivarlo dagli indirizzi AMQP; supporta prefissi di percorso di reverse proxy come https://proxy/rmq", + mqRabbitmqAdminUrlHint: "Lascia vuoto solo con la porta AMQP predefinita; le porte AMQP personalizzate richiedono l'URL dell'API di gestione configurato separatamente. Supporta percorsi reverse proxy come https://proxy/rmq", mqRabbitmqUsernamePlaceholder: "Predefinito guest", mqRabbitmqPasswordPlaceholder: "Predefinito guest", mqSecurity: "Security", diff --git a/apps/desktop/src/i18n/locales/ja.ts b/apps/desktop/src/i18n/locales/ja.ts index fd9dc14a0..98e088559 100644 --- a/apps/desktop/src/i18n/locales/ja.ts +++ b/apps/desktop/src/i18n/locales/ja.ts @@ -601,7 +601,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "空欄の場合は AMQP アドレスから派生します。https://proxy/rmq のようなリバースプロキシのパスプレフィックスに対応しています", + mqRabbitmqAdminUrlHint: "空欄から派生できるのは既定の AMQP ポートだけです。カスタム AMQP ポートでは個別に設定した Management API URL が必要です。https://proxy/rmq のようなリバースプロキシパスにも対応しています", mqRabbitmqUsernamePlaceholder: "デフォルトは guest", mqRabbitmqPasswordPlaceholder: "デフォルトは guest", mqSecurity: "Security", diff --git a/apps/desktop/src/i18n/locales/ko.ts b/apps/desktop/src/i18n/locales/ko.ts index 3d7b3e518..c3b66bcc9 100644 --- a/apps/desktop/src/i18n/locales/ko.ts +++ b/apps/desktop/src/i18n/locales/ko.ts @@ -437,7 +437,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "관리 URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "AMQP 주소에서 자동 유추하려면 비워 두세요. https://proxy/rmq 같은 역방향 프록시 경로 접두사를 지원합니다", + mqRabbitmqAdminUrlHint: "기본 AMQP 포트에서만 비워 둘 수 있습니다. 사용자 지정 AMQP 포트는 별도로 구성된 Management API URL이 필요합니다. https://proxy/rmq 같은 역방향 프록시 경로도 지원합니다", mqRabbitmqUsernamePlaceholder: "기본값 guest", mqRabbitmqPasswordPlaceholder: "기본값 guest", mqSecurity: "보안", diff --git a/apps/desktop/src/i18n/locales/pt-BR.ts b/apps/desktop/src/i18n/locales/pt-BR.ts index 0a1c5988e..255d423e0 100644 --- a/apps/desktop/src/i18n/locales/pt-BR.ts +++ b/apps/desktop/src/i18n/locales/pt-BR.ts @@ -602,7 +602,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "Deixe vazio para derivar dos endereços AMQP; suporta prefixos de caminho de proxy reverso como https://proxy/rmq", + mqRabbitmqAdminUrlHint: "Deixe vazio apenas com a porta AMQP padrão; portas AMQP personalizadas exigem a URL da API de gerenciamento configurada separadamente. Suporta caminhos de proxy reverso como https://proxy/rmq", mqRabbitmqUsernamePlaceholder: "Padrão guest", mqRabbitmqPasswordPlaceholder: "Padrão guest", mqSecurity: "Security", diff --git a/apps/desktop/src/i18n/locales/zh-CN.ts b/apps/desktop/src/i18n/locales/zh-CN.ts index 7d93deafa..e2821da78 100644 --- a/apps/desktop/src/i18n/locales/zh-CN.ts +++ b/apps/desktop/src/i18n/locales/zh-CN.ts @@ -442,7 +442,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "留空则按 AMQP 地址派生;支持反代路径前缀,如 https://proxy/rmq", + mqRabbitmqAdminUrlHint: "仅默认 AMQP 端口可留空派生;非默认 AMQP 端口必须填写独立配置的管理 API 地址。支持反代路径前缀,如 https://proxy/rmq", mqRabbitmqUsernamePlaceholder: "缺省 guest", mqRabbitmqPasswordPlaceholder: "缺省 guest", mqSecurity: "Security", diff --git a/apps/desktop/src/i18n/locales/zh-TW.ts b/apps/desktop/src/i18n/locales/zh-TW.ts index 5bd0cf575..807e07deb 100644 --- a/apps/desktop/src/i18n/locales/zh-TW.ts +++ b/apps/desktop/src/i18n/locales/zh-TW.ts @@ -601,7 +601,7 @@ export default withEnglishFallback({ mqVirtualHostPlaceholder: "/", mqRabbitmqAdminUrl: "Management URL", mqRabbitmqAdminUrlPlaceholder: "http://192.168.1.1:15672", - mqRabbitmqAdminUrlHint: "留空則按 AMQP 位址派生;支援反代路徑前綴,如 https://proxy/rmq", + mqRabbitmqAdminUrlHint: "僅預設 AMQP 連接埠可留空派生;非預設 AMQP 連接埠必須填寫獨立設定的管理 API 位址。支援反代路徑前綴,如 https://proxy/rmq", mqRabbitmqUsernamePlaceholder: "預設 guest", mqRabbitmqPasswordPlaceholder: "預設 guest", mqSecurity: "Security", diff --git a/crates/dbx-core/src/connection.rs b/crates/dbx-core/src/connection.rs index 18c30843b..4e2497132 100644 --- a/crates/dbx-core/src/connection.rs +++ b/crates/dbx-core/src/connection.rs @@ -2213,6 +2213,60 @@ impl AppState { return Ok(mqc); } + if mqc.system_kind == crate::mq::types::MqSystemKind::RabbitMq { + let transport_layers = self.resolved_transport_layers(config).await?; + let (amqp_host, amqp_port) = crate::mq::adapters::rabbitmq::primary_amqp_endpoint(&mqc)?; + let management_endpoint = crate::mq::adapters::rabbitmq::management_endpoint(&mqc)?; + let amqp_local_port = db::transport_layer_tunnel::start_transport_layers( + connection_id, + &transport_layers, + &amqp_host, + amqp_port, + &self.tunnels, + &self.proxy_tunnels, + &self.http_tunnels, + ) + .await?; + let mut mqc = mqc.with_connect_override("127.0.0.1", amqp_local_port); + if let Some((management_host, management_port)) = management_endpoint { + let management_transport_id = rabbitmq_management_transport_id(connection_id); + let management_local_port = match db::transport_layer_tunnel::start_transport_layers( + &management_transport_id, + &transport_layers, + &management_host, + management_port, + &self.tunnels, + &self.proxy_tunnels, + &self.http_tunnels, + ) + .await + { + Ok(port) => port, + Err(error) => { + db::transport_layer_tunnel::stop_transport_layers( + &management_transport_id, + transport_layers.len(), + &self.tunnels, + &self.proxy_tunnels, + &self.http_tunnels, + ) + .await; + db::transport_layer_tunnel::stop_transport_layers( + connection_id, + transport_layers.len(), + &self.tunnels, + &self.proxy_tunnels, + &self.http_tunnels, + ) + .await; + return Err(error); + } + }; + mqc = mqc.with_management_connect_override("127.0.0.1", management_local_port); + } + return Ok(mqc); + } + let (host, port) = self.connection_host_port(connection_id, config).await?; Ok(mqc.with_connect_override(&host, port)) } @@ -3052,6 +3106,14 @@ impl AppState { &self.http_tunnels, ) .await; + db::transport_layer_tunnel::stop_transport_layers( + &rabbitmq_management_transport_id(connection_id), + layer_count, + &self.tunnels, + &self.proxy_tunnels, + &self.http_tunnels, + ) + .await; self.tunnels.stop_tunnel(connection_id).await; self.proxy_tunnels.stop_tunnel(connection_id).await; self.http_tunnels.stop_tunnel(connection_id).await; @@ -3618,6 +3680,10 @@ fn rnacos_console_transport_id(connection_id: &str) -> String { format!("{connection_id}:rnacos-console") } +fn rabbitmq_management_transport_id(connection_id: &str) -> String { + format!("{connection_id}:rabbitmq-management") +} + fn parse_mq_admin_host_port(config: &ConnectionConfig) -> Option<(String, u16)> { let value = config .external_config @@ -6122,6 +6188,55 @@ for line in sys.stdin: let _ = std::fs::remove_dir_all(dir); } + #[cfg(feature = "mq-admin")] + #[tokio::test] + async fn rabbitmq_transport_uses_separate_amqp_and_management_tunnels() { + let (state, dir) = test_app_state().await; + let mut config = mysql_config(None); + config.id = "proxied-rabbitmq".to_string(); + config.db_type = DatabaseType::MessageQueue; + config.host = "rabbit.internal".to_string(); + config.port = 5672; + config.external_config = Some(serde_json::json!({ + "systemKind": "rabbitmq", + "adminUrl": "http://management.internal:15672/rmq", + "auth": { "kind": "none" }, + "extra": { + "addresses": "rabbit.internal:5672", + "virtualHost": "/" + } + })); + config.transport_layers = vec![TransportLayerConfig::Proxy(ProxyTunnelConfig { + profile_id: String::new(), + id: "proxy".to_string(), + name: String::new(), + enabled: true, + proxy_type: ProxyType::Socks5, + host: "127.0.0.1".to_string(), + port: 65000, + username: String::new(), + password: String::new(), + test_target: None, + })]; + + let mqc = state.mq_admin_config_for_connection("proxied-rabbitmq", &config).await.unwrap(); + let amqp_override = mqc.connect_override.expect("RabbitMQ AMQP tunnel override"); + let management_override = mqc.management_connect_override.expect("RabbitMQ Management tunnel override"); + assert_eq!(amqp_override.host, "127.0.0.1"); + assert_eq!(management_override.host, "127.0.0.1"); + assert_ne!(amqp_override.port, management_override.port); + assert_eq!(state.proxy_tunnels.local_port("proxied-rabbitmq:transport:0").await, Some(amqp_override.port)); + assert_eq!( + state.proxy_tunnels.local_port("proxied-rabbitmq:rabbitmq-management:transport:0").await, + Some(management_override.port) + ); + + state.reset_connection_transport_for_config("proxied-rabbitmq", &config).await; + assert!(state.proxy_tunnels.local_port("proxied-rabbitmq:transport:0").await.is_none()); + assert!(state.proxy_tunnels.local_port("proxied-rabbitmq:rabbitmq-management:transport:0").await.is_none()); + let _ = std::fs::remove_dir_all(dir); + } + #[tokio::test] async fn nacos_admin_config_allows_domain_server_addr_without_transport_override() { let (state, dir) = test_app_state().await; diff --git a/crates/dbx-core/src/mq/adapters/kafka.rs b/crates/dbx-core/src/mq/adapters/kafka.rs index 0ff008a7b..64edaeb6a 100644 --- a/crates/dbx-core/src/mq/adapters/kafka.rs +++ b/crates/dbx-core/src/mq/adapters/kafka.rs @@ -766,6 +766,7 @@ mod tests { pinned_version: None, token_signing: None, connect_override: None, + management_connect_override: None, extra, } } diff --git a/crates/dbx-core/src/mq/adapters/pulsar.rs b/crates/dbx-core/src/mq/adapters/pulsar.rs index 704450469..e2f1e7b44 100644 --- a/crates/dbx-core/src/mq/adapters/pulsar.rs +++ b/crates/dbx-core/src/mq/adapters/pulsar.rs @@ -1150,6 +1150,7 @@ mod tests { pinned_version: Some("3.1".to_string()), token_signing: None, connect_override: None, + management_connect_override: None, extra: serde_json::Value::Null, }) .await diff --git a/crates/dbx-core/src/mq/adapters/rabbitmq.rs b/crates/dbx-core/src/mq/adapters/rabbitmq.rs index fb2d2bbbf..513a10d55 100644 --- a/crates/dbx-core/src/mq/adapters/rabbitmq.rs +++ b/crates/dbx-core/src/mq/adapters/rabbitmq.rs @@ -861,7 +861,73 @@ fn extra_port(extra: &serde_json::Value) -> Result { } } -/// Build the connection params JSON from MqAdminConfig for the Java agent. +fn configured_management_port(cfg: &MqAdminConfig) -> Result, String> { + let Some(value) = cfg.extra.get("properties").and_then(|value| value.get("management_port")) else { + return Ok(None); + }; + let port = if let Some(port) = value.as_u64() { + u16::try_from(port).map_err(|_| format!("RabbitMQ management port {port} is out of range (1-65535)"))? + } else if let Some(port) = value.as_str() { + port.trim().parse::().map_err(|_| format!("invalid RabbitMQ management port '{port}'"))? + } else { + return Err("RabbitMQ management port must be a number or a numeric string".to_string()); + }; + if port == 0 { + return Err("RabbitMQ management port must be between 1 and 65535".to_string()); + } + Ok(Some(port)) +} + +pub(crate) fn primary_amqp_endpoint(cfg: &MqAdminConfig) -> Result<(String, u16), String> { + let address_list = addresses(cfg); + let first = address_list + .split(',') + .map(str::trim) + .find(|address| !address.is_empty()) + .ok_or("RabbitMQ addresses are empty")?; + let parsed = reqwest::Url::parse(&format!("amqp://{first}")) + .map_err(|error| format!("RabbitMQ address '{first}' is invalid: {error}"))?; + let host = parsed + .host_str() + .filter(|host| !host.is_empty()) + .ok_or_else(|| format!("RabbitMQ address '{first}' has no host"))?; + Ok((host.to_string(), parsed.port().unwrap_or(extra_port(&cfg.extra)?))) +} + +/// Resolve the independently configured RabbitMQ Management HTTP endpoint. +/// +/// A blank admin URL is only safe to derive for RabbitMQ's default AMQP ports. +/// Custom AMQP listeners do not define the Management listener port, so callers +/// must not silently fall back to another broker's default port 15672/15671. +pub(crate) fn management_endpoint(cfg: &MqAdminConfig) -> Result, String> { + if !cfg.admin_url.trim().is_empty() { + let parsed = reqwest::Url::parse(cfg.admin_url.trim()) + .map_err(|error| format!("RabbitMQ Management API URL is invalid: {error}"))?; + let host = parsed + .host_str() + .filter(|host| !host.is_empty()) + .ok_or("RabbitMQ Management API URL does not include a host")?; + let port = parsed.port_or_known_default().ok_or("RabbitMQ Management API URL does not include a port")?; + return Ok(Some((host.to_string(), port))); + } + + let (host, amqp_port) = primary_amqp_endpoint(cfg)?; + if let Some(port) = configured_management_port(cfg)? { + return Ok(Some((host, port))); + } + let tls_enabled = + cfg.extra.get("properties").and_then(|properties| properties.as_object()).is_some_and(|properties| { + properties.get("ssl").and_then(|value| value.as_bool()).unwrap_or(false) + || properties.get("tls").and_then(|value| value.as_bool()).unwrap_or(false) + }); + let is_default_amqp_port = amqp_port == 5672 || (tls_enabled && amqp_port == 5671); + if !is_default_amqp_port { + return Ok(None); + } + Ok(Some((host, if tls_enabled { 15671 } else { 15672 }))) +} + +/// Build the connection params JSON from MqAdminConfig for the native agent. /// Blank credentials are omitted so the agent falls back to its guest/guest /// default instead of authenticating as `:`. A non-empty `admin_url` is /// forwarded as `management_url`; otherwise the agent derives the management @@ -894,6 +960,18 @@ fn build_connection_params(cfg: &MqAdminConfig) -> Result, + /// Runtime-only TCP endpoint override for a secondary management endpoint. + /// RabbitMQ uses this in addition to the AMQP `connect_override` because + /// its Management HTTP API listens on an independently configured port. + #[serde(skip)] + pub management_connect_override: Option, /// System-specific extension fields (e.g. Kafka bootstrap servers). #[serde(default, skip_serializing_if = "serde_json::Value::is_null")] pub extra: serde_json::Value, @@ -87,6 +92,11 @@ impl MqAdminConfig { self.connect_override = Some(MqConnectOverride { host: host.to_string(), port }); self } + + pub fn with_management_connect_override(mut self, host: &str, port: u16) -> Self { + self.management_connect_override = Some(MqConnectOverride { host: host.to_string(), port }); + self + } } pub fn admin_url_with_endpoint(admin_url: &str, host: &str, port: u16) -> Result { diff --git a/docs/mq-quick-start.md b/docs/mq-quick-start.md index eaced23db..e1e15ed8e 100644 --- a/docs/mq-quick-start.md +++ b/docs/mq-quick-start.md @@ -457,6 +457,8 @@ docker run -d --name dbx-rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-managem 说明:RabbitMQ 适配器将 topic 映射为队列(queue),支持队列列表/声明/删除、清空队列(purge)、消费者列表、消息预览(basic.get + requeue)与发送(basic.publish);vhost 映射为 namespace,可在控制台查看/创建/删除。AMQP 操作按 vhost 透传(agent 侧 per-vhost 通道缓存),tenant 语义不适用(固定合成 `_rabbitmq`)。 +RabbitMQ 的 AMQP 监听端口与 Management HTTP 监听端口是独立配置。使用默认 AMQP 端口时可将 `adminUrl` 留空,由 DBX 使用默认 Management 端口;使用非默认 AMQP 端口时必须填写实际的 Management API URL,例如 AMQP 为 `127.0.0.1:5673`、Management 为 `http://127.0.0.1:15673`。SSH、SOCKS 或 HTTP 隧道会分别转发 AMQP 与 Management 两个端点。 + ### Kafka 适配器参考 1. 实现 `MessageQueueAdmin` trait: