This is an automated email from the ASF dual-hosted git repository.
anton-vinogradov pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ignite.git
The following commit(s) were added to refs/heads/master by this push:
new 4b92f71c2bb IGNITE-28901 Split GridEventStorageMessage into a request
and a response (#13428)
4b92f71c2bb is described below
commit 4b92f71c2bbb8865e8572f6a7ac8cc22e7cc614e
Author: Anton Vinogradov <[email protected]>
AuthorDate: Tue Aug 4 23:21:11 2026 +0300
IGNITE-28901 Split GridEventStorageMessage into a request and a response
(#13428)
---
.../ignite/internal/CoreMessagesProvider.java | 6 +-
.../managers/communication/GridIoManager.java | 1 +
.../eventstorage/GridEventStorageManager.java | 29 ++--
...geMessage.java => GridEventStorageRequest.java} | 146 ++++-----------------
.../eventstorage/GridEventStorageResponse.java | 76 +++++++++++
.../main/resources/META-INF/classnames.properties | 1 -
6 files changed, 122 insertions(+), 137 deletions(-)
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 8b3fe7bbd4f..cd03c5a2f2f 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.GroupKeyEncrypted;
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 class CoreMessagesProvider extends
AbstractMarshallableMessageFactoryProv
// [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 3eb86d13a2b..bba3c4c08bd 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 @@ public class GridIoManager extends
GridManagerAdapter<CommunicationSpi<Object>>
}
/** */
+ // 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 618ad6bb05e..1c883d1474a 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.cluster.ClusterTopologyCheckedException;
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.LT;
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<EventStorageSpi>
/** 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<EventStorageSpi>
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 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
assert nodeId != null;
assert msg != null;
- if (!(msg instanceof GridEventStorageMessage)) {
+ if (!(msg instanceof GridEventStorageResponse)) {
U.error(log, "Received unknown message: " + msg);
return;
}
- GridEventStorageMessage res = (GridEventStorageMessage)msg;
+ GridEventStorageResponse res = (GridEventStorageResponse)msg;
synchronized (qryMux) {
if (uids.remove(nodeId)) {
@@ -1058,7 +1053,9 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
}
};
- Object resTopic =
TOPIC_EVENT.topic(IgniteUuid.fromUuid(ctx.localNodeId()));
+ IgniteUuid resTopicId = IgniteUuid.fromUuid(ctx.localNodeId());
+
+ Object resTopic = TOPIC_EVENT.topic(resTopicId);
try {
addLocalEventListener(evtLsnr, new int[] {
@@ -1073,8 +1070,8 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
if (dep == null)
throw new IgniteDeploymentCheckedException("Failed to deploy
event filter: " + p);
- GridEventStorageMessage msg = new GridEventStorageMessage(
- resTopic,
+ GridEventStorageRequest msg = new GridEventStorageRequest(
+ resTopicId,
p,
dep.classLoaderId(),
dep.deployMode(),
@@ -1144,7 +1141,7 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
* @throws IgniteCheckedException If sending failed.
*/
private void sendMessage(Collection<? extends ClusterNode> 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<? extends ClusterNode> rmtNodes = F.view(nodes,
remoteNodes(ctx.localNodeId()));
@@ -1225,13 +1222,13 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
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 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
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<Event>)req.filter();
@@ -1295,7 +1292,7 @@ public class GridEventStorageManager extends
GridManagerAdapter<EventStorageSpi>
}
// 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 7ba5215611e..a2fd6a0c26c 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<Event> 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<UUID, IgniteUuid> 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<UUID, IgniteUuid> 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<Event> 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<Event> 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<UUID, IgniteUuid> 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 00000000000..aa296d6c82d
--- /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<Event> 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<Event> evts, @Nullable Throwable ex) {
+ this.evts = evts;
+
+ if (ex != null)
+ errMsg = new ErrorMessage(ex);
+ }
+
+ /** @return Events. */
+ @Nullable Collection<Event> 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 931395d4b7e..987db34c4e8 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