This is an automated email from the ASF dual-hosted git repository. squakez pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/camel-k.git
commit 5aa092050467fdd390975c86a81a60a105a257fd Author: Keerthan <[email protected]> AuthorDate: Thu Sep 24 22:44:35 2026 +0530 Fix #5935: Use AMQP component and ApplicationProperties for ArkMQ binding --- pkg/controller/pipe/initialize_test.go | 12 +++- pkg/util/bindings/arkmq.go | 116 ++++++++++++++++++++++----------- pkg/util/bindings/arkmq_test.go | 36 ++++++++-- 3 files changed, 117 insertions(+), 47 deletions(-) diff --git a/pkg/controller/pipe/initialize_test.go b/pkg/controller/pipe/initialize_test.go index 9d4abcbd8..6a2920812 100644 --- a/pkg/controller/pipe/initialize_test.go +++ b/pkg/controller/pipe/initialize_test.go @@ -519,8 +519,10 @@ func TestNewPipeArkMQBinding(t *testing.T) { assert.Equal(t, "my-pipe", expectedIT.Labels[kubernetes.CamelCreatorLabelName]) flow, err := json.Marshal(expectedIT.Spec.Flows[0].RawMessage) require.NoError(t, err) - assert.Equal(t, "{\"route\":{\"from\":{\"steps\":[{\"to\":\"jms:queue:my-queue?brokerURL=tcp%3A%2F%2Fmy-broker-hdls-svc%3A61616\"}],"+ + assert.Equal(t, "{\"route\":{\"from\":{\"steps\":[{\"to\":\"amqp:queue:my-queue\"}],"+ "\"uri\":\"direct:something\"},\"id\":\"binding\"}}", string(flow)) + assert.Equal(t, "amqp://my-broker-hdls-svc:61616", expectedIT.Spec.GetConfigurationProperty("quarkus.qpid-jms.url")) + assert.Equal(t, "amqp://my-broker-hdls-svc:61616", expectedIT.Spec.GetConfigurationProperty("camel.component.amqp.broker-url")) } func TestNewPipeArkMQSourceBinding(t *testing.T) { @@ -571,7 +573,9 @@ func TestNewPipeArkMQSourceBinding(t *testing.T) { flow, err := json.Marshal(expectedIT.Spec.Flows[0].RawMessage) require.NoError(t, err) assert.Equal(t, "{\"route\":{\"from\":{\"steps\":[{\"to\":\"log:info\"}],"+ - "\"uri\":\"jms:queue:my-queue?brokerURL=tcp%3A%2F%2Fmy-broker-hdls-svc%3A61616\"},\"id\":\"binding\"}}", string(flow)) + "\"uri\":\"amqp:queue:my-queue\"},\"id\":\"binding\"}}", string(flow)) + assert.Equal(t, "amqp://my-broker-hdls-svc:61616", expectedIT.Spec.GetConfigurationProperty("quarkus.qpid-jms.url")) + assert.Equal(t, "amqp://my-broker-hdls-svc:61616", expectedIT.Spec.GetConfigurationProperty("camel.component.amqp.broker-url")) } func TestNewPipeArkMQBrokerBinding(t *testing.T) { @@ -622,8 +626,10 @@ func TestNewPipeArkMQBrokerBinding(t *testing.T) { assert.Equal(t, "my-pipe-broker", expectedIT.Labels[kubernetes.CamelCreatorLabelName]) flow, err := json.Marshal(expectedIT.Spec.Flows[0].RawMessage) require.NoError(t, err) - assert.Equal(t, "{\"route\":{\"from\":{\"steps\":[{\"to\":\"jms:queue:my-orders?brokerURL=tcp%3A%2F%2Fmy-broker-hdls-svc%3A61616\"}],"+ + assert.Equal(t, "{\"route\":{\"from\":{\"steps\":[{\"to\":\"amqp:queue:my-orders\"}],"+ "\"uri\":\"direct:something\"},\"id\":\"binding\"}}", string(flow)) + assert.Equal(t, "amqp://my-broker-hdls-svc:61616", expectedIT.Spec.GetConfigurationProperty("quarkus.qpid-jms.url")) + assert.Equal(t, "amqp://my-broker-hdls-svc:61616", expectedIT.Spec.GetConfigurationProperty("camel.component.amqp.broker-url")) } func asEndpointProperties(props map[string]string) *v1.EndpointProperties { diff --git a/pkg/util/bindings/arkmq.go b/pkg/util/bindings/arkmq.go index df397bbb7..10f70e709 100644 --- a/pkg/util/bindings/arkmq.go +++ b/pkg/util/bindings/arkmq.go @@ -21,6 +21,7 @@ import ( "fmt" "net" "strconv" + "strings" camelv1 "github.com/apache/camel-k/v2/pkg/apis/camel/v1" arkmqv1beta1 "github.com/apache/camel-k/v2/pkg/apis/duck/arkmq/v1beta1" @@ -34,13 +35,14 @@ func init() { RegisterBindingProvider(ArkMQBindingProvider{}) } -// camelArkMQ represents the configuration required by Camel JMS component for ArkMQ. +// camelArkMQ represents the configuration required by Camel AMQP component for ArkMQ. type camelArkMQ struct { queueName string + brokerURL string properties map[string]string } -// defaultArtemisPort is the default core port for ActiveMQ Artemis. +// defaultArtemisPort is the default port for ActiveMQ Artemis. const defaultArtemisPort = int32(61616) // ArkMQBindingProvider allows connecting to an ArkMQ queue via Binding. @@ -68,11 +70,18 @@ func (a ArkMQBindingProvider) Translate(ctx BindingContext, _ EndpointContext, e if err != nil { return nil, err } - arkmqURI := "jms:queue:" + camelArkMQ.queueName + arkmqURI := "amqp:queue:" + camelArkMQ.queueName arkmqURI = uri.AppendParameters(arkmqURI, camelArkMQ.properties) + appProps := make(map[string]string) + if camelArkMQ.brokerURL != "" { + appProps["quarkus.qpid-jms.url"] = camelArkMQ.brokerURL + appProps["camel.component.amqp.broker-url"] = camelArkMQ.brokerURL + } + return &Binding{ - URI: arkmqURI, + URI: arkmqURI, + ApplicationProperties: appProps, }, nil } @@ -89,7 +98,7 @@ func (a ArkMQBindingProvider) toCamelArkMQ(ctx BindingContext, endpoint camelv1. endpoint.Ref.Kind, arkmqv1beta1.ArkMQKindBroker, arkmqv1beta1.ArkMQKindAddress) } -// Verify and transform an ActiveMQArtemis broker resource to Camel JMS queue endpoint parameters. +// Verify and transform an ActiveMQArtemis broker resource to Camel AMQP queue endpoint parameters. func (a ArkMQBindingProvider) fromBrokerToCamel(ctx BindingContext, endpoint camelv1.Endpoint) (*camelArkMQ, error) { props, err := endpoint.Properties.GetPropertyMap() if err != nil { @@ -110,27 +119,39 @@ func (a ArkMQBindingProvider) fromBrokerToCamel(ctx BindingContext, endpoint cam delete(props, "destination") delete(props, "queue") - if props["brokerURL"] == "" { - namespace := endpoint.Ref.Namespace - if namespace == "" { - namespace = ctx.Namespace - } + brokerURL := props["brokerURL"] + if brokerURL == "" { + brokerURL = props["brokerUrl"] + } + delete(props, "brokerURL") + delete(props, "brokerUrl") + + if brokerURL != "" { + return &camelArkMQ{ + queueName: queueName, + brokerURL: normalizeBrokerURL(brokerURL), + properties: props, + }, nil + } - brokerURL, err := a.getBrokerURL(ctx, endpoint.Ref.Name, namespace) - if err != nil { - return nil, err - } + namespace := endpoint.Ref.Namespace + if namespace == "" { + namespace = ctx.Namespace + } - props["brokerURL"] = brokerURL + brokerURL, err = a.getBrokerURL(ctx, endpoint.Ref.Name, namespace) + if err != nil { + return nil, err } return &camelArkMQ{ queueName: queueName, + brokerURL: brokerURL, properties: props, }, nil } -// Verify and transform an ActiveMQArtemisAddress resource to Camel JMS queue endpoint parameters. +// Verify and transform an ActiveMQArtemisAddress resource to Camel AMQP queue endpoint parameters. func (a ArkMQBindingProvider) fromAddressToCamel(ctx BindingContext, endpoint camelv1.Endpoint) (*camelArkMQ, error) { props, err := endpoint.Properties.GetPropertyMap() if err != nil { @@ -141,37 +162,56 @@ func (a ArkMQBindingProvider) fromAddressToCamel(ctx BindingContext, endpoint ca } queueName := endpoint.Ref.Name - //nolint:nestif - if props["brokerURL"] == "" { - address, err := a.lookupAddress(ctx, endpoint) - if err != nil { - return nil, err - } + brokerURL := props["brokerURL"] + if brokerURL == "" { + brokerURL = props["brokerUrl"] + } + delete(props, "brokerURL") + delete(props, "brokerUrl") + + if brokerURL != "" { + return &camelArkMQ{ + queueName: queueName, + brokerURL: normalizeBrokerURL(brokerURL), + properties: props, + }, nil + } - if address.Spec.RoutingType == "multicast" { - return nil, fmt.Errorf("multicast addresses (topics) are not supported on queue binding %s", endpoint.Ref.Name) - } + address, err := a.lookupAddress(ctx, endpoint) + if err != nil { + return nil, err + } - if address.Spec.QueueName != "" { - queueName = address.Spec.QueueName - } else if address.Spec.AddressName != "" { - queueName = address.Spec.AddressName - } + if address.Spec.RoutingType == "multicast" { + return nil, fmt.Errorf("multicast addresses (topics) are not supported on queue binding %s", endpoint.Ref.Name) + } - brokerURL, err := a.lookupBrokerURL(ctx, address, endpoint) - if err != nil { - return nil, err - } + if address.Spec.QueueName != "" { + queueName = address.Spec.QueueName + } else if address.Spec.AddressName != "" { + queueName = address.Spec.AddressName + } - props["brokerURL"] = brokerURL + brokerURL, err = a.lookupBrokerURL(ctx, address, endpoint) + if err != nil { + return nil, err } return &camelArkMQ{ queueName: queueName, + brokerURL: brokerURL, properties: props, }, nil } +func normalizeBrokerURL(url string) string { + if after, ok := strings.CutPrefix(url, "tcp://"); ok { + return "amqp://" + after + } + + return url +} + func (a ArkMQBindingProvider) lookupBrokerURL(ctx BindingContext, address *arkmqv1beta1.ActiveMQArtemisAddress, endpoint camelv1.Endpoint) (string, error) { clusterName := address.Spec.ApplyTo if clusterName == "" && address.Labels != nil { @@ -224,7 +264,7 @@ func (a ArkMQBindingProvider) getBrokerURL(ctx BindingContext, clusterName, name port := defaultArtemisPort for _, p := range broker.Status.PortStatus { - if p.Name == "core" || p.Name == "all" || p.Name == "openwire" { + if p.Name == "amqp" || p.Name == "all" || p.Name == "core" || p.Name == "openwire" { if p.Port > 0 { port = p.Port @@ -251,7 +291,7 @@ func (a ArkMQBindingProvider) getBrokerURL(ctx BindingContext, clusterName, name } for _, p := range svc.Spec.Ports { - if p.Name == "core" || p.Name == "all" || p.Name == "openwire" || p.Port == defaultArtemisPort { + if p.Name == "amqp" || p.Name == "all" || p.Name == "core" || p.Name == "openwire" || p.Port == defaultArtemisPort { port = p.Port break @@ -263,7 +303,7 @@ func (a ArkMQBindingProvider) getBrokerURL(ctx BindingContext, clusterName, name host = svc.Name + "." + svc.Namespace + ".svc" } - return "tcp://" + net.JoinHostPort(host, strconv.Itoa(int(port))), nil + return "amqp://" + net.JoinHostPort(host, strconv.Itoa(int(port))), nil } func (a ArkMQBindingProvider) lookupAddress(ctx BindingContext, endpoint camelv1.Endpoint) (*arkmqv1beta1.ActiveMQArtemisAddress, error) { diff --git a/pkg/util/bindings/arkmq_test.go b/pkg/util/bindings/arkmq_test.go index 6cb610db3..1fac85d00 100644 --- a/pkg/util/bindings/arkmq_test.go +++ b/pkg/util/bindings/arkmq_test.go @@ -61,7 +61,11 @@ func TestArkMQDirect(t *testing.T) { }, endpoint) require.NoError(t, err) assert.NotNil(t, binding) - assert.Equal(t, "jms:queue:myqueue?brokerURL=tcp%3A%2F%2Fcustom-broker%3A61616", binding.URI) + assert.Equal(t, "amqp:queue:myqueue", binding.URI) + assert.Equal(t, map[string]string{ + "quarkus.qpid-jms.url": "amqp://custom-broker:61616", + "camel.component.amqp.broker-url": "amqp://custom-broker:61616", + }, binding.ApplicationProperties) } func TestArkMQLookupAddress(t *testing.T) { @@ -134,7 +138,11 @@ func TestArkMQLookupAddress(t *testing.T) { }, endpoint) require.NoError(t, err) assert.NotNil(t, binding) - assert.Equal(t, "jms:queue:resolved-queue?brokerURL=tcp%3A%2F%2Fmybroker-hdls-svc.test.svc%3A61616", binding.URI) + assert.Equal(t, "amqp:queue:resolved-queue", binding.URI) + assert.Equal(t, map[string]string{ + "quarkus.qpid-jms.url": "amqp://mybroker-hdls-svc.test.svc:61616", + "camel.component.amqp.broker-url": "amqp://mybroker-hdls-svc.test.svc:61616", + }, binding.ApplicationProperties) } func TestArkMQMulticastUnsupported(t *testing.T) { @@ -243,7 +251,11 @@ func TestArkMQLookupAddressByName(t *testing.T) { }, endpoint) require.NoError(t, err) assert.NotNil(t, binding) - assert.Equal(t, "jms:queue:events-queue?brokerURL=tcp%3A%2F%2Fmybroker-hdls-svc.test.svc%3A61616", binding.URI) + assert.Equal(t, "amqp:queue:events-queue", binding.URI) + assert.Equal(t, map[string]string{ + "quarkus.qpid-jms.url": "amqp://mybroker-hdls-svc.test.svc:61616", + "camel.component.amqp.broker-url": "amqp://mybroker-hdls-svc.test.svc:61616", + }, binding.ApplicationProperties) } func TestArkMQBrokerDirect(t *testing.T) { @@ -309,7 +321,11 @@ func TestArkMQBrokerDirect(t *testing.T) { }, endpoint) require.NoError(t, err) assert.NotNil(t, binding) - assert.Equal(t, "jms:queue:orders?brokerURL=tcp%3A%2F%2Fmybroker-hdls-svc.test.svc%3A61616", binding.URI) + assert.Equal(t, "amqp:queue:orders", binding.URI) + assert.Equal(t, map[string]string{ + "quarkus.qpid-jms.url": "amqp://mybroker-hdls-svc.test.svc:61616", + "camel.component.amqp.broker-url": "amqp://mybroker-hdls-svc.test.svc:61616", + }, binding.ApplicationProperties) // 2. With "queue" property fallback endpointQueue := camelv1.Endpoint{ @@ -327,7 +343,11 @@ func TestArkMQBrokerDirect(t *testing.T) { }, endpointQueue) require.NoError(t, err) assert.NotNil(t, binding) - assert.Equal(t, "jms:queue:invoices?brokerURL=tcp%3A%2F%2Fmybroker-hdls-svc.test.svc%3A61616", binding.URI) + assert.Equal(t, "amqp:queue:invoices", binding.URI) + assert.Equal(t, map[string]string{ + "quarkus.qpid-jms.url": "amqp://mybroker-hdls-svc.test.svc:61616", + "camel.component.amqp.broker-url": "amqp://mybroker-hdls-svc.test.svc:61616", + }, binding.ApplicationProperties) // 3. Missing destination/queue property returns error endpointMissing := camelv1.Endpoint{ @@ -524,7 +544,11 @@ func TestArkMQBrokerFallback(t *testing.T) { }, endpoint) require.NoError(t, err) assert.NotNil(t, binding) - assert.Equal(t, "jms:queue:my-fallback-queue?brokerURL=tcp%3A%2F%2Fdefault-broker-hdls-svc.test.svc%3A61616", binding.URI) + assert.Equal(t, "amqp:queue:my-fallback-queue", binding.URI) + assert.Equal(t, map[string]string{ + "quarkus.qpid-jms.url": "amqp://default-broker-hdls-svc.test.svc:61616", + "camel.component.amqp.broker-url": "amqp://default-broker-hdls-svc.test.svc:61616", + }, binding.ApplicationProperties) } func TestArkMQMultipleBrokersFallback(t *testing.T) {
