From 974c481b1d5d56720e4d003952ea85fcb716d380 Mon Sep 17 00:00:00 2001 From: "gaoran_10@126.com" Date: Thu, 9 Jul 2026 18:00:49 +0800 Subject: [PATCH 1/3] Fix deliveryTag not found issue --- .../streamnative/pulsar/handlers/amqp/AmqpChannel.java | 5 +++-- .../pulsar/handlers/amqp/AmqpPulsarConsumer.java | 10 ++++++---- .../pulsar/handlers/amqp/MessageFetchContext.java | 8 +++++--- 3 files changed, 14 insertions(+), 9 deletions(-) diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpChannel.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpChannel.java index 15ad4ec8a..395561d44 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpChannel.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpChannel.java @@ -124,8 +124,9 @@ public class AmqpChannel implements ServerChannelMethodProcessor { /** * The delivery tag is unique per channel. This is pre-incremented before putting into the deliver frame so that * value of this represents the last tag sent out. + * AtomicLong is required because multi-bundle mode may deliver from multiple consumers concurrently. */ - protected volatile long deliveryTag = 0; + protected final AtomicLong deliveryTag = new AtomicLong(0); protected final AmqpFlowCreditManager creditManager; protected final AtomicBoolean blockedOnCredit = new AtomicBoolean(false); public static final int DEFAULT_CONSUMER_PERMIT = 1000; @@ -905,7 +906,7 @@ public void closeChannel(int cause, final String message) { } public long getNextDeliveryTag() { - return ++deliveryTag; + return deliveryTag.incrementAndGet(); } public AmqpConnection getConnection() { diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpPulsarConsumer.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpPulsarConsumer.java index 5cb7087bc..280fc43b0 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpPulsarConsumer.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpPulsarConsumer.java @@ -75,6 +75,12 @@ private void consume() { MessageIdImpl messageId = (MessageIdImpl) message.getMessageId(); long deliveryIndex = this.amqpChannel.getNextDeliveryTag(); + // Register unacked before deliver so a fast client ack cannot miss the tag. + if (!this.autoAck) { + this.amqpChannel.getUnacknowledgedMessageMap().add( + deliveryIndex, PositionFactory.create(messageId.getLedgerId(), messageId.getEntryId()), + AmqpPulsarConsumer.this, message.size()); + } this.amqpChannel.getConnection().getAmqpOutputConverter().writeDeliver( MessageConvertUtils.messageToAmqpBody(message), this.amqpChannel.getChannelId(), @@ -87,10 +93,6 @@ private void consume() { messageId, consumer.getTopic(), t); return null; }); - } else { - this.amqpChannel.getUnacknowledgedMessageMap().add( - deliveryIndex, PositionFactory.create(messageId.getLedgerId(), messageId.getEntryId()), - AmqpPulsarConsumer.this, message.size()); } consumeBackoff.reset(); this.consume(); diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/MessageFetchContext.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/MessageFetchContext.java index dd0bf32a7..a782c5c19 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/MessageFetchContext.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/MessageFetchContext.java @@ -71,13 +71,15 @@ public static void handleFetch(AmqpChannel channel, AmqpConsumer consumer, boole long deliveryTag = channel.getNextDeliveryTag(); boolean isRedelivery = consumer.getRedeliveryTracker() .getRedeliveryCount(pair.getLeft().getLedgerId(), pair.getLeft().getEntryId()) > 0; + // Register unacked before get-ok so a fast client ack cannot miss the tag. + if (!autoAck) { + channel.getUnacknowledgedMessageMap().add(deliveryTag, pair.getLeft(), consumer, 0); + channel.getCreditManager().useCreditForMessages(1, 0); + } channel.getConnection().getAmqpOutputConverter().writeGetOk(pair.getRight(), channel.getChannelId(), isRedelivery, deliveryTag, 0); if (autoAck) { consumer.messageAck(pair.getLeft()); - } else { - channel.getUnacknowledgedMessageMap().add(deliveryTag, pair.getLeft(), consumer, 0); - channel.getCreditManager().useCreditForMessages(1, 0); } } else { if (pair != null && pair.getLeft() != null) { From 8ff609f4f3a0129b478f9bc9a5e3831c39da67e3 Mon Sep 17 00:00:00 2001 From: "gaoran_10@126.com" Date: Fri, 10 Jul 2026 03:58:44 +0800 Subject: [PATCH 2/3] Add the AMQP admin port in the protocol data --- .../handlers/amqp/AmqpBrokerService.java | 25 ++++++++ .../handlers/amqp/AmqpProtocolHandler.java | 57 ++++++++++++++++++- .../amqp/admin/impl/BaseResources.java | 56 ++++++++++++++++-- .../proxy/PulsarServiceLookupHandler.java | 2 +- .../amqp/test/AmqpProtocolDataTest.java | 50 ++++++++++++++++ 5 files changed, 181 insertions(+), 9 deletions(-) create mode 100644 amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/AmqpProtocolDataTest.java diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpBrokerService.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpBrokerService.java index 19cfac8e3..a6a8a79b7 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpBrokerService.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpBrokerService.java @@ -16,15 +16,19 @@ import io.netty.util.concurrent.DefaultThreadFactory; import io.streamnative.pulsar.handlers.amqp.admin.AmqpAdmin; +import java.io.IOException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import lombok.Getter; +import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.authentication.AuthenticationService; +import org.apache.pulsar.broker.resources.MetadataStoreCacheLoader; /** * AMQP broker related. */ +@Slf4j public class AmqpBrokerService { @Getter private AmqpTopicManager amqpTopicManager; @@ -42,6 +46,8 @@ public class AmqpBrokerService { private PulsarService pulsarService; @Getter private AmqpAdmin amqpAdmin; + @Getter + private MetadataStoreCacheLoader metadataStoreCacheLoader; public AmqpBrokerService(PulsarService pulsarService, AmqpServiceConfiguration config) { this.pulsarService = pulsarService; @@ -53,6 +59,15 @@ public AmqpBrokerService(PulsarService pulsarService, AmqpServiceConfiguration c this.queueService = new QueueServiceImpl(exchangeContainer, queueContainer); this.connectionContainer = new ConnectionContainer(pulsarService, exchangeContainer, queueContainer); this.amqpAdmin = new AmqpAdmin("localhost", config.getAmqpAdminPort()); + try { + // Used by admin ownership redirects to resolve the owner broker's amqpAdminPort. + this.metadataStoreCacheLoader = new MetadataStoreCacheLoader(pulsarService.getPulsarResources(), + 30_000); + } catch (Exception e) { + // Unit tests may mock PulsarResources without load-report store; keep service usable. + log.warn("Failed to init MetadataStoreCacheLoader for AoP, admin redirects may fail", e); + this.metadataStoreCacheLoader = null; + } } private ExecutorService initRouteExecutor(AmqpServiceConfiguration config) { @@ -67,4 +82,14 @@ public boolean isAuthenticationEnabled() { public AuthenticationService getAuthenticationService() { return pulsarService.getBrokerService().getAuthenticationService(); } + + public void close() { + if (metadataStoreCacheLoader != null) { + try { + metadataStoreCacheLoader.close(); + } catch (IOException e) { + log.warn("Failed to close MetadataStoreCacheLoader", e); + } + } + } } diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpProtocolHandler.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpProtocolHandler.java index 095a5ea72..9d23b647d 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpProtocolHandler.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpProtocolHandler.java @@ -24,6 +24,7 @@ import io.streamnative.pulsar.handlers.amqp.utils.ConfigurationUtils; import java.net.InetSocketAddress; import java.util.Map; +import java.util.Optional; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.ServiceConfiguration; @@ -46,6 +47,7 @@ public class AmqpProtocolHandler implements ProtocolHandler { public static final String SSL_PREFIX = "amqp+ssl://"; public static final String PLAINTEXT_PREFIX = "amqp://"; public static final String LISTENER_DEL = ","; + public static final String PROTOCOL_DATA_SEP = "|"; public static final String LISTENER_PATTEN = "^(amqp)://[-a-zA-Z0-9+&@#/%?=~_|!:,.;]*[-0-9+]"; @Getter @@ -84,10 +86,13 @@ public void initialize(ServiceConfiguration conf) throws Exception { // This method is called after initialize @Override public String getProtocolDataToAdvertise() { + // Format: | + // Admin port is required for multi-broker ownership redirects of AoP admin REST calls. + String protocolData = getAppliedAmqpListeners(amqpConfig) + PROTOCOL_DATA_SEP + amqpConfig.getAmqpAdminPort(); if (log.isDebugEnabled()) { - log.debug("Get configured listeners: {}", getAppliedAmqpListeners(amqpConfig)); + log.debug("Get protocol data to advertise: {}", protocolData); } - return getAppliedAmqpListeners(amqpConfig); + return protocolData; } @Override @@ -182,10 +187,15 @@ public Map> newChannelIniti @Override public void close() { try { - webServer.stop(); + if (webServer != null) { + webServer.stop(); + } } catch (Exception e) { log.error("Failed to stop web server for aop", e); } + if (amqpBrokerService != null) { + amqpBrokerService.close(); + } } public static int getListenerPort(String listener) { @@ -207,4 +217,45 @@ public static String getAppliedAmqpListeners(AmqpServiceConfiguration configurat public static String amqpUrl(String host, int port) { return String.format("amqp://%s:%d", host, port); } + + /** + * Extract AMQP listeners from protocol advertise data. + * Compatible with both legacy format (`amqp://host:port`) and + * new format (`amqp://host:port|adminPort`). + */ + public static String extractAmqpListeners(String protocolData) { + if (protocolData == null) { + return null; + } + int sep = protocolData.lastIndexOf(PROTOCOL_DATA_SEP); + if (sep < 0) { + return protocolData; + } + String maybeAdminPort = protocolData.substring(sep + 1); + try { + Integer.parseInt(maybeAdminPort); + return protocolData.substring(0, sep); + } catch (NumberFormatException e) { + return protocolData; + } + } + + /** + * Extract AoP admin port from protocol advertise data. + * Returns empty if the data uses the legacy format without admin port. + */ + public static Optional extractAmqpAdminPort(String protocolData) { + if (protocolData == null) { + return Optional.empty(); + } + int sep = protocolData.lastIndexOf(PROTOCOL_DATA_SEP); + if (sep < 0) { + return Optional.empty(); + } + try { + return Optional.of(Integer.parseInt(protocolData.substring(sep + 1))); + } catch (NumberFormatException e) { + return Optional.empty(); + } + } } diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/admin/impl/BaseResources.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/admin/impl/BaseResources.java index 2522972d9..8580a23f2 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/admin/impl/BaseResources.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/admin/impl/BaseResources.java @@ -30,19 +30,25 @@ import java.net.URI; import java.util.ArrayList; import java.util.List; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.lookup.LookupResult; import org.apache.pulsar.broker.namespace.LookupOptions; import org.apache.pulsar.broker.namespace.NamespaceService; +import org.apache.pulsar.broker.resources.MetadataStoreCacheLoader; import org.apache.pulsar.broker.resources.NamespaceResources; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.common.lookup.data.LookupData; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.metadata.api.MetadataStoreException; +import org.apache.pulsar.policies.data.loadbalancer.LoadManagerReport; /** * Base resources. @@ -161,18 +167,23 @@ protected CompletableFuture validateTopicOwnershipAsync(TopicName topicNam "Failed to find ownership for topic:" + topicName); } return lookupResult.get(); - }).thenCompose(webUrl -> nsService.isServiceUnitOwnedAsync(topicName) - .thenApply(isTopicOwned -> Pair.of(webUrl, isTopicOwned)) + }).thenCompose(lookupResult -> nsService.isServiceUnitOwnedAsync(topicName) + .thenApply(isTopicOwned -> Pair.of(lookupResult, isTopicOwned)) ).thenAccept(pair -> { - URI webUri = pair.getLeft().toLookupRedirectUri(uri.getRequestUri()); + LookupResult lookupResult = pair.getLeft(); + URI webUri = lookupResult.toLookupRedirectUri(uri.getRequestUri()); boolean isTopicOwned = pair.getRight(); if (!isTopicOwned) { boolean newAuthoritative = isLeaderBroker(pulsar()); - // Replace the host and port of the current request and redirect + int adminPort = resolveOwnerAmqpAdminPort(lookupResult) + .orElseThrow(() -> new RestException(Response.Status.PRECONDITION_FAILED, + "Failed to resolve amqp admin port for topic:" + topicName)); + // Redirect to the owner broker's AoP admin endpoint. + // Host comes from lookup; port must be the owner admin port (not local). URI redirect = UriBuilder.fromUri(uri.getRequestUri()) .host(webUri.getHost()) - .port(aop().getAmqpConfig().getAmqpAdminPort()) + .port(adminPort) .replaceQueryParam("authoritative", newAuthoritative) .build(); // Redirect @@ -197,6 +208,41 @@ protected CompletableFuture validateTopicOwnershipAsync(TopicName topicNam }); } + private Optional resolveOwnerAmqpAdminPort(LookupResult lookupResult) { + LookupData lookupData = lookupResult.getLookupData(); + MetadataStoreCacheLoader cacheLoader = aop().getAmqpBrokerService().getMetadataStoreCacheLoader(); + if (cacheLoader == null) { + log.warn("MetadataStoreCacheLoader is unavailable, cannot resolve owner amqp admin port"); + return Optional.empty(); + } + List brokers = cacheLoader.getAvailableBrokers(); + Optional owner = brokers.stream() + .filter(report -> matchesOwnerBroker(report, lookupData)) + .findFirst(); + if (owner.isEmpty()) { + log.warn("Unable to locate load report for owner broker. httpUrl={}, brokerUrl={}, available={}", + lookupData.getHttpUrl(), lookupData.getBrokerUrl(), brokers.size()); + return Optional.empty(); + } + Optional protocolData = owner.get().getProtocol(AmqpProtocolHandler.PROTOCOL_NAME); + if (protocolData.isEmpty()) { + log.warn("Owner broker has no amqp protocol data. webServiceUrl={}", owner.get().getWebServiceUrl()); + return Optional.empty(); + } + Optional adminPort = AmqpProtocolHandler.extractAmqpAdminPort(protocolData.get()); + if (adminPort.isEmpty()) { + log.warn("Owner broker amqp protocol data has no admin port: {}", protocolData.get()); + } + return adminPort; + } + + private static boolean matchesOwnerBroker(LoadManagerReport report, LookupData lookupData) { + return StringUtils.equals(report.getWebServiceUrl(), lookupData.getHttpUrl()) + || StringUtils.equals(report.getWebServiceUrlTls(), lookupData.getHttpUrlTls()) + || StringUtils.equals(report.getPulsarServiceUrl(), lookupData.getBrokerUrl()) + || StringUtils.equals(report.getPulsarServiceUrlTls(), lookupData.getBrokerUrlTls()); + } + protected static boolean isLeaderBroker(PulsarService pulsar) { return pulsar.getLeaderElectionService().isLeader(); } diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/proxy/PulsarServiceLookupHandler.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/proxy/PulsarServiceLookupHandler.java index fcea60569..86a342438 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/proxy/PulsarServiceLookupHandler.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/proxy/PulsarServiceLookupHandler.java @@ -84,7 +84,7 @@ public CompletableFuture> findBroker(TopicName topicName, return; } - String amqpBrokerAddress = protocolData.get(); + String amqpBrokerAddress = AmqpProtocolHandler.extractAmqpListeners(protocolData.get()); if (!StringUtils.startsWith(amqpBrokerAddress, AmqpProtocolHandler.PLAINTEXT_PREFIX) && !StringUtils.startsWith(amqpBrokerAddress, AmqpProtocolHandler.SSL_PREFIX)) { amqpBrokerAddress = AmqpProtocolHandler.PLAINTEXT_PREFIX + amqpBrokerAddress; diff --git a/amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/AmqpProtocolDataTest.java b/amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/AmqpProtocolDataTest.java new file mode 100644 index 000000000..43e4db8d4 --- /dev/null +++ b/amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/AmqpProtocolDataTest.java @@ -0,0 +1,50 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.amqp.test; + +import io.streamnative.pulsar.handlers.amqp.AmqpProtocolHandler; +import io.streamnative.pulsar.handlers.amqp.AmqpServiceConfiguration; +import org.testng.Assert; +import org.testng.annotations.Test; + +/** + * Protocol advertise data format tests. + */ +public class AmqpProtocolDataTest { + + @Test + public void testAdvertiseDataIncludesAdminPort() { + AmqpServiceConfiguration conf = new AmqpServiceConfiguration(); + conf.setAmqpListeners("amqp://127.0.0.1:5672"); + conf.setAmqpAdminPort(15673); + + AmqpProtocolHandler handler = new AmqpProtocolHandler(); + try { + handler.initialize(conf); + } catch (Exception e) { + throw new RuntimeException(e); + } + String protocolData = handler.getProtocolDataToAdvertise(); + Assert.assertEquals(protocolData, "amqp://127.0.0.1:5672|15673"); + Assert.assertEquals(AmqpProtocolHandler.extractAmqpListeners(protocolData), "amqp://127.0.0.1:5672"); + Assert.assertEquals(AmqpProtocolHandler.extractAmqpAdminPort(protocolData).orElse(-1).intValue(), 15673); + } + + @Test + public void testLegacyProtocolDataCompatible() { + String legacy = "amqp://127.0.0.1:5672"; + Assert.assertEquals(AmqpProtocolHandler.extractAmqpListeners(legacy), legacy); + Assert.assertFalse(AmqpProtocolHandler.extractAmqpAdminPort(legacy).isPresent()); + } +} From 4d7c78e08c1aee44f35bbef167bbd6fa0f00350d Mon Sep 17 00:00:00 2001 From: "gaoran_10@126.com" Date: Fri, 10 Jul 2026 04:23:42 +0800 Subject: [PATCH 3/3] fix CI --- .github/workflows/pr-test.yml | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/.github/workflows/pr-test.yml b/.github/workflows/pr-test.yml index 28890fcef..71fb3c43e 100644 --- a/.github/workflows/pr-test.yml +++ b/.github/workflows/pr-test.yml @@ -2,13 +2,17 @@ name: aop mvn build check and ut on: pull_request: + # PR target branches (e.g. master, branch-4.0) branches: - master - branch-* push: + # Only maintenance / main lines. Exclude PR head branches like branch-4.0.10.2 + # so "branch-4.0.10.2 -> branch-4.0" PRs are not double-triggered by push. branches: - master - branch-* + - '!branch-*.*.*' jobs: build: