From 41ba05f5daa08980840ac959658fc7d3d61edeef Mon Sep 17 00:00:00 2001 From: Julien Viet Date: Fri, 24 Jul 2026 09:01:29 +0200 Subject: [PATCH] Support event-bus interceptor chain failure. Motivation: The event-bus interceptor chain does not provide a clean way to fail a message interception, some case can be achieved by failing the message itself but that works for limited cases and it does not handle the interception in progress. Cleanly failing a message delivery can be useful in some scenario, specially for testing purposes. Changes: Add a DeliveryContext#fail method that performs a clean fail of the message by interrupting the interception chain and properly failing the message: The message will not make progress anymore and the sender will receive a failure notification whether it is the write future for outbound message writes or the reply handler. --- .../vertx/core/eventbus/DeliveryContext.java | 15 +++ .../eventbus/impl/DeliveryContextImpl.java | 28 ++++- .../eventbus/impl/HandlerRegistration.java | 8 +- .../vertx/core/eventbus/impl/SendContext.java | 43 +++++--- .../eventbus/EventBusInterceptorTest.java | 103 +++++++++++++++++- 5 files changed, 174 insertions(+), 23 deletions(-) diff --git a/vertx-core/src/main/java/io/vertx/core/eventbus/DeliveryContext.java b/vertx-core/src/main/java/io/vertx/core/eventbus/DeliveryContext.java index b3e49edc02f..076e3193ab6 100644 --- a/vertx-core/src/main/java/io/vertx/core/eventbus/DeliveryContext.java +++ b/vertx-core/src/main/java/io/vertx/core/eventbus/DeliveryContext.java @@ -34,6 +34,21 @@ public interface DeliveryContext { */ void next(); + /** + * {@link #fail(ReplyFailure, int, String)} with {@link ReplyFailure#RECIPIENT_FAILURE} + */ + boolean fail(int failureCode, String message); + + /** + * Terminate the interception chain immediately and fail the message. + * + * @param failure the failure + * @param failureCode the failure code + * @param message the message + * @return true if the operation succeeded + */ + boolean fail(ReplyFailure failure, int failureCode, String message); + /** * @return true if the message is being sent (point to point) or False if the message is being published */ diff --git a/vertx-core/src/main/java/io/vertx/core/eventbus/impl/DeliveryContextImpl.java b/vertx-core/src/main/java/io/vertx/core/eventbus/impl/DeliveryContextImpl.java index 31867a92257..47f631a796d 100644 --- a/vertx-core/src/main/java/io/vertx/core/eventbus/impl/DeliveryContextImpl.java +++ b/vertx-core/src/main/java/io/vertx/core/eventbus/impl/DeliveryContextImpl.java @@ -13,6 +13,8 @@ import io.vertx.core.Handler; import io.vertx.core.eventbus.DeliveryContext; import io.vertx.core.eventbus.Message; +import io.vertx.core.eventbus.ReplyException; +import io.vertx.core.eventbus.ReplyFailure; import io.vertx.core.internal.ContextInternal; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; @@ -24,12 +26,12 @@ class DeliveryContextImpl implements DeliveryContext { private final MessageImpl message; private final ContextInternal context; private final Object body; - private final Runnable dispatch; + private final Handler dispatch; private final Handler>[] interceptors; private volatile int interceptorIdx; protected DeliveryContextImpl(MessageImpl message, Handler>[] interceptors, - ContextInternal context, Object body, Runnable dispatch) { + ContextInternal context, Object body, Handler dispatch) { this.message = message; this.interceptors = interceptors; this.context = context; @@ -53,6 +55,26 @@ public Object body() { return body; } + @Override + public boolean fail(int failureCode, String message) { + return fail(ReplyFailure.RECIPIENT_FAILURE, failureCode, message); + } + + @Override + public boolean fail(ReplyFailure failure, int failureCode, String message) { + while (true) { + int idx = UPDATER.get(this); + if (idx > interceptors.length) { + return false; + } + if (UPDATER.compareAndSet(this, idx, interceptors.length + 1)) { + break; + } + } + dispatch.handle(new ReplyException(failure, failureCode, message)); + return true; + } + @Override public void next() { int idx = UPDATER.getAndIncrement(this); @@ -68,7 +90,7 @@ public void next() { } } } else if (idx == interceptors.length) { - dispatch.run(); + dispatch.handle(null); } else { throw new IllegalStateException(); } diff --git a/vertx-core/src/main/java/io/vertx/core/eventbus/impl/HandlerRegistration.java b/vertx-core/src/main/java/io/vertx/core/eventbus/impl/HandlerRegistration.java index 30458e942e6..44cae57a5bb 100644 --- a/vertx-core/src/main/java/io/vertx/core/eventbus/impl/HandlerRegistration.java +++ b/vertx-core/src/main/java/io/vertx/core/eventbus/impl/HandlerRegistration.java @@ -95,7 +95,13 @@ public Future unregister() { void dispatchMessage(Function, Future> processor, MessageImpl message, ContextInternal context) { Handler>[] interceptors = message.bus.inboundInterceptors(); if (interceptors.length > 0) { - Runnable dispatch = () -> dispatch(context, message, processor); + Handler dispatch = failure -> { + if (failure == null) { + dispatch(context, message, processor); + } else { + message.reply(failure); + } + }; DeliveryContextImpl deliveryCtx = new DeliveryContextImpl<>(message, interceptors, context, message.receivedBody, dispatch); deliveryCtx.next(); } else { diff --git a/vertx-core/src/main/java/io/vertx/core/eventbus/impl/SendContext.java b/vertx-core/src/main/java/io/vertx/core/eventbus/impl/SendContext.java index 6caab55fb22..43e068b5455 100644 --- a/vertx-core/src/main/java/io/vertx/core/eventbus/impl/SendContext.java +++ b/vertx-core/src/main/java/io/vertx/core/eventbus/impl/SendContext.java @@ -49,10 +49,10 @@ public class SendContext implements Promise { void send() { Handler>[] interceptors = message.bus.outboundInterceptors(); if (interceptors.length > 0) { - DeliveryContextImpl deliveryContext = new DeliveryContextImpl<>(message, interceptors, ctx, message.sentBody, this::sendOrPub); + DeliveryContextImpl deliveryContext = new DeliveryContextImpl<>(message, interceptors, ctx, message.sentBody, this::sendOrPubOrFail); deliveryContext.next(); } else { - sendOrPub(); + sendOrPubOrFail(null); } } @@ -113,23 +113,32 @@ private void written(Throwable failure) { } } - private void sendOrPub() { - VertxTracer tracer = ctx.tracer(); - if (tracer != null) { - if (message.trace == null) { - src = true; - BiConsumer biConsumer = (String key, String val) -> message.headers().set(key, val); - TracingPolicy tracingPolicy = options.getTracingPolicy(); - if (tracingPolicy == null) { - tracingPolicy = TracingPolicy.PROPAGATE; + private void sendOrPubOrFail(ReplyException failure) { + if (failure == null) { + VertxTracer tracer = ctx.tracer(); + if (tracer != null) { + if (message.trace == null) { + src = true; + BiConsumer biConsumer = (String key, String val) -> message.headers().set(key, val); + TracingPolicy tracingPolicy = options.getTracingPolicy(); + if (tracingPolicy == null) { + tracingPolicy = TracingPolicy.PROPAGATE; + } + message.trace = tracer.sendRequest(ctx, SpanKind.RPC, tracingPolicy, message, message.send ? "send" : "publish", biConsumer, MessageTagExtractor.INSTANCE); + } else { + // Handle failure here + tracer.sendResponse(ctx, null, message.trace, null, TagExtractor.empty()); } - message.trace = tracer.sendRequest(ctx, SpanKind.RPC, tracingPolicy, message, message.send ? "send" : "publish", biConsumer, MessageTagExtractor.INSTANCE); - } else { - // Handle failure here - tracer.sendResponse(ctx, null, message.trace, null, TagExtractor.empty()); + } + bus.sendOrPub(this); + } else { + // Fail write promise + writePromise.fail(failure); + + // Fail reply handler when it exists + if (replyHandler != null) { + replyHandler.fail(failure); } } - bus.sendOrPub(this); } - } diff --git a/vertx-core/src/test/java/io/vertx/tests/eventbus/EventBusInterceptorTest.java b/vertx-core/src/test/java/io/vertx/tests/eventbus/EventBusInterceptorTest.java index 6f429c4dae9..48b4cc46747 100644 --- a/vertx-core/src/test/java/io/vertx/tests/eventbus/EventBusInterceptorTest.java +++ b/vertx-core/src/test/java/io/vertx/tests/eventbus/EventBusInterceptorTest.java @@ -12,10 +12,10 @@ package io.vertx.tests.eventbus; import io.vertx.core.Context; +import io.vertx.core.Future; import io.vertx.core.Handler; import io.vertx.core.Vertx; -import io.vertx.core.eventbus.DeliveryContext; -import io.vertx.core.eventbus.EventBus; +import io.vertx.core.eventbus.*; import io.vertx.test.core.VertxTestBase; import org.junit.Test; @@ -425,6 +425,105 @@ public void testInboundInterceptorFromNonVertxThreadFailure() { assertSame(expected, caught.get()); } + @Test + public void testOutboundReplySendFailure() { + eb.addOutboundInterceptor(sc -> { + assertTrue(sc.fail(3, "test")); + try { + sc.next(); + fail(); + } catch (IllegalStateException expected) { + } + }); + eb.consumer("some-address", msg -> { + fail(); + }); + MessageProducer sender = eb.sender("some-address"); + Future fut = sender.write("armadillo"); + fut.onComplete(onFailure(expected -> { + testComplete(); + })); + await(); + } + + @Test + public void testOutboundReplyRequestFailure() { + eb.addOutboundInterceptor(sc -> { + assertTrue(sc.fail(ReplyFailure.NO_HANDLERS, 3, "test")); + try { + sc.next(); + fail(); + } catch (IllegalStateException expected) { + } + }); + eb.consumer("some-address", msg -> { + fail(); + }); + Future> fut = eb.request("some-address", "armadillo"); + fut.onComplete(onFailure(expected -> { + ReplyException replyException = (ReplyException) expected; + assertEquals(3, replyException.failureCode()); + assertEquals(ReplyFailure.NO_HANDLERS, replyException.failureType()); + testComplete(); + })); + await(); + } + + @Test + public void testInboundReplySendFailure() { + eb.addInboundInterceptor(sc -> { + assertTrue(sc.fail(3, "test")); + try { + sc.next(); + fail(); + } catch (IllegalStateException expected) { + } + }); + eb.consumer("some-address", msg -> { + fail(); + }); + MessageProducer sender = eb.sender("some-address"); + Future fut = sender.write("armadillo"); + fut.onComplete(onSuccess(expected -> { + testComplete(); + })); + await(); + } + + @Test + public void testInboundReplyRequestFailure() { + eb.addInboundInterceptor(dc -> { + if (dc.message().address().equals("some-address")) { + assertTrue(dc.fail(3, "test")); + try { + dc.next(); + fail(); + } catch (IllegalStateException expected) { + } + } else { + dc.next(); + } + }); + eb.addInboundInterceptor(dc -> { + if (dc.message().address().equals("some-address")) { + fail(); + } else { + dc.next(); + } + }); + eb.consumer("some-address", msg -> { + fail(); + }); + Future> fut = eb.request("some-address", "armadillo"); + fut.onComplete(onFailure(expected -> { + ReplyException replyException = (ReplyException) expected; + assertEquals(3, replyException.failureCode()); + assertEquals(ReplyFailure.RECIPIENT_FAILURE, replyException.failureType()); + testComplete(); + })); + await(); + } + @Override public void setUp() throws Exception { super.setUp();