fix(rabbitmq): honor tunnels and management endpoints
This commit is contained in:
parent
dd913b5ab2
commit
9057b746a6
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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: "보안",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -766,6 +766,7 @@ mod tests {
|
|||
pinned_version: None,
|
||||
token_signing: None,
|
||||
connect_override: None,
|
||||
management_connect_override: None,
|
||||
extra,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -861,7 +861,73 @@ fn extra_port(extra: &serde_json::Value) -> Result<u16, String> {
|
|||
}
|
||||
}
|
||||
|
||||
/// Build the connection params JSON from MqAdminConfig for the Java agent.
|
||||
fn configured_management_port(cfg: &MqAdminConfig) -> Result<Option<u16>, 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::<u16>().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<Option<(String, u16)>, 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<serde_json::Value, Str
|
|||
if !cfg.admin_url.trim().is_empty() {
|
||||
params["management_url"] = serde_json::json!(cfg.admin_url);
|
||||
}
|
||||
if let Some(connect_override) = &cfg.connect_override {
|
||||
params["connect_override"] = serde_json::json!({
|
||||
"host": connect_override.host,
|
||||
"port": connect_override.port,
|
||||
});
|
||||
}
|
||||
if let Some(connect_override) = &cfg.management_connect_override {
|
||||
params["management_connect_override"] = serde_json::json!({
|
||||
"host": connect_override.host,
|
||||
"port": connect_override.port,
|
||||
});
|
||||
}
|
||||
Ok(params)
|
||||
}
|
||||
|
||||
|
|
@ -946,6 +1024,7 @@ mod tests {
|
|||
pinned_version: None,
|
||||
token_signing: None,
|
||||
connect_override: None,
|
||||
management_connect_override: None,
|
||||
extra,
|
||||
}
|
||||
}
|
||||
|
|
@ -1046,10 +1125,36 @@ mod tests {
|
|||
false,
|
||||
);
|
||||
cfg.admin_url = "http://rabbit.internal:15672/proxy".to_string();
|
||||
cfg = cfg.with_connect_override("127.0.0.1", 45672).with_management_connect_override("127.0.0.1", 45673);
|
||||
|
||||
let params = build_connection_params(&cfg).expect("connection params");
|
||||
|
||||
assert_eq!(params.get("management_url").and_then(|v| v.as_str()), Some("http://rabbit.internal:15672/proxy"));
|
||||
assert_eq!(params.pointer("/connect_override/host").and_then(|v| v.as_str()), Some("127.0.0.1"));
|
||||
assert_eq!(params.pointer("/connect_override/port").and_then(|v| v.as_u64()), Some(45672));
|
||||
assert_eq!(params.pointer("/management_connect_override/host").and_then(|v| v.as_str()), Some("127.0.0.1"));
|
||||
assert_eq!(params.pointer("/management_connect_override/port").and_then(|v| v.as_u64()), Some(45673));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn management_endpoint_does_not_guess_for_custom_amqp_port() {
|
||||
let custom = rabbitmq_config(serde_json::json!({ "addresses": "rabbit.internal:5673" }), MqAuth::None, false);
|
||||
assert_eq!(primary_amqp_endpoint(&custom).unwrap(), ("rabbit.internal".to_string(), 5673));
|
||||
assert_eq!(management_endpoint(&custom).unwrap(), None);
|
||||
|
||||
let configured = rabbitmq_config(
|
||||
serde_json::json!({
|
||||
"addresses": "rabbit.internal:5673",
|
||||
"properties": { "management_port": 15673 }
|
||||
}),
|
||||
MqAuth::None,
|
||||
false,
|
||||
);
|
||||
assert_eq!(management_endpoint(&configured).unwrap(), Some(("rabbit.internal".to_string(), 15673)));
|
||||
|
||||
let mut explicit = custom;
|
||||
explicit.admin_url = "https://management.internal:8443/rmq".to_string();
|
||||
assert_eq!(management_endpoint(&explicit).unwrap(), Some(("management.internal".to_string(), 8443)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
|
|
@ -966,6 +966,7 @@ mod tests {
|
|||
pinned_version: None,
|
||||
token_signing: None,
|
||||
connect_override: None,
|
||||
management_connect_override: None,
|
||||
extra,
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -10,12 +10,12 @@ use crate::models::connection::ConnectionConfig;
|
|||
use crate::mq::auth::MqAuth;
|
||||
use crate::mq::types::{MqSystemKind, MqTokenSigningConfig};
|
||||
|
||||
/// Runtime TCP endpoint override for MQ admin requests.
|
||||
/// Runtime TCP endpoint override for an MQ transport.
|
||||
///
|
||||
/// The public admin URL remains unchanged so TLS hostname verification, SNI and
|
||||
/// the HTTP Host header continue to target the broker name. The HTTP client uses
|
||||
/// this endpoint only for the underlying TCP connection, e.g. after an SSH/proxy
|
||||
/// tunnel has mapped the broker to a local port.
|
||||
/// The logical broker endpoint remains unchanged so TLS hostname verification,
|
||||
/// SNI and protocol-level host names continue to target the broker. The client
|
||||
/// uses this endpoint only for the underlying TCP connection, e.g. after an
|
||||
/// SSH/proxy tunnel has mapped the broker to a local port.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct MqConnectOverride {
|
||||
|
|
@ -48,6 +48,11 @@ pub struct MqAdminConfig {
|
|||
/// Runtime-only TCP endpoint override used by transport layers.
|
||||
#[serde(skip)]
|
||||
pub connect_override: Option<MqConnectOverride>,
|
||||
/// 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<MqConnectOverride>,
|
||||
/// 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<String, String> {
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
Loading…
Reference in New Issue