This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 8a789d09a60 [Pipe] Allow startup when a plugin JAR is unavailable
(#18496)
8a789d09a60 is described below
commit 8a789d09a6046691b3b6003740c48334a0b79aa4
Author: Caideyipi <[email protected]>
AuthorDate: Tue Sep 8 09:19:15 2026 +0800
[Pipe] Allow startup when a plugin JAR is unavailable (#18496)
* fix(pipe): allow startup when plugin jar is unavailable
* fix(pipe): isolate unavailable plugin jars during startup
* Update PipeAgentLauncherTest.java
* test(pipe): move plugin load failure test out of PowerMock
* fix(pipe): address review comments
---
.../apache/iotdb/db/i18n/DataNodePipeMessages.java | 19 +++
.../apache/iotdb/db/i18n/DataNodePipeMessages.java | 17 ++
.../db/pipe/agent/runtime/PipeAgentLauncher.java | 167 +++++++++++++++---
.../agent/plugin/PipeDataNodePluginAgentTest.java | 24 +++
.../pipe/agent/runtime/PipeAgentLauncherTest.java | 186 +++++++++++++++++++++
5 files changed, 394 insertions(+), 19 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
index 520a300c8e1..dd0e205ddf5 100644
---
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
+++
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
@@ -144,6 +144,25 @@ public final class DataNodePipeMessages {
"Failed to get pipe metas, will be synced by configNode later...";
public static final String FAILED_TO_GET_PIPE_PLUGIN_JAR_FROM =
"Failed to get pipe plugin jar from config node.";
+ public static final String
+
LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JARS_FROM_CONFIGNODE_PLUGINS_ARG_JARS_ARG_STATUS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_574C0077
=
+ "Failed to fetch pipe plugin jars from ConfigNode. Plugins: {},
jars: {}, status: {}. "
+ + "Retrying each plugin individually.";
+ public static final String
+
LOG_CONFIGNODE_RETURNED_ARG_PIPE_PLUGIN_JARS_FOR_ARG_REQUESTED_PLUGINS_PLUGINS_ARG_JARS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_27E32FDE
=
+ "ConfigNode returned {} pipe plugin jars for {} requested plugins.
Plugins: {}, jars: {}. "
+ + "Retrying each plugin individually.";
+ public static final String
+
EXCEPTION_FAILED_TO_FETCH_PIPE_PLUGIN_JAR_FROM_CONFIGNODE_FOR_PLUGIN_ARG_JAR_ARG_STATUS_ARG_B7C7FDE5
=
+ "Failed to fetch pipe plugin jar from ConfigNode for plugin %s (jar
%s). Status: %s.";
+ public static final String
+
EXCEPTION_CONFIGNODE_RETURNED_ARG_JARS_FOR_PIPE_PLUGIN_ARG_WHILE_ONE_WAS_REQUESTED_A724E582
=
+ "ConfigNode returned %d jars for pipe plugin %s while one was
requested.";
+ public static final String
+
LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JAR_ARG_FOR_PIPE_PLUGIN_ARG_FROM_CONFIGNODE_4929C5D9
=
+ "Failed to fetch pipe plugin jar {} for pipe plugin {} from
ConfigNode.";
+ public static final String
LOG_FAILED_TO_SAVE_JAR_ARG_FOR_PIPE_PLUGIN_ARG_A64D1530 =
+ "Failed to save jar {} for pipe plugin {}.";
public static final String FAILED_TO_GET_PIPE_TASK_META_FROM =
"Failed to get pipe task meta from config node. Ignore the exception,
because config "
+ "node may not be ready yet, and meta will be pushed by config node
later.";
diff --git
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
index 2a43385342f..61290db1c25 100644
---
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
+++
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
@@ -133,6 +133,23 @@ public final class DataNodePipeMessages {
"获取 pipe metas 失败,稍后会从 CN 处再次同步。";
public static final String FAILED_TO_GET_PIPE_PLUGIN_JAR_FROM =
"从 CN 处获取 pipe 插件 jar 包失败。";
+ public static final String
+
LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JARS_FROM_CONFIGNODE_PLUGINS_ARG_JARS_ARG_STATUS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_574C0077
=
+ "从 ConfigNode 获取 pipe plugin jars 失败。插件:{},jars:{},状态:{}。将逐个重试每个插件。";
+ public static final String
+
LOG_CONFIGNODE_RETURNED_ARG_PIPE_PLUGIN_JARS_FOR_ARG_REQUESTED_PLUGINS_PLUGINS_ARG_JARS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_27E32FDE
=
+ "ConfigNode 为 {} 个请求的插件返回了 {} 个 pipe plugin
jars。插件:{},jars:{}。将逐个重试每个插件。";
+ public static final String
+
EXCEPTION_FAILED_TO_FETCH_PIPE_PLUGIN_JAR_FROM_CONFIGNODE_FOR_PLUGIN_ARG_JAR_ARG_STATUS_ARG_B7C7FDE5
=
+ "从 ConfigNode 获取插件 %s(jar %s)的 pipe plugin jar 失败。状态:%s。";
+ public static final String
+
EXCEPTION_CONFIGNODE_RETURNED_ARG_JARS_FOR_PIPE_PLUGIN_ARG_WHILE_ONE_WAS_REQUESTED_A724E582
=
+ "ConfigNode 返回了 %d 个 jars,而 pipe plugin %s 只请求了一个。";
+ public static final String
+
LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JAR_ARG_FOR_PIPE_PLUGIN_ARG_FROM_CONFIGNODE_4929C5D9
=
+ "从 ConfigNode 获取 pipe plugin jar {}(pipe plugin {})失败。";
+ public static final String
LOG_FAILED_TO_SAVE_JAR_ARG_FOR_PIPE_PLUGIN_ARG_A64D1530 =
+ "保存 jar {}(pipe plugin {})失败。";
public static final String FAILED_TO_GET_PIPE_TASK_META_FROM =
"获取 pipe task meta from config node. Ignore the exception 失败,原因:config
node may not be "
+ "ready yet, and meta will be pushed by config node later.";
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncher.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncher.java
index cde729a2e51..1b3e210b1e8 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncher.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncher.java
@@ -19,7 +19,6 @@
package org.apache.iotdb.db.pipe.agent.runtime;
-import org.apache.iotdb.commons.client.exception.ClientManagerException;
import org.apache.iotdb.commons.exception.StartupException;
import org.apache.iotdb.commons.pipe.agent.plugin.meta.PipePluginMeta;
import
org.apache.iotdb.commons.pipe.agent.plugin.service.PipePluginClassLoaderManager;
@@ -39,14 +38,16 @@ import
org.apache.iotdb.db.service.ResourcesInformationHolder;
import org.apache.iotdb.pipe.api.exception.PipeException;
import org.apache.iotdb.rpc.TSStatusCode;
-import org.apache.thrift.TException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashSet;
import java.util.List;
+import java.util.Set;
import java.util.stream.Collectors;
class PipeAgentLauncher {
@@ -70,6 +71,7 @@ class PipeAgentLauncher {
final List<PipePluginMeta> uninstalledOrConflictedPipePluginMetaList =
getUninstalledOrConflictedPipePluginMetaList(resourcesInformationHolder);
+ final Set<String> unavailablePipePluginNameSet = new HashSet<>();
int index = 0;
while (index < uninstalledOrConflictedPipePluginMetaList.size()) {
List<PipePluginMeta> curList = new ArrayList<>();
@@ -80,7 +82,7 @@ class PipeAgentLauncher {
offset++;
}
index += offset;
- fetchAndSavePipePluginJars(curList);
+ unavailablePipePluginNameSet.addAll(fetchAndSavePipePluginJars(curList));
}
// create instances of pipe plugins and do registration
@@ -88,12 +90,15 @@ class PipeAgentLauncher {
if (meta.isBuiltin()) {
continue;
}
+ if (unavailablePipePluginNameSet.contains(meta.getPluginName())) {
+ continue;
+ }
try {
PipeDataNodeAgent.plugin().doRegister(meta);
} catch (Throwable e) {
PipeDataNodeAgent.plugin().markPluginLoadFailure(meta, e);
// Ignore a single broken plugin and continue startup.
- LOGGER.warn(
+ LOGGER.error(
DataNodePipeMessages.FAILURE_WHEN_REGISTER_PIPE_PLUGIN_SKIP_THIS,
meta.getPluginName(),
e);
@@ -137,28 +142,152 @@ class PipeAgentLauncher {
return pipePluginMetaList;
}
- private static void fetchAndSavePipePluginJars(List<PipePluginMeta>
pipePluginMetaList)
- throws StartupException {
+ static Set<String> fetchAndSavePipePluginJars(List<PipePluginMeta>
pipePluginMetaList) {
+ if (pipePluginMetaList.isEmpty()) {
+ return Collections.emptySet();
+ }
+
+ final List<String> pluginNameList =
+
pipePluginMetaList.stream().map(PipePluginMeta::getPluginName).collect(Collectors.toList());
+ final List<String> jarNameList =
+
pipePluginMetaList.stream().map(PipePluginMeta::getJarName).collect(Collectors.toList());
+ final TGetJarInListResp resp;
+
try (ConfigNodeClient configNodeClient =
ConfigNodeClientManager.getInstance().borrowClient(ConfigNodeInfo.CONFIG_REGION_ID))
{
- final List<String> jarNameList =
-
pipePluginMetaList.stream().map(PipePluginMeta::getJarName).collect(Collectors.toList());
- final TGetJarInListResp resp =
- configNodeClient.getPipePluginJar(new TGetJarInListReq(jarNameList));
- if (resp.getStatus().getCode() ==
TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()) {
- throw new
StartupException(DataNodePipeMessages.FAILED_TO_GET_PIPE_PLUGIN_JAR_FROM);
+ resp = configNodeClient.getPipePluginJar(new
TGetJarInListReq(jarNameList));
+ } catch (Exception e) {
+ LOGGER.error(
+ DataNodePipeMessages
+
.LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JARS_FROM_CONFIGNODE_PLUGINS_ARG_JARS_ARG_STATUS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_574C0077,
+ pluginNameList,
+ jarNameList,
+ null,
+ e);
+ return fetchAndSavePipePluginJarsIndividually(pipePluginMetaList);
+ }
+
+ if (resp == null
+ || resp.getStatus() == null
+ || resp.getStatus().getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ LOGGER.error(
+ DataNodePipeMessages
+
.LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JARS_FROM_CONFIGNODE_PLUGINS_ARG_JARS_ARG_STATUS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_574C0077,
+ pluginNameList,
+ jarNameList,
+ resp == null ? null : resp.getStatus());
+ return fetchAndSavePipePluginJarsIndividually(pipePluginMetaList);
+ }
+
+ final List<ByteBuffer> jarList = resp.getJarList();
+ if (jarList == null || jarList.size() != pipePluginMetaList.size()) {
+ LOGGER.error(
+ DataNodePipeMessages
+
.LOG_CONFIGNODE_RETURNED_ARG_PIPE_PLUGIN_JARS_FOR_ARG_REQUESTED_PLUGINS_PLUGINS_ARG_JARS_ARG_RETRYING_EACH_PLUGIN_INDIVIDUALLY_27E32FDE,
+ jarList == null ? 0 : jarList.size(),
+ pipePluginMetaList.size(),
+ pluginNameList,
+ jarNameList);
+ return fetchAndSavePipePluginJarsIndividually(pipePluginMetaList);
+ }
+
+ return savePipePluginJars(pipePluginMetaList, jarList);
+ }
+
+ private static Set<String> fetchAndSavePipePluginJarsIndividually(
+ List<PipePluginMeta> pipePluginMetaList) {
+ final Set<String> unavailablePipePluginNameSet = new HashSet<>();
+ for (PipePluginMeta pipePluginMeta : pipePluginMetaList) {
+ if (!fetchAndSavePipePluginJarIndividually(pipePluginMeta)) {
+ unavailablePipePluginNameSet.add(pipePluginMeta.getPluginName());
}
- final List<ByteBuffer> jarList = resp.getJarList();
- for (int i = 0; i < pipePluginMetaList.size(); i++) {
+ }
+ return unavailablePipePluginNameSet;
+ }
+
+ private static boolean fetchAndSavePipePluginJarIndividually(PipePluginMeta
pipePluginMeta) {
+ final String pluginName = pipePluginMeta.getPluginName();
+ final String jarName = pipePluginMeta.getJarName();
+ final TGetJarInListResp resp;
+ try (ConfigNodeClient configNodeClient =
+
ConfigNodeClientManager.getInstance().borrowClient(ConfigNodeInfo.CONFIG_REGION_ID))
{
+ resp = configNodeClient.getPipePluginJar(new
TGetJarInListReq(List.of(jarName)));
+ } catch (Exception e) {
+ PipeDataNodeAgent.plugin().markPluginLoadFailure(pipePluginMeta, e);
+ LOGGER.error(
+ DataNodePipeMessages
+
.LOG_FAILED_TO_FETCH_PIPE_PLUGIN_JAR_ARG_FOR_PIPE_PLUGIN_ARG_FROM_CONFIGNODE_4929C5D9,
+ jarName,
+ pluginName,
+ e);
+ return false;
+ }
+
+ if (resp == null
+ || resp.getStatus() == null
+ || resp.getStatus().getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ final PipeException exception =
+ new PipeException(
+ String.format(
+ DataNodePipeMessages
+
.EXCEPTION_FAILED_TO_FETCH_PIPE_PLUGIN_JAR_FROM_CONFIGNODE_FOR_PLUGIN_ARG_JAR_ARG_STATUS_ARG_B7C7FDE5,
+ pluginName,
+ jarName,
+ resp == null ? null : resp.getStatus()));
+ PipeDataNodeAgent.plugin().markPluginLoadFailure(pipePluginMeta,
exception);
+ LOGGER.error(exception.getMessage(), exception);
+ return false;
+ }
+
+ final List<ByteBuffer> jarList = resp.getJarList();
+ if (jarList == null || jarList.size() != 1) {
+ final PipeException exception =
+ new PipeException(
+ String.format(
+ DataNodePipeMessages
+
.EXCEPTION_CONFIGNODE_RETURNED_ARG_JARS_FOR_PIPE_PLUGIN_ARG_WHILE_ONE_WAS_REQUESTED_A724E582,
+ jarList == null ? 0 : jarList.size(),
+ pluginName));
+ PipeDataNodeAgent.plugin().markPluginLoadFailure(pipePluginMeta,
exception);
+ LOGGER.error(exception.getMessage(), exception);
+ return false;
+ }
+
+ try {
+ PipePluginExecutableManager.getInstance()
+ .savePluginToInstallDir(jarList.get(0), pluginName, jarName);
+ return true;
+ } catch (Exception e) {
+ PipeDataNodeAgent.plugin().markPluginLoadFailure(pipePluginMeta, e);
+ LOGGER.error(
+
DataNodePipeMessages.LOG_FAILED_TO_SAVE_JAR_ARG_FOR_PIPE_PLUGIN_ARG_A64D1530,
+ jarName,
+ pluginName,
+ e);
+ return false;
+ }
+ }
+
+ private static Set<String> savePipePluginJars(
+ List<PipePluginMeta> pipePluginMetaList, List<ByteBuffer> jarList) {
+ final Set<String> unavailablePipePluginNameSet = new HashSet<>();
+ for (int i = 0; i < pipePluginMetaList.size(); i++) {
+ final PipePluginMeta pipePluginMeta = pipePluginMetaList.get(i);
+ try {
PipePluginExecutableManager.getInstance()
.savePluginToInstallDir(
- jarList.get(i),
- pipePluginMetaList.get(i).getPluginName(),
- pipePluginMetaList.get(i).getJarName());
+ jarList.get(i), pipePluginMeta.getPluginName(),
pipePluginMeta.getJarName());
+ } catch (Exception e) {
+ PipeDataNodeAgent.plugin().markPluginLoadFailure(pipePluginMeta, e);
+ LOGGER.error(
+
DataNodePipeMessages.LOG_FAILED_TO_SAVE_JAR_ARG_FOR_PIPE_PLUGIN_ARG_A64D1530,
+ pipePluginMeta.getJarName(),
+ pipePluginMeta.getPluginName(),
+ e);
+ unavailablePipePluginNameSet.add(pipePluginMeta.getPluginName());
}
- } catch (IOException | TException | ClientManagerException e) {
- throw new StartupException(e);
}
+ return unavailablePipePluginNameSet;
}
public static synchronized void launchPipeTaskAgent() {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/plugin/PipeDataNodePluginAgentTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/plugin/PipeDataNodePluginAgentTest.java
index 1fbb615ef7a..7b065a010df 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/plugin/PipeDataNodePluginAgentTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/plugin/PipeDataNodePluginAgentTest.java
@@ -20,12 +20,14 @@
package org.apache.iotdb.db.pipe.agent.plugin;
import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
+import
org.apache.iotdb.commons.pipe.agent.plugin.meta.DataNodePipePluginMetaKeeper;
import org.apache.iotdb.commons.pipe.agent.plugin.meta.PipePluginMeta;
import
org.apache.iotdb.commons.pipe.agent.plugin.service.PipePluginClassLoaderManager;
import
org.apache.iotdb.commons.pipe.agent.plugin.service.PipePluginExecutableManager;
import org.apache.iotdb.commons.pipe.config.constant.PipeProcessorConstant;
import org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant;
import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
+import org.apache.iotdb.commons.pipe.datastructure.visibility.Visibility;
import
org.apache.iotdb.db.pipe.processor.iotconsensusv2.IoTConsensusV2Processor;
import
org.apache.iotdb.db.pipe.sink.protocol.iotconsensusv2.IoTConsensusV2AsyncSink;
import
org.apache.iotdb.db.pipe.sink.protocol.thrift.async.IoTDBDataRegionAsyncSink;
@@ -40,6 +42,7 @@ import org.junit.Before;
import org.junit.Test;
import java.io.IOException;
+import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.HashMap;
@@ -180,4 +183,25 @@ public class PipeDataNodePluginAgentTest {
}))
.getClass());
}
+
+ @Test
+ public void testPluginLoadFailureIsRecordedForShowPipePlugins() throws
Exception {
+ final PipeDataNodePluginAgent agent = new PipeDataNodePluginAgent();
+ final PipePluginMeta plugin =
+ new PipePluginMeta("failed", "test.class", false, "failed.jar",
"test-md5");
+
+ agent.markPluginLoadFailure(plugin, new IOException("missing jar"));
+
+ final Field metaKeeperField =
+ PipeDataNodePluginAgent.class.getDeclaredField("pipePluginMetaKeeper");
+ metaKeeperField.setAccessible(true);
+ final DataNodePipePluginMetaKeeper metaKeeper =
+ (DataNodePipePluginMetaKeeper) metaKeeperField.get(agent);
+ final PipePluginMeta recordedPlugin =
metaKeeper.getPipePluginMeta("FAILED");
+
+ Assert.assertEquals(
+ "IOException: missing jar",
recordedPlugin.getPluginLoadingExceptionMessage());
+ Assert.assertEquals(
+ Visibility.BOTH,
metaKeeper.getPipePluginNameToVisibilityMap().get("FAILED"));
+ }
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncherTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncherTest.java
new file mode 100644
index 00000000000..c2ec80f0735
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/runtime/PipeAgentLauncherTest.java
@@ -0,0 +1,186 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.pipe.agent.runtime;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.client.IClientManager;
+import org.apache.iotdb.commons.consensus.ConfigRegionId;
+import org.apache.iotdb.commons.pipe.agent.plugin.meta.PipePluginMeta;
+import
org.apache.iotdb.commons.pipe.agent.plugin.service.PipePluginExecutableManager;
+import org.apache.iotdb.confignode.rpc.thrift.TGetJarInListReq;
+import org.apache.iotdb.confignode.rpc.thrift.TGetJarInListResp;
+import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
+import org.apache.iotdb.db.pipe.agent.plugin.PipeDataNodePluginAgent;
+import org.apache.iotdb.db.protocol.client.ConfigNodeClient;
+import org.apache.iotdb.db.protocol.client.ConfigNodeClientManager;
+import org.apache.iotdb.db.protocol.client.ConfigNodeInfo;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
+import org.powermock.api.mockito.PowerMockito;
+import org.powermock.core.classloader.annotations.PowerMockIgnore;
+import org.powermock.core.classloader.annotations.PrepareForTest;
+import org.powermock.modules.junit4.PowerMockRunner;
+
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.Set;
+
+@PowerMockIgnore({"com.sun.org.apache.xerces.*", "javax.xml.*", "org.xml.*",
"javax.management.*"})
+@RunWith(PowerMockRunner.class)
+@PrepareForTest({
+ ConfigNodeClientManager.class,
+ PipePluginExecutableManager.class,
+ PipeDataNodeAgent.class
+})
+public class PipeAgentLauncherTest {
+
+ private IClientManager<ConfigRegionId, ConfigNodeClient>
configNodeClientManager;
+ private ConfigNodeClient configNodeClient;
+ private PipePluginExecutableManager pipePluginExecutableManager;
+ private PipeDataNodePluginAgent pipeDataNodePluginAgent;
+
+ @Before
+ public void setUp() throws Exception {
+ configNodeClientManager = Mockito.mock(IClientManager.class);
+ configNodeClient = Mockito.mock(ConfigNodeClient.class);
+ pipePluginExecutableManager =
Mockito.mock(PipePluginExecutableManager.class);
+ pipeDataNodePluginAgent = Mockito.mock(PipeDataNodePluginAgent.class);
+
+ PowerMockito.mockStatic(ConfigNodeClientManager.class);
+ PowerMockito.mockStatic(PipePluginExecutableManager.class);
+ PowerMockito.mockStatic(PipeDataNodeAgent.class);
+
PowerMockito.when(ConfigNodeClientManager.getInstance()).thenReturn(configNodeClientManager);
+ PowerMockito.when(PipePluginExecutableManager.getInstance())
+ .thenReturn(pipePluginExecutableManager);
+
PowerMockito.when(PipeDataNodeAgent.plugin()).thenReturn(pipeDataNodePluginAgent);
+
Mockito.when(configNodeClientManager.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID))
+ .thenReturn(configNodeClient);
+ }
+
+ @Test
+ public void
testBatchFailureRetriesIndividuallyAndOnlyMarksFailedPluginUnavailable()
+ throws Exception {
+ final PipePluginMeta healthyPlugin = pipePluginMeta("healthy",
"healthy.jar");
+ final PipePluginMeta missingPlugin = pipePluginMeta("missing",
"missing.jar");
+ final TGetJarInListResp failedResponse = failureResponse();
+
+
Mockito.when(configNodeClient.getPipePluginJar(Mockito.any(TGetJarInListReq.class)))
+ .thenAnswer(
+ invocation -> {
+ final List<String> jarNames =
+ ((TGetJarInListReq)
invocation.getArgument(0)).getJarNameList();
+ if (jarNames.size() == 2 ||
"missing.jar".equals(jarNames.get(0))) {
+ return failedResponse;
+ }
+ return
successResponse(Collections.singletonList(ByteBuffer.wrap(new byte[] {1})));
+ });
+
+ final Set<String> unavailablePlugins =
+
PipeAgentLauncher.fetchAndSavePipePluginJars(Arrays.asList(healthyPlugin,
missingPlugin));
+
+ Assert.assertEquals(Collections.singleton("MISSING"), unavailablePlugins);
+ Mockito.verify(pipePluginExecutableManager)
+ .savePluginToInstallDir(
+ Mockito.any(ByteBuffer.class), Mockito.eq("HEALTHY"),
Mockito.eq("healthy.jar"));
+ Mockito.verify(pipeDataNodePluginAgent)
+ .markPluginLoadFailure(Mockito.eq(missingPlugin), Mockito.any());
+ Mockito.verify(pipeDataNodePluginAgent, Mockito.never())
+ .markPluginLoadFailure(Mockito.eq(healthyPlugin), Mockito.any());
+ }
+
+ @Test
+ public void testBatchJarCountMismatchRetriesEachPluginIndividually() throws
Exception {
+ final PipePluginMeta firstPlugin = pipePluginMeta("first", "first.jar");
+ final PipePluginMeta secondPlugin = pipePluginMeta("second", "second.jar");
+
+
Mockito.when(configNodeClient.getPipePluginJar(Mockito.any(TGetJarInListReq.class)))
+
.thenReturn(successResponse(Collections.singletonList(ByteBuffer.wrap(new
byte[] {1}))));
+
+ final Set<String> unavailablePlugins =
+
PipeAgentLauncher.fetchAndSavePipePluginJars(Arrays.asList(firstPlugin,
secondPlugin));
+
+ Assert.assertTrue(unavailablePlugins.isEmpty());
+ final ArgumentCaptor<TGetJarInListReq> requestCaptor =
+ ArgumentCaptor.forClass(TGetJarInListReq.class);
+ Mockito.verify(configNodeClient,
Mockito.times(3)).getPipePluginJar(requestCaptor.capture());
+ Assert.assertEquals(
+ Arrays.asList("first.jar", "second.jar"),
+ requestCaptor.getAllValues().get(0).getJarNameList());
+ Assert.assertEquals(
+ Collections.singletonList("first.jar"),
+ requestCaptor.getAllValues().get(1).getJarNameList());
+ Assert.assertEquals(
+ Collections.singletonList("second.jar"),
+ requestCaptor.getAllValues().get(2).getJarNameList());
+ Mockito.verify(pipePluginExecutableManager, Mockito.times(2))
+ .savePluginToInstallDir(
+ Mockito.any(ByteBuffer.class), Mockito.anyString(),
Mockito.anyString());
+ }
+
+ @Test
+ public void testSaveFailureDoesNotBlockSubsequentPlugin() throws Exception {
+ final PipePluginMeta failedPlugin = pipePluginMeta("failed", "failed.jar");
+ final PipePluginMeta healthyPlugin = pipePluginMeta("healthy",
"healthy.jar");
+
Mockito.when(configNodeClient.getPipePluginJar(Mockito.any(TGetJarInListReq.class)))
+ .thenReturn(
+ successResponse(
+ Arrays.asList(ByteBuffer.wrap(new byte[] {1}),
ByteBuffer.wrap(new byte[] {2}))));
+ Mockito.doThrow(new IOException("injected save failure"))
+ .when(pipePluginExecutableManager)
+ .savePluginToInstallDir(
+ Mockito.any(ByteBuffer.class), Mockito.eq("FAILED"),
Mockito.eq("failed.jar"));
+
+ final Set<String> unavailablePlugins =
+
PipeAgentLauncher.fetchAndSavePipePluginJars(Arrays.asList(failedPlugin,
healthyPlugin));
+
+ Assert.assertEquals(Collections.singleton("FAILED"), unavailablePlugins);
+ Mockito.verify(pipePluginExecutableManager)
+ .savePluginToInstallDir(
+ Mockito.any(ByteBuffer.class), Mockito.eq("HEALTHY"),
Mockito.eq("healthy.jar"));
+ Mockito.verify(pipeDataNodePluginAgent)
+ .markPluginLoadFailure(Mockito.eq(failedPlugin),
Mockito.any(IOException.class));
+ Mockito.verify(pipeDataNodePluginAgent, Mockito.never())
+ .markPluginLoadFailure(Mockito.eq(healthyPlugin), Mockito.any());
+ }
+
+ private static PipePluginMeta pipePluginMeta(final String pluginName, final
String jarName) {
+ return new PipePluginMeta(pluginName, "test.class", false, jarName,
"test-md5");
+ }
+
+ private static TGetJarInListResp successResponse(final List<ByteBuffer>
jarList) {
+ return new TGetJarInListResp(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()), jarList);
+ }
+
+ private static TGetJarInListResp failureResponse() {
+ return new TGetJarInListResp(
+ new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()),
+ Collections.emptyList());
+ }
+}