diff --git a/infra/common/src/main/java/cn/hippo4j/common/toolkit/ContentUtil.java b/infra/common/src/main/java/cn/hippo4j/common/toolkit/ContentUtil.java index ae5b6dfe..b3137407 100644 --- a/infra/common/src/main/java/cn/hippo4j/common/toolkit/ContentUtil.java +++ b/infra/common/src/main/java/cn/hippo4j/common/toolkit/ContentUtil.java @@ -37,8 +37,12 @@ public class ContentUtil { threadPoolParameterInfo.setTenantId(parameter.getTenantId()) .setItemId(parameter.getItemId()) .setTpId(parameter.getTpId()) - .setCoreSize(parameter.getCoreSize()) - .setMaxSize(parameter.getMaxSize()) + .setCoreSize((parameter instanceof ThreadPoolParameterInfo) + ? ((ThreadPoolParameterInfo) parameter).corePoolSizeAdapt() + : parameter.getCoreSize()) + .setMaxSize((parameter instanceof ThreadPoolParameterInfo) + ? ((ThreadPoolParameterInfo) parameter).maximumPoolSizeAdapt() + : parameter.getMaxSize()) .setQueueType(parameter.getQueueType()) .setCapacity(parameter.getCapacity()) .setKeepAliveTime(parameter.getKeepAliveTime()) diff --git a/infra/common/src/main/java/cn/hippo4j/common/toolkit/IncrementalContentUtil.java b/infra/common/src/main/java/cn/hippo4j/common/toolkit/IncrementalContentUtil.java new file mode 100644 index 00000000..8246c355 --- /dev/null +++ b/infra/common/src/main/java/cn/hippo4j/common/toolkit/IncrementalContentUtil.java @@ -0,0 +1,195 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package cn.hippo4j.common.toolkit; + +import cn.hippo4j.common.model.ThreadPoolParameter; +import cn.hippo4j.common.model.ThreadPoolParameterInfo; +import lombok.extern.slf4j.Slf4j; + +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; + +/** + * Incremental content util for thread pool parameter comparison. + * Supports version compatibility and incremental updates. + */ +@Slf4j +public class IncrementalContentUtil { + + /** + * Version of the incremental protocol + */ + public static final int PROTOCOL_VERSION = 2; + + /** + * Core parameters that affect thread pool behavior + */ + private static final String[] CORE_PARAMETERS = { + "coreSize", "maxSize", "queueType", "capacity", + "keepAliveTime", "rejectedType", "allowCoreThreadTimeOut" + }; + + /** + * Extended parameters that don't affect core behavior + */ + private static final String[] EXTENDED_PARAMETERS = { + "executeTimeOut", "isAlarm", "capacityAlarm", "livenessAlarm" + }; + + /** + * Get core content for MD5 calculation (only essential parameters) + * + * @param parameter thread-pool parameter + * @return core content string for MD5 + */ + public static String getCoreContent(ThreadPoolParameter parameter) { + ThreadPoolParameterInfo threadPoolParameterInfo = new ThreadPoolParameterInfo(); + threadPoolParameterInfo.setTenantId(parameter.getTenantId()) + .setItemId(parameter.getItemId()) + .setTpId(parameter.getTpId()); + if (parameter instanceof ThreadPoolParameterInfo) { + ThreadPoolParameterInfo info = (ThreadPoolParameterInfo) parameter; + threadPoolParameterInfo.setCorePoolSize(info.corePoolSizeAdapt()) + .setMaximumPoolSize(info.maximumPoolSizeAdapt()); + } else { + // Fallback to deprecated methods for non-ThreadPoolParameterInfo implementations + threadPoolParameterInfo.setCorePoolSize(parameter.getCoreSize()) + .setMaximumPoolSize(parameter.getMaxSize()); + } + + threadPoolParameterInfo.setQueueType(parameter.getQueueType()) + .setCapacity(parameter.getCapacity()) + .setKeepAliveTime(parameter.getKeepAliveTime()) + .setRejectedType(parameter.getRejectedType()) + .setAllowCoreThreadTimeOut(parameter.getAllowCoreThreadTimeOut()); + return JSONUtil.toJSONString(threadPoolParameterInfo); + } + + /** + * Get full content for MD5 calculation (all parameters) + * + * @param parameter thread-pool parameter + * @return full content string for MD5 + */ + public static String getFullContent(ThreadPoolParameter parameter) { + return ContentUtil.getPoolContent(parameter); + } + + /** + * Get incremental content for version compatibility + * + * @param parameter thread-pool parameter + * @param version client version + * @return incremental content string + */ + public static String getIncrementalContent(ThreadPoolParameter parameter, int version) { + if (version >= PROTOCOL_VERSION) { + return getCoreContent(parameter); + } else { + return getFullContent(parameter); + } + } + + /** + * Check if parameters have core changes that require thread pool refresh + * + * @param oldParameter old parameter + * @param newParameter new parameter + * @return true if core parameters changed + */ + public static boolean hasCoreChanges(ThreadPoolParameter oldParameter, ThreadPoolParameter newParameter) { + if (oldParameter == null || newParameter == null) { + return true; + } + // Use adapt methods for ThreadPoolParameterInfo, fallback to deprecated methods for other implementations + Integer oldCoreSize = (oldParameter instanceof ThreadPoolParameterInfo) ? ((ThreadPoolParameterInfo) oldParameter).corePoolSizeAdapt() : oldParameter.getCoreSize(); + Integer newCoreSize = (newParameter instanceof ThreadPoolParameterInfo) ? ((ThreadPoolParameterInfo) newParameter).corePoolSizeAdapt() : newParameter.getCoreSize(); + Integer oldMaxSize = (oldParameter instanceof ThreadPoolParameterInfo) ? ((ThreadPoolParameterInfo) oldParameter).maximumPoolSizeAdapt() : oldParameter.getMaxSize(); + Integer newMaxSize = (newParameter instanceof ThreadPoolParameterInfo) ? ((ThreadPoolParameterInfo) newParameter).maximumPoolSizeAdapt() : newParameter.getMaxSize(); + return !Objects.equals(oldCoreSize, newCoreSize) || + !Objects.equals(oldMaxSize, newMaxSize) || + !Objects.equals(oldParameter.getQueueType(), newParameter.getQueueType()) || + !Objects.equals(oldParameter.getCapacity(), newParameter.getCapacity()) || + !Objects.equals(oldParameter.getKeepAliveTime(), newParameter.getKeepAliveTime()) || + !Objects.equals(oldParameter.getRejectedType(), newParameter.getRejectedType()) || + !Objects.equals(oldParameter.getAllowCoreThreadTimeOut(), newParameter.getAllowCoreThreadTimeOut()); + } + + /** + * Check if parameters have extended changes (non-core) + * + * @param oldParameter old parameter + * @param newParameter new parameter + * @return true if extended parameters changed + */ + public static boolean hasExtendedChanges(ThreadPoolParameter oldParameter, ThreadPoolParameter newParameter) { + if (oldParameter == null || newParameter == null) { + return true; + } + return !Objects.equals(oldParameter.getExecuteTimeOut(), newParameter.getExecuteTimeOut()) || + !Objects.equals(oldParameter.getIsAlarm(), newParameter.getIsAlarm()) || + !Objects.equals(oldParameter.getCapacityAlarm(), newParameter.getCapacityAlarm()) || + !Objects.equals(oldParameter.getLivenessAlarm(), newParameter.getLivenessAlarm()); + } + + /** + * Get parameter changes summary + * + * @param oldParameter old parameter + * @param newParameter new parameter + * @return changes summary map + */ + public static Map getChangesSummary(ThreadPoolParameter oldParameter, ThreadPoolParameter newParameter) { + Map changes = new HashMap<>(); + if (oldParameter == null || newParameter == null) { + changes.put("type", "full"); + changes.put("reason", "initial_load"); + return changes; + } + boolean coreChanges = hasCoreChanges(oldParameter, newParameter); + boolean extendedChanges = hasExtendedChanges(oldParameter, newParameter); + if (coreChanges) { + changes.put("type", "core"); + changes.put("reason", "core_parameters_changed"); + } else if (extendedChanges) { + changes.put("type", "extended"); + changes.put("reason", "extended_parameters_changed"); + } else { + changes.put("type", "none"); + changes.put("reason", "no_changes"); + } + return changes; + } + + /** + * Create versioned content for backward compatibility + * + * @param parameter thread-pool parameter + * @param clientVersion client protocol version + * @return versioned content + */ + public static String createVersionedContent(ThreadPoolParameter parameter, int clientVersion) { + Map versionedContent = new HashMap<>(); + versionedContent.put("version", PROTOCOL_VERSION); + versionedContent.put("clientVersion", clientVersion); + versionedContent.put("content", getIncrementalContent(parameter, clientVersion)); + versionedContent.put("changes", getChangesSummary(null, parameter)); + return JSONUtil.toJSONString(versionedContent); + } +} diff --git a/infra/common/src/main/java/cn/hippo4j/common/toolkit/IncrementalMd5Util.java b/infra/common/src/main/java/cn/hippo4j/common/toolkit/IncrementalMd5Util.java new file mode 100644 index 00000000..7e7e5fc1 --- /dev/null +++ b/infra/common/src/main/java/cn/hippo4j/common/toolkit/IncrementalMd5Util.java @@ -0,0 +1,133 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package cn.hippo4j.common.toolkit; + +import cn.hippo4j.common.model.ThreadPoolParameter; +import lombok.extern.slf4j.Slf4j; + +/** + * Incremental MD5 util for thread pool parameter comparison. + * Supports version compatibility and reduces unnecessary refreshes. + */ +@Slf4j +public class IncrementalMd5Util { + + /** + * Get core MD5 for essential parameters only + * + * @param config thread pool parameter + * @return core MD5 hash + */ + public static String getCoreMd5(ThreadPoolParameter config) { + String coreContent = IncrementalContentUtil.getCoreContent(config); + return Md5Util.md5Hex(coreContent, "UTF-8"); + } + + /** + * Get full MD5 for all parameters (legacy compatibility) + * + * @param config thread pool parameter + * @return full MD5 hash + */ + public static String getFullMd5(ThreadPoolParameter config) { + return Md5Util.getTpContentMd5(config); + } + + /** + * Get versioned MD5 based on client version + * + * @param config thread pool parameter + * @param clientVersion client protocol version + * @return versioned MD5 hash + */ + public static String getVersionedMd5(ThreadPoolParameter config, int clientVersion) { + if (clientVersion >= IncrementalContentUtil.PROTOCOL_VERSION) { + String coreMd5 = getCoreMd5(config); + if (log.isDebugEnabled()) { + log.debug("Protocol v{}: Using incremental MD5 (core parameters only), MD5={}", clientVersion, coreMd5); + } + return coreMd5; + } else { + String fullMd5 = getFullMd5(config); + if (log.isDebugEnabled()) { + log.debug("Protocol v{}: Using full MD5 (all parameters), MD5={}", clientVersion, fullMd5); + } + return fullMd5; + } + } + + /** + * Compare MD5 with version support + * + * @param oldConfig old configuration + * @param newConfig new configuration + * @param clientVersion client version + * @return true if configurations are different + */ + public static boolean isDifferent(ThreadPoolParameter oldConfig, ThreadPoolParameter newConfig, int clientVersion) { + if (oldConfig == null || newConfig == null) { + return true; + } + String oldMd5 = getVersionedMd5(oldConfig, clientVersion); + String newMd5 = getVersionedMd5(newConfig, clientVersion); + boolean different = !oldMd5.equals(newMd5); + if (different) { + log.debug("Configuration changed - Old MD5: {}, New MD5: {}, Client Version: {}", + oldMd5, newMd5, clientVersion); + } + return different; + } + + /** + * Check if only extended parameters changed (non-core) + * + * @param oldConfig old configuration + * @param newConfig new configuration + * @return true if only extended parameters changed + */ + public static boolean onlyExtendedChanged(ThreadPoolParameter oldConfig, ThreadPoolParameter newConfig) { + if (oldConfig == null || newConfig == null) { + return false; + } + boolean coreChanged = IncrementalContentUtil.hasCoreChanges(oldConfig, newConfig); + boolean extendedChanged = IncrementalContentUtil.hasExtendedChanges(oldConfig, newConfig); + return !coreChanged && extendedChanged; + } + + /** + * Get change type for logging and monitoring + * + * @param oldConfig old configuration + * @param newConfig new configuration + * @return change type string + */ + public static String getChangeType(ThreadPoolParameter oldConfig, ThreadPoolParameter newConfig) { + if (oldConfig == null || newConfig == null) { + return "INITIAL"; + } + boolean coreChanged = IncrementalContentUtil.hasCoreChanges(oldConfig, newConfig); + boolean extendedChanged = IncrementalContentUtil.hasExtendedChanges(oldConfig, newConfig); + if (coreChanged) { + return "CORE"; + } else if (extendedChanged) { + return "EXTENDED"; + } else { + return "NONE"; + } + } +}