This is an automated email from the ASF dual-hosted git repository.
rdhabalia pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 358b68d add admin-api to update function with url (#1987)
358b68d is described below
commit 358b68dd560640ed025f91b13971367ff345a591
Author: Rajan Dhabalia <[email protected]>
AuthorDate: Thu Jun 21 23:42:33 2018 -0700
add admin-api to update function with url (#1987)
* add admin-api to update function with url
* add cli option for jar-url
---
.../pulsar/broker/admin/impl/FunctionsBase.java | 3 +-
.../org/apache/pulsar/io/PulsarSinkE2ETest.java | 7 ++--
.../org/apache/pulsar/client/admin/Functions.java | 22 ++++++++++++
.../client/admin/internal/FunctionsImpl.java | 17 ++++++++++
.../org/apache/pulsar/admin/cli/CmdFunctions.java | 6 +++-
.../java/org/apache/pulsar/admin/cli/CmdSinks.java | 6 +++-
.../org/apache/pulsar/admin/cli/CmdSources.java | 12 +++++--
.../functions/worker/rest/api/FunctionsImpl.java | 21 ++++++++----
.../worker/rest/api/v2/FunctionApiV2Resource.java | 3 +-
.../rest/api/v2/FunctionApiV2ResourceTest.java | 39 ++++++++++++++++++++++
10 files changed, 120 insertions(+), 16 deletions(-)
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java
index bacade6..6955d5f 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/FunctionsBase.java
@@ -95,10 +95,11 @@ public class FunctionsBase extends AdminResource implements
Supplier<WorkerServi
final @PathParam("functionName") String
functionName,
final @FormDataParam("data") InputStream
uploadedInputStream,
final @FormDataParam("data")
FormDataContentDisposition fileDetail,
+ final @FormDataParam("url") String
functionPkgUrl,
final @FormDataParam("functionDetails")
String functionDetailsJson) {
return functions.updateFunction(
- tenant, namespace, functionName, uploadedInputStream, fileDetail,
functionDetailsJson);
+ tenant, namespace, functionName, uploadedInputStream, fileDetail,
functionPkgUrl, functionDetailsJson);
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java
index 4987a32..7b08f6b 100644
--- a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java
+++ b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarSinkE2ETest.java
@@ -182,9 +182,7 @@ public class PulsarSinkE2ETest {
workerConfig.getClientAuthenticationParameters());
}
pulsarClient = clientBuilder.build();
- // pulsarClient =
PulsarClient.builder().serviceUrl(urlTls.toString()).statsInterval(0,
- // TimeUnit.SECONDS).build();
-
+
TenantInfo propAdmin = new TenantInfo();
propAdmin.setAllowedClusters(Sets.newHashSet(Lists.newArrayList("use")));
admin.tenants().updateTenant(tenant, propAdmin);
@@ -257,6 +255,9 @@ public class PulsarSinkE2ETest {
+
PulsarSink.class.getProtectionDomain().getCodeSource().getLocation().getPath();
FunctionDetails functionDetails = createSinkConfig(jarFilePathUrl,
tenant, namespacePortion, "PulsarSink-test");
admin.functions().createFunctionWithUrl(functionDetails,
jarFilePathUrl);
+
+ // try to update function to test: update-function functionality
+ admin.functions().updateFunctionWithUrl(functionDetails,
jarFilePathUrl);
retryStrategically((test) -> {
try {
diff --git
a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Functions.java
b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Functions.java
index a86036c..135c337 100644
---
a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Functions.java
+++
b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Functions.java
@@ -117,6 +117,28 @@ public interface Functions {
* Unexpected error
*/
void updateFunction(FunctionDetails functionDetails, String fileName)
throws PulsarAdminException;
+
+ /**
+ * Update the configuration for a function.
+ * <pre>
+ * Update a function by providing url from which fun-pkg can be
downloaded. supported url: http/file
+ * eg:
+ * File: file:/dir/fileName.jar
+ * Http: http://www.repo.com/fileName.jar
+ * </pre>
+ *
+ * @param functionDetails
+ * the function configuration object
+ * @param pkgUrl
+ * url from which pkg can be downloaded
+ * @throws NotAuthorizedException
+ * You don't have admin permission to create the cluster
+ * @throws NotFoundException
+ * Cluster doesn't exist
+ * @throws PulsarAdminException
+ * Unexpected error
+ */
+ void updateFunctionWithUrl(FunctionDetails functionDetails, String pkgUrl)
throws PulsarAdminException;
/**
* Delete an existing function
diff --git
a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/FunctionsImpl.java
b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/FunctionsImpl.java
index 68d3046..e7e39b1 100644
---
a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/FunctionsImpl.java
+++
b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/FunctionsImpl.java
@@ -166,6 +166,23 @@ public class FunctionsImpl extends BaseResource implements
Functions {
}
@Override
+ public void updateFunctionWithUrl(FunctionDetails functionDetails, String
pkgUrl) throws PulsarAdminException {
+ try {
+ final FormDataMultiPart mp = new FormDataMultiPart();
+
+ mp.bodyPart(new FormDataBodyPart("url", pkgUrl,
MediaType.TEXT_PLAIN_TYPE));
+
+ mp.bodyPart(new FormDataBodyPart("functionDetails",
printJson(functionDetails),
+ MediaType.APPLICATION_JSON_TYPE));
+
request(functions.path(functionDetails.getTenant()).path(functionDetails.getNamespace())
+ .path(functionDetails.getName())).put(Entity.entity(mp,
MediaType.MULTIPART_FORM_DATA),
+ ErrorData.class);
+ } catch (Exception e) {
+ throw getApiException(e);
+ }
+ }
+
+ @Override
public String triggerFunction(String tenant, String namespace, String
functionName, String topic, String triggerValue, String triggerFile) throws
PulsarAdminException {
try {
final FormDataMultiPart mp = new FormDataMultiPart();
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
index 0850172..5900295 100644
---
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
+++
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java
@@ -758,7 +758,11 @@ public class CmdFunctions extends CmdBase {
class UpdateFunction extends FunctionDetailsCommand {
@Override
void runCmd() throws Exception {
- admin.functions().updateFunction(convert(functionConfig),
userCodeFile);
+ if (Utils.isFunctionPackageUrlSupported(jarFile)) {
+
admin.functions().updateFunctionWithUrl(convert(functionConfig), jarFile);
+ } else {
+ admin.functions().updateFunction(convert(functionConfig),
userCodeFile);
+ }
print("Updated successfully");
}
}
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
index 8d05ec8..b5993e3 100644
---
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
+++
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSinks.java
@@ -156,7 +156,11 @@ public class CmdSinks extends CmdBase {
class UpdateSink extends SinkCommand {
@Override
void runCmd() throws Exception {
- admin.functions().updateFunction(createSinkConfig(sinkConfig),
jarFile);
+ if (Utils.isFunctionPackageUrlSupported(jarFile)) {
+
admin.functions().updateFunctionWithUrl(createSinkConfig(sinkConfig), jarFile);
+ } else {
+ admin.functions().updateFunction(createSinkConfig(sinkConfig),
jarFile);
+ }
print("Updated successfully");
}
}
diff --git
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
index e49357c..b3e49fd 100644
---
a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
+++
b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdSources.java
@@ -142,7 +142,11 @@ public class CmdSources extends CmdBase {
public class CreateSource extends SourceCommand {
@Override
void runCmd() throws Exception {
- admin.functions().createFunction(createSourceConfig(sourceConfig),
jarFile);
+ if (Utils.isFunctionPackageUrlSupported(jarFile)) {
+
admin.functions().createFunctionWithUrl(createSourceConfig(sourceConfig),
jarFile);
+ } else {
+
admin.functions().createFunction(createSourceConfig(sourceConfig), jarFile);
+ }
print("Created successfully");
}
}
@@ -151,7 +155,11 @@ public class CmdSources extends CmdBase {
public class UpdateSource extends SourceCommand {
@Override
void runCmd() throws Exception {
- admin.functions().updateFunction(createSourceConfig(sourceConfig),
jarFile);
+ if (Utils.isFunctionPackageUrlSupported(jarFile)) {
+
admin.functions().updateFunctionWithUrl(createSourceConfig(sourceConfig),
jarFile);
+ } else {
+
admin.functions().updateFunction(createSourceConfig(sourceConfig), jarFile);
+ }
print("Updated successfully");
}
}
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 9d475b2..9da5e6d 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
@@ -176,6 +176,7 @@ public class FunctionsImpl {
final @PathParam("functionName") String
functionName,
final @FormDataParam("data") InputStream
uploadedInputStream,
final @FormDataParam("data")
FormDataContentDisposition fileDetail,
+ final @FormDataParam("url") String
functionPkgUrl,
final @FormDataParam("functionDetails")
String functionDetailsJson) {
if (!isWorkerServiceAvailable()) {
@@ -183,13 +184,19 @@ public class FunctionsImpl {
}
FunctionDetails functionDetails;
+ boolean isPkgUrlProvided = StringUtils.isNotBlank(functionPkgUrl);
// validate parameters
try {
- functionDetails = validateUpdateRequestParams(tenant, namespace,
functionName,
- uploadedInputStream, fileDetail, functionDetailsJson);
- } catch (IllegalArgumentException e) {
- log.error("Invalid update function request @ /{}/{}/{}",
- tenant, namespace, functionName, e);
+ if(isPkgUrlProvided) {
+ functionDetails =
validateUpdateRequestParamsWithPkgUrl(tenant, namespace, functionName,
+ functionPkgUrl, functionDetailsJson);
+ }else {
+ functionDetails = validateUpdateRequestParams(tenant,
namespace, functionName,
+ uploadedInputStream, fileDetail, functionDetailsJson);
+ }
+ } catch (Exception e) {
+ log.error("Invalid register function request @ /{}/{}/{}",
+ tenant, namespace, functionName, e);
return Response.status(Status.BAD_REQUEST)
.type(MediaType.APPLICATION_JSON)
.entity(new ErrorData(e.getMessage())).build();
@@ -210,10 +217,10 @@ public class FunctionsImpl {
.setVersion(0);
PackageLocationMetaData.Builder packageLocationMetaDataBuilder =
PackageLocationMetaData.newBuilder()
- .setPackagePath(createPackagePath(tenant, namespace,
functionName, fileDetail.getFileName()));
+ .setPackagePath(isPkgUrlProvided ? functionPkgUrl :
createPackagePath(tenant, namespace, functionName, fileDetail.getFileName()));
functionMetaDataBuilder.setPackageLocation(packageLocationMetaDataBuilder);
- return updateRequest(functionMetaDataBuilder.build(),
uploadedInputStream);
+ return isPkgUrlProvided ?
updateRequest(functionMetaDataBuilder.build()) :
updateRequest(functionMetaDataBuilder.build(), uploadedInputStream);
}
diff --git
a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2Resource.java
b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2Resource.java
index 778eaab..c382bd6 100644
---
a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2Resource.java
+++
b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2Resource.java
@@ -64,10 +64,11 @@ public class FunctionApiV2Resource extends
FunctionApiResource {
final @PathParam("functionName") String
functionName,
final @FormDataParam("data") InputStream
uploadedInputStream,
final @FormDataParam("data")
FormDataContentDisposition fileDetail,
+ final @FormDataParam("url") String
functionPkgUrl,
final @FormDataParam("functionDetails")
String functionDetailsJson) {
return functions.updateFunction(
- tenant, namespace, functionName, uploadedInputStream, fileDetail,
functionDetailsJson);
+ tenant, namespace, functionName, uploadedInputStream, fileDetail,
functionPkgUrl, functionDetailsJson);
}
diff --git
a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2ResourceTest.java
b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2ResourceTest.java
index 49dede5..599a30e 100644
---
a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2ResourceTest.java
+++
b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v2/FunctionApiV2ResourceTest.java
@@ -585,6 +585,7 @@ public class FunctionApiV2ResourceTest {
function,
inputStream,
details,
+ null,
org.apache.pulsar.functions.utils.Utils.printJson(functionDetails));
assertEquals(Status.BAD_REQUEST.getStatusCode(), response.getStatus());
@@ -612,6 +613,7 @@ public class FunctionApiV2ResourceTest {
function,
mockedInputStream,
mockedFormData,
+ null,
org.apache.pulsar.functions.utils.Utils.printJson(functionDetails));
}
@@ -662,6 +664,43 @@ public class FunctionApiV2ResourceTest {
}
@Test
+ public void testUpdateFunctionWithUrl() throws IOException {
+ Configurator.setRootLevel(Level.DEBUG);
+
+ String fileLocation =
FutureUtil.class.getProtectionDomain().getCodeSource().getLocation().getPath();
+ String filePackageUrl = "file://" + fileLocation;
+
+ SinkSpec sinkSpec = SinkSpec.newBuilder()
+ .setTopic(outputTopic)
+ .setSerDeClassName(outputSerdeClassName).build();
+ FunctionDetails functionDetails = FunctionDetails.newBuilder()
+ .setTenant(tenant).setNamespace(namespace).setName(function)
+ .setSink(sinkSpec)
+ .setClassName(className)
+ .setParallelism(parallelism)
+
.setSource(SourceSpec.newBuilder().setSubscriptionType(subscriptionType)
+
.putAllTopicsToSerDeClassName(topicsToSerDeClassName)).build();
+
+ when(mockedManager.containsFunction(eq(tenant), eq(namespace),
eq(function))).thenReturn(true);
+ RequestResult rr = new RequestResult()
+ .setSuccess(true)
+ .setMessage("function registered");
+ CompletableFuture<RequestResult> requestResult =
CompletableFuture.completedFuture(rr);
+
when(mockedManager.updateFunction(any(FunctionMetaData.class))).thenReturn(requestResult);
+
+ Response response = resource.updateFunction(
+ tenant,
+ namespace,
+ function,
+ null,
+ null,
+ filePackageUrl,
+
org.apache.pulsar.functions.utils.Utils.printJson(functionDetails));
+
+ assertEquals(Status.OK.getStatusCode(), response.getStatus());
+ }
+
+ @Test
public void testUpdateFunctionFailure() throws Exception {
mockStatic(Utils.class);
doNothing().when(Utils.class);