Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -1032,13 +1027,13 @@ private <T extends Event> List<T> query(IgnitePredicate<T> p, Collection<? exten
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)) {
Expand All @@ -1058,7 +1053,9 @@ private <T extends Event> List<T> query(IgnitePredicate<T> p, Collection<? exten
}
};

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[] {
Expand All @@ -1073,8 +1070,8 @@ private <T extends Event> List<T> query(IgnitePredicate<T> p, Collection<? exten
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(),
Expand Down Expand Up @@ -1144,7 +1141,7 @@ private <T extends Event> List<T> query(IgnitePredicate<T> p, Collection<? exten
* @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()));
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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<Event>)req.filter();

Expand Down Expand Up @@ -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())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Loading
Loading