Skip to content

Fix #5935: Add ArkMQ operator queue Pipe binding - #6845

Open
Thundercloud12 wants to merge 18 commits into
apache:mainfrom
Thundercloud12:feature/5935-arkmq-pipe-binding
Open

Thundercloud12 wants to merge 18 commits into
apache:mainfrom
Thundercloud12:feature/5935-arkmq-pipe-binding

Conversation

@Thundercloud12

@Thundercloud12 Thundercloud12 commented Sep 23, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #6837

Motivation & Context

Camel K Pipe resources currently support Strimzi Kafka, Knative, and Kamelet resources as declarative endpoints. However, there was no native binding provider for messaging queues managed by the ArkMQ / ActiveMQ Artemis operator (broker.amq.io).

This PR adds a minimal, queue-only Pipe binding provider for ArkMQ operator-managed queues, following the architectural pattern of the existing Strimzi binding.


Architecture & Design Decisions

Following maintainer guidance on the issue, this implementation focuses strictly on a minimal queue use case:

ActiveMQArtemisAddress CR
(spec.routingType: anycast)
        ↓
Camel K duck type
        ↓
ArkMQ binding provider
        ↓
JMS queue endpoint
(jms:queue:<queueName>?brokerURL=...)

1. Queue-Only Resource Target

  • Binds exclusively to ActiveMQArtemisAddress (broker.amq.io/v1beta1).
  • Rejects non-address kinds with an informative error.
  • Actively rejects spec.routingType: multicast (topics) to ensure only queue bindings are configured.

2. Automated Broker Discovery

If brokerURL is not manually provided in the Pipe endpoint properties, the provider:

  • Resolves the referenced ActiveMQArtemis cluster from spec.applyTo or the ActiveMQArtemis label.
  • Discovers the broker port from status (core, all, or openwire), defaulting to the standard 61616.
  • Resolves the in-cluster headless service:
tcp://<cluster>-hdls-svc.<namespace>.svc:<port>

3. Queue Name Resolution

The queue name is resolved using the following priority:

spec.queueName
    ↓
spec.addressName
    ↓
metadata.name

4. Deterministic Component URI

The provider always generates:

jms:queue:<queueName>?brokerURL=...

Any user-supplied Pipe endpoint properties are appended to the generated URI.


Key Changes

  • Duck Types

    • Added minimal ActiveMQArtemisAddress and ActiveMQArtemis duck types in pkg/apis/duck/arkmq/v1beta1/.
    • Registered the duck types with the scheme.
  • Client Generation

    • Added client-gen configuration in script/gen_client.sh.
    • Generated typed clientsets under pkg/client/arkmq/.
  • Binding Provider

    • Implemented ArkMQBindingProvider in pkg/util/bindings/arkmq.go.
  • RBAC

    • Added namespaced and descoped roles and bindings under pkg/resources/config/rbac/.

    • Grants read-only (get, list, watch) permissions for:

      • activemqartemises
      • activemqartemisaddresses
  • Unit Tests

    • Added pkg/util/bindings/arkmq_test.go covering:

      • Direct brokerURL property override
      • Cluster headless service resolution
      • Address-by-name fallback
      • Multicast topic rejection
      • Unsupported resource kinds
      • Pass-through behavior
    • Added pkg/controller/pipe/initialize_test.go covering bidirectional support:

      • Sink: Timer → ArkMQ Queue
      • Source: ArkMQ Queue → Log
  • Documentation

    • Added ArkMQ queue Pipe documentation and a sample to:
      docs/modules/ROOT/pages/pipes/pipes.adoc
  • E2E Infrastructure

    • Added the end-to-end test suite under e2e/arkmq/.
    • Added the test-arkmq Makefile target.
    • Added the GitHub Actions workflow:
      .github/workflows/arkmq.yml
    • Uses Gomega Eventually for asynchronous assertions.

How Has This Been Tested?

The following checks were run successfully:

Unit Tests

go test -v ./pkg/util/bindings -run TestArkMQ
go test -v ./pkg/controller/pipe -run TestNewPipeArkMQ

Static Analysis

go vet ./...

go vet passes across all packages with zero issues.

Formatting

make fmt goimport

Both formatting and import checks pass successfully.

@squakez squakez left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice work. However I think we need to consider the following thread: arkmq-org/arkmq-org-broker-operator#1121 - which it seems to point to a deprecation of certain resources used here. Maybe we can only use the newer approach instead and making it easier to maintain in the long term.

Also, when complete, we should add a short documentation as provided for the other bindings in docs/modules/ROOT/pages/pipes/pipes.adoc#more-advanced-examples

Comment thread e2e/arkmq/setup/setup.sh Outdated
Comment thread e2e/arkmq/setup/setup.sh
Comment thread pkg/util/bindings/arkmq.go Outdated
Comment thread pkg/util/bindings/arkmq.go Outdated
Comment thread pkg/util/bindings/arkmq.go Outdated
@squakez

squakez commented Sep 23, 2026

Copy link
Copy Markdown
Contributor

BTW, also mind #6846 - while reviewing the PR I realized the Strimzi binding which you may have used as a reference is flawed. We can avoid that problem reported there if possible.

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez so the issues you mentioned are resolved, actually i wanted to follow kafka to the point so ended up following that doc( was recommended by llm when i said i am following kafka implementation should have atleast checked that), also after this i am planning to raise a pr to fix issue #6846. Thank you!

@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 63.9% (+0.1%)

@squakez squakez left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM, though there are a few points to clarify.

Comment thread pkg/util/bindings/arkmq.go Outdated
Comment thread pkg/util/bindings/arkmq.go Outdated
Comment thread pkg/util/bindings/arkmq.go Outdated
Comment thread script/gen_client.sh
@squakez

squakez commented Sep 24, 2026

Copy link
Copy Markdown
Contributor

after this i am planning to raise a pr to fix issue #6846

There may be already somebody else working on it. Please, verify the assignment is clear before starting to avoid duplicate works and add a comment to let some maintainer assign it as well if it's free.

@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 64% (+0.2%)

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez yeah at the time I checked it wasn't assigned I can see someone is working on it, I wouldn't pick that up

@Thundercloud12

Thundercloud12 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

@squakez all the changes you mentioned are reconcile that rm -rf was a leftover line during previous refactoring that part you raised about the search that was a minconception because in an activemqartemisaddress does not match underlying queue address name. But while researching for it i got to know that kubernetes semantics and camel k conventions would never allow it so i have ensure that too, also a lint ci was failing, fixes that too

Comment thread pkg/util/bindings/arkmq.go Outdated
@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 64% (+0.2%)

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez please have a look at the updated code

@squakez squakez left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think it's more or less fine. Now it's a matter to verify the proper configuration that Camel expects. I think this step requires some manual testing against an existing ArkMQ broker in order to find the proper configuration params.

@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 64.1% (+0.3%)

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

ArkMQ / ActiveMQ Artemis AMQP Verification

Issue: #5935 – Support ArkMQ in Pipe Binding
Date: 2026-09-26

Objective

Verify the AMQP configuration used by the ArkMQ Pipe binding against an ActiveMQ Artemis broker.

Verification

Tested with:

  • Camel: 4.22.1 LTS
  • Qpid JMS: 2.11.0
  • ActiveMQ Artemis: 2.44.0
  • JDK: 21

An Artemis broker was started with both AMQP acceptors:

  • 61616 — multi-protocol (AMQP, CORE, MQTT, etc.)
  • 5672 — dedicated AMQP

Both were tested using:

amqp://localhost:<port>

with the Camel endpoint:

amqp:queue:<destination>

Results

Port 61616: SUCCESS
Camel connected, created/used the queue, produced a message, and consumed it successfully.

Port 5672: SUCCESS
Identical behavior; AMQP 1.0 connection and message flow worked successfully.

Configuration

For Camel K, the following properties were verified:

quarkus.qpid-jms.url=amqp://<host>:<port>
camel.component.amqp.broker-url=amqp://<host>:<port>

The Pipe binding can therefore continue using:

amqp:queue:<destination>

without additional custom connection/bean configuration.

Conclusion

The current ArkMQ binding configuration is compatible with ActiveMQ Artemis AMQP.

  • 61616 supports AMQP 1.0 despite being a multi-protocol port.
  • 5672 also works as the dedicated AMQP port.
  • No changes to the current binding translation are required based on these tests.

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez the above verification was done using docker because as you mentioned you need to check the correct config params and this does it well, is it okay or should i look towards setting up a cluster and checking?

@squakez

squakez commented Sep 26, 2026

Copy link
Copy Markdown
Contributor

The only failing test is the new one that should prove the feature. This is what we get:

namespace/arkmq created
error: unable to read URL "https://github.com/arkmq-org/arkmq-org-broker-operator/releases/latest/download/arkmq-org-broker-operator.yaml", server reported 404 Not Found, status code=404
Error from server (NotFound): deployments.apps "arkmq-org-broker-operator" not found
Error from server (NotFound): customresourcedefinitions.apiextensions.k8s.io "activemqartemises.broker.amq.io" not found
Error from server (NotFound): customresourcedefinitions.apiextensions.k8s.io "activemqartemisaddresses.broker.amq.io" not found
error: resource mapping not found for name: "my-broker" namespace: "" from "/home/runner/work/camel-k/camel-k/e2e/arkmq/setup/broker.yaml": no matches for kind "ActiveMQArtemis" in version "broker.amq.io/v1beta1"
ensure CRDs are installed first
error: the server doesn't have a resource type "activemqartemis"
error: resource mapping not found for name: "my-queue" namespace: "" from "/home/runner/work/camel-k/camel-k/e2e/arkmq/setup/queue.yaml": no matches for kind "ActiveMQArtemisAddress" in version "broker.amq.io/v1beta1"
ensure CRDs are installed first
error: the server doesn't have a resource type "activemqartemisaddress"

As suggested during the review, let's use https://github.com/arkmq-org/arkmq-org-broker-operator/releases/download/v2.2.2/activemq-artemis-operator.yaml instead of latest and let's add a parameter so we can control each release update easily.

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez Updated e2e/arkmq/setup/setup.sh to pin and parameterize the ArkMQ operator version

@squakez squakez left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Starting the checks

@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 64.1% (+0.3%)

@squakez squakez left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 64.1% (+0.3%)

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez so it has been frustrating on me and i think on you as well i think i was lax in my work earlier i thought fixxinf this one failing test issue would solve the pipeline but it was not the case i should have been more thorough with my methods so to actually confirm i ran tests in the below manner:

  1. Installed local tooling: Installed kubectl (v1.31.1) and kind (v0.24.0).
  2. Created a live cluster: Spun up a real Kubernetes cluster container (kind-camel-k) with control plane in Ready
    state.
  3. Executed setup.sh directly against the live cluster:
    • Operator deployment rolled out and reached 1/1 Running.
    • CRDs established.
    • Broker StatefulSet my-broker-ss-0 pulled images, initialized, and reached 1/1 Running.
    • Headless service my-broker-hdls-svc was created with port 61616.
    • Broker CR reached Ready: True.
    • Address CR my-queue created cleanly.
    • Script finished with exit code 0 without any timeout.
  4. Inspected live cluster state:
    pod/my-broker-ss-0 1/1 Running
    service/my-broker-hdls-svc 61616/TCP
    activemqartemis.broker.amq.io/my-broker Ready: True
    activemqartemisaddress.broker.amq.io/my-queue Created

only aftyer this was verified i have pushed the change if you think the testing methodology should be changed please let me know because that would avoid us to go into again this ci failing loop

@github-actions

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 63.8% to 64.1% (+0.3%)

@squakez

squakez commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

The new worflow runs, but we have some issue to fix:

{"level":"error","ts":"2026-09-29T07:05:32Z","logger":"controller-runtime.cache.UnhandledError","msg":"Failed to watch","reflector":"k8s.io/client-go@v0.37.0/tools/cache/reflector.go:343","type":"*v1beta1.ActiveMQArtemisAddress","error":"activemqartemisaddresses.broker.amq.io is forbidden: User \"system:serviceaccount:camel-k:camel-k-operator\" cannot watch resource \"activemqartemisaddresses\" in API group \"broker.amq.io\" at the cluster scope"

-> we must include watch and list rbac to the resource

on the queue side:

2026-09-29 07:05:04,898 WARN  [org.apache.activemq.artemis.core.client] AMQ212037: CORE connection failure to 10.244.0.7:33796 has been detected: Frame size exceeded: 1095586128 [code=GENERIC_EXCEPTION]

which indicates that maybe we need to configure the protocol of the acceptor. I wonder if we should try to configure jms on the camel conversion side instead of amqp. However, that would require to understand how we can include a connection factory via camel properties (which was the main reason to try amqp component).

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez please check the changes, let me know for any changes!

@github-actions

github-actions Bot commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 64.4% to 64.7% (+0.3%)

@squakez

squakez commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

You may need to rebase with main and remove any registry action previously used as we're now running a dev mode with an internal embedded registry

- Support binding directly to ActiveMQArtemis cluster with destination/queue property
    - Avoid reliance on deprecated ActiveMQArtemisAddress while keeping backward compatibility
    - Update unit tests and pipes.adoc documentation
@Thundercloud12
Thundercloud12 force-pushed the feature/5935-arkmq-pipe-binding branch from 3b77bac to bb089e9 Compare October 2, 2026 16:46
@github-actions

github-actions Bot commented Oct 3, 2026

Copy link
Copy Markdown
Contributor

✔️ Unit test coverage report - coverage increased from 64.4% to 64.7% (+0.3%)

@Thundercloud12

Copy link
Copy Markdown
Contributor Author

@squakez rebase with main!

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Provide Artemis Queue Pipe binding support

2 participants