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);

Reply via email to