From ad7833b3729a5e6457abab4240e50c80bb6e7349 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Fri, 13 Aug 2021 16:55:54 +0800 Subject: [PATCH 1/7] add package url support in pulsar-admin --- .../pulsar/packages/management/core/common/PackageType.java | 0 .../main/java/org/apache/pulsar/common/functions/Utils.java | 6 +++++- .../pulsar/functions/worker/rest/api/FunctionsImpl.java | 3 ++- 3 files changed, 7 insertions(+), 2 deletions(-) rename {pulsar-package-management/core => pulsar-client-admin-api}/src/main/java/org/apache/pulsar/packages/management/core/common/PackageType.java (100%) diff --git a/pulsar-package-management/core/src/main/java/org/apache/pulsar/packages/management/core/common/PackageType.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/packages/management/core/common/PackageType.java similarity index 100% rename from pulsar-package-management/core/src/main/java/org/apache/pulsar/packages/management/core/common/PackageType.java rename to pulsar-client-admin-api/src/main/java/org/apache/pulsar/packages/management/core/common/PackageType.java diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java b/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java index 6f6a6840aba02..726426ffaf419 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java @@ -21,8 +21,10 @@ import static org.apache.commons.lang3.StringUtils.isNotBlank; import static org.apache.pulsar.common.naming.TopicName.DEFAULT_NAMESPACE; import static org.apache.pulsar.common.naming.TopicName.PUBLIC_TENANT; +import java.util.Arrays; import org.apache.pulsar.common.io.SinkConfig; import org.apache.pulsar.common.io.SourceConfig; +import org.apache.pulsar.packages.management.core.common.PackageType; /** * Helper class to work with configuration. @@ -34,7 +36,9 @@ public class Utils { public static boolean isFunctionPackageUrlSupported(String functionPkgUrl) { return isNotBlank(functionPkgUrl) && (functionPkgUrl.startsWith(HTTP) - || functionPkgUrl.startsWith(FILE)); + || functionPkgUrl.startsWith(FILE) + || Arrays.stream(PackageType.values()).anyMatch(type -> functionPkgUrl.startsWith(type.toString())) + && functionPkgUrl.contains("://")); } public static void inferMissingFunctionName(FunctionConfig functionConfig) { diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java index 5aa438922db8b..69f014317a4b7 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java @@ -760,7 +760,8 @@ private Function.FunctionDetails validateUpdateRequestParams(final String tenant } private static boolean hasPackageTypePrefix(String destPkgUrl) { - return Arrays.stream(PackageType.values()).anyMatch(type -> destPkgUrl.startsWith(type.toString())); + return Arrays.stream(PackageType.values()).anyMatch(type -> destPkgUrl.startsWith(type.toString()) + && destPkgUrl.contains("://")); } private File downloadPackageFile(String packageName) throws IOException, PulsarAdminException { From 714c8688da2ab32bad978c9b1929b17bed358d3b Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Fri, 13 Aug 2021 17:31:17 +0800 Subject: [PATCH 2/7] add tests --- .../pulsar/admin/cli/CmdFunctionsTest.java | 71 ++++++++++++++++++- .../integration/cli/PackagesCliTest.java | 40 +++++++++++ 2 files changed, 108 insertions(+), 3 deletions(-) diff --git a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java index d7f68a4b790f0..e4284bdeebdf9 100644 --- a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java +++ b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java @@ -82,14 +82,17 @@ public IObjectFactory getObjectFactory() { private static final String JAR_NAME = CmdFunctionsTest.class.getClassLoader().getResource("dummyexamples.jar").getFile(); private static final String GO_EXEC_FILE_NAME = "test-go-function-with-url"; private static final String PYTHON_FILE_NAME = "test-go-function-with-url"; - private static final String URL ="file:" + JAR_NAME; - private static final String URL_WITH_GO ="file:" + GO_EXEC_FILE_NAME; - private static final String URL_WITH_PY ="file:" + PYTHON_FILE_NAME; + private static final String URL = "file:" + JAR_NAME; + private static final String URL_WITH_GO = "file:" + GO_EXEC_FILE_NAME; + private static final String URL_WITH_PY = "file:" + PYTHON_FILE_NAME; private static final String FN_NAME = TEST_NAME + "-function"; private static final String INPUT_TOPIC_NAME = TEST_NAME + "-input-topic"; private static final String OUTPUT_TOPIC_NAME = TEST_NAME + "-output-topic"; private static final String TENANT = TEST_NAME + "-tenant"; private static final String NAMESPACE = TEST_NAME + "-namespace"; + private static final String PACKAGE_URL = "function://sample/ns1/jardummyexamples@1"; + private static final String PACKAGE_GO_URL = "function://sample/ns1/godummyexamples@1"; + private static final String PACKAGE_PY_URL = "function://sample/ns1/pydummyexamples@1"; private PulsarAdmin admin; private Functions functions; @@ -361,6 +364,68 @@ public void testCreatePyFunctionWithFileUrl() throws Exception { verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); } + @Test + public void testCreateFunctionWithPackageUrl() throws Exception { + cmd.run(new String[] { + "create", + "--name", FN_NAME, + "--inputs", INPUT_TOPIC_NAME, + "--output", OUTPUT_TOPIC_NAME, + "--jar", PACKAGE_URL, + "--tenant", "sample", + "--namespace", "ns1", + "--className", DummyFunction.class.getName(), + }); + + CreateFunction creater = cmd.getCreater(); + + assertEquals(FN_NAME, creater.getFunctionName()); + assertEquals(INPUT_TOPIC_NAME, creater.getInputs()); + assertEquals(OUTPUT_TOPIC_NAME, creater.getOutput()); + verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); + } + + @Test + public void testCreateGoFunctionWithPackageUrl() throws Exception { + cmd.run(new String[] { + "create", + "--name", "test-go-function", + "--inputs", INPUT_TOPIC_NAME, + "--output", OUTPUT_TOPIC_NAME, + "--go", PACKAGE_GO_URL, + "--tenant", "sample", + "--namespace", "ns1", + }); + + CreateFunction creater = cmd.getCreater(); + + assertEquals("test-go-function", creater.getFunctionName()); + assertEquals(INPUT_TOPIC_NAME, creater.getInputs()); + assertEquals(OUTPUT_TOPIC_NAME, creater.getOutput()); + verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); + } + + @Test + public void testCreatePyFunctionWithPackageUrl() throws Exception { + cmd.run(new String[] { + "create", + "--name", "test-py-function", + "--inputs", INPUT_TOPIC_NAME, + "--output", OUTPUT_TOPIC_NAME, + "--py", PACKAGE_PY_URL, + "--tenant", "sample", + "--namespace", "ns1", + "--className", "process_python_function", + }); + + CreateFunction creater = cmd.getCreater(); + + assertEquals("test-py-function", creater.getFunctionName()); + assertEquals(INPUT_TOPIC_NAME, creater.getInputs()); + assertEquals(OUTPUT_TOPIC_NAME, creater.getOutput()); + verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); + } + @Test public void testCreateFunctionWithoutBasicArguments() throws Exception { cmd.run(new String[] { diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java index 2f213a0b50d44..1ae7155054efd 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java @@ -137,10 +137,50 @@ public void testPackagesOperationsWithUploadingPackages() throws Exception { result.assertNoStdout(); } + @Test(timeOut = 60000 * 5) + public void testCreateFunctionFromPackagesURLWithoutUploadingPackages() throws Exception { + ContainerExecResult result = runPackagesCommand("list", "--type", "function", "public/default"); + assertEquals(result.getExitCode(), 0); + + result = runPackagesCommand("list-versions", "function://public/default/nonexist"); + assertEquals(result.getExitCode(), 0); + + try { + result = runFunctionsCommand("create", + "--jar", "function://public/default/nonexist@v1", + "--inputs", "nonexist-input-topic", + "--className", "org.DummyFunction"); + fail("this command should be failed"); + } catch (Exception e) { + // expected exception + } + } + + @Test(timeOut = 60000 * 5) + public void testCreateFunctionFromPackagesURLWithUploadingPackages() throws Exception { + String testPackageName = "function://public/default/uploaded@v1"; + ContainerExecResult result = runPackagesCommand("upload", "--description", "a test package", + "--path", PulsarCluster.ADMIN_SCRIPT, testPackageName); + assertEquals(result.getExitCode(), 0); + + result = runFunctionsCommand("create", + "--jar", testPackageName, + "--inputs", "uploaded-input-topic", + "--className", "org.DummyFunction"); + assertEquals(result.getExitCode(), 0); + } + private ContainerExecResult runPackagesCommand(String... commands) throws Exception { String[] cmds = new String[commands.length + 1]; cmds[0] = "packages"; System.arraycopy(commands, 0, cmds, 1, commands.length); return pulsarCluster.runAdminCommandOnAnyBroker(cmds); } + + private ContainerExecResult runFunctionsCommand(String... commands) throws Exception { + String[] cmds = new String[commands.length + 1]; + cmds[0] = "functions"; + System.arraycopy(commands, 0, cmds, 1, commands.length); + return pulsarCluster.runAdminCommandOnAnyBroker(cmds); + } } From f58d07d7f02530a4c5b2685ae70425b49d964a7d Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Fri, 13 Aug 2021 21:19:45 +0800 Subject: [PATCH 3/7] add functions worker to package integration test --- .../integration/cli/PackagesCliTest.java | 24 +++++++++++++++---- 1 file changed, 20 insertions(+), 4 deletions(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java index 1ae7155054efd..a229a1c36e98d 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java @@ -21,6 +21,7 @@ import org.apache.pulsar.tests.TestRetrySupport; import org.apache.pulsar.tests.integration.containers.BrokerContainer; import org.apache.pulsar.tests.integration.docker.ContainerExecResult; +import org.apache.pulsar.tests.integration.topologies.FunctionRuntimeType; import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.apache.pulsar.tests.integration.topologies.PulsarClusterSpec; import org.testcontainers.shaded.org.apache.commons.lang.RandomStringUtils; @@ -49,15 +50,20 @@ public final void setup() throws Exception { .brokerEnvs(getPackagesManagementServiceEnvs()) .build(); pulsarCluster = PulsarCluster.forSpec(spec); + setupFunctionWorkers(); pulsarCluster.start(); } @AfterClass(alwaysRun = true) public final void cleanup() { - markCurrentSetupNumberCleaned(); - if (pulsarCluster != null) { - pulsarCluster.stop(); - pulsarCluster = null; + try { + teardownFunctionWorkers(); + } finally { + markCurrentSetupNumberCleaned(); + if (pulsarCluster != null) { + pulsarCluster.stop(); + pulsarCluster = null; + } } } @@ -183,4 +189,14 @@ private ContainerExecResult runFunctionsCommand(String... commands) throws Excep System.arraycopy(commands, 0, cmds, 1, commands.length); return pulsarCluster.runAdminCommandOnAnyBroker(cmds); } + + private void setupFunctionWorkers() { + final int numFunctionWorkers = 2; + pulsarCluster.setupFunctionWorkers("", FunctionRuntimeType.THREAD, numFunctionWorkers); + } + + private void teardownFunctionWorkers() { + pulsarCluster.stopWorkers(); + } } + From bd95b3f19d05de9d9800fe48c503c3cb4baffa0e Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Sun, 15 Aug 2021 10:40:00 +0800 Subject: [PATCH 4/7] revert changes --- .../integration/cli/PackagesCliTest.java | 76 +++---------------- 1 file changed, 10 insertions(+), 66 deletions(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java index a229a1c36e98d..f6e5db7a6c101 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/PackagesCliTest.java @@ -21,7 +21,6 @@ import org.apache.pulsar.tests.TestRetrySupport; import org.apache.pulsar.tests.integration.containers.BrokerContainer; import org.apache.pulsar.tests.integration.docker.ContainerExecResult; -import org.apache.pulsar.tests.integration.topologies.FunctionRuntimeType; import org.apache.pulsar.tests.integration.topologies.PulsarCluster; import org.apache.pulsar.tests.integration.topologies.PulsarClusterSpec; import org.testcontainers.shaded.org.apache.commons.lang.RandomStringUtils; @@ -46,24 +45,19 @@ public class PackagesCliTest extends TestRetrySupport { public final void setup() throws Exception { incrementSetupNumber(); PulsarClusterSpec spec = PulsarClusterSpec.builder() - .clusterName(String.format("%s-%s", clusterNamePrefix, RandomStringUtils.randomAlphabetic(6))) - .brokerEnvs(getPackagesManagementServiceEnvs()) - .build(); + .clusterName(String.format("%s-%s", clusterNamePrefix, RandomStringUtils.randomAlphabetic(6))) + .brokerEnvs(getPackagesManagementServiceEnvs()) + .build(); pulsarCluster = PulsarCluster.forSpec(spec); - setupFunctionWorkers(); pulsarCluster.start(); } @AfterClass(alwaysRun = true) public final void cleanup() { - try { - teardownFunctionWorkers(); - } finally { - markCurrentSetupNumberCleaned(); - if (pulsarCluster != null) { - pulsarCluster.stop(); - pulsarCluster = null; - } + markCurrentSetupNumberCleaned(); + if (pulsarCluster != null) { + pulsarCluster.stop(); + pulsarCluster = null; } } @@ -94,13 +88,13 @@ public void testPackagesOperationsWithoutUploadingPackages() throws Exception { public void testPackagesOperationsWithUploadingPackages() throws Exception { String testPackageName = "function://public/default/test@v1"; ContainerExecResult result = runPackagesCommand("upload", "--description", "a test package", - "--path", PulsarCluster.ADMIN_SCRIPT, testPackageName); + "--path", PulsarCluster.ADMIN_SCRIPT, testPackageName); assertEquals(result.getExitCode(), 0); BrokerContainer container = pulsarCluster.getBroker(0); String downloadFile = "tmp-file-" + RandomStringUtils.randomAlphabetic(8); String[] downloadCmd = new String[]{PulsarCluster.ADMIN_SCRIPT, "packages", "download", - "--path", downloadFile, testPackageName}; + "--path", downloadFile, testPackageName}; result = container.execCmd(downloadCmd); assertEquals(result.getExitCode(), 0); @@ -125,7 +119,7 @@ public void testPackagesOperationsWithUploadingPackages() throws Exception { String contact = "test@apache.org"; result = runPackagesCommand("update-metadata", "--description", "a test package", - "--contact", contact, "-PpropertyA=A", testPackageName); + "--contact", contact, "-PpropertyA=A", testPackageName); assertEquals(result.getExitCode(), 0); result = runPackagesCommand("get-metadata", testPackageName); @@ -143,60 +137,10 @@ public void testPackagesOperationsWithUploadingPackages() throws Exception { result.assertNoStdout(); } - @Test(timeOut = 60000 * 5) - public void testCreateFunctionFromPackagesURLWithoutUploadingPackages() throws Exception { - ContainerExecResult result = runPackagesCommand("list", "--type", "function", "public/default"); - assertEquals(result.getExitCode(), 0); - - result = runPackagesCommand("list-versions", "function://public/default/nonexist"); - assertEquals(result.getExitCode(), 0); - - try { - result = runFunctionsCommand("create", - "--jar", "function://public/default/nonexist@v1", - "--inputs", "nonexist-input-topic", - "--className", "org.DummyFunction"); - fail("this command should be failed"); - } catch (Exception e) { - // expected exception - } - } - - @Test(timeOut = 60000 * 5) - public void testCreateFunctionFromPackagesURLWithUploadingPackages() throws Exception { - String testPackageName = "function://public/default/uploaded@v1"; - ContainerExecResult result = runPackagesCommand("upload", "--description", "a test package", - "--path", PulsarCluster.ADMIN_SCRIPT, testPackageName); - assertEquals(result.getExitCode(), 0); - - result = runFunctionsCommand("create", - "--jar", testPackageName, - "--inputs", "uploaded-input-topic", - "--className", "org.DummyFunction"); - assertEquals(result.getExitCode(), 0); - } - private ContainerExecResult runPackagesCommand(String... commands) throws Exception { String[] cmds = new String[commands.length + 1]; cmds[0] = "packages"; System.arraycopy(commands, 0, cmds, 1, commands.length); return pulsarCluster.runAdminCommandOnAnyBroker(cmds); } - - private ContainerExecResult runFunctionsCommand(String... commands) throws Exception { - String[] cmds = new String[commands.length + 1]; - cmds[0] = "functions"; - System.arraycopy(commands, 0, cmds, 1, commands.length); - return pulsarCluster.runAdminCommandOnAnyBroker(cmds); - } - - private void setupFunctionWorkers() { - final int numFunctionWorkers = 2; - pulsarCluster.setupFunctionWorkers("", FunctionRuntimeType.THREAD, numFunctionWorkers); - } - - private void teardownFunctionWorkers() { - pulsarCluster.stopWorkers(); - } } - From 3c487982702838cd46efdecd70defed936d59016 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Sun, 15 Aug 2021 13:48:23 +0800 Subject: [PATCH 5/7] add more test --- .../pulsar/admin/cli/CmdFunctionsTest.java | 22 +++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java index e4284bdeebdf9..a441d2f9e478b 100644 --- a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java +++ b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java @@ -93,6 +93,7 @@ public IObjectFactory getObjectFactory() { private static final String PACKAGE_URL = "function://sample/ns1/jardummyexamples@1"; private static final String PACKAGE_GO_URL = "function://sample/ns1/godummyexamples@1"; private static final String PACKAGE_PY_URL = "function://sample/ns1/pydummyexamples@1"; + private static final String PACKAGE_INVALID_URL = "functionsample.jar"; private PulsarAdmin admin; private Functions functions; @@ -426,6 +427,27 @@ public void testCreatePyFunctionWithPackageUrl() throws Exception { verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); } + @Test + public void testCreateFunctionWithInvalidPackageUrl() throws Exception { + cmd.run(new String[] { + "create", + "--name", FN_NAME, + "--inputs", INPUT_TOPIC_NAME, + "--output", OUTPUT_TOPIC_NAME, + "--jar", PACKAGE_INVALID_URL, + "--tenant", "sample", + "--namespace", "ns1", + "--className", DummyFunction.class.getName(), + }); + + CreateFunction creater = cmd.getCreater(); + + assertEquals(FN_NAME, creater.getFunctionName()); + assertEquals(INPUT_TOPIC_NAME, creater.getInputs()); + assertEquals(OUTPUT_TOPIC_NAME, creater.getOutput()); + verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); + } + @Test public void testCreateFunctionWithoutBasicArguments() throws Exception { cmd.run(new String[] { From 4615fb3d18807a209407c796e5566e7e0919c205 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Mon, 16 Aug 2021 09:11:47 +0800 Subject: [PATCH 6/7] fix CI --- .../test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java index a441d2f9e478b..3e32961d8040f 100644 --- a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java +++ b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/CmdFunctionsTest.java @@ -445,7 +445,7 @@ public void testCreateFunctionWithInvalidPackageUrl() throws Exception { assertEquals(FN_NAME, creater.getFunctionName()); assertEquals(INPUT_TOPIC_NAME, creater.getInputs()); assertEquals(OUTPUT_TOPIC_NAME, creater.getOutput()); - verify(functions, times(1)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); + verify(functions, times(0)).createFunctionWithUrl(any(FunctionConfig.class), anyString()); } @Test From e703fa46df341fe2fa1a122e711610356287e717 Mon Sep 17 00:00:00 2001 From: Rui Fu Date: Mon, 16 Aug 2021 21:23:40 +0800 Subject: [PATCH 7/7] reduce duplicate code --- .../java/org/apache/pulsar/common/functions/Utils.java | 8 ++++++-- .../pulsar/functions/worker/rest/api/FunctionsImpl.java | 9 ++------- .../pulsar/functions/worker/rest/api/SinksImpl.java | 8 ++------ .../pulsar/functions/worker/rest/api/SourcesImpl.java | 8 ++------ 4 files changed, 12 insertions(+), 21 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java b/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java index 726426ffaf419..2de38bbe02fbc 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/functions/Utils.java @@ -37,8 +37,12 @@ public class Utils { public static boolean isFunctionPackageUrlSupported(String functionPkgUrl) { return isNotBlank(functionPkgUrl) && (functionPkgUrl.startsWith(HTTP) || functionPkgUrl.startsWith(FILE) - || Arrays.stream(PackageType.values()).anyMatch(type -> functionPkgUrl.startsWith(type.toString())) - && functionPkgUrl.contains("://")); + || hasPackageTypePrefix(functionPkgUrl)); + } + + public static boolean hasPackageTypePrefix(String destPkgUrl) { + return Arrays.stream(PackageType.values()).anyMatch(type -> destPkgUrl.startsWith(type.toString()) + && destPkgUrl.contains("://")); } public static void inferMissingFunctionName(FunctionConfig functionConfig) { diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java index 69f014317a4b7..8c3585100cb95 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImpl.java @@ -154,7 +154,7 @@ public void registerFunction(final String tenant, // validate parameters try { if (isPkgUrlProvided) { - if (hasPackageTypePrefix(functionPkgUrl)) { + if (Utils.hasPackageTypePrefix(functionPkgUrl)) { componentPackageFile = downloadPackageFile(functionPkgUrl); } else { if (!Utils.isFunctionPackageUrlSupported(functionPkgUrl)) { @@ -323,7 +323,7 @@ public void updateFunction(final String tenant, // validate parameters try { if (isNotBlank(functionPkgUrl)) { - if (hasPackageTypePrefix(functionPkgUrl)) { + if (Utils.hasPackageTypePrefix(functionPkgUrl)) { componentPackageFile = downloadPackageFile(functionName); } else { try { @@ -759,11 +759,6 @@ private Function.FunctionDetails validateUpdateRequestParams(final String tenant } - private static boolean hasPackageTypePrefix(String destPkgUrl) { - return Arrays.stream(PackageType.values()).anyMatch(type -> destPkgUrl.startsWith(type.toString()) - && destPkgUrl.contains("://")); - } - private File downloadPackageFile(String packageName) throws IOException, PulsarAdminException { return downloadPackageFile(worker(), packageName); } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java index 6bd80d5cc92bc..31a72347de04c 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SinksImpl.java @@ -153,7 +153,7 @@ public void registerSink(final String tenant, // validate parameters try { if (isPkgUrlProvided) { - if (hasPackageTypePrefix(sinkPkgUrl)) { + if (Utils.hasPackageTypePrefix(sinkPkgUrl)) { componentPackageFile = downloadPackageFile(sinkPkgUrl); } else { if (!Utils.isFunctionPackageUrlSupported(sinkPkgUrl)) { @@ -323,7 +323,7 @@ public void updateSink(final String tenant, // validate parameters try { if (isNotBlank(sinkPkgUrl)) { - if (hasPackageTypePrefix(sinkPkgUrl)) { + if (Utils.hasPackageTypePrefix(sinkPkgUrl)) { componentPackageFile = downloadPackageFile(sinkPkgUrl); } else { try { @@ -749,10 +749,6 @@ private Function.FunctionDetails validateUpdateRequestParams(final String tenant return SinkConfigUtils.convert(sinkConfig, sinkDetails); } - private static boolean hasPackageTypePrefix(String destPkgUrl) { - return Arrays.stream(PackageType.values()).anyMatch(type -> destPkgUrl.startsWith(type.toString())); - } - private File downloadPackageFile(String packageName) throws IOException, PulsarAdminException { return FunctionsImpl.downloadPackageFile(worker(), packageName); } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java index 2a0edb187e989..1e9148bd658ae 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/SourcesImpl.java @@ -153,7 +153,7 @@ public void registerSource(final String tenant, // validate parameters try { if (isPkgUrlProvided) { - if (hasPackageTypePrefix(sourcePkgUrl)) { + if (Utils.hasPackageTypePrefix(sourcePkgUrl)) { componentPackageFile = downloadPackageFile(sourcePkgUrl); } else { if (!Utils.isFunctionPackageUrlSupported(sourcePkgUrl)) { @@ -321,7 +321,7 @@ public void updateSource(final String tenant, // validate parameters try { if (isNotBlank(sourcePkgUrl)) { - if (hasPackageTypePrefix(sourcePkgUrl)) { + if (Utils.hasPackageTypePrefix(sourcePkgUrl)) { componentPackageFile = downloadPackageFile(sourcePkgUrl); } else { try { @@ -746,10 +746,6 @@ private Function.FunctionDetails validateUpdateRequestParams(final String tenant return SourceConfigUtils.convert(sourceConfig, sourceDetails); } - private static boolean hasPackageTypePrefix(String destPkgUrl) { - return Arrays.stream(PackageType.values()).anyMatch(type -> destPkgUrl.startsWith(type.toString())); - } - private File downloadPackageFile(String packageName) throws IOException, PulsarAdminException { return FunctionsImpl.downloadPackageFile(worker(), packageName); }