From e65847bd8a78ddb0f94d63e87b6cfabdf24a81f3 Mon Sep 17 00:00:00 2001 From: Anton Vinogradov Date: Tue, 4 Aug 2026 00:27:00 +0300 Subject: [PATCH] IGNITE-28901 Split GridEventStorageMessage into a request and a response Co-Authored-By: Claude Opus 5 --- .../ignite/internal/CoreMessagesProvider.java | 6 +- .../managers/communication/GridIoManager.java | 1 + .../eventstorage/GridEventStorageManager.java | 29 ++-- ...sage.java => GridEventStorageRequest.java} | 146 ++++-------------- .../GridEventStorageResponse.java | 76 +++++++++ .../resources/META-INF/classnames.properties | 1 - 6 files changed, 122 insertions(+), 137 deletions(-) rename modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/{GridEventStorageMessage.java => GridEventStorageRequest.java} (50%) create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageResponse.java diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java index 8b3fe7bbd4f00..cd03c5a2f2f6f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java @@ -42,7 +42,8 @@ import org.apache.ignite.internal.managers.encryption.MasterKeyChangeRequest; import org.apache.ignite.internal.managers.encryption.NodeEncryptionKeys; import org.apache.ignite.internal.managers.eventstorage.EventsDataBagItem; -import org.apache.ignite.internal.managers.eventstorage.GridEventStorageMessage; +import org.apache.ignite.internal.managers.eventstorage.GridEventStorageRequest; +import org.apache.ignite.internal.managers.eventstorage.GridEventStorageResponse; import org.apache.ignite.internal.plugin.AbstractMarshallableMessageFactoryProvider; import org.apache.ignite.internal.processors.authentication.AuthentificationDataBagItem; import org.apache.ignite.internal.processors.authentication.User; @@ -703,7 +704,8 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { // [13000 - 13300]: Control, configuration, diagnostics and other messages. msgIdx = 13000; - register(GridEventStorageMessage.class); + register(GridEventStorageRequest.class); + register(GridEventStorageResponse.class); register(ChangeGlobalStateMessage.class); register(GridChangeGlobalStateMessageResponse.class); register(IgniteDiagnosticRequest.class); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java index 3eb86d13a2b83..bba3c4c08bd5a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java @@ -1460,6 +1460,7 @@ private void processRegularMessage0(GridIoMessage msg, UUID nodeId) { } /** */ + // TODO IGNITE-28950: the regular path drops the message without a trace, unlike the ordered one. private void unmarshalPayload(GridIoMessage msg) { if (msg.message() instanceof DeferredUnmarshalMessage) return; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java index 618ad6bb05eed..1c883d1474aa3 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageManager.java @@ -50,6 +50,7 @@ import org.apache.ignite.internal.managers.GridManagerAdapter; import org.apache.ignite.internal.managers.communication.GridIoManager; import org.apache.ignite.internal.managers.communication.GridMessageListener; +import org.apache.ignite.internal.managers.communication.MessageMarshalling; import org.apache.ignite.internal.managers.deployment.GridDeployment; import org.apache.ignite.internal.managers.discovery.DiscoCache; import org.apache.ignite.internal.processors.platform.PlatformEventFilterListener; @@ -63,7 +64,6 @@ import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgnitePredicate; import org.apache.ignite.lang.IgniteUuid; -import org.apache.ignite.marshaller.Marshaller; import org.apache.ignite.plugin.security.SecurityPermission; import org.apache.ignite.spi.IgniteSpiException; import org.apache.ignite.spi.discovery.DiscoveryDataBag; @@ -105,9 +105,6 @@ public class GridEventStorageManager extends GridManagerAdapter /** Recordable events arrays length. */ private final int len; - /** Marshaller. */ - private final Marshaller marsh; - /** Request listener. */ private RequestListener msgLsnr; @@ -142,8 +139,6 @@ public class GridEventStorageManager extends GridManagerAdapter public GridEventStorageManager(GridKernalContext ctx) { super(ctx, ctx.config().getEventStorageSpi()); - marsh = ctx.marshaller(); - int[] cfgInclEvtTypes0 = ctx.config().getIncludeEventTypes(); if (F.isEmpty(cfgInclEvtTypes0)) @@ -1032,13 +1027,13 @@ private List query(IgnitePredicate p, Collection List query(IgnitePredicate p, Collection List query(IgnitePredicate p, Collection List query(IgnitePredicate p, Collection nodes, GridTopic topic, - GridEventStorageMessage msg, byte plc) throws IgniteCheckedException { + GridEventStorageRequest msg, byte plc) throws IgniteCheckedException { ClusterNode locNode = F.find(nodes, null, localNode(ctx.localNodeId())); Collection rmtNodes = F.view(nodes, remoteNodes(ctx.localNodeId())); @@ -1225,13 +1222,13 @@ private class RequestListener implements GridMessageListener { return; try { - if (!(msg instanceof GridEventStorageMessage)) { + if (!(msg instanceof GridEventStorageRequest)) { U.warn(log, "Received unknown message: " + msg); return; } - GridEventStorageMessage req = (GridEventStorageMessage)msg; + GridEventStorageRequest req = (GridEventStorageRequest)msg; ClusterNode node = ctx.discovery().node(nodeId); @@ -1265,7 +1262,7 @@ private class RequestListener implements GridMessageListener { throw new IgniteDeploymentCheckedException("Failed to obtain deployment for event filter " + "(is peer class loading turned on?): " + req); - req.finishUnmarshalFilters(marsh, U.resolveClassLoader(dep.classLoader(), ctx.config())); + MessageMarshalling.unmarshal(req, ctx, null, U.resolveClassLoader(dep.classLoader(), ctx.config())); filter = (IgnitePredicate)req.filter(); @@ -1295,7 +1292,7 @@ private class RequestListener implements GridMessageListener { } // Response message. - GridEventStorageMessage res = new GridEventStorageMessage(evts, ex); + GridEventStorageResponse res = new GridEventStorageResponse(evts, ex); try { if (log.isDebugEnabled()) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java similarity index 50% rename from modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageMessage.java rename to modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java index 7ba5215611e8f..a2fd6a0c26c3c 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageRequest.java @@ -17,215 +17,125 @@ package org.apache.ignite.internal.managers.eventstorage; -import java.util.Collection; import java.util.Collections; import java.util.Map; import java.util.UUID; -import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.configuration.DeploymentMode; -import org.apache.ignite.events.Event; -import org.apache.ignite.internal.GridTopicMessage; -import org.apache.ignite.internal.MarshallableMessage; +import org.apache.ignite.internal.DeferredUnmarshalMessage; import org.apache.ignite.internal.Marshalled; import org.apache.ignite.internal.Order; import org.apache.ignite.internal.UseBinaryMarshaller; -import org.apache.ignite.internal.managers.communication.ErrorMessage; import org.apache.ignite.internal.util.tostring.GridToStringInclude; import org.apache.ignite.internal.util.typedef.internal.S; -import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgnitePredicate; import org.apache.ignite.lang.IgniteUuid; -import org.apache.ignite.marshaller.Marshaller; import org.jetbrains.annotations.Nullable; -/** - * Event storage message. - */ +import static org.apache.ignite.internal.GridTopic.TOPIC_EVENT; + +/** Remote event query. The filter is a user class, hence the deferred unmarshalling. */ @UseBinaryMarshaller -public class GridEventStorageMessage implements MarshallableMessage { +public class GridEventStorageRequest implements DeferredUnmarshalMessage { /** */ @Order(0) - GridTopicMessage resTopicMsg; + IgniteUuid resTopicId; /** */ - private IgnitePredicate filter; + @Marshalled("filterBytes") + IgnitePredicate filter; /** */ @Order(1) byte[] filterBytes; - /** */ - @Marshalled("evtsBytes") - Collection evts; - /** */ @Order(2) - byte[] evtsBytes; - - /** */ - @Order(3) - ErrorMessage errMsg; - - /** */ - @Order(4) IgniteUuid clsLdrId; /** */ - @Order(5) + @Order(3) DeploymentMode depMode; /** */ - @Order(6) + @Order(4) String filterClsName; /** */ - @Order(7) + @Order(5) String userVer; /** Node class loader participants. */ @GridToStringInclude - @Order(8) + @Order(6) Map ldrParties; /** */ - public GridEventStorageMessage() { + public GridEventStorageRequest() { // No-op. } /** - * @param resTopic Response topic. + * @param resTopicId Id of the node waiting for the response. * @param filter Query filter. * @param clsLdrId Class loader ID. * @param depMode Deployment mode. * @param userVer User version. * @param ldrParties Node loader participant map. */ - GridEventStorageMessage( - Object resTopic, + GridEventStorageRequest( + IgniteUuid resTopicId, IgnitePredicate filter, IgniteUuid clsLdrId, DeploymentMode depMode, String userVer, Map ldrParties) { - resTopicMsg = new GridTopicMessage(resTopic); + this.resTopicId = resTopicId; this.filter = filter; - filterClsName = filter.getClass().getName(); - this.depMode = depMode; this.clsLdrId = clsLdrId; + this.depMode = depMode; this.userVer = userVer; this.ldrParties = ldrParties; - evts = null; - errMsg = null; - } - - /** - * @param evts Grid events. - * @param ex Exception occurred during processing. - */ - GridEventStorageMessage(Collection evts, Throwable ex) { - this.evts = evts; - - if (ex != null) - errMsg = new ErrorMessage(ex); - - resTopicMsg = null; - filter = null; - filterClsName = null; - depMode = null; - clsLdrId = null; - userVer = null; + filterClsName = filter.getClass().getName(); } - /** - * @return Response topic. - */ + /** @return Topic to answer to. */ Object responseTopic() { - return GridTopicMessage.topic(resTopicMsg); + return TOPIC_EVENT.topic(resTopicId); } - /** - * @return Filter. - */ + /** @return Filter. */ IgnitePredicate filter() { return filter; } - /** - * @return Events. - */ - @Nullable Collection events() { - return evts != null ? Collections.unmodifiableCollection(evts) : null; - } - - /** - * @return the Class loader ID. - */ + /** @return Class loader ID. */ IgniteUuid classLoaderId() { return clsLdrId; } - /** - * @return Deployment mode. - */ + /** @return Deployment mode. */ DeploymentMode deploymentMode() { return depMode; } - /** - * @return Filter class name. - */ + /** @return Filter class name. */ String filterClassName() { return filterClsName; } - /** - * @return User version. - */ + /** @return User version. */ String userVersion() { return userVer; } - /** - * @return Node class loader participant map. - */ + /** @return Node class loader participant map. */ @Nullable Map loaderParticipants() { return ldrParties != null ? Collections.unmodifiableMap(ldrParties) : null; } - /** - * @return Exception. - */ - @Nullable Throwable exception() { - return ErrorMessage.error(errMsg); - } - - /** {@inheritDoc} */ - @Override public void marshal(Marshaller marsh) throws IgniteCheckedException { - if (filter != null) - filterBytes = U.marshal(marsh, filter); - } - - /** {@inheritDoc} */ - @Override public void unmarshal(Marshaller marsh, ClassLoader ldr) throws IgniteCheckedException { - // No-op. - } - - /** - * @param marsh Marshaller. - * @param filterClsLdr Class loader for filter. - */ - // TODO IGNITE-28901: revise the filters marshalling. - public void finishUnmarshalFilters(Marshaller marsh, ClassLoader filterClsLdr) throws IgniteCheckedException { - if (filterBytes != null && filter == null) { - filter = U.unmarshal(marsh, filterBytes, filterClsLdr); - - filterBytes = null; - } - } - /** {@inheritDoc} */ @Override public String toString() { - return S.toString(GridEventStorageMessage.class, this); + return S.toString(GridEventStorageRequest.class, this); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageResponse.java b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageResponse.java new file mode 100644 index 0000000000000..aa296d6c82d1c --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/managers/eventstorage/GridEventStorageResponse.java @@ -0,0 +1,76 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.ignite.internal.managers.eventstorage; + +import java.util.Collection; +import java.util.Collections; +import org.apache.ignite.events.Event; +import org.apache.ignite.internal.Marshalled; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.UseBinaryMarshaller; +import org.apache.ignite.internal.managers.communication.ErrorMessage; +import org.apache.ignite.internal.util.typedef.internal.S; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; + +/** Events collected for a {@link GridEventStorageRequest}, or the failure that prevented it. */ +@UseBinaryMarshaller +public class GridEventStorageResponse implements Message { + /** */ + @Marshalled("evtsBytes") + Collection evts; + + /** */ + @Order(0) + byte[] evtsBytes; + + /** */ + @Order(1) + ErrorMessage errMsg; + + /** */ + public GridEventStorageResponse() { + // No-op. + } + + /** + * @param evts Grid events. + * @param ex Exception occurred during processing. + */ + GridEventStorageResponse(Collection evts, @Nullable Throwable ex) { + this.evts = evts; + + if (ex != null) + errMsg = new ErrorMessage(ex); + } + + /** @return Events. */ + @Nullable Collection events() { + return evts != null ? Collections.unmodifiableCollection(evts) : null; + } + + /** @return Exception. */ + @Nullable Throwable exception() { + return ErrorMessage.error(errMsg); + } + + /** {@inheritDoc} */ + @Override public String toString() { + return S.toString(GridEventStorageResponse.class, this); + } +} diff --git a/modules/core/src/main/resources/META-INF/classnames.properties b/modules/core/src/main/resources/META-INF/classnames.properties index 931395d4b7e35..987db34c4e8a2 100644 --- a/modules/core/src/main/resources/META-INF/classnames.properties +++ b/modules/core/src/main/resources/META-INF/classnames.properties @@ -735,7 +735,6 @@ org.apache.ignite.internal.managers.encryption.GridEncryptionManager$EmptyResult org.apache.ignite.internal.managers.encryption.GridEncryptionManager$MasterKeyChangeRequest org.apache.ignite.internal.managers.encryption.NodeEncryptionKeys org.apache.ignite.internal.managers.encryption.GroupKeyEncrypted -org.apache.ignite.internal.managers.eventstorage.GridEventStorageMessage org.apache.ignite.internal.managers.indexing.GridIndexingManager$1 org.apache.ignite.internal.managers.loadbalancer.GridLoadBalancerAdapter org.apache.ignite.internal.managers.loadbalancer.GridLoadBalancerManager$1