From 7e61596fc01e30474ad6fcf04ff5e2cd77df9e84 Mon Sep 17 00:00:00 2001 From: Anton Vinogradov Date: Fri, 31 Jul 2026 16:24:51 +0300 Subject: [PATCH 1/4] IGNITE-28937 Register core messages uniformly: drop withSchema/withNoSchema After IGNITE-28929 the marshaller of a message is set by @UseBinaryMarshaller, and the two registration helpers became the same code: they differed only by an assert that the annotation matches the helper. The same fact was stored twice, and the assert only checked that the two copies agree. Replaced the 299 withSchema/withNoSchema calls with register(factory, cls, msgIdx++) and deleted both helpers, so registration in core now reads like CalciteMessageFactory. Also stopped passing a Marshaller into the serializer and deployer lookups: neither takes one - 0 of 197 generated serializers and 0 of 31 generated deployers declare such a constructor, and the generators cannot emit one. The marshaller now reaches only the generated marshaller companion. Co-Authored-By: Claude Opus 5 --- .../ignite/internal/CoreMessagesProvider.java | 615 +++++++++--------- ...actMarshallableMessageFactoryProvider.java | 56 +- 2 files changed, 342 insertions(+), 329 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 ffabd9754894b..4e07d4a25d366 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 @@ -269,7 +269,6 @@ import org.apache.ignite.internal.util.distributed.SingleNodeMessage; import org.apache.ignite.marshaller.Marshaller; import org.apache.ignite.marshaller.jdk.JdkMarshaller; -import org.apache.ignite.plugin.extensions.communication.Message; import org.apache.ignite.plugin.security.SecurityBasicPermissionSet; import org.apache.ignite.spi.collision.jobstealing.JobStealingRequest; import org.apache.ignite.spi.communication.tcp.internal.TcpConnectionRequestDiscoveryMessage; @@ -366,380 +365,364 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { // [5000 - 5500]: Utility messages. Most of them originally come from Discovery. msgIdx = 5000; - withNoSchema(CompressedMessage.class); - withNoSchema(ErrorMessage.class); - withNoSchema(InetSocketAddressMessage.class); - withNoSchema(InetAddressMessage.class); - withNoSchema(TcpDiscoveryNode.class); - withNoSchema(DiscoveryDataPacket.class); - withNoSchema(GridByteArrayList.class); - withNoSchema(CacheVersionedValue.class); - withNoSchema(KeyedVersionedValue.class); - withNoSchema(WALPointer.class); - withNoSchema(SerializableDataBagItemWrapper.class); - withSchema(GridTopicMessage.class); - withNoSchema(GridIntList.class); + register(factory, CompressedMessage.class, msgIdx++); + register(factory, ErrorMessage.class, msgIdx++); + register(factory, InetSocketAddressMessage.class, msgIdx++); + register(factory, InetAddressMessage.class, msgIdx++); + register(factory, TcpDiscoveryNode.class, msgIdx++); + register(factory, DiscoveryDataPacket.class, msgIdx++); + register(factory, GridByteArrayList.class, msgIdx++); + register(factory, CacheVersionedValue.class, msgIdx++); + register(factory, KeyedVersionedValue.class, msgIdx++); + register(factory, WALPointer.class, msgIdx++); + register(factory, SerializableDataBagItemWrapper.class, msgIdx++); + register(factory, GridTopicMessage.class, msgIdx++); + register(factory, GridIntList.class, msgIdx++); // [5700 - 5900]: Discovery originated messages. msgIdx = 5700; - withNoSchema(TcpDiscoveryHandshakeRequest.class); - withNoSchema(TcpDiscoveryHandshakeResponse.class); - withNoSchema(TcpDiscoveryJoinRequestMessage.class); - withNoSchema(TcpDiscoveryNodeAddedMessage.class); - withNoSchema(TcpDiscoveryNodeAddFinishedMessage.class); - withNoSchema(TcpDiscoveryNodeLeftMessage.class); - withNoSchema(TcpDiscoveryNodeFailedMessage.class); - withNoSchema(TcpDiscoveryConnectionCheckMessage.class); - withNoSchema(TcpDiscoveryPingRequest.class); - withNoSchema(TcpDiscoveryPingResponse.class); - withNoSchema(TcpDiscoveryClientPingRequest.class); - withNoSchema(TcpDiscoveryClientPingResponse.class); - withNoSchema(TcpDiscoveryClientAckResponse.class); - withNoSchema(TcpDiscoveryClientReconnectMessage.class); - withNoSchema(TcpDiscoveryDiscardMessage.class); - withNoSchema(TcpDiscoveryCheckFailedMessage.class); - withNoSchema(TcpDiscoveryLoopbackProblemMessage.class); - withNoSchema(TcpDiscoveryRingLatencyCheckMessage.class); - withNoSchema(TcpDiscoveryDuplicateIdMessage.class); - withNoSchema(TcpDiscoveryCustomEventMessage.class); - withNoSchema(TcpDiscoveryServerOnlyCustomEventMessage.class); + register(factory, TcpDiscoveryHandshakeRequest.class, msgIdx++); + register(factory, TcpDiscoveryHandshakeResponse.class, msgIdx++); + register(factory, TcpDiscoveryJoinRequestMessage.class, msgIdx++); + register(factory, TcpDiscoveryNodeAddedMessage.class, msgIdx++); + register(factory, TcpDiscoveryNodeAddFinishedMessage.class, msgIdx++); + register(factory, TcpDiscoveryNodeLeftMessage.class, msgIdx++); + register(factory, TcpDiscoveryNodeFailedMessage.class, msgIdx++); + register(factory, TcpDiscoveryConnectionCheckMessage.class, msgIdx++); + register(factory, TcpDiscoveryPingRequest.class, msgIdx++); + register(factory, TcpDiscoveryPingResponse.class, msgIdx++); + register(factory, TcpDiscoveryClientPingRequest.class, msgIdx++); + register(factory, TcpDiscoveryClientPingResponse.class, msgIdx++); + register(factory, TcpDiscoveryClientAckResponse.class, msgIdx++); + register(factory, TcpDiscoveryClientReconnectMessage.class, msgIdx++); + register(factory, TcpDiscoveryDiscardMessage.class, msgIdx++); + register(factory, TcpDiscoveryCheckFailedMessage.class, msgIdx++); + register(factory, TcpDiscoveryLoopbackProblemMessage.class, msgIdx++); + register(factory, TcpDiscoveryRingLatencyCheckMessage.class, msgIdx++); + register(factory, TcpDiscoveryDuplicateIdMessage.class, msgIdx++); + register(factory, TcpDiscoveryCustomEventMessage.class, msgIdx++); + register(factory, TcpDiscoveryServerOnlyCustomEventMessage.class, msgIdx++); msgIdx = 5900; - withNoSchema(TcpDiscoveryStatusCheckMessage.class); + register(factory, TcpDiscoveryStatusCheckMessage.class, msgIdx++); // [6000 - 6200]: Snapshot operation messages. Most of them originally come from Discovery. msgIdx = 6000; - withNoSchema(SnapshotStartDiscoveryMessage.class); - withNoSchema(SnapshotCheckProcessRequest.class); - withNoSchema(SnapshotOperationRequest.class); - withNoSchema(SnapshotOperationEndRequest.class); - withNoSchema(SnapshotRestoreStartRequest.class); - withNoSchema(SnapshotOperationResponse.class); - withNoSchema(SnapshotHandlerResult.class); - withNoSchema(SnapshotCheckResponse.class); - withNoSchema(SnapshotPartitionsVerifyHandlerResponse.class); - withNoSchema(SnapshotRestoreOperationResponse.class); - withNoSchema(SnapshotMetadataResponse.class); - withNoSchema(SnapshotMetadata.class); - withNoSchema(SnapshotCheckPartitionHashesResponse.class); - withNoSchema(SnapshotCheckHandlersResponse.class); - withNoSchema(SnapshotFilesRequestMessage.class); - withNoSchema(SnapshotFilesFailureMessage.class); - withNoSchema(IncrementalSnapshotVerifyResult.class); - withNoSchema(IncrementalSnapshotAwareMessage.class); + register(factory, SnapshotStartDiscoveryMessage.class, msgIdx++); + register(factory, SnapshotCheckProcessRequest.class, msgIdx++); + register(factory, SnapshotOperationRequest.class, msgIdx++); + register(factory, SnapshotOperationEndRequest.class, msgIdx++); + register(factory, SnapshotRestoreStartRequest.class, msgIdx++); + register(factory, SnapshotOperationResponse.class, msgIdx++); + register(factory, SnapshotHandlerResult.class, msgIdx++); + register(factory, SnapshotCheckResponse.class, msgIdx++); + register(factory, SnapshotPartitionsVerifyHandlerResponse.class, msgIdx++); + register(factory, SnapshotRestoreOperationResponse.class, msgIdx++); + register(factory, SnapshotMetadataResponse.class, msgIdx++); + register(factory, SnapshotMetadata.class, msgIdx++); + register(factory, SnapshotCheckPartitionHashesResponse.class, msgIdx++); + register(factory, SnapshotCheckHandlersResponse.class, msgIdx++); + register(factory, SnapshotFilesRequestMessage.class, msgIdx++); + register(factory, SnapshotFilesFailureMessage.class, msgIdx++); + register(factory, IncrementalSnapshotVerifyResult.class, msgIdx++); + register(factory, IncrementalSnapshotAwareMessage.class, msgIdx++); // [6300 - 6400]: Services messages. Most of them originally come from Discovery. msgIdx = 6300; - withNoSchema(ServiceDeploymentProcessId.class); - withSchema(ServiceSingleNodeDeploymentResult.class); - withNoSchema(ServiceClusterDeploymentResult.class); - withNoSchema(ServiceDeploymentRequest.class); - withNoSchema(ServiceUndeploymentRequest.class); - withNoSchema(ServiceClusterDeploymentResultBatch.class); - withNoSchema(ServiceChangeBatchRequest.class); - withNoSchema(ServiceSingleNodeDeploymentResultBatch.class); - withNoSchema(ServiceProcessorCommonDiscoveryData.class); - withNoSchema(ServiceProcessorJoinNodeDiscoveryData.class); - withNoSchema(ServiceInfo.class); - withNoSchema(ServiceTopology.class); - withNoSchema(LazyServiceConfigurationMessage.class); + register(factory, ServiceDeploymentProcessId.class, msgIdx++); + register(factory, ServiceSingleNodeDeploymentResult.class, msgIdx++); + register(factory, ServiceClusterDeploymentResult.class, msgIdx++); + register(factory, ServiceDeploymentRequest.class, msgIdx++); + register(factory, ServiceUndeploymentRequest.class, msgIdx++); + register(factory, ServiceClusterDeploymentResultBatch.class, msgIdx++); + register(factory, ServiceChangeBatchRequest.class, msgIdx++); + register(factory, ServiceSingleNodeDeploymentResultBatch.class, msgIdx++); + register(factory, ServiceProcessorCommonDiscoveryData.class, msgIdx++); + register(factory, ServiceProcessorJoinNodeDiscoveryData.class, msgIdx++); + register(factory, ServiceInfo.class, msgIdx++); + register(factory, ServiceTopology.class, msgIdx++); + register(factory, LazyServiceConfigurationMessage.class, msgIdx++); // [6500 - 6700]: DiscoveryCustomMessage msgIdx = 6500; - withNoSchema(TcpConnectionRequestDiscoveryMessage.class); - withNoSchema(DistributedMetaStorageUpdateMessage.class); - withNoSchema(DistributedMetaStorageUpdateAckMessage.class); - withNoSchema(DistributedMetaStorageCasMessage.class); - withNoSchema(DistributedMetaStorageCasAckMessage.class); - withNoSchema(FullMessage.class); - withNoSchema(InitMessage.class); - withNoSchema(CacheStatisticsModeChangeMessage.class); - withNoSchema(MetadataRemoveAcceptedMessage.class); - withNoSchema(MetadataRemoveProposedMessage.class); - withNoSchema(WalStateFinishMessage.class); - withNoSchema(WalStateProposeMessage.class); - withNoSchema(MetadataUpdateAcceptedMessage.class); - withNoSchema(MetadataUpdateProposedMessage.class); - withNoSchema(TxTimeoutOnPartitionMapExchangeChangeMessage.class); - withNoSchema(UserAcceptedMessage.class); - withNoSchema(UserProposedMessage.class); - withNoSchema(ChangeGlobalStateFinishMessage.class); - withNoSchema(StopRoutineAckDiscoveryMessage.class); - withNoSchema(StopRoutineDiscoveryMessage.class); - withNoSchema(CacheAffinityChangeMessage.class); - withNoSchema(ClientCacheChangeDiscoveryMessage.class); - withNoSchema(MappingAcceptedMessage.class); - withNoSchema(MappingProposedMessage.class); - withNoSchema(ExchangeFailureMessage.class); - withNoSchema(CacheStatisticsClearMessage.class); - withNoSchema(ClientCacheChangeDummyDiscoveryMessage.class); - withNoSchema(DynamicCacheChangeBatch.class); - withNoSchema(CacheClientReconnectDiscoveryData.class); - withNoSchema(CacheGroupRecoveryState.class); - withNoSchema(CacheJoinInfo.class); - withNoSchema(CacheJoinNodeDiscoveryData.class); - withNoSchema(CacheReconnectInfo.class); - withNoSchema(ClusterCacheGroupRecoveryData.class); + register(factory, TcpConnectionRequestDiscoveryMessage.class, msgIdx++); + register(factory, DistributedMetaStorageUpdateMessage.class, msgIdx++); + register(factory, DistributedMetaStorageUpdateAckMessage.class, msgIdx++); + register(factory, DistributedMetaStorageCasMessage.class, msgIdx++); + register(factory, DistributedMetaStorageCasAckMessage.class, msgIdx++); + register(factory, FullMessage.class, msgIdx++); + register(factory, InitMessage.class, msgIdx++); + register(factory, CacheStatisticsModeChangeMessage.class, msgIdx++); + register(factory, MetadataRemoveAcceptedMessage.class, msgIdx++); + register(factory, MetadataRemoveProposedMessage.class, msgIdx++); + register(factory, WalStateFinishMessage.class, msgIdx++); + register(factory, WalStateProposeMessage.class, msgIdx++); + register(factory, MetadataUpdateAcceptedMessage.class, msgIdx++); + register(factory, MetadataUpdateProposedMessage.class, msgIdx++); + register(factory, TxTimeoutOnPartitionMapExchangeChangeMessage.class, msgIdx++); + register(factory, UserAcceptedMessage.class, msgIdx++); + register(factory, UserProposedMessage.class, msgIdx++); + register(factory, ChangeGlobalStateFinishMessage.class, msgIdx++); + register(factory, StopRoutineAckDiscoveryMessage.class, msgIdx++); + register(factory, StopRoutineDiscoveryMessage.class, msgIdx++); + register(factory, CacheAffinityChangeMessage.class, msgIdx++); + register(factory, ClientCacheChangeDiscoveryMessage.class, msgIdx++); + register(factory, MappingAcceptedMessage.class, msgIdx++); + register(factory, MappingProposedMessage.class, msgIdx++); + register(factory, ExchangeFailureMessage.class, msgIdx++); + register(factory, CacheStatisticsClearMessage.class, msgIdx++); + register(factory, ClientCacheChangeDummyDiscoveryMessage.class, msgIdx++); + register(factory, DynamicCacheChangeBatch.class, msgIdx++); + register(factory, CacheClientReconnectDiscoveryData.class, msgIdx++); + register(factory, CacheGroupRecoveryState.class, msgIdx++); + register(factory, CacheJoinInfo.class, msgIdx++); + register(factory, CacheJoinNodeDiscoveryData.class, msgIdx++); + register(factory, CacheReconnectInfo.class, msgIdx++); + register(factory, ClusterCacheGroupRecoveryData.class, msgIdx++); // [10000 - 10200]: Transaction and lock related messages. Most of them originally comes from Communication. msgIdx = 10000; - withNoSchema(TxInfo.class); - withSchema(TxEntriesInfo.class); - withNoSchema(TxLock.class); - withSchema(TxLocksRequest.class); - withSchema(TxLocksResponse.class); - withSchema(IgniteTxKey.class); - withSchema(IgniteTxEntry.class); - withSchema(TxEntryValueHolder.class); - withNoSchema(GridCacheTxRecoveryRequest.class); - withNoSchema(GridCacheTxRecoveryResponse.class); - withNoSchema(GridDistributedTxFinishRequest.class); - withNoSchema(GridDistributedTxFinishResponse.class); - withSchema(GridDistributedTxPrepareRequest.class); - withNoSchema(GridDistributedTxPrepareResponse.class); - withNoSchema(GridDhtTxFinishRequest.class); - withSchema(GridDhtTxFinishResponse.class); - withSchema(GridDhtTxPrepareRequest.class); - withSchema(GridDhtTxPrepareResponse.class); - withNoSchema(GridNearTxFinishRequest.class); - withNoSchema(GridNearTxFinishResponse.class); - withNoSchema(GridNearTxPrepareRequest.class); - withSchema(GridNearTxPrepareResponse.class); - withSchema(GridDhtLockRequest.class); - withSchema(GridDhtLockResponse.class); - withSchema(GridDhtUnlockRequest.class); - withNoSchema(GridNearLockRequest.class); - withNoSchema(GridNearLockResponse.class); - withSchema(GridNearUnlockRequest.class); - withSchema(GridDistributedLockRequest.class); - withSchema(GridDistributedLockResponse.class); - withNoSchema(GridDhtTxOnePhaseCommitAckRequest.class); - withSchema(TransactionAttributesAwareRequest.class); + register(factory, TxInfo.class, msgIdx++); + register(factory, TxEntriesInfo.class, msgIdx++); + register(factory, TxLock.class, msgIdx++); + register(factory, TxLocksRequest.class, msgIdx++); + register(factory, TxLocksResponse.class, msgIdx++); + register(factory, IgniteTxKey.class, msgIdx++); + register(factory, IgniteTxEntry.class, msgIdx++); + register(factory, TxEntryValueHolder.class, msgIdx++); + register(factory, GridCacheTxRecoveryRequest.class, msgIdx++); + register(factory, GridCacheTxRecoveryResponse.class, msgIdx++); + register(factory, GridDistributedTxFinishRequest.class, msgIdx++); + register(factory, GridDistributedTxFinishResponse.class, msgIdx++); + register(factory, GridDistributedTxPrepareRequest.class, msgIdx++); + register(factory, GridDistributedTxPrepareResponse.class, msgIdx++); + register(factory, GridDhtTxFinishRequest.class, msgIdx++); + register(factory, GridDhtTxFinishResponse.class, msgIdx++); + register(factory, GridDhtTxPrepareRequest.class, msgIdx++); + register(factory, GridDhtTxPrepareResponse.class, msgIdx++); + register(factory, GridNearTxFinishRequest.class, msgIdx++); + register(factory, GridNearTxFinishResponse.class, msgIdx++); + register(factory, GridNearTxPrepareRequest.class, msgIdx++); + register(factory, GridNearTxPrepareResponse.class, msgIdx++); + register(factory, GridDhtLockRequest.class, msgIdx++); + register(factory, GridDhtLockResponse.class, msgIdx++); + register(factory, GridDhtUnlockRequest.class, msgIdx++); + register(factory, GridNearLockRequest.class, msgIdx++); + register(factory, GridNearLockResponse.class, msgIdx++); + register(factory, GridNearUnlockRequest.class, msgIdx++); + register(factory, GridDistributedLockRequest.class, msgIdx++); + register(factory, GridDistributedLockResponse.class, msgIdx++); + register(factory, GridDhtTxOnePhaseCommitAckRequest.class, msgIdx++); + register(factory, TransactionAttributesAwareRequest.class, msgIdx++); // [10300 - 10500]: Cache, DHT messages. msgIdx = 10300; - withSchema(GridDhtForceKeysRequest.class); - withSchema(GridDhtForceKeysResponse.class); - withNoSchema(GridDhtAtomicDeferredUpdateResponse.class); - withNoSchema(GridDhtAtomicUpdateRequest.class); - withSchema(GridDhtAtomicUpdateResponse.class); - withSchema(GridNearAtomicFullUpdateRequest.class); - withSchema(GridDhtAtomicSingleUpdateRequest.class); - withSchema(GridNearAtomicUpdateResponse.class); - withSchema(GridNearAtomicSingleUpdateRequest.class); - withSchema(GridNearAtomicSingleUpdateInvokeRequest.class); - withSchema(GridNearAtomicSingleUpdateFilterRequest.class); - withNoSchema(GridNearAtomicCheckUpdateRequest.class); - withSchema(NearCacheUpdates.class); - withSchema(GridNearGetRequest.class); - withSchema(GridNearGetResponse.class); - withSchema(GridNearSingleGetRequest.class); - withSchema(GridNearSingleGetResponse.class); - withNoSchema(GridDhtAtomicNearResponse.class); - withSchema(GridCacheTtlUpdateRequest.class); - withSchema(GridCacheReturn.class); - withSchema(GridCacheEntryInfo.class); - withSchema(CacheInvokeDirectResult.class); - withNoSchema(GridCacheRawVersionedEntry.class); - withSchema(CacheEvictionEntry.class); - withSchema(CacheEntryPredicateAdapter.class); - withNoSchema(GridContinuousMessage.class); - withNoSchema(ContinuousRoutineStartResultMessage.class); - withSchema(UpdateErrors.class); - withNoSchema(LatchAckMessage.class); - withSchema(AtomicApplicationAttributesAwareRequest.class); - withNoSchema(StartRequestData.class); - withNoSchema(StartRoutineAckDiscoveryMessage.class); - withNoSchema(StartRoutineDiscoveryMessage.class); - withNoSchema(StoredCacheData.class); + register(factory, GridDhtForceKeysRequest.class, msgIdx++); + register(factory, GridDhtForceKeysResponse.class, msgIdx++); + register(factory, GridDhtAtomicDeferredUpdateResponse.class, msgIdx++); + register(factory, GridDhtAtomicUpdateRequest.class, msgIdx++); + register(factory, GridDhtAtomicUpdateResponse.class, msgIdx++); + register(factory, GridNearAtomicFullUpdateRequest.class, msgIdx++); + register(factory, GridDhtAtomicSingleUpdateRequest.class, msgIdx++); + register(factory, GridNearAtomicUpdateResponse.class, msgIdx++); + register(factory, GridNearAtomicSingleUpdateRequest.class, msgIdx++); + register(factory, GridNearAtomicSingleUpdateInvokeRequest.class, msgIdx++); + register(factory, GridNearAtomicSingleUpdateFilterRequest.class, msgIdx++); + register(factory, GridNearAtomicCheckUpdateRequest.class, msgIdx++); + register(factory, NearCacheUpdates.class, msgIdx++); + register(factory, GridNearGetRequest.class, msgIdx++); + register(factory, GridNearGetResponse.class, msgIdx++); + register(factory, GridNearSingleGetRequest.class, msgIdx++); + register(factory, GridNearSingleGetResponse.class, msgIdx++); + register(factory, GridDhtAtomicNearResponse.class, msgIdx++); + register(factory, GridCacheTtlUpdateRequest.class, msgIdx++); + register(factory, GridCacheReturn.class, msgIdx++); + register(factory, GridCacheEntryInfo.class, msgIdx++); + register(factory, CacheInvokeDirectResult.class, msgIdx++); + register(factory, GridCacheRawVersionedEntry.class, msgIdx++); + register(factory, CacheEvictionEntry.class, msgIdx++); + register(factory, CacheEntryPredicateAdapter.class, msgIdx++); + register(factory, GridContinuousMessage.class, msgIdx++); + register(factory, ContinuousRoutineStartResultMessage.class, msgIdx++); + register(factory, UpdateErrors.class, msgIdx++); + register(factory, LatchAckMessage.class, msgIdx++); + register(factory, AtomicApplicationAttributesAwareRequest.class, msgIdx++); + register(factory, StartRequestData.class, msgIdx++); + register(factory, StartRoutineAckDiscoveryMessage.class, msgIdx++); + register(factory, StartRoutineDiscoveryMessage.class, msgIdx++); + register(factory, StoredCacheData.class, msgIdx++); // [10600-10800]: Affinity & partition maps. msgIdx = 10600; - withNoSchema(GridDhtAffinityAssignmentRequest.class); - withNoSchema(GridDhtAffinityAssignmentResponse.class); - withNoSchema(CacheGroupAffinityMessage.class); - withNoSchema(ExchangeInfo.class); - withNoSchema(PartitionUpdateCountersMessage.class); - withNoSchema(CachePartitionPartialCountersMap.class); - withNoSchema(IgniteDhtDemandedPartitionsMap.class); - withNoSchema(CachePartitionFullCountersMap.class); - withNoSchema(GroupPartitionIdPair.class); - withNoSchema(GridPartitionStateMap.class); - withNoSchema(GridDhtPartitionMap.class); - withNoSchema(GridDhtPartitionFullMap.class); - withNoSchema(GridDhtPartitionExchangeId.class); - withNoSchema(GridCheckpointRequest.class); - withNoSchema(GridDhtPartitionDemandMessage.class); - withSchema(GridDhtPartitionSupplyMessage.class); - withNoSchema(GridDhtPartitionsFullMessage.class); - withNoSchema(GridDhtPartitionsSingleMessage.class); - withNoSchema(GridDhtPartitionsSingleRequest.class); - withNoSchema(PartitionKey.class); + register(factory, GridDhtAffinityAssignmentRequest.class, msgIdx++); + register(factory, GridDhtAffinityAssignmentResponse.class, msgIdx++); + register(factory, CacheGroupAffinityMessage.class, msgIdx++); + register(factory, ExchangeInfo.class, msgIdx++); + register(factory, PartitionUpdateCountersMessage.class, msgIdx++); + register(factory, CachePartitionPartialCountersMap.class, msgIdx++); + register(factory, IgniteDhtDemandedPartitionsMap.class, msgIdx++); + register(factory, CachePartitionFullCountersMap.class, msgIdx++); + register(factory, GroupPartitionIdPair.class, msgIdx++); + register(factory, GridPartitionStateMap.class, msgIdx++); + register(factory, GridDhtPartitionMap.class, msgIdx++); + register(factory, GridDhtPartitionFullMap.class, msgIdx++); + register(factory, GridDhtPartitionExchangeId.class, msgIdx++); + register(factory, GridCheckpointRequest.class, msgIdx++); + register(factory, GridDhtPartitionDemandMessage.class, msgIdx++); + register(factory, GridDhtPartitionSupplyMessage.class, msgIdx++); + register(factory, GridDhtPartitionsFullMessage.class, msgIdx++); + register(factory, GridDhtPartitionsSingleMessage.class, msgIdx++); + register(factory, GridDhtPartitionsSingleRequest.class, msgIdx++); + register(factory, PartitionKey.class, msgIdx++); // [10900-11100]: Query, schema and SQL related messages. msgIdx = 10900; - withNoSchema(SchemaAlterTableAddColumnOperation.class); - withNoSchema(SchemaIndexCreateOperation.class); - withNoSchema(SchemaIndexDropOperation.class); - withNoSchema(SchemaAlterTableDropColumnOperation.class); - withNoSchema(SchemaAddQueryEntityOperation.class); - withNoSchema(SchemaOperationStatusMessage.class); - withNoSchema(SchemaProposeDiscoveryMessage.class); - withNoSchema(SchemaFinishDiscoveryMessage.class); - withNoSchema(QueryField.class); - withNoSchema(QueryIndexMessage.class); - withNoSchema(GridCacheSqlQuery.class); - withSchema(GridCacheQueryRequest.class); - withSchema(GridCacheQueryResponse.class); - withNoSchema(GridQueryCancelRequest.class); - withNoSchema(GridQueryFailResponse.class); - withNoSchema(GridQueryNextPageRequest.class); - withNoSchema(GridQueryNextPageResponse.class); - withNoSchema(GridQueryKillRequest.class); - withNoSchema(GridQueryKillResponse.class); - withNoSchema(IndexKeyDefinition.class); - withNoSchema(IndexKeyTypeSettings.class); - withNoSchema(IndexQueryResultMeta.class); - withNoSchema(StatisticsKeyMessage.class); - withNoSchema(StatisticsDecimalMessage.class); - withNoSchema(StatisticsObjectData.class); - withNoSchema(StatisticsColumnData.class); - withNoSchema(StatisticsRequest.class); - withNoSchema(StatisticsResponse.class); - withNoSchema(CacheContinuousQueryBatchAck.class); - withNoSchema(GridDhtTxSalvageMessage.class); - withSchema(CacheContinuousQueryEntry.class); - withNoSchema(QueryInlineSizesDataBagItem.class); - withNoSchema(QueryProposalsDataBagItem.class); - withNoSchema(QueryEntityMessage.class); - withNoSchema(QueryEntityExMessage.class); + register(factory, SchemaAlterTableAddColumnOperation.class, msgIdx++); + register(factory, SchemaIndexCreateOperation.class, msgIdx++); + register(factory, SchemaIndexDropOperation.class, msgIdx++); + register(factory, SchemaAlterTableDropColumnOperation.class, msgIdx++); + register(factory, SchemaAddQueryEntityOperation.class, msgIdx++); + register(factory, SchemaOperationStatusMessage.class, msgIdx++); + register(factory, SchemaProposeDiscoveryMessage.class, msgIdx++); + register(factory, SchemaFinishDiscoveryMessage.class, msgIdx++); + register(factory, QueryField.class, msgIdx++); + register(factory, QueryIndexMessage.class, msgIdx++); + register(factory, GridCacheSqlQuery.class, msgIdx++); + register(factory, GridCacheQueryRequest.class, msgIdx++); + register(factory, GridCacheQueryResponse.class, msgIdx++); + register(factory, GridQueryCancelRequest.class, msgIdx++); + register(factory, GridQueryFailResponse.class, msgIdx++); + register(factory, GridQueryNextPageRequest.class, msgIdx++); + register(factory, GridQueryNextPageResponse.class, msgIdx++); + register(factory, GridQueryKillRequest.class, msgIdx++); + register(factory, GridQueryKillResponse.class, msgIdx++); + register(factory, IndexKeyDefinition.class, msgIdx++); + register(factory, IndexKeyTypeSettings.class, msgIdx++); + register(factory, IndexQueryResultMeta.class, msgIdx++); + register(factory, StatisticsKeyMessage.class, msgIdx++); + register(factory, StatisticsDecimalMessage.class, msgIdx++); + register(factory, StatisticsObjectData.class, msgIdx++); + register(factory, StatisticsColumnData.class, msgIdx++); + register(factory, StatisticsRequest.class, msgIdx++); + register(factory, StatisticsResponse.class, msgIdx++); + register(factory, CacheContinuousQueryBatchAck.class, msgIdx++); + register(factory, GridDhtTxSalvageMessage.class, msgIdx++); + register(factory, CacheContinuousQueryEntry.class, msgIdx++); + register(factory, QueryInlineSizesDataBagItem.class, msgIdx++); + register(factory, QueryProposalsDataBagItem.class, msgIdx++); + register(factory, QueryEntityMessage.class, msgIdx++); + register(factory, QueryEntityExMessage.class, msgIdx++); // [11200 - 11300]: Compute, distributed process messages. msgIdx = 11200; - withNoSchema(GridJobCancelRequest.class); - withSchema(GridJobExecuteRequest.class); - withSchema(GridJobExecuteResponse.class); - withNoSchema(GridJobSiblingsRequest.class); - withSchema(GridJobSiblingsResponse.class); - withNoSchema(GridTaskCancelRequest.class); - withSchema(GridTaskSessionRequest.class); - withNoSchema(GridTaskResultRequest.class); - withSchema(GridTaskResultResponse.class); - withNoSchema(JobStealingRequest.class); - withNoSchema(SingleNodeMessage.class); + register(factory, GridJobCancelRequest.class, msgIdx++); + register(factory, GridJobExecuteRequest.class, msgIdx++); + register(factory, GridJobExecuteResponse.class, msgIdx++); + register(factory, GridJobSiblingsRequest.class, msgIdx++); + register(factory, GridJobSiblingsResponse.class, msgIdx++); + register(factory, GridTaskCancelRequest.class, msgIdx++); + register(factory, GridTaskSessionRequest.class, msgIdx++); + register(factory, GridTaskResultRequest.class, msgIdx++); + register(factory, GridTaskResultResponse.class, msgIdx++); + register(factory, JobStealingRequest.class, msgIdx++); + register(factory, SingleNodeMessage.class, msgIdx++); // [11500 - 11600]: IO, networking messages. msgIdx = NODE_ID_MSG_TYPE; - withNoSchema(NodeIdMessage.class); + register(factory, NodeIdMessage.class, msgIdx++); msgIdx = HANDSHAKE_MSG_TYPE; - withNoSchema(HandshakeMessage.class); + register(factory, HandshakeMessage.class, msgIdx++); msgIdx = HANDSHAKE_WAIT_MSG_TYPE; - withNoSchema(HandshakeWaitMessage.class); - withNoSchema(GridIoMessage.class); - withNoSchema(IgniteIoTestMessage.class); - withSchema(GridIoUserMessage.class); - withNoSchema(RecoveryLastReceivedMessage.class); - withNoSchema(TcpInverseConnectionResponseMessage.class); - withNoSchema(SessionChannelMessage.class); + register(factory, HandshakeWaitMessage.class, msgIdx++); + register(factory, GridIoMessage.class, msgIdx++); + register(factory, IgniteIoTestMessage.class, msgIdx++); + register(factory, GridIoUserMessage.class, msgIdx++); + register(factory, RecoveryLastReceivedMessage.class, msgIdx++); + register(factory, TcpInverseConnectionResponseMessage.class, msgIdx++); + register(factory, SessionChannelMessage.class, msgIdx++); // [11700 - 11800]: Datastreamer messages. msgIdx = 11700; - withNoSchema(DataStreamerUpdatesHandlerResult.class); - withSchema(DataStreamerEntry.class); - withNoSchema(DataStreamerRequest.class); - withNoSchema(DataStreamerResponse.class); + register(factory, DataStreamerUpdatesHandlerResult.class, msgIdx++); + register(factory, DataStreamerEntry.class, msgIdx++); + register(factory, DataStreamerRequest.class, msgIdx++); + register(factory, DataStreamerResponse.class, msgIdx++); // [11900 - 12000]: Metrics, monitoring messages. msgIdx = 11900; - withNoSchema(CacheMetricsMessage.class); - withNoSchema(NodeMetricsMessage.class); - withNoSchema(NodeFullMetricsMessage.class); - withNoSchema(ClusterMetricsUpdateMessage.class); - withNoSchema(TcpDiscoveryClientNodesMetricsMessage.class); - withNoSchema(TcpDiscoveryMetricsUpdateMessage.class); - withNoSchema(TcpDiscoveryClientMetricsUpdateMessage.class); + register(factory, CacheMetricsMessage.class, msgIdx++); + register(factory, NodeMetricsMessage.class, msgIdx++); + register(factory, NodeFullMetricsMessage.class, msgIdx++); + register(factory, ClusterMetricsUpdateMessage.class, msgIdx++); + register(factory, TcpDiscoveryClientNodesMetricsMessage.class, msgIdx++); + register(factory, TcpDiscoveryMetricsUpdateMessage.class, msgIdx++); + register(factory, TcpDiscoveryClientMetricsUpdateMessage.class, msgIdx++); // [12000 - 12100]: Authentication, security messages. msgIdx = 12000; - withNoSchema(User.class); - withNoSchema(UserManagementOperation.class); - withNoSchema(UserManagementOperationFinishedMessage.class); - withNoSchema(UserAuthenticateRequestMessage.class); - withNoSchema(UserAuthenticateResponseMessage.class); - withNoSchema(TcpDiscoveryAuthFailedMessage.class); - withNoSchema(AuthentificationDataBagItem.class); - withNoSchema(SecurityBasicPermissionSet.class); + register(factory, User.class, msgIdx++); + register(factory, UserManagementOperation.class, msgIdx++); + register(factory, UserManagementOperationFinishedMessage.class, msgIdx++); + register(factory, UserAuthenticateRequestMessage.class, msgIdx++); + register(factory, UserAuthenticateResponseMessage.class, msgIdx++); + register(factory, TcpDiscoveryAuthFailedMessage.class, msgIdx++); + register(factory, AuthentificationDataBagItem.class, msgIdx++); + register(factory, SecurityBasicPermissionSet.class, msgIdx++); // [12200 - 12300]: Binary, classloading and marshalling messages. msgIdx = 12200; - withNoSchema(GridDeploymentInfoBean.class); - withNoSchema(GridDeploymentRequest.class); - withNoSchema(GridDeploymentResponse.class); - withNoSchema(MissingMappingRequestMessage.class); - withNoSchema(MissingMappingResponseMessage.class); - withNoSchema(MetadataRequestMessage.class); - withNoSchema(MetadataResponseMessage.class); - withNoSchema(MarshallerMappingItem.class); - withSchema(BinaryMetadataVersionInfo.class); - withNoSchema(CacheBinaryDataBagItem.class); - withNoSchema(MappedName.class); - withNoSchema(MarshallerDataBagItem.class); + register(factory, GridDeploymentInfoBean.class, msgIdx++); + register(factory, GridDeploymentRequest.class, msgIdx++); + register(factory, GridDeploymentResponse.class, msgIdx++); + register(factory, MissingMappingRequestMessage.class, msgIdx++); + register(factory, MissingMappingResponseMessage.class, msgIdx++); + register(factory, MetadataRequestMessage.class, msgIdx++); + register(factory, MetadataResponseMessage.class, msgIdx++); + register(factory, MarshallerMappingItem.class, msgIdx++); + register(factory, BinaryMetadataVersionInfo.class, msgIdx++); + register(factory, CacheBinaryDataBagItem.class, msgIdx++); + register(factory, MappedName.class, msgIdx++); + register(factory, MarshallerDataBagItem.class, msgIdx++); // [12400 - 12500]: Encryption messages. msgIdx = 12400; - withNoSchema(GenerateEncryptionKeyRequest.class); - withNoSchema(GenerateEncryptionKeyResponse.class); - withNoSchema(ChangeCacheEncryptionRequest.class); - withNoSchema(MasterKeyChangeRequest.class); - withNoSchema(GroupKeyEncrypted.class); - withNoSchema(NodeEncryptionKeys.class); - withNoSchema(EncryptionDataBagItem.class); + register(factory, GenerateEncryptionKeyRequest.class, msgIdx++); + register(factory, GenerateEncryptionKeyResponse.class, msgIdx++); + register(factory, ChangeCacheEncryptionRequest.class, msgIdx++); + register(factory, MasterKeyChangeRequest.class, msgIdx++); + register(factory, GroupKeyEncrypted.class, msgIdx++); + register(factory, NodeEncryptionKeys.class, msgIdx++); + register(factory, EncryptionDataBagItem.class, msgIdx++); // [13000 - 13300]: Control, configuration, diagnostics and other messages. msgIdx = 13000; - withSchema(GridEventStorageMessage.class); - withNoSchema(ChangeGlobalStateMessage.class); - withNoSchema(GridChangeGlobalStateMessageResponse.class); - withSchema(IgniteDiagnosticRequest.class); - withNoSchema(IgniteDiagnosticResponse.class); - withNoSchema(WalStateAckMessage.class); - withNoSchema(CacheConfigurationEnrichment.class); - withNoSchema(DynamicCacheChangeRequest.class); - withNoSchema(PartitionHashRecord.class); - withNoSchema(TransactionsHashRecord.class); - withNoSchema(ClusterIdAndTag.class); - withNoSchema(ClusterUpdateNotifierDataBagItem.class); - withNoSchema(PluginsDataBagItem.class); - withSchema(EventsDataBagItem.class); + register(factory, GridEventStorageMessage.class, msgIdx++); + register(factory, ChangeGlobalStateMessage.class, msgIdx++); + register(factory, GridChangeGlobalStateMessageResponse.class, msgIdx++); + register(factory, IgniteDiagnosticRequest.class, msgIdx++); + register(factory, IgniteDiagnosticResponse.class, msgIdx++); + register(factory, WalStateAckMessage.class, msgIdx++); + register(factory, CacheConfigurationEnrichment.class, msgIdx++); + register(factory, DynamicCacheChangeRequest.class, msgIdx++); + register(factory, PartitionHashRecord.class, msgIdx++); + register(factory, TransactionsHashRecord.class, msgIdx++); + register(factory, ClusterIdAndTag.class, msgIdx++); + register(factory, ClusterUpdateNotifierDataBagItem.class, msgIdx++); + register(factory, PluginsDataBagItem.class, msgIdx++); + register(factory, EventsDataBagItem.class, msgIdx++); // [13400 - 13500]: Operation context messages. msgIdx = 13400; - withNoSchema(OperationContextSnapshotMessage.class); - withNoSchema(SecurityContextWrapper.class); + register(factory, OperationContextSnapshotMessage.class, msgIdx++); + register(factory, SecurityContextWrapper.class, msgIdx++); // [13600 - 13700]: Rolling Upgrade messages. msgIdx = 13600; - withNoSchema(IgniteFeatureSet.class); - withNoSchema(IgniteCoreFeatureSet.class); - withNoSchema(IgnitePluginFeatureSet.class); - withNoSchema(RollingUpgradeClusterData.class); + register(factory, IgniteFeatureSet.class, msgIdx++); + register(factory, IgniteCoreFeatureSet.class, msgIdx++); + register(factory, IgnitePluginFeatureSet.class, msgIdx++); + register(factory, RollingUpgradeClusterData.class, msgIdx++); assert msgIdx <= MAX_MESSAGE_ID; } - - /** Registers message using {@link #dfltMarsh}. */ - private void withNoSchema(Class cls) { - assert cls.getAnnotation(UseBinaryMarshaller.class) == null : - "Remove @" + UseBinaryMarshaller.class.getSimpleName() + " for class: " + cls.getSimpleName(); - - register(factory, cls, msgIdx++); - } - - /** Registers message using {@link #schemaAwareMarsh}. */ - private void withSchema(Class cls) { - assert cls.getAnnotation(UseBinaryMarshaller.class) != null : - "Add @" + UseBinaryMarshaller.class.getSimpleName() + " for class: " + cls.getSimpleName(); - - register(factory, cls, msgIdx++); - } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java index 05299f7b438f1..2cdcdb91b0912 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java @@ -68,7 +68,7 @@ protected void register(IgniteMessageFactory factory, Class< /** */ private static void register(IgniteMessageFactory factory, Class cls, short id, Marshaller marsh) { - MessageSerializer serializer = requireGenerated(cls, "Serializer", marsh); + MessageSerializer serializer = requireGenerated(cls, "Serializer"); // A MarshallableMessage always gets a generated marshaller (the hook call alone is a statement), so its // absence is a build problem. For the rest the generator skips statement-free marshallers, so absence @@ -79,52 +79,82 @@ private static void register(IgniteMessageFactory factory, C if (NonMarshallableMessage.class.isAssignableFrom(cls)) marshaller = null; else if (MarshallableMessage.class.isAssignableFrom(cls)) - marshaller = requireGenerated(cls, "Marshaller", marsh); + marshaller = requireMarshaller(cls, marsh); else - marshaller = loadGenerated(cls, "Marshaller", marsh); + marshaller = loadMarshaller(cls, marsh); // Deployers are generated for GridCacheMessage subclasses only, so the class lookup is skipped for the rest; // a DeployableMessage left without a deployer is then rejected at registration. GridCacheMessageDeployer deployer = GridCacheMessage.class.isAssignableFrom(cls) - ? loadGenerated(cls, "Deployer", marsh) + ? loadGenerated(cls, "Deployer") : null; factory.register(id, serializer, marshaller, deployer); } /** Loads the generated companion like {@link #loadGenerated}, failing fast when it is missing. */ - private static T requireGenerated(Class cls, String suffix, Marshaller marsh) { - T res = loadGenerated(cls, suffix, marsh); + private static T requireGenerated(Class cls, String suffix) { + return require(loadGenerated(cls, suffix), cls, suffix); + } + + /** Loads the generated marshaller like {@link #loadMarshaller}, failing fast when it is missing. */ + private static T requireMarshaller(Class cls, Marshaller marsh) { + return require(loadMarshaller(cls, marsh), cls, "Marshaller"); + } - if (res == null) { + /** */ + private static T require(@Nullable T companion, Class cls, String suffix) { + if (companion == null) { throw new IgniteException("No " + cls.getSimpleName() + suffix + " found for " + cls.getName() + ". Either the class is not processed by codegen or the generated sources are stale," + " try 'mvn clean install'."); } - return res; + return companion; } /** - * Instantiates the generated companion class {@code Serializer/Marshaller/Deployer}, or returns - * {@code null} when it does not exist. The sole declared constructor is used, passing {@code marsh} when it takes - * one. Constructor lookups, including missing companions, are cached per message class in {@link #COMPANIONS}. + * Instantiates the generated companion class {@code Serializer/Deployer}, or returns {@code null} when it + * does not exist. Neither takes a marshaller: the serializer writes the wire fields as they are, and the deployer + * only walks cache objects. */ @SuppressWarnings("unchecked") - private static @Nullable T loadGenerated(Class cls, String suffix, Marshaller marsh) { + private static @Nullable T loadGenerated(Class cls, String suffix) { Constructor ctor = COMPANIONS.get(cls).ctor(suffix); if (ctor == null) return null; + assert ctor.getParameterCount() == 0 : cls.getSimpleName() + suffix + " must have a no-arg constructor"; + try { - return (T)(ctor.getParameterCount() == 0 ? ctor.newInstance() : ctor.newInstance(marsh)); + return (T)ctor.newInstance(); } catch (Exception e) { throw new IgniteException("Failed to instantiate " + cls.getSimpleName() + suffix, e); } } + /** + * Instantiates the generated {@code Marshaller}, or returns {@code null} when it does not exist. The + * generator gives it a {@code Marshaller} constructor only when the message has fields to marshal with one; + * otherwise the companion just walks nested messages and cache objects, and takes no arguments. + */ + @SuppressWarnings("unchecked") + private static @Nullable T loadMarshaller(Class cls, Marshaller marsh) { + Constructor ctor = COMPANIONS.get(cls).ctor("Marshaller"); + + if (ctor == null) + return null; + + try { + return (T)(ctor.getParameterCount() == 0 ? ctor.newInstance() : ctor.newInstance(marsh)); + } + catch (Exception e) { + throw new IgniteException("Failed to instantiate " + cls.getSimpleName() + "Marshaller", e); + } + } + /** @return the sole public constructor of the generated companion {@code }, or {@code null} when it does not exist. */ private static @Nullable Constructor companionCtor(Class cls, String suffix) { try { From 849ff55e6eb52835d8165b1a36e4bc1bee00a313 Mon Sep 17 00:00:00 2001 From: Anton Vinogradov Date: Fri, 31 Jul 2026 19:33:22 +0300 Subject: [PATCH 2/4] IGNITE-28937 Review fixes: keep a short register(cls), inline single-use helpers Co-Authored-By: Claude Opus 5 --- .../ignite/internal/CoreMessagesProvider.java | 604 +++++++++--------- ...actMarshallableMessageFactoryProvider.java | 16 +- 2 files changed, 308 insertions(+), 312 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 4e07d4a25d366..df35e9ce099cc 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 @@ -269,6 +269,7 @@ import org.apache.ignite.internal.util.distributed.SingleNodeMessage; import org.apache.ignite.marshaller.Marshaller; import org.apache.ignite.marshaller.jdk.JdkMarshaller; +import org.apache.ignite.plugin.extensions.communication.Message; import org.apache.ignite.plugin.security.SecurityBasicPermissionSet; import org.apache.ignite.spi.collision.jobstealing.JobStealingRequest; import org.apache.ignite.spi.communication.tcp.internal.TcpConnectionRequestDiscoveryMessage; @@ -365,364 +366,369 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { // [5000 - 5500]: Utility messages. Most of them originally come from Discovery. msgIdx = 5000; - register(factory, CompressedMessage.class, msgIdx++); - register(factory, ErrorMessage.class, msgIdx++); - register(factory, InetSocketAddressMessage.class, msgIdx++); - register(factory, InetAddressMessage.class, msgIdx++); - register(factory, TcpDiscoveryNode.class, msgIdx++); - register(factory, DiscoveryDataPacket.class, msgIdx++); - register(factory, GridByteArrayList.class, msgIdx++); - register(factory, CacheVersionedValue.class, msgIdx++); - register(factory, KeyedVersionedValue.class, msgIdx++); - register(factory, WALPointer.class, msgIdx++); - register(factory, SerializableDataBagItemWrapper.class, msgIdx++); - register(factory, GridTopicMessage.class, msgIdx++); - register(factory, GridIntList.class, msgIdx++); + register(CompressedMessage.class); + register(ErrorMessage.class); + register(InetSocketAddressMessage.class); + register(InetAddressMessage.class); + register(TcpDiscoveryNode.class); + register(DiscoveryDataPacket.class); + register(GridByteArrayList.class); + register(CacheVersionedValue.class); + register(KeyedVersionedValue.class); + register(WALPointer.class); + register(SerializableDataBagItemWrapper.class); + register(GridTopicMessage.class); + register(GridIntList.class); // [5700 - 5900]: Discovery originated messages. msgIdx = 5700; - register(factory, TcpDiscoveryHandshakeRequest.class, msgIdx++); - register(factory, TcpDiscoveryHandshakeResponse.class, msgIdx++); - register(factory, TcpDiscoveryJoinRequestMessage.class, msgIdx++); - register(factory, TcpDiscoveryNodeAddedMessage.class, msgIdx++); - register(factory, TcpDiscoveryNodeAddFinishedMessage.class, msgIdx++); - register(factory, TcpDiscoveryNodeLeftMessage.class, msgIdx++); - register(factory, TcpDiscoveryNodeFailedMessage.class, msgIdx++); - register(factory, TcpDiscoveryConnectionCheckMessage.class, msgIdx++); - register(factory, TcpDiscoveryPingRequest.class, msgIdx++); - register(factory, TcpDiscoveryPingResponse.class, msgIdx++); - register(factory, TcpDiscoveryClientPingRequest.class, msgIdx++); - register(factory, TcpDiscoveryClientPingResponse.class, msgIdx++); - register(factory, TcpDiscoveryClientAckResponse.class, msgIdx++); - register(factory, TcpDiscoveryClientReconnectMessage.class, msgIdx++); - register(factory, TcpDiscoveryDiscardMessage.class, msgIdx++); - register(factory, TcpDiscoveryCheckFailedMessage.class, msgIdx++); - register(factory, TcpDiscoveryLoopbackProblemMessage.class, msgIdx++); - register(factory, TcpDiscoveryRingLatencyCheckMessage.class, msgIdx++); - register(factory, TcpDiscoveryDuplicateIdMessage.class, msgIdx++); - register(factory, TcpDiscoveryCustomEventMessage.class, msgIdx++); - register(factory, TcpDiscoveryServerOnlyCustomEventMessage.class, msgIdx++); + register(TcpDiscoveryHandshakeRequest.class); + register(TcpDiscoveryHandshakeResponse.class); + register(TcpDiscoveryJoinRequestMessage.class); + register(TcpDiscoveryNodeAddedMessage.class); + register(TcpDiscoveryNodeAddFinishedMessage.class); + register(TcpDiscoveryNodeLeftMessage.class); + register(TcpDiscoveryNodeFailedMessage.class); + register(TcpDiscoveryConnectionCheckMessage.class); + register(TcpDiscoveryPingRequest.class); + register(TcpDiscoveryPingResponse.class); + register(TcpDiscoveryClientPingRequest.class); + register(TcpDiscoveryClientPingResponse.class); + register(TcpDiscoveryClientAckResponse.class); + register(TcpDiscoveryClientReconnectMessage.class); + register(TcpDiscoveryDiscardMessage.class); + register(TcpDiscoveryCheckFailedMessage.class); + register(TcpDiscoveryLoopbackProblemMessage.class); + register(TcpDiscoveryRingLatencyCheckMessage.class); + register(TcpDiscoveryDuplicateIdMessage.class); + register(TcpDiscoveryCustomEventMessage.class); + register(TcpDiscoveryServerOnlyCustomEventMessage.class); msgIdx = 5900; - register(factory, TcpDiscoveryStatusCheckMessage.class, msgIdx++); + register(TcpDiscoveryStatusCheckMessage.class); // [6000 - 6200]: Snapshot operation messages. Most of them originally come from Discovery. msgIdx = 6000; - register(factory, SnapshotStartDiscoveryMessage.class, msgIdx++); - register(factory, SnapshotCheckProcessRequest.class, msgIdx++); - register(factory, SnapshotOperationRequest.class, msgIdx++); - register(factory, SnapshotOperationEndRequest.class, msgIdx++); - register(factory, SnapshotRestoreStartRequest.class, msgIdx++); - register(factory, SnapshotOperationResponse.class, msgIdx++); - register(factory, SnapshotHandlerResult.class, msgIdx++); - register(factory, SnapshotCheckResponse.class, msgIdx++); - register(factory, SnapshotPartitionsVerifyHandlerResponse.class, msgIdx++); - register(factory, SnapshotRestoreOperationResponse.class, msgIdx++); - register(factory, SnapshotMetadataResponse.class, msgIdx++); - register(factory, SnapshotMetadata.class, msgIdx++); - register(factory, SnapshotCheckPartitionHashesResponse.class, msgIdx++); - register(factory, SnapshotCheckHandlersResponse.class, msgIdx++); - register(factory, SnapshotFilesRequestMessage.class, msgIdx++); - register(factory, SnapshotFilesFailureMessage.class, msgIdx++); - register(factory, IncrementalSnapshotVerifyResult.class, msgIdx++); - register(factory, IncrementalSnapshotAwareMessage.class, msgIdx++); + register(SnapshotStartDiscoveryMessage.class); + register(SnapshotCheckProcessRequest.class); + register(SnapshotOperationRequest.class); + register(SnapshotOperationEndRequest.class); + register(SnapshotRestoreStartRequest.class); + register(SnapshotOperationResponse.class); + register(SnapshotHandlerResult.class); + register(SnapshotCheckResponse.class); + register(SnapshotPartitionsVerifyHandlerResponse.class); + register(SnapshotRestoreOperationResponse.class); + register(SnapshotMetadataResponse.class); + register(SnapshotMetadata.class); + register(SnapshotCheckPartitionHashesResponse.class); + register(SnapshotCheckHandlersResponse.class); + register(SnapshotFilesRequestMessage.class); + register(SnapshotFilesFailureMessage.class); + register(IncrementalSnapshotVerifyResult.class); + register(IncrementalSnapshotAwareMessage.class); // [6300 - 6400]: Services messages. Most of them originally come from Discovery. msgIdx = 6300; - register(factory, ServiceDeploymentProcessId.class, msgIdx++); - register(factory, ServiceSingleNodeDeploymentResult.class, msgIdx++); - register(factory, ServiceClusterDeploymentResult.class, msgIdx++); - register(factory, ServiceDeploymentRequest.class, msgIdx++); - register(factory, ServiceUndeploymentRequest.class, msgIdx++); - register(factory, ServiceClusterDeploymentResultBatch.class, msgIdx++); - register(factory, ServiceChangeBatchRequest.class, msgIdx++); - register(factory, ServiceSingleNodeDeploymentResultBatch.class, msgIdx++); - register(factory, ServiceProcessorCommonDiscoveryData.class, msgIdx++); - register(factory, ServiceProcessorJoinNodeDiscoveryData.class, msgIdx++); - register(factory, ServiceInfo.class, msgIdx++); - register(factory, ServiceTopology.class, msgIdx++); - register(factory, LazyServiceConfigurationMessage.class, msgIdx++); + register(ServiceDeploymentProcessId.class); + register(ServiceSingleNodeDeploymentResult.class); + register(ServiceClusterDeploymentResult.class); + register(ServiceDeploymentRequest.class); + register(ServiceUndeploymentRequest.class); + register(ServiceClusterDeploymentResultBatch.class); + register(ServiceChangeBatchRequest.class); + register(ServiceSingleNodeDeploymentResultBatch.class); + register(ServiceProcessorCommonDiscoveryData.class); + register(ServiceProcessorJoinNodeDiscoveryData.class); + register(ServiceInfo.class); + register(ServiceTopology.class); + register(LazyServiceConfigurationMessage.class); // [6500 - 6700]: DiscoveryCustomMessage msgIdx = 6500; - register(factory, TcpConnectionRequestDiscoveryMessage.class, msgIdx++); - register(factory, DistributedMetaStorageUpdateMessage.class, msgIdx++); - register(factory, DistributedMetaStorageUpdateAckMessage.class, msgIdx++); - register(factory, DistributedMetaStorageCasMessage.class, msgIdx++); - register(factory, DistributedMetaStorageCasAckMessage.class, msgIdx++); - register(factory, FullMessage.class, msgIdx++); - register(factory, InitMessage.class, msgIdx++); - register(factory, CacheStatisticsModeChangeMessage.class, msgIdx++); - register(factory, MetadataRemoveAcceptedMessage.class, msgIdx++); - register(factory, MetadataRemoveProposedMessage.class, msgIdx++); - register(factory, WalStateFinishMessage.class, msgIdx++); - register(factory, WalStateProposeMessage.class, msgIdx++); - register(factory, MetadataUpdateAcceptedMessage.class, msgIdx++); - register(factory, MetadataUpdateProposedMessage.class, msgIdx++); - register(factory, TxTimeoutOnPartitionMapExchangeChangeMessage.class, msgIdx++); - register(factory, UserAcceptedMessage.class, msgIdx++); - register(factory, UserProposedMessage.class, msgIdx++); - register(factory, ChangeGlobalStateFinishMessage.class, msgIdx++); - register(factory, StopRoutineAckDiscoveryMessage.class, msgIdx++); - register(factory, StopRoutineDiscoveryMessage.class, msgIdx++); - register(factory, CacheAffinityChangeMessage.class, msgIdx++); - register(factory, ClientCacheChangeDiscoveryMessage.class, msgIdx++); - register(factory, MappingAcceptedMessage.class, msgIdx++); - register(factory, MappingProposedMessage.class, msgIdx++); - register(factory, ExchangeFailureMessage.class, msgIdx++); - register(factory, CacheStatisticsClearMessage.class, msgIdx++); - register(factory, ClientCacheChangeDummyDiscoveryMessage.class, msgIdx++); - register(factory, DynamicCacheChangeBatch.class, msgIdx++); - register(factory, CacheClientReconnectDiscoveryData.class, msgIdx++); - register(factory, CacheGroupRecoveryState.class, msgIdx++); - register(factory, CacheJoinInfo.class, msgIdx++); - register(factory, CacheJoinNodeDiscoveryData.class, msgIdx++); - register(factory, CacheReconnectInfo.class, msgIdx++); - register(factory, ClusterCacheGroupRecoveryData.class, msgIdx++); + register(TcpConnectionRequestDiscoveryMessage.class); + register(DistributedMetaStorageUpdateMessage.class); + register(DistributedMetaStorageUpdateAckMessage.class); + register(DistributedMetaStorageCasMessage.class); + register(DistributedMetaStorageCasAckMessage.class); + register(FullMessage.class); + register(InitMessage.class); + register(CacheStatisticsModeChangeMessage.class); + register(MetadataRemoveAcceptedMessage.class); + register(MetadataRemoveProposedMessage.class); + register(WalStateFinishMessage.class); + register(WalStateProposeMessage.class); + register(MetadataUpdateAcceptedMessage.class); + register(MetadataUpdateProposedMessage.class); + register(TxTimeoutOnPartitionMapExchangeChangeMessage.class); + register(UserAcceptedMessage.class); + register(UserProposedMessage.class); + register(ChangeGlobalStateFinishMessage.class); + register(StopRoutineAckDiscoveryMessage.class); + register(StopRoutineDiscoveryMessage.class); + register(CacheAffinityChangeMessage.class); + register(ClientCacheChangeDiscoveryMessage.class); + register(MappingAcceptedMessage.class); + register(MappingProposedMessage.class); + register(ExchangeFailureMessage.class); + register(CacheStatisticsClearMessage.class); + register(ClientCacheChangeDummyDiscoveryMessage.class); + register(DynamicCacheChangeBatch.class); + register(CacheClientReconnectDiscoveryData.class); + register(CacheGroupRecoveryState.class); + register(CacheJoinInfo.class); + register(CacheJoinNodeDiscoveryData.class); + register(CacheReconnectInfo.class); + register(ClusterCacheGroupRecoveryData.class); // [10000 - 10200]: Transaction and lock related messages. Most of them originally comes from Communication. msgIdx = 10000; - register(factory, TxInfo.class, msgIdx++); - register(factory, TxEntriesInfo.class, msgIdx++); - register(factory, TxLock.class, msgIdx++); - register(factory, TxLocksRequest.class, msgIdx++); - register(factory, TxLocksResponse.class, msgIdx++); - register(factory, IgniteTxKey.class, msgIdx++); - register(factory, IgniteTxEntry.class, msgIdx++); - register(factory, TxEntryValueHolder.class, msgIdx++); - register(factory, GridCacheTxRecoveryRequest.class, msgIdx++); - register(factory, GridCacheTxRecoveryResponse.class, msgIdx++); - register(factory, GridDistributedTxFinishRequest.class, msgIdx++); - register(factory, GridDistributedTxFinishResponse.class, msgIdx++); - register(factory, GridDistributedTxPrepareRequest.class, msgIdx++); - register(factory, GridDistributedTxPrepareResponse.class, msgIdx++); - register(factory, GridDhtTxFinishRequest.class, msgIdx++); - register(factory, GridDhtTxFinishResponse.class, msgIdx++); - register(factory, GridDhtTxPrepareRequest.class, msgIdx++); - register(factory, GridDhtTxPrepareResponse.class, msgIdx++); - register(factory, GridNearTxFinishRequest.class, msgIdx++); - register(factory, GridNearTxFinishResponse.class, msgIdx++); - register(factory, GridNearTxPrepareRequest.class, msgIdx++); - register(factory, GridNearTxPrepareResponse.class, msgIdx++); - register(factory, GridDhtLockRequest.class, msgIdx++); - register(factory, GridDhtLockResponse.class, msgIdx++); - register(factory, GridDhtUnlockRequest.class, msgIdx++); - register(factory, GridNearLockRequest.class, msgIdx++); - register(factory, GridNearLockResponse.class, msgIdx++); - register(factory, GridNearUnlockRequest.class, msgIdx++); - register(factory, GridDistributedLockRequest.class, msgIdx++); - register(factory, GridDistributedLockResponse.class, msgIdx++); - register(factory, GridDhtTxOnePhaseCommitAckRequest.class, msgIdx++); - register(factory, TransactionAttributesAwareRequest.class, msgIdx++); + register(TxInfo.class); + register(TxEntriesInfo.class); + register(TxLock.class); + register(TxLocksRequest.class); + register(TxLocksResponse.class); + register(IgniteTxKey.class); + register(IgniteTxEntry.class); + register(TxEntryValueHolder.class); + register(GridCacheTxRecoveryRequest.class); + register(GridCacheTxRecoveryResponse.class); + register(GridDistributedTxFinishRequest.class); + register(GridDistributedTxFinishResponse.class); + register(GridDistributedTxPrepareRequest.class); + register(GridDistributedTxPrepareResponse.class); + register(GridDhtTxFinishRequest.class); + register(GridDhtTxFinishResponse.class); + register(GridDhtTxPrepareRequest.class); + register(GridDhtTxPrepareResponse.class); + register(GridNearTxFinishRequest.class); + register(GridNearTxFinishResponse.class); + register(GridNearTxPrepareRequest.class); + register(GridNearTxPrepareResponse.class); + register(GridDhtLockRequest.class); + register(GridDhtLockResponse.class); + register(GridDhtUnlockRequest.class); + register(GridNearLockRequest.class); + register(GridNearLockResponse.class); + register(GridNearUnlockRequest.class); + register(GridDistributedLockRequest.class); + register(GridDistributedLockResponse.class); + register(GridDhtTxOnePhaseCommitAckRequest.class); + register(TransactionAttributesAwareRequest.class); // [10300 - 10500]: Cache, DHT messages. msgIdx = 10300; - register(factory, GridDhtForceKeysRequest.class, msgIdx++); - register(factory, GridDhtForceKeysResponse.class, msgIdx++); - register(factory, GridDhtAtomicDeferredUpdateResponse.class, msgIdx++); - register(factory, GridDhtAtomicUpdateRequest.class, msgIdx++); - register(factory, GridDhtAtomicUpdateResponse.class, msgIdx++); - register(factory, GridNearAtomicFullUpdateRequest.class, msgIdx++); - register(factory, GridDhtAtomicSingleUpdateRequest.class, msgIdx++); - register(factory, GridNearAtomicUpdateResponse.class, msgIdx++); - register(factory, GridNearAtomicSingleUpdateRequest.class, msgIdx++); - register(factory, GridNearAtomicSingleUpdateInvokeRequest.class, msgIdx++); - register(factory, GridNearAtomicSingleUpdateFilterRequest.class, msgIdx++); - register(factory, GridNearAtomicCheckUpdateRequest.class, msgIdx++); - register(factory, NearCacheUpdates.class, msgIdx++); - register(factory, GridNearGetRequest.class, msgIdx++); - register(factory, GridNearGetResponse.class, msgIdx++); - register(factory, GridNearSingleGetRequest.class, msgIdx++); - register(factory, GridNearSingleGetResponse.class, msgIdx++); - register(factory, GridDhtAtomicNearResponse.class, msgIdx++); - register(factory, GridCacheTtlUpdateRequest.class, msgIdx++); - register(factory, GridCacheReturn.class, msgIdx++); - register(factory, GridCacheEntryInfo.class, msgIdx++); - register(factory, CacheInvokeDirectResult.class, msgIdx++); - register(factory, GridCacheRawVersionedEntry.class, msgIdx++); - register(factory, CacheEvictionEntry.class, msgIdx++); - register(factory, CacheEntryPredicateAdapter.class, msgIdx++); - register(factory, GridContinuousMessage.class, msgIdx++); - register(factory, ContinuousRoutineStartResultMessage.class, msgIdx++); - register(factory, UpdateErrors.class, msgIdx++); - register(factory, LatchAckMessage.class, msgIdx++); - register(factory, AtomicApplicationAttributesAwareRequest.class, msgIdx++); - register(factory, StartRequestData.class, msgIdx++); - register(factory, StartRoutineAckDiscoveryMessage.class, msgIdx++); - register(factory, StartRoutineDiscoveryMessage.class, msgIdx++); - register(factory, StoredCacheData.class, msgIdx++); + register(GridDhtForceKeysRequest.class); + register(GridDhtForceKeysResponse.class); + register(GridDhtAtomicDeferredUpdateResponse.class); + register(GridDhtAtomicUpdateRequest.class); + register(GridDhtAtomicUpdateResponse.class); + register(GridNearAtomicFullUpdateRequest.class); + register(GridDhtAtomicSingleUpdateRequest.class); + register(GridNearAtomicUpdateResponse.class); + register(GridNearAtomicSingleUpdateRequest.class); + register(GridNearAtomicSingleUpdateInvokeRequest.class); + register(GridNearAtomicSingleUpdateFilterRequest.class); + register(GridNearAtomicCheckUpdateRequest.class); + register(NearCacheUpdates.class); + register(GridNearGetRequest.class); + register(GridNearGetResponse.class); + register(GridNearSingleGetRequest.class); + register(GridNearSingleGetResponse.class); + register(GridDhtAtomicNearResponse.class); + register(GridCacheTtlUpdateRequest.class); + register(GridCacheReturn.class); + register(GridCacheEntryInfo.class); + register(CacheInvokeDirectResult.class); + register(GridCacheRawVersionedEntry.class); + register(CacheEvictionEntry.class); + register(CacheEntryPredicateAdapter.class); + register(GridContinuousMessage.class); + register(ContinuousRoutineStartResultMessage.class); + register(UpdateErrors.class); + register(LatchAckMessage.class); + register(AtomicApplicationAttributesAwareRequest.class); + register(StartRequestData.class); + register(StartRoutineAckDiscoveryMessage.class); + register(StartRoutineDiscoveryMessage.class); + register(StoredCacheData.class); // [10600-10800]: Affinity & partition maps. msgIdx = 10600; - register(factory, GridDhtAffinityAssignmentRequest.class, msgIdx++); - register(factory, GridDhtAffinityAssignmentResponse.class, msgIdx++); - register(factory, CacheGroupAffinityMessage.class, msgIdx++); - register(factory, ExchangeInfo.class, msgIdx++); - register(factory, PartitionUpdateCountersMessage.class, msgIdx++); - register(factory, CachePartitionPartialCountersMap.class, msgIdx++); - register(factory, IgniteDhtDemandedPartitionsMap.class, msgIdx++); - register(factory, CachePartitionFullCountersMap.class, msgIdx++); - register(factory, GroupPartitionIdPair.class, msgIdx++); - register(factory, GridPartitionStateMap.class, msgIdx++); - register(factory, GridDhtPartitionMap.class, msgIdx++); - register(factory, GridDhtPartitionFullMap.class, msgIdx++); - register(factory, GridDhtPartitionExchangeId.class, msgIdx++); - register(factory, GridCheckpointRequest.class, msgIdx++); - register(factory, GridDhtPartitionDemandMessage.class, msgIdx++); - register(factory, GridDhtPartitionSupplyMessage.class, msgIdx++); - register(factory, GridDhtPartitionsFullMessage.class, msgIdx++); - register(factory, GridDhtPartitionsSingleMessage.class, msgIdx++); - register(factory, GridDhtPartitionsSingleRequest.class, msgIdx++); - register(factory, PartitionKey.class, msgIdx++); + register(GridDhtAffinityAssignmentRequest.class); + register(GridDhtAffinityAssignmentResponse.class); + register(CacheGroupAffinityMessage.class); + register(ExchangeInfo.class); + register(PartitionUpdateCountersMessage.class); + register(CachePartitionPartialCountersMap.class); + register(IgniteDhtDemandedPartitionsMap.class); + register(CachePartitionFullCountersMap.class); + register(GroupPartitionIdPair.class); + register(GridPartitionStateMap.class); + register(GridDhtPartitionMap.class); + register(GridDhtPartitionFullMap.class); + register(GridDhtPartitionExchangeId.class); + register(GridCheckpointRequest.class); + register(GridDhtPartitionDemandMessage.class); + register(GridDhtPartitionSupplyMessage.class); + register(GridDhtPartitionsFullMessage.class); + register(GridDhtPartitionsSingleMessage.class); + register(GridDhtPartitionsSingleRequest.class); + register(PartitionKey.class); // [10900-11100]: Query, schema and SQL related messages. msgIdx = 10900; - register(factory, SchemaAlterTableAddColumnOperation.class, msgIdx++); - register(factory, SchemaIndexCreateOperation.class, msgIdx++); - register(factory, SchemaIndexDropOperation.class, msgIdx++); - register(factory, SchemaAlterTableDropColumnOperation.class, msgIdx++); - register(factory, SchemaAddQueryEntityOperation.class, msgIdx++); - register(factory, SchemaOperationStatusMessage.class, msgIdx++); - register(factory, SchemaProposeDiscoveryMessage.class, msgIdx++); - register(factory, SchemaFinishDiscoveryMessage.class, msgIdx++); - register(factory, QueryField.class, msgIdx++); - register(factory, QueryIndexMessage.class, msgIdx++); - register(factory, GridCacheSqlQuery.class, msgIdx++); - register(factory, GridCacheQueryRequest.class, msgIdx++); - register(factory, GridCacheQueryResponse.class, msgIdx++); - register(factory, GridQueryCancelRequest.class, msgIdx++); - register(factory, GridQueryFailResponse.class, msgIdx++); - register(factory, GridQueryNextPageRequest.class, msgIdx++); - register(factory, GridQueryNextPageResponse.class, msgIdx++); - register(factory, GridQueryKillRequest.class, msgIdx++); - register(factory, GridQueryKillResponse.class, msgIdx++); - register(factory, IndexKeyDefinition.class, msgIdx++); - register(factory, IndexKeyTypeSettings.class, msgIdx++); - register(factory, IndexQueryResultMeta.class, msgIdx++); - register(factory, StatisticsKeyMessage.class, msgIdx++); - register(factory, StatisticsDecimalMessage.class, msgIdx++); - register(factory, StatisticsObjectData.class, msgIdx++); - register(factory, StatisticsColumnData.class, msgIdx++); - register(factory, StatisticsRequest.class, msgIdx++); - register(factory, StatisticsResponse.class, msgIdx++); - register(factory, CacheContinuousQueryBatchAck.class, msgIdx++); - register(factory, GridDhtTxSalvageMessage.class, msgIdx++); - register(factory, CacheContinuousQueryEntry.class, msgIdx++); - register(factory, QueryInlineSizesDataBagItem.class, msgIdx++); - register(factory, QueryProposalsDataBagItem.class, msgIdx++); - register(factory, QueryEntityMessage.class, msgIdx++); - register(factory, QueryEntityExMessage.class, msgIdx++); + register(SchemaAlterTableAddColumnOperation.class); + register(SchemaIndexCreateOperation.class); + register(SchemaIndexDropOperation.class); + register(SchemaAlterTableDropColumnOperation.class); + register(SchemaAddQueryEntityOperation.class); + register(SchemaOperationStatusMessage.class); + register(SchemaProposeDiscoveryMessage.class); + register(SchemaFinishDiscoveryMessage.class); + register(QueryField.class); + register(QueryIndexMessage.class); + register(GridCacheSqlQuery.class); + register(GridCacheQueryRequest.class); + register(GridCacheQueryResponse.class); + register(GridQueryCancelRequest.class); + register(GridQueryFailResponse.class); + register(GridQueryNextPageRequest.class); + register(GridQueryNextPageResponse.class); + register(GridQueryKillRequest.class); + register(GridQueryKillResponse.class); + register(IndexKeyDefinition.class); + register(IndexKeyTypeSettings.class); + register(IndexQueryResultMeta.class); + register(StatisticsKeyMessage.class); + register(StatisticsDecimalMessage.class); + register(StatisticsObjectData.class); + register(StatisticsColumnData.class); + register(StatisticsRequest.class); + register(StatisticsResponse.class); + register(CacheContinuousQueryBatchAck.class); + register(GridDhtTxSalvageMessage.class); + register(CacheContinuousQueryEntry.class); + register(QueryInlineSizesDataBagItem.class); + register(QueryProposalsDataBagItem.class); + register(QueryEntityMessage.class); + register(QueryEntityExMessage.class); // [11200 - 11300]: Compute, distributed process messages. msgIdx = 11200; - register(factory, GridJobCancelRequest.class, msgIdx++); - register(factory, GridJobExecuteRequest.class, msgIdx++); - register(factory, GridJobExecuteResponse.class, msgIdx++); - register(factory, GridJobSiblingsRequest.class, msgIdx++); - register(factory, GridJobSiblingsResponse.class, msgIdx++); - register(factory, GridTaskCancelRequest.class, msgIdx++); - register(factory, GridTaskSessionRequest.class, msgIdx++); - register(factory, GridTaskResultRequest.class, msgIdx++); - register(factory, GridTaskResultResponse.class, msgIdx++); - register(factory, JobStealingRequest.class, msgIdx++); - register(factory, SingleNodeMessage.class, msgIdx++); + register(GridJobCancelRequest.class); + register(GridJobExecuteRequest.class); + register(GridJobExecuteResponse.class); + register(GridJobSiblingsRequest.class); + register(GridJobSiblingsResponse.class); + register(GridTaskCancelRequest.class); + register(GridTaskSessionRequest.class); + register(GridTaskResultRequest.class); + register(GridTaskResultResponse.class); + register(JobStealingRequest.class); + register(SingleNodeMessage.class); // [11500 - 11600]: IO, networking messages. msgIdx = NODE_ID_MSG_TYPE; - register(factory, NodeIdMessage.class, msgIdx++); + register(NodeIdMessage.class); msgIdx = HANDSHAKE_MSG_TYPE; - register(factory, HandshakeMessage.class, msgIdx++); + register(HandshakeMessage.class); msgIdx = HANDSHAKE_WAIT_MSG_TYPE; - register(factory, HandshakeWaitMessage.class, msgIdx++); - register(factory, GridIoMessage.class, msgIdx++); - register(factory, IgniteIoTestMessage.class, msgIdx++); - register(factory, GridIoUserMessage.class, msgIdx++); - register(factory, RecoveryLastReceivedMessage.class, msgIdx++); - register(factory, TcpInverseConnectionResponseMessage.class, msgIdx++); - register(factory, SessionChannelMessage.class, msgIdx++); + register(HandshakeWaitMessage.class); + register(GridIoMessage.class); + register(IgniteIoTestMessage.class); + register(GridIoUserMessage.class); + register(RecoveryLastReceivedMessage.class); + register(TcpInverseConnectionResponseMessage.class); + register(SessionChannelMessage.class); // [11700 - 11800]: Datastreamer messages. msgIdx = 11700; - register(factory, DataStreamerUpdatesHandlerResult.class, msgIdx++); - register(factory, DataStreamerEntry.class, msgIdx++); - register(factory, DataStreamerRequest.class, msgIdx++); - register(factory, DataStreamerResponse.class, msgIdx++); + register(DataStreamerUpdatesHandlerResult.class); + register(DataStreamerEntry.class); + register(DataStreamerRequest.class); + register(DataStreamerResponse.class); // [11900 - 12000]: Metrics, monitoring messages. msgIdx = 11900; - register(factory, CacheMetricsMessage.class, msgIdx++); - register(factory, NodeMetricsMessage.class, msgIdx++); - register(factory, NodeFullMetricsMessage.class, msgIdx++); - register(factory, ClusterMetricsUpdateMessage.class, msgIdx++); - register(factory, TcpDiscoveryClientNodesMetricsMessage.class, msgIdx++); - register(factory, TcpDiscoveryMetricsUpdateMessage.class, msgIdx++); - register(factory, TcpDiscoveryClientMetricsUpdateMessage.class, msgIdx++); + register(CacheMetricsMessage.class); + register(NodeMetricsMessage.class); + register(NodeFullMetricsMessage.class); + register(ClusterMetricsUpdateMessage.class); + register(TcpDiscoveryClientNodesMetricsMessage.class); + register(TcpDiscoveryMetricsUpdateMessage.class); + register(TcpDiscoveryClientMetricsUpdateMessage.class); // [12000 - 12100]: Authentication, security messages. msgIdx = 12000; - register(factory, User.class, msgIdx++); - register(factory, UserManagementOperation.class, msgIdx++); - register(factory, UserManagementOperationFinishedMessage.class, msgIdx++); - register(factory, UserAuthenticateRequestMessage.class, msgIdx++); - register(factory, UserAuthenticateResponseMessage.class, msgIdx++); - register(factory, TcpDiscoveryAuthFailedMessage.class, msgIdx++); - register(factory, AuthentificationDataBagItem.class, msgIdx++); - register(factory, SecurityBasicPermissionSet.class, msgIdx++); + register(User.class); + register(UserManagementOperation.class); + register(UserManagementOperationFinishedMessage.class); + register(UserAuthenticateRequestMessage.class); + register(UserAuthenticateResponseMessage.class); + register(TcpDiscoveryAuthFailedMessage.class); + register(AuthentificationDataBagItem.class); + register(SecurityBasicPermissionSet.class); // [12200 - 12300]: Binary, classloading and marshalling messages. msgIdx = 12200; - register(factory, GridDeploymentInfoBean.class, msgIdx++); - register(factory, GridDeploymentRequest.class, msgIdx++); - register(factory, GridDeploymentResponse.class, msgIdx++); - register(factory, MissingMappingRequestMessage.class, msgIdx++); - register(factory, MissingMappingResponseMessage.class, msgIdx++); - register(factory, MetadataRequestMessage.class, msgIdx++); - register(factory, MetadataResponseMessage.class, msgIdx++); - register(factory, MarshallerMappingItem.class, msgIdx++); - register(factory, BinaryMetadataVersionInfo.class, msgIdx++); - register(factory, CacheBinaryDataBagItem.class, msgIdx++); - register(factory, MappedName.class, msgIdx++); - register(factory, MarshallerDataBagItem.class, msgIdx++); + register(GridDeploymentInfoBean.class); + register(GridDeploymentRequest.class); + register(GridDeploymentResponse.class); + register(MissingMappingRequestMessage.class); + register(MissingMappingResponseMessage.class); + register(MetadataRequestMessage.class); + register(MetadataResponseMessage.class); + register(MarshallerMappingItem.class); + register(BinaryMetadataVersionInfo.class); + register(CacheBinaryDataBagItem.class); + register(MappedName.class); + register(MarshallerDataBagItem.class); // [12400 - 12500]: Encryption messages. msgIdx = 12400; - register(factory, GenerateEncryptionKeyRequest.class, msgIdx++); - register(factory, GenerateEncryptionKeyResponse.class, msgIdx++); - register(factory, ChangeCacheEncryptionRequest.class, msgIdx++); - register(factory, MasterKeyChangeRequest.class, msgIdx++); - register(factory, GroupKeyEncrypted.class, msgIdx++); - register(factory, NodeEncryptionKeys.class, msgIdx++); - register(factory, EncryptionDataBagItem.class, msgIdx++); + register(GenerateEncryptionKeyRequest.class); + register(GenerateEncryptionKeyResponse.class); + register(ChangeCacheEncryptionRequest.class); + register(MasterKeyChangeRequest.class); + register(GroupKeyEncrypted.class); + register(NodeEncryptionKeys.class); + register(EncryptionDataBagItem.class); // [13000 - 13300]: Control, configuration, diagnostics and other messages. msgIdx = 13000; - register(factory, GridEventStorageMessage.class, msgIdx++); - register(factory, ChangeGlobalStateMessage.class, msgIdx++); - register(factory, GridChangeGlobalStateMessageResponse.class, msgIdx++); - register(factory, IgniteDiagnosticRequest.class, msgIdx++); - register(factory, IgniteDiagnosticResponse.class, msgIdx++); - register(factory, WalStateAckMessage.class, msgIdx++); - register(factory, CacheConfigurationEnrichment.class, msgIdx++); - register(factory, DynamicCacheChangeRequest.class, msgIdx++); - register(factory, PartitionHashRecord.class, msgIdx++); - register(factory, TransactionsHashRecord.class, msgIdx++); - register(factory, ClusterIdAndTag.class, msgIdx++); - register(factory, ClusterUpdateNotifierDataBagItem.class, msgIdx++); - register(factory, PluginsDataBagItem.class, msgIdx++); - register(factory, EventsDataBagItem.class, msgIdx++); + register(GridEventStorageMessage.class); + register(ChangeGlobalStateMessage.class); + register(GridChangeGlobalStateMessageResponse.class); + register(IgniteDiagnosticRequest.class); + register(IgniteDiagnosticResponse.class); + register(WalStateAckMessage.class); + register(CacheConfigurationEnrichment.class); + register(DynamicCacheChangeRequest.class); + register(PartitionHashRecord.class); + register(TransactionsHashRecord.class); + register(ClusterIdAndTag.class); + register(ClusterUpdateNotifierDataBagItem.class); + register(PluginsDataBagItem.class); + register(EventsDataBagItem.class); // [13400 - 13500]: Operation context messages. msgIdx = 13400; - register(factory, OperationContextSnapshotMessage.class, msgIdx++); - register(factory, SecurityContextWrapper.class, msgIdx++); + register(OperationContextSnapshotMessage.class); + register(SecurityContextWrapper.class); // [13600 - 13700]: Rolling Upgrade messages. msgIdx = 13600; - register(factory, IgniteFeatureSet.class, msgIdx++); - register(factory, IgniteCoreFeatureSet.class, msgIdx++); - register(factory, IgnitePluginFeatureSet.class, msgIdx++); - register(factory, RollingUpgradeClusterData.class, msgIdx++); + register(IgniteFeatureSet.class); + register(IgniteCoreFeatureSet.class); + register(IgnitePluginFeatureSet.class); + register(RollingUpgradeClusterData.class); assert msgIdx <= MAX_MESSAGE_ID; } + + /** Registers the message under the next id. */ + private void register(Class cls) { + register(factory, cls, msgIdx++); + } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java index 2cdcdb91b0912..ceaaeeafe0b92 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java @@ -68,7 +68,7 @@ protected void register(IgniteMessageFactory factory, Class< /** */ private static void register(IgniteMessageFactory factory, Class cls, short id, Marshaller marsh) { - MessageSerializer serializer = requireGenerated(cls, "Serializer"); + MessageSerializer serializer = require(loadGenerated(cls, "Serializer"), cls, "Serializer"); // A MarshallableMessage always gets a generated marshaller (the hook call alone is a statement), so its // absence is a build problem. For the rest the generator skips statement-free marshallers, so absence @@ -79,7 +79,7 @@ private static void register(IgniteMessageFactory factory, C if (NonMarshallableMessage.class.isAssignableFrom(cls)) marshaller = null; else if (MarshallableMessage.class.isAssignableFrom(cls)) - marshaller = requireMarshaller(cls, marsh); + marshaller = require(loadMarshaller(cls, marsh), cls, "Marshaller"); else marshaller = loadMarshaller(cls, marsh); @@ -92,17 +92,7 @@ else if (MarshallableMessage.class.isAssignableFrom(cls)) factory.register(id, serializer, marshaller, deployer); } - /** Loads the generated companion like {@link #loadGenerated}, failing fast when it is missing. */ - private static T requireGenerated(Class cls, String suffix) { - return require(loadGenerated(cls, suffix), cls, suffix); - } - - /** Loads the generated marshaller like {@link #loadMarshaller}, failing fast when it is missing. */ - private static T requireMarshaller(Class cls, Marshaller marsh) { - return require(loadMarshaller(cls, marsh), cls, "Marshaller"); - } - - /** */ + /** @return {@code companion}, failing fast when it is missing. */ private static T require(@Nullable T companion, Class cls, String suffix) { if (companion == null) { throw new IgniteException("No " + cls.getSimpleName() + suffix + " found for " + cls.getName() + From 4daa565d061d912a9fb3b75a479e81e384b51e66 Mon Sep 17 00:00:00 2001 From: Anton Vinogradov Date: Fri, 31 Jul 2026 20:03:21 +0300 Subject: [PATCH 3/4] IGNITE-28937 Review fix: join the companion loaders back into one method Co-Authored-By: Claude Opus 5 --- ...actMarshallableMessageFactoryProvider.java | 42 ++++++------------- 1 file changed, 12 insertions(+), 30 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java index ceaaeeafe0b92..1b589dba546b9 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java @@ -68,7 +68,7 @@ protected void register(IgniteMessageFactory factory, Class< /** */ private static void register(IgniteMessageFactory factory, Class cls, short id, Marshaller marsh) { - MessageSerializer serializer = require(loadGenerated(cls, "Serializer"), cls, "Serializer"); + MessageSerializer serializer = require(loadGenerated(cls, "Serializer", null), cls, "Serializer"); // A MarshallableMessage always gets a generated marshaller (the hook call alone is a statement), so its // absence is a build problem. For the rest the generator skips statement-free marshallers, so absence @@ -79,14 +79,14 @@ private static void register(IgniteMessageFactory factory, C if (NonMarshallableMessage.class.isAssignableFrom(cls)) marshaller = null; else if (MarshallableMessage.class.isAssignableFrom(cls)) - marshaller = require(loadMarshaller(cls, marsh), cls, "Marshaller"); + marshaller = require(loadGenerated(cls, "Marshaller", marsh), cls, "Marshaller"); else - marshaller = loadMarshaller(cls, marsh); + marshaller = loadGenerated(cls, "Marshaller", marsh); // Deployers are generated for GridCacheMessage subclasses only, so the class lookup is skipped for the rest; // a DeployableMessage left without a deployer is then rejected at registration. GridCacheMessageDeployer deployer = GridCacheMessage.class.isAssignableFrom(cls) - ? loadGenerated(cls, "Deployer") + ? loadGenerated(cls, "Deployer", null) : null; factory.register(id, serializer, marshaller, deployer); @@ -104,44 +104,26 @@ private static T require(@Nullable T companion, Class cls, String suffix) } /** - * Instantiates the generated companion class {@code Serializer/Deployer}, or returns {@code null} when it - * does not exist. Neither takes a marshaller: the serializer writes the wire fields as they are, and the deployer - * only walks cache objects. + * Instantiates the generated companion class {@code Serializer/Marshaller/Deployer}, or returns + * {@code null} when it does not exist. Only the marshaller companion ever takes a {@code Marshaller}, and only + * when the message has fields to marshal with one, so {@code marsh} is {@code null} for the other two. + * Constructor lookups, including missing companions, are cached per message class in {@link #COMPANIONS}. */ @SuppressWarnings("unchecked") - private static @Nullable T loadGenerated(Class cls, String suffix) { + private static @Nullable T loadGenerated(Class cls, String suffix, @Nullable Marshaller marsh) { Constructor ctor = COMPANIONS.get(cls).ctor(suffix); if (ctor == null) return null; - assert ctor.getParameterCount() == 0 : cls.getSimpleName() + suffix + " must have a no-arg constructor"; - - try { - return (T)ctor.newInstance(); - } - catch (Exception e) { - throw new IgniteException("Failed to instantiate " + cls.getSimpleName() + suffix, e); - } - } - - /** - * Instantiates the generated {@code Marshaller}, or returns {@code null} when it does not exist. The - * generator gives it a {@code Marshaller} constructor only when the message has fields to marshal with one; - * otherwise the companion just walks nested messages and cache objects, and takes no arguments. - */ - @SuppressWarnings("unchecked") - private static @Nullable T loadMarshaller(Class cls, Marshaller marsh) { - Constructor ctor = COMPANIONS.get(cls).ctor("Marshaller"); - - if (ctor == null) - return null; + assert ctor.getParameterCount() == 0 || marsh != null : + cls.getSimpleName() + suffix + " takes a marshaller, but none was provided"; try { return (T)(ctor.getParameterCount() == 0 ? ctor.newInstance() : ctor.newInstance(marsh)); } catch (Exception e) { - throw new IgniteException("Failed to instantiate " + cls.getSimpleName() + "Marshaller", e); + throw new IgniteException("Failed to instantiate " + cls.getSimpleName() + suffix, e); } } From 3fdf446d3d7231c0d08cdcb23d446879bb3aecfe Mon Sep 17 00:00:00 2001 From: Anton Vinogradov Date: Fri, 31 Jul 2026 20:17:51 +0300 Subject: [PATCH 4/4] IGNITE-28937 Review fix: fold the required check into loadGenerated Co-Authored-By: Claude Opus 5 --- ...actMarshallableMessageFactoryProvider.java | 40 +++++++++---------- 1 file changed, 19 insertions(+), 21 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java index 1b589dba546b9..74e8768167cb0 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/plugin/AbstractMarshallableMessageFactoryProvider.java @@ -68,7 +68,7 @@ protected void register(IgniteMessageFactory factory, Class< /** */ private static void register(IgniteMessageFactory factory, Class cls, short id, Marshaller marsh) { - MessageSerializer serializer = require(loadGenerated(cls, "Serializer", null), cls, "Serializer"); + MessageSerializer serializer = loadGenerated(cls, "Serializer", null, true); // A MarshallableMessage always gets a generated marshaller (the hook call alone is a statement), so its // absence is a build problem. For the rest the generator skips statement-free marshallers, so absence @@ -79,42 +79,40 @@ private static void register(IgniteMessageFactory factory, C if (NonMarshallableMessage.class.isAssignableFrom(cls)) marshaller = null; else if (MarshallableMessage.class.isAssignableFrom(cls)) - marshaller = require(loadGenerated(cls, "Marshaller", marsh), cls, "Marshaller"); + marshaller = loadGenerated(cls, "Marshaller", marsh, true); else - marshaller = loadGenerated(cls, "Marshaller", marsh); + marshaller = loadGenerated(cls, "Marshaller", marsh, false); // Deployers are generated for GridCacheMessage subclasses only, so the class lookup is skipped for the rest; // a DeployableMessage left without a deployer is then rejected at registration. GridCacheMessageDeployer deployer = GridCacheMessage.class.isAssignableFrom(cls) - ? loadGenerated(cls, "Deployer", null) + ? loadGenerated(cls, "Deployer", null, false) : null; factory.register(id, serializer, marshaller, deployer); } - /** @return {@code companion}, failing fast when it is missing. */ - private static T require(@Nullable T companion, Class cls, String suffix) { - if (companion == null) { - throw new IgniteException("No " + cls.getSimpleName() + suffix + " found for " + cls.getName() + - ". Either the class is not processed by codegen or the generated sources are stale," + - " try 'mvn clean install'."); - } - - return companion; - } - /** - * Instantiates the generated companion class {@code Serializer/Marshaller/Deployer}, or returns - * {@code null} when it does not exist. Only the marshaller companion ever takes a {@code Marshaller}, and only - * when the message has fields to marshal with one, so {@code marsh} is {@code null} for the other two. - * Constructor lookups, including missing companions, are cached per message class in {@link #COMPANIONS}. + * Instantiates the generated companion class {@code Serializer/Marshaller/Deployer}. Only the marshaller + * companion ever takes a {@code Marshaller}, and only when the message has fields to marshal with one, so + * {@code marsh} is {@code null} for the other two. Constructor lookups, including missing companions, are cached + * per message class in {@link #COMPANIONS}. + * + * @return the companion, or {@code null} when it is not generated and {@code required} is {@code false}. */ @SuppressWarnings("unchecked") - private static @Nullable T loadGenerated(Class cls, String suffix, @Nullable Marshaller marsh) { + private static @Nullable T loadGenerated(Class cls, String suffix, @Nullable Marshaller marsh, boolean required) { Constructor ctor = COMPANIONS.get(cls).ctor(suffix); - if (ctor == null) + if (ctor == null) { + if (required) { + throw new IgniteException("No " + cls.getSimpleName() + suffix + " found for " + cls.getName() + + ". Either the class is not processed by codegen or the generated sources are stale," + + " try 'mvn clean install'."); + } + return null; + } assert ctor.getParameterCount() == 0 || marsh != null : cls.getSimpleName() + suffix + " takes a marshaller, but none was provided";