rdhabalia closed pull request #1987: add admin-api to update function with url
URL: https://github.com/apache/incubator-pulsar/pull/1987
This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:
As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):
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 bacade64b6..6955d5f6ab 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 Response updateFunction(final @PathParam("tenant")
String tenant,
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 4987a32c1c..7b08f6b05b 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 @@ void setup(Method method) throws Exception {
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 void testE2EPulsarSink() throws Exception {
+
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 a86036c194..135c33749a 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 @@
* 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 68d3046250..e7e39b1577 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
@@ -165,6 +165,23 @@ public void updateFunction(FunctionDetails
functionDetails, String fileName) thr
}
}
+ @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 {
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 1dbe50b0ba..47a06bdaaf 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 @@ void runCmd() throws Exception {
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 83eb7f577f..d175833c77 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 @@ void runCmd() throws Exception {
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 32d9a6ca4c..1d564525ba 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 @@ void runCmd() throws Exception {
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 @@ void runCmd() throws Exception {
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 9d475b24cf..9da5e6d07c 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 Response updateFunction(final @PathParam("tenant")
String tenant,
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 Response updateFunction(final @PathParam("tenant")
String tenant,
}
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 Response updateFunction(final @PathParam("tenant")
String tenant,
.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 778eaaba6f..c382bd6cb0 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 Response updateFunction(final @PathParam("tenant")
String tenant,
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 49dede5265..599a30ef95 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 @@ private void testUpdateFunctionMissingArguments(
function,
inputStream,
details,
+ null,
org.apache.pulsar.functions.utils.Utils.printJson(functionDetails));
assertEquals(Status.BAD_REQUEST.getStatusCode(), response.getStatus());
@@ -612,6 +613,7 @@ private Response updateDefaultFunction() throws IOException
{
function,
mockedInputStream,
mockedFormData,
+ null,
org.apache.pulsar.functions.utils.Utils.printJson(functionDetails));
}
@@ -661,6 +663,43 @@ public void testUpdateFunctionSuccess() throws Exception {
assertEquals(Status.OK.getStatusCode(), response.getStatus());
}
+ @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);
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services