diff --git a/conf/functions_worker.yml b/conf/functions_worker.yml index 6f995576ebd64..d362887537fe9 100644 --- a/conf/functions_worker.yml +++ b/conf/functions_worker.yml @@ -129,6 +129,8 @@ clusterCoordinationTopicName: coordinate useCompactedMetadataTopic: false # Number of threads to use for HTTP requests processing. Default is set to 8 numHttpServerThreads: 8 +# Maximum parallelism used by function status-summary batch queries. Default is 4. +functionsStatusSummaryMaxParallelism: 4 # function assignment and scheduler diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Functions.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Functions.java index 3932ccf6c9180..3c5a223f6bf47 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Functions.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Functions.java @@ -18,6 +18,8 @@ */ package org.apache.pulsar.client.admin; +import java.util.ArrayList; +import java.util.Comparator; import java.util.List; import java.util.Set; import java.util.concurrent.CompletableFuture; @@ -32,6 +34,8 @@ import org.apache.pulsar.common.policies.data.FunctionInstanceStatsData; import org.apache.pulsar.common.policies.data.FunctionStats; import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; /** * Admin interface for function management. @@ -69,6 +73,148 @@ public interface Functions { */ CompletableFuture> getFunctionsAsync(String tenant, String namespace); + /** + * Get a batch status summary for all functions in a namespace. + *

+ * Returns a lightweight summary (name, state, instance counts) for every function + * under the given tenant/namespace in a single API call. This is not equivalent to + * calling {@link #getFunctionStatus} for each function - it returns only aggregated + * counts, not per-instance details. + *

+ * Individual function failures are isolated: a function whose status cannot be + * retrieved will appear with {@code state = UNKNOWN} and a non-null {@code error}. + * + * @param tenant + * Tenant name + * @param namespace + * Namespace name + * + * @return list of status summaries, one per function + * + * @throws NotAuthorizedException + * Don't have admin permission + * @throws PulsarAdminException + * Unexpected error + */ + default FunctionStatusPage getFunctionsWithStatus(String tenant, String namespace) + throws PulsarAdminException { + return getFunctionsWithStatus(tenant, namespace, null, null); + } + + /** + * Get a batch status summary for all functions in a namespace asynchronously. + *

+ * Async version of {@link #getFunctionsWithStatus(String, String)}. + * + * @param tenant + * Tenant name + * @param namespace + * Namespace name + * + * @return a future that completes with the list of status summaries + */ + default CompletableFuture getFunctionsWithStatusAsync(String tenant, + String namespace) { + return getFunctionsWithStatusAsync(tenant, namespace, null, null); + } + + /** + * Get a paginated batch status summary for functions in a namespace. + *

+ * The {@code startAfter} is an exclusive cursor based on function name + * in lexicographical order. + * + * @param tenant + * Tenant name + * @param namespace + * Namespace name + * @param limit + * Maximum number of functions to return; must be greater than 0 when provided + * @param startAfter + * Exclusive continuation token from previous page; null means from beginning + * @return list of status summaries for the requested page + * @throws PulsarAdminException + * Unexpected error + */ + default FunctionStatusPage getFunctionsWithStatus( + String tenant, String namespace, Integer limit, String startAfter) + throws PulsarAdminException { + if (limit != null && limit <= 0) { + throw new IllegalArgumentException("limit must be greater than 0"); + } + + List functionNames = getFunctions(tenant, namespace); + List pagedNames = new ArrayList<>(functionNames); + pagedNames.sort(String::compareTo); + + int startIndex = 0; + if (startAfter != null && !startAfter.isEmpty()) { + while (startIndex < pagedNames.size() && pagedNames.get(startIndex).compareTo(startAfter) <= 0) { + startIndex++; + } + } + + int endIndex = limit == null ? pagedNames.size() : Math.min(pagedNames.size(), startIndex + limit); + List summaries = new ArrayList<>(Math.max(0, endIndex - startIndex)); + for (int index = startIndex; index < endIndex; index++) { + String functionName = pagedNames.get(index); + try { + FunctionStatus status = getFunctionStatus(tenant, namespace, functionName); + FunctionStatusSummary.SummaryState state = status.getNumInstances() <= 0 + ? FunctionStatusSummary.SummaryState.UNKNOWN + : status.getNumRunning() == status.getNumInstances() + ? FunctionStatusSummary.SummaryState.RUNNING + : status.getNumRunning() == 0 + ? FunctionStatusSummary.SummaryState.STOPPED + : FunctionStatusSummary.SummaryState.PARTIAL; + summaries.add(FunctionStatusSummary.builder() + .name(functionName) + .state(state) + .numInstances(status.getNumInstances()) + .numRunning(status.getNumRunning()) + .build()); + } catch (PulsarAdminException e) { + summaries.add(FunctionStatusSummary.builder() + .name(functionName) + .state(FunctionStatusSummary.SummaryState.UNKNOWN) + .error(e.getMessage()) + .build()); + } + } + summaries.sort(Comparator.comparing(FunctionStatusSummary::getName)); + + String nextStartAfter = endIndex < pagedNames.size() ? pagedNames.get(endIndex - 1) : null; + return FunctionStatusPage.builder() + .summaries(summaries) + .nextStartAfter(nextStartAfter) + .build(); + } + + /** + * Async paginated version of {@link #getFunctionsWithStatus(String, String, Integer, String)}. + * + * @param tenant + * Tenant name + * @param namespace + * Namespace name + * @param limit + * Maximum number of functions to return; must be greater than 0 when provided + * @param startAfter + * Exclusive cursor (function name) from previous page; null means from beginning + * @return a future that completes with the paginated response + */ + default CompletableFuture getFunctionsWithStatusAsync( + String tenant, String namespace, Integer limit, String startAfter) { + try { + return CompletableFuture.completedFuture( + getFunctionsWithStatus(tenant, namespace, limit, startAfter)); + } catch (Exception e) { + CompletableFuture future = new CompletableFuture<>(); + future.completeExceptionally(e); + return future; + } + } + /** * Get the configuration for the specified function. *

diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/FunctionStatusPage.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/FunctionStatusPage.java new file mode 100644 index 0000000000000..546d436bb4688 --- /dev/null +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/FunctionStatusPage.java @@ -0,0 +1,34 @@ +/* + * 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.pulsar.common.policies.data; + +import java.util.List; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class FunctionStatusPage { + private List summaries; + private String nextStartAfter; +} diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/FunctionStatusSummary.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/FunctionStatusSummary.java new file mode 100644 index 0000000000000..836b466f7c018 --- /dev/null +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/FunctionStatusSummary.java @@ -0,0 +1,67 @@ +/* + * 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.pulsar.common.policies.data; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + +/** + * Summary status of a single Pulsar Function, used by the batch + * status-summary endpoint to avoid N+1 API calls. + */ +@Data +@Builder +@NoArgsConstructor +@AllArgsConstructor +public class FunctionStatusSummary { + + /** + * Runtime state derived from instance counts. + */ + public enum SummaryState { + RUNNING, + STOPPED, + PARTIAL, + UNKNOWN + } + + /** + * Classification of the status retrieval failure. + */ + public enum ErrorType { + AUTHENTICATION_FAILED, + FUNCTION_NOT_FOUND, + NETWORK_ERROR, + INTERNAL_ERROR + } + + private String name; + private SummaryState state; + private int numInstances; + private int numRunning; + + /** + * Non-null when the status query for this function failed; + * the function will have {@code state = UNKNOWN} in that case. + */ + private String error; + private ErrorType errorType; +} 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 bfcc3fe39a444..b0668acdc2e22 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 @@ -25,7 +25,13 @@ import java.io.File; import java.io.FileOutputStream; import java.io.IOException; +import java.net.ConnectException; +import java.net.SocketTimeoutException; +import java.net.UnknownHostException; import java.nio.channels.FileChannel; +import java.nio.channels.UnresolvedAddressException; +import java.util.ArrayList; +import java.util.Comparator; import java.util.List; import java.util.Set; import java.util.concurrent.CompletableFuture; @@ -54,6 +60,9 @@ import org.apache.pulsar.common.policies.data.FunctionStats; import org.apache.pulsar.common.policies.data.FunctionStatsImpl; import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; +import org.apache.pulsar.common.util.FutureUtil; import org.asynchttpclient.AsyncCompletionHandlerBase; import org.asynchttpclient.HttpResponseBodyPart; import org.asynchttpclient.RequestBuilder; @@ -89,6 +98,183 @@ public CompletableFuture> getFunctionsAsync(String tenant, String n return asyncGetRequest(path, new GenericType>() {}); } + @Override + public FunctionStatusPage getFunctionsWithStatus(String tenant, String namespace) + throws PulsarAdminException { + return sync(() -> getFunctionsWithStatusAsync(tenant, namespace, null, null)); + } + + @Override + public CompletableFuture getFunctionsWithStatusAsync( + String tenant, String namespace) { + return getFunctionsWithStatusAsync(tenant, namespace, null, null); + } + + @Override + public FunctionStatusPage getFunctionsWithStatus( + String tenant, String namespace, Integer limit, String startAfter) + throws PulsarAdminException { + return sync(() -> getFunctionsWithStatusAsync(tenant, namespace, limit, startAfter)); + } + + @Override + public CompletableFuture getFunctionsWithStatusAsync( + String tenant, String namespace, Integer limit, String startAfter) { + WebTarget path = functions.path(tenant).path(namespace).path("status").path("summary"); + if (limit != null) { + path = path.queryParam("limit", limit); + } + if (startAfter != null && !startAfter.isEmpty()) { + path = path.queryParam("startAfter", startAfter); + } + CompletableFuture result = new CompletableFuture<>(); + asyncGetRequest(path, new GenericType() {}) + .whenComplete((summaries, error) -> { + if (error == null) { + result.complete(summaries); + return; + } + + Throwable cause = FutureUtil.unwrapCompletionException(error); + if (isUnsupportedStatusSummaryEndpoint(cause)) { + log.debug( + "Falling back to legacy functions status queries for {}/{}", + tenant, namespace, cause); + getFunctionsWithStatusLegacyAsync(tenant, namespace, limit, startAfter) + .whenComplete((fallbackSummaries, fallbackError) -> { + if (fallbackError == null) { + result.complete(fallbackSummaries); + } else { + result.completeExceptionally( + FutureUtil.unwrapCompletionException(fallbackError)); + } + }); + return; + } + + result.completeExceptionally(cause); + }); + return result; + } + + private CompletableFuture getFunctionsWithStatusLegacyAsync( + String tenant, String namespace, Integer limit, String startAfter) { + return getFunctionsAsync(tenant, namespace).thenCompose(functionNames -> { + List sorted = new ArrayList<>(functionNames); + sorted.sort(String::compareTo); + List pagedNames = pageFunctionNames(functionNames, limit, startAfter); + List> summaryFutures = pagedNames.stream() + .map(functionName -> getFunctionStatusAsync(tenant, namespace, functionName) + .handle((status, error) -> buildStatusSummary(functionName, status, error))) + .collect(Collectors.toList()); + return FutureUtil.waitForAll(new ArrayList<>(summaryFutures)) + .thenApply(__ -> { + List summaries = summaryFutures.stream() + .map(CompletableFuture::join) + .sorted(Comparator.comparing(FunctionStatusSummary::getName)) + .collect(Collectors.toList()); + + String nextStartAfter = null; + if (limit != null && !pagedNames.isEmpty()) { + String lastReturned = pagedNames.get(pagedNames.size() - 1); + int lastIndex = sorted.indexOf(lastReturned); + if (lastIndex >= 0 && lastIndex < sorted.size() - 1) { + nextStartAfter = lastReturned; + } + } + + return FunctionStatusPage.builder() + .summaries(summaries) + .nextStartAfter(nextStartAfter) + .build(); + }); + }); + } + + private static FunctionStatusSummary buildStatusSummary(String functionName, + FunctionStatus status, + Throwable error) { + if (error == null) { + return FunctionStatusSummary.builder() + .name(functionName) + .state(deriveState(status.getNumInstances(), status.getNumRunning())) + .numInstances(status.getNumInstances()) + .numRunning(status.getNumRunning()) + .build(); + } + + Throwable cause = FutureUtil.unwrapCompletionException(error); + return FunctionStatusSummary.builder() + .name(functionName) + .state(FunctionStatusSummary.SummaryState.UNKNOWN) + .error(cause.getMessage()) + .errorType(classifyError(cause)) + .build(); + } + + private static List pageFunctionNames(List functionNames, Integer limit, String startAfter) { + if (limit != null && limit <= 0) { + throw new IllegalArgumentException("limit must be greater than 0"); + } + + List sorted = new ArrayList<>(functionNames); + sorted.sort(String::compareTo); + int startIndex = 0; + if (startAfter != null && !startAfter.isEmpty()) { + while (startIndex < sorted.size() && sorted.get(startIndex).compareTo(startAfter) <= 0) { + startIndex++; + } + } + int endIndex = limit == null ? sorted.size() : Math.min(sorted.size(), startIndex + limit); + return startIndex >= sorted.size() ? List.of() : sorted.subList(startIndex, endIndex); + } + + private static FunctionStatusSummary.SummaryState deriveState(int numInstances, int numRunning) { + if (numInstances <= 0) { + return FunctionStatusSummary.SummaryState.UNKNOWN; + } + if (numRunning == numInstances) { + return FunctionStatusSummary.SummaryState.RUNNING; + } + if (numRunning == 0) { + return FunctionStatusSummary.SummaryState.STOPPED; + } + return FunctionStatusSummary.SummaryState.PARTIAL; + } + + private static boolean isUnsupportedStatusSummaryEndpoint(Throwable cause) { + return cause instanceof PulsarAdminException + && (((PulsarAdminException) cause).getStatusCode() + == Response.Status.NOT_FOUND.getStatusCode() + || ((PulsarAdminException) cause).getStatusCode() + == Response.Status.METHOD_NOT_ALLOWED.getStatusCode()); + } + + private static FunctionStatusSummary.ErrorType classifyError(Throwable error) { + if (error instanceof PulsarAdminException) { + int statusCode = ((PulsarAdminException) error).getStatusCode(); + if (statusCode == Response.Status.UNAUTHORIZED.getStatusCode() + || statusCode == Response.Status.FORBIDDEN.getStatusCode()) { + return FunctionStatusSummary.ErrorType.AUTHENTICATION_FAILED; + } + if (statusCode == Response.Status.NOT_FOUND.getStatusCode()) { + return FunctionStatusSummary.ErrorType.FUNCTION_NOT_FOUND; + } + } + + Throwable current = error; + while (current != null) { + if (current instanceof ConnectException + || current instanceof SocketTimeoutException + || current instanceof UnknownHostException + || current instanceof UnresolvedAddressException) { + return FunctionStatusSummary.ErrorType.NETWORK_ERROR; + } + current = current.getCause(); + } + return FunctionStatusSummary.ErrorType.INTERNAL_ERROR; + } + @Override public FunctionConfig getFunction(String tenant, String namespace, String function) throws PulsarAdminException { return sync(() -> getFunctionAsync(tenant, namespace, function)); diff --git a/pulsar-client-admin/src/test/java/org/apache/pulsar/client/admin/internal/FunctionsImplTest.java b/pulsar-client-admin/src/test/java/org/apache/pulsar/client/admin/internal/FunctionsImplTest.java new file mode 100644 index 0000000000000..0e7b5edc9d1f0 --- /dev/null +++ b/pulsar-client-admin/src/test/java/org/apache/pulsar/client/admin/internal/FunctionsImplTest.java @@ -0,0 +1,192 @@ +/* + * 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.pulsar.client.admin.internal; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import javax.ws.rs.client.WebTarget; +import javax.ws.rs.core.GenericType; +import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; +import org.testng.annotations.Test; + +public class FunctionsImplTest { + + @Test + public void testGetFunctionsWithStatusAsyncBuildsExpectedPath() throws Exception { + WebTarget root = mock(WebTarget.class); + WebTarget adminV3Functions = mock(WebTarget.class); + WebTarget tenantTarget = mock(WebTarget.class); + WebTarget namespaceTarget = mock(WebTarget.class); + WebTarget statusTarget = mock(WebTarget.class); + WebTarget summaryTarget = mock(WebTarget.class); + + when(root.path("/admin/v3/functions")).thenReturn(adminV3Functions); + when(adminV3Functions.path("tenant-a")).thenReturn(tenantTarget); + when(tenantTarget.path("namespace-a")).thenReturn(namespaceTarget); + when(namespaceTarget.path("status")).thenReturn(statusTarget); + when(statusTarget.path("summary")).thenReturn(summaryTarget); + + FunctionsImpl functions = org.mockito.Mockito.spy(new FunctionsImpl(root, null, null, 0)); + FunctionStatusPage expected = FunctionStatusPage.builder() + .summaries(Collections.singletonList( + FunctionStatusSummary.builder() + .name("fn-1") + .state(FunctionStatusSummary.SummaryState.RUNNING) + .build())) + .build(); + CompletableFuture response = CompletableFuture.completedFuture(expected); + doReturn(response).when(functions).asyncGetRequest(eq(summaryTarget), any(GenericType.class)); + + FunctionStatusPage actual = + functions.getFunctionsWithStatusAsync("tenant-a", "namespace-a").get(); + + verify(adminV3Functions).path("tenant-a"); + verify(tenantTarget).path("namespace-a"); + verify(namespaceTarget).path("status"); + verify(statusTarget).path("summary"); + assertEquals(actual, expected); + } + + @Test + public void testGetFunctionsWithStatusSyncDelegatesToAsync() throws Exception { + WebTarget root = mock(WebTarget.class); + WebTarget adminV3Functions = mock(WebTarget.class); + WebTarget tenantTarget = mock(WebTarget.class); + WebTarget namespaceTarget = mock(WebTarget.class); + WebTarget statusTarget = mock(WebTarget.class); + WebTarget summaryTarget = mock(WebTarget.class); + + when(root.path("/admin/v3/functions")).thenReturn(adminV3Functions); + when(adminV3Functions.path("tenant-b")).thenReturn(tenantTarget); + when(tenantTarget.path("namespace-b")).thenReturn(namespaceTarget); + when(namespaceTarget.path("status")).thenReturn(statusTarget); + when(statusTarget.path("summary")).thenReturn(summaryTarget); + + FunctionsImpl functions = org.mockito.Mockito.spy(new FunctionsImpl(root, null, null, 0)); + FunctionStatusPage expected = FunctionStatusPage.builder() + .summaries(Collections.singletonList( + FunctionStatusSummary.builder() + .name("fn-2") + .state(FunctionStatusSummary.SummaryState.STOPPED) + .build())) + .build(); + CompletableFuture response = CompletableFuture.completedFuture(expected); + doReturn(response).when(functions).asyncGetRequest(eq(summaryTarget), any(GenericType.class)); + + FunctionStatusPage actual = functions.getFunctionsWithStatus("tenant-b", "namespace-b"); + assertEquals(actual, expected); + } + + @Test + public void testGetFunctionsWithStatusAsyncFallsBackForLegacyBroker() throws Exception { + WebTarget root = mock(WebTarget.class); + WebTarget adminV3Functions = mock(WebTarget.class); + WebTarget tenantTarget = mock(WebTarget.class); + WebTarget namespaceTarget = mock(WebTarget.class); + WebTarget statusTarget = mock(WebTarget.class); + WebTarget summaryTarget = mock(WebTarget.class); + + when(root.path("/admin/v3/functions")).thenReturn(adminV3Functions); + when(adminV3Functions.path("tenant-c")).thenReturn(tenantTarget); + when(tenantTarget.path("namespace-c")).thenReturn(namespaceTarget); + when(namespaceTarget.path("status")).thenReturn(statusTarget); + when(statusTarget.path("summary")).thenReturn(summaryTarget); + + FunctionsImpl functions = org.mockito.Mockito.spy(new FunctionsImpl(root, null, null, 0)); + CompletableFuture failed = new CompletableFuture<>(); + failed.completeExceptionally(new PulsarAdminException.NotFoundException(null, "Not Found", 404)); + doReturn(failed).when(functions).asyncGetRequest(eq(summaryTarget), any(GenericType.class)); + doReturn(CompletableFuture.completedFuture(List.of("fn-b", "fn-a"))) + .when(functions).getFunctionsAsync("tenant-c", "namespace-c"); + + FunctionStatus runningStatus = new FunctionStatus(); + runningStatus.setNumInstances(1); + runningStatus.setNumRunning(1); + FunctionStatus stoppedStatus = new FunctionStatus(); + stoppedStatus.setNumInstances(1); + stoppedStatus.setNumRunning(0); + doReturn(CompletableFuture.completedFuture(stoppedStatus)) + .when(functions).getFunctionStatusAsync("tenant-c", "namespace-c", "fn-b"); + doReturn(CompletableFuture.completedFuture(runningStatus)) + .when(functions).getFunctionStatusAsync("tenant-c", "namespace-c", "fn-a"); + + FunctionStatusPage actual = + functions.getFunctionsWithStatusAsync("tenant-c", "namespace-c").get(); + + assertEquals(actual.getSummaries().size(), 2); + assertEquals(actual.getSummaries().get(0).getName(), "fn-a"); + assertEquals(actual.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.RUNNING); + assertEquals(actual.getSummaries().get(1).getName(), "fn-b"); + assertEquals(actual.getSummaries().get(1).getState(), FunctionStatusSummary.SummaryState.STOPPED); + } + + @Test + public void testGetFunctionsWithStatusAsyncFallbackPagesBeforeQueryingStatus() throws Exception { + WebTarget root = mock(WebTarget.class); + WebTarget adminV3Functions = mock(WebTarget.class); + WebTarget tenantTarget = mock(WebTarget.class); + WebTarget namespaceTarget = mock(WebTarget.class); + WebTarget statusTarget = mock(WebTarget.class); + WebTarget summaryTarget = mock(WebTarget.class); + WebTarget limitTarget = mock(WebTarget.class); + WebTarget continuationTarget = mock(WebTarget.class); + + when(root.path("/admin/v3/functions")).thenReturn(adminV3Functions); + when(adminV3Functions.path("tenant-d")).thenReturn(tenantTarget); + when(tenantTarget.path("namespace-d")).thenReturn(namespaceTarget); + when(namespaceTarget.path("status")).thenReturn(statusTarget); + when(statusTarget.path("summary")).thenReturn(summaryTarget); + when(summaryTarget.queryParam("limit", 1)).thenReturn(limitTarget); + when(limitTarget.queryParam("startAfter", "fn-a")).thenReturn(continuationTarget); + + FunctionsImpl functions = org.mockito.Mockito.spy(new FunctionsImpl(root, null, null, 0)); + CompletableFuture failed = new CompletableFuture<>(); + failed.completeExceptionally(new PulsarAdminException.NotFoundException(null, "Not Found", 404)); + doReturn(failed).when(functions).asyncGetRequest(eq(continuationTarget), any(GenericType.class)); + doReturn(CompletableFuture.completedFuture(List.of("fn-c", "fn-a", "fn-b"))) + .when(functions).getFunctionsAsync("tenant-d", "namespace-d"); + + FunctionStatus runningStatus = new FunctionStatus(); + runningStatus.setNumInstances(1); + runningStatus.setNumRunning(1); + doReturn(CompletableFuture.completedFuture(runningStatus)) + .when(functions).getFunctionStatusAsync("tenant-d", "namespace-d", "fn-b"); + + FunctionStatusPage actual = + functions.getFunctionsWithStatusAsync("tenant-d", "namespace-d", 1, "fn-a").get(); + + assertEquals(actual.getSummaries().size(), 1); + assertEquals(actual.getSummaries().get(0).getName(), "fn-b"); + verify(functions).getFunctionStatusAsync("tenant-d", "namespace-d", "fn-b"); + verify(functions, never()).getFunctionStatusAsync("tenant-d", "namespace-d", "fn-a"); + verify(functions, never()).getFunctionStatusAsync("tenant-d", "namespace-d", "fn-c"); + } +} 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 e5323fcd7fa0a..f8006df9f106a 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 @@ -27,10 +27,12 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; import java.io.PrintWriter; import java.io.StringWriter; +import java.util.List; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.admin.cli.CmdFunctions.CreateFunction; @@ -46,6 +48,8 @@ import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.common.functions.FunctionConfig; import org.apache.pulsar.common.functions.UpdateOptionsImpl; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; import org.apache.pulsar.functions.api.Context; import org.apache.pulsar.functions.api.Function; import org.apache.pulsar.functions.api.utils.IdentityFunction; @@ -630,6 +634,7 @@ public void testListFunctions() throws Exception { assertEquals(NAMESPACE, lister.getNamespace()); verify(functions, times(1)).getFunctions(eq(TENANT), eq(NAMESPACE)); + verify(functions, times(0)).getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE)); } @Test @@ -643,6 +648,7 @@ public void testListFunctionsWithDefaultValue() throws Exception { assertEquals("default", lister.getNamespace()); verify(functions, times(1)).getFunctions(eq("public"), eq("default")); + verify(functions, times(0)).getFunctionsWithStatus(eq("public"), eq("default")); } @Test @@ -914,4 +920,190 @@ public void testDownloadTransformFunction() throws Exception { verify(functions, times(1)) .downloadFunction(JAR_NAME, TENANT, NAMESPACE, FN_NAME, true); } + + @Test + public void testListFunctionsLongFormat() throws Exception { + FunctionStatusPage statusPage = FunctionStatusPage.builder() + .summaries(List.of( + FunctionStatusSummary.builder() + .name("fn-a") + .state(FunctionStatusSummary.SummaryState.RUNNING) + .numInstances(2).numRunning(2).build() + )) + .build(); + when(functions.getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE))) + .thenReturn(statusPage); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "-l" + }); + + verify(functions, times(1)) + .getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE)); + verify(functions, times(0)) + .getFunctions(eq(TENANT), eq(NAMESPACE)); + } + + @Test + public void testListFunctionsWithStateFilter() throws Exception { + FunctionStatusPage statusPage = FunctionStatusPage.builder() + .summaries(List.of( + FunctionStatusSummary.builder() + .name("fn-running") + .state(FunctionStatusSummary.SummaryState.RUNNING) + .numInstances(1).numRunning(1).build(), + FunctionStatusSummary.builder() + .name("fn-stopped") + .state(FunctionStatusSummary.SummaryState.STOPPED) + .numInstances(1).numRunning(0).build() + )) + .build(); + when(functions.getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE))) + .thenReturn(statusPage); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "--state", "RUNNING" + }); + + verify(functions, times(1)) + .getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE)); + verify(functions, times(0)) + .getFunctions(eq(TENANT), eq(NAMESPACE)); + } + + @Test + public void testListFunctionsWithUnknownStateFilter() throws Exception { + FunctionStatusPage statusPage = FunctionStatusPage.builder() + .summaries(List.of( + FunctionStatusSummary.builder() + .name("fn-unknown") + .state(FunctionStatusSummary.SummaryState.UNKNOWN) + .error("status unavailable") + .build(), + FunctionStatusSummary.builder() + .name("fn-running") + .state(FunctionStatusSummary.SummaryState.RUNNING) + .numInstances(1).numRunning(1).build() + )) + .build(); + when(functions.getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE))) + .thenReturn(statusPage); + + @Cleanup + StringWriter stringWriter = new StringWriter(); + @Cleanup + PrintWriter printWriter = new PrintWriter(stringWriter); + cmd.getCommander().setOut(printWriter); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "--state", "UNKNOWN" + }); + + verify(functions, times(1)) + .getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE)); + verify(functions, times(0)) + .getFunctions(eq(TENANT), eq(NAMESPACE)); + assertTrue(stringWriter.toString().contains("fn-unknown")); + assertFalse(stringWriter.toString().contains("fn-running")); + } + + @Test + public void testListFunctionsWithInvalidStateValue() throws Exception { + @Cleanup + StringWriter stringWriter = new StringWriter(); + @Cleanup + PrintWriter printWriter = new PrintWriter(stringWriter); + cmd.getCommander().setErr(printWriter); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "--state", "INVALID_STATE" + }); + + assertTrue(stringWriter.toString().contains("--state")); + assertTrue(stringWriter.toString().contains("INVALID_STATE")); + verify(functions, times(0)).getFunctionsWithStatus(anyString(), anyString()); + verify(functions, times(0)).getFunctions(anyString(), anyString()); + } + + @Test + public void testListFunctionsWithPaginationParams() throws Exception { + FunctionStatusPage statusPage = FunctionStatusPage.builder() + .summaries(List.of( + FunctionStatusSummary.builder() + .name("fn-b") + .state(FunctionStatusSummary.SummaryState.RUNNING) + .numInstances(1) + .numRunning(1) + .build() + )) + .build(); + when(functions.getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE), eq(1), eq("fn-a"))) + .thenReturn(statusPage); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "--limit", "1", + "--continuation-token", "fn-a" + }); + + verify(functions, times(1)).getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE), eq(1), eq("fn-a")); + verify(functions, times(0)).getFunctions(eq(TENANT), eq(NAMESPACE)); + verify(functions, times(0)).getFunctionsWithStatus(eq(TENANT), eq(NAMESPACE)); + } + + @Test + public void testListFunctionsWithInvalidLimitValue() throws Exception { + @Cleanup + StringWriter stringWriter = new StringWriter(); + @Cleanup + PrintWriter printWriter = new PrintWriter(stringWriter); + cmd.getCommander().setErr(printWriter); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "--limit", "0" + }); + + assertTrue(stringWriter.toString().contains("--limit")); + verify(functions, times(0)).getFunctions(anyString(), anyString()); + verify(functions, times(0)).getFunctionsWithStatus(anyString(), anyString()); + verify(functions, times(0)).getFunctionsWithStatus(anyString(), anyString(), any(), anyString()); + } + + @Test + public void testListFunctionsRejectsStateWithPagination() throws Exception { + @Cleanup + StringWriter stringWriter = new StringWriter(); + @Cleanup + PrintWriter printWriter = new PrintWriter(stringWriter); + cmd.getCommander().setErr(printWriter); + + cmd.run(new String[] { + "list", + "--tenant", TENANT, + "--namespace", NAMESPACE, + "--state", "RUNNING", + "--limit", "1" + }); + + verify(functions, times(0)).getFunctions(anyString(), anyString()); + verify(functions, times(0)).getFunctionsWithStatus(anyString(), anyString()); + verify(functions, times(0)).getFunctionsWithStatus(anyString(), anyString(), any(), anyString()); + } } 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 2a3b660d9b47b..feb28c66eb0bd 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 @@ -36,6 +36,7 @@ import java.util.List; import java.util.Map; import java.util.function.Supplier; +import java.util.stream.Collectors; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; @@ -53,6 +54,8 @@ import org.apache.pulsar.common.functions.UpdateOptionsImpl; import org.apache.pulsar.common.functions.Utils; import org.apache.pulsar.common.functions.WindowConfig; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; import org.apache.pulsar.common.util.ObjectMapperFactory; import picocli.CommandLine.Command; import picocli.CommandLine.Option; @@ -700,16 +703,16 @@ protected void validateFunctionConfigs(FunctionConfig functionConfig) { } } if (StringUtils.isEmpty(functionConfig.getName())) { - org.apache.pulsar.common.functions.Utils.inferMissingFunctionName(functionConfig); + Utils.inferMissingFunctionName(functionConfig); } if (StringUtils.isEmpty(functionConfig.getName())) { throw new IllegalArgumentException("No Function name specified"); } if (StringUtils.isEmpty(functionConfig.getTenant())) { - org.apache.pulsar.common.functions.Utils.inferMissingTenant(functionConfig); + Utils.inferMissingTenant(functionConfig); } if (StringUtils.isEmpty(functionConfig.getNamespace())) { - org.apache.pulsar.common.functions.Utils.inferMissingNamespace(functionConfig); + Utils.inferMissingNamespace(functionConfig); } if (isNotBlank(functionConfig.getJar()) && isNotBlank(functionConfig.getPy()) @@ -1017,16 +1020,16 @@ class UpdateFunction extends FunctionDetailsCommand { @Override protected void validateFunctionConfigs(FunctionConfig functionConfig) { if (StringUtils.isEmpty(functionConfig.getName())) { - org.apache.pulsar.common.functions.Utils.inferMissingFunctionName(functionConfig); + Utils.inferMissingFunctionName(functionConfig); } if (StringUtils.isEmpty(functionConfig.getName())) { throw new ParameterException("Function Name not provided"); } if (StringUtils.isEmpty(functionConfig.getTenant())) { - org.apache.pulsar.common.functions.Utils.inferMissingTenant(functionConfig); + Utils.inferMissingTenant(functionConfig); } if (StringUtils.isEmpty(functionConfig.getNamespace())) { - org.apache.pulsar.common.functions.Utils.inferMissingNamespace(functionConfig); + Utils.inferMissingNamespace(functionConfig); } } @@ -1050,9 +1053,77 @@ void runCmd() throws Exception { @Command(description = "List all Pulsar Functions running under a specific tenant and namespace") class ListFunctions extends NamespaceCommand { + + @Option(names = "--state", + description = "Filter by runtime state: RUNNING, STOPPED, PARTIAL, UNKNOWN; cannot be combined" + + " with --limit or --start-after") + private FunctionStatusSummary.SummaryState state; + + @Option(names = {"-l", "--long"}, + description = "Show extended output with state and instance counts") + private boolean longFormat; + + @Option(names = "--limit", + description = "Limit the number of status summaries returned (only with status-summary path)") + private Integer limit; + + @Option(names = "--start-after", + description = "Exclusive cursor (function name) for status-summary pagination") + private String startAfter; + @Override void runCmd() throws Exception { - print(getAdmin().functions().getFunctions(tenant, namespace)); + if (limit != null && limit <= 0) { + throw new ParameterException("--limit must be greater than 0"); + } + + // Prevent ambiguity in semantics + if (state != null && (limit != null || startAfter != null)) { + throw new ParameterException("--state cannot be combined with --limit or --start-after"); + } + + if (state == null && !longFormat && limit == null && startAfter == null) { + print(getAdmin().functions().getFunctions(tenant, namespace)); + return; + } + + FunctionStatusPage page = limit == null && startAfter == null + ? getAdmin().functions().getFunctionsWithStatus(tenant, namespace) + : getAdmin().functions().getFunctionsWithStatus(tenant, namespace, limit, startAfter); + + List summaries = page.getSummaries(); + + if (state != null) { + summaries = summaries.stream() + .filter(s -> s.getState() == state) + .collect(Collectors.toList()); + } + + if (longFormat) { + printLongFormat(summaries); + } else { + for (FunctionStatusSummary s : summaries) { + print(s.getName()); + } + } + + if (page.getNextStartAfter() != null) { + print("\nNext page: --start-after " + page.getNextStartAfter()); + } + } + + private void printLongFormat(List summaries) { + String header = String.format("%-40s %-10s %s", "NAME", "STATE", "RUNNING/INSTANCES"); + print(header); + for (FunctionStatusSummary s : summaries) { + String running = s.getState() == FunctionStatusSummary.SummaryState.UNKNOWN + ? "?" : String.valueOf(s.getNumRunning()); + String instances = s.getState() == FunctionStatusSummary.SummaryState.UNKNOWN + ? "?" : String.valueOf(s.getNumInstances()); + String line = String.format("%-40s %-10s %s/%s", + s.getName(), s.getState(), running, instances); + print(line); + } } } diff --git a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java index 036311ea13230..d7f4980b811f2 100644 --- a/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java +++ b/pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/worker/WorkerConfig.java @@ -129,6 +129,12 @@ public class WorkerConfig implements Serializable, PulsarConfiguration { ) private int numHttpServerThreads = 8; + @FieldContext( + category = CATEGORY_WORKER, + doc = "Maximum parallelism for function status-summary batch queries" + ) + private int functionsStatusSummaryMaxParallelism = 4; + @FieldContext( category = CATEGORY_WORKER, doc = "Enable the enforcement of limits on the incoming HTTP requests" diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerStatsManager.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerStatsManager.java index d746f867d8ece..ecc17287e0aba 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerStatsManager.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerStatsManager.java @@ -50,6 +50,7 @@ public class WorkerStatsManager { private static final String STOPPING_INSTANCE_PROCESS_TIME = "stop_instance_process_time_ms"; private static final String STARTING_INSTANCE_PROCESS_TIME = "start_instance_process_time_ms"; private static final String DRAIN_TOTAL_EXEC_TIME = "drain_execution_time_total_ms"; + private static final String FUNCTIONS_STATUS_SUMMARY_QUERY_TIME = "functions_status_summary_query_time_ms"; private static final String IS_LEADER = "is_leader"; @@ -79,6 +80,7 @@ public class WorkerStatsManager { private final Summary stopInstanceProcessTime; private final Summary startInstanceProcessTime; private final Summary drainTotalExecutionTime; + private final Summary functionsStatusSummaryQueryTime; // As an optimization private final Summary.Child statWorkerStartupTimeChild; @@ -90,6 +92,7 @@ public class WorkerStatsManager { private final Summary.Child stopInstanceProcessTimeChild; private final Summary.Child startInstanceProcessTimeChild; private final Summary.Child drainTotalExecutionTimeChild; + private final Summary.Child functionsStatusSummaryQueryTimeChild; public WorkerStatsManager(WorkerConfig workerConfig, boolean runAsStandalone) { @@ -179,6 +182,16 @@ public WorkerStatsManager(WorkerConfig workerConfig, boolean runAsStandalone) { .register(collectorRegistry); drainTotalExecutionTimeChild = drainTotalExecutionTime.labels(metricsLabels); + functionsStatusSummaryQueryTime = Summary.build() + .name(PULSAR_FUNCTION_WORKER_METRICS_PREFIX + FUNCTIONS_STATUS_SUMMARY_QUERY_TIME) + .help("Execution time of functions status summary batch query in milliseconds.") + .labelNames(metricsLabelNames) + .quantile(0.5, 0.01) + .quantile(0.9, 0.01) + .quantile(1, 0.01) + .register(collectorRegistry); + functionsStatusSummaryQueryTimeChild = functionsStatusSummaryQueryTime.labels(metricsLabels); + if (runAsStandalone) { Gauge.build("jvm_memory_direct_bytes_used", "-").create().setChild(new Gauge.Child() { @Override @@ -292,6 +305,10 @@ public void startInstanceProcessTimeEnd() { } } + public void observeFunctionsStatusSummaryQueryTime(double elapsedMs) { + functionsStatusSummaryQueryTimeChild.observe(elapsedMs); + } + public String getStatsAsString() throws IOException { statNumInstancesChild.set(functionRuntimeManager.getMyInstances()); 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 4cbd7c8cbcb12..dcd4bd4c43157 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 @@ -27,18 +27,30 @@ import java.io.File; import java.io.IOException; import java.io.InputStream; +import java.net.ConnectException; +import java.net.SocketTimeoutException; import java.net.URI; +import java.net.UnknownHostException; +import java.nio.channels.UnresolvedAddressException; +import java.util.ArrayList; import java.util.Collection; +import java.util.Collections; import java.util.LinkedList; import java.util.List; import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.function.Supplier; +import java.util.stream.Collectors; import javax.ws.rs.WebApplicationException; import javax.ws.rs.core.Response; import javax.ws.rs.core.UriBuilder; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.authentication.AuthenticationParameters; +import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.functions.FunctionConfig; import org.apache.pulsar.common.functions.FunctionDefinition; @@ -47,6 +59,8 @@ import org.apache.pulsar.common.functions.WorkerInfo; import org.apache.pulsar.common.policies.data.ExceptionInformation; import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; import org.apache.pulsar.common.util.RestException; import org.apache.pulsar.functions.auth.FunctionAuthData; import org.apache.pulsar.functions.instance.InstanceUtils; @@ -743,7 +757,7 @@ private Function.FunctionDetails validateUpdateRequestParams(final String tenant ValidatableFunctionPackage functionPackage = null; // check if function is builtin and extract classloader if (!StringUtils.isEmpty(archive)) { - if (archive.startsWith(org.apache.pulsar.common.functions.Utils.BUILTIN)) { + if (archive.startsWith(Utils.BUILTIN)) { archive = archive.replaceFirst("^builtin://", ""); FunctionsManager functionsManager = worker().getFunctionsManager(); @@ -788,4 +802,241 @@ private Function.FunctionDetails validateUpdateRequestParams(final String tenant } } } + + @Override + public FunctionStatusPage listFunctionsWithStatus( + final String tenant, + final String namespace, + final AuthenticationParameters authParams) { + return listFunctionsWithStatus(tenant, namespace, null, null, authParams); + } + + @Override + public FunctionStatusPage listFunctionsWithStatus( + final String tenant, + final String namespace, + final Integer limit, + final String startAfter, + final AuthenticationParameters authParams) { + if (!isWorkerServiceAvailable()) { + throwUnavailableException(); + } + if (limit != null && limit <= 0) { + throw new RestException(Response.Status.BAD_REQUEST, "limit must be greater than 0"); + } + + long startNs = System.nanoTime(); + ExecutorService summaryExecutor = null; + try { + // listFunctions already handles auth check and parameter validation + List functionNames = listFunctions(tenant, namespace, authParams); + List sorted = new ArrayList<>(functionNames); + sorted.sort(String::compareTo); + List pagedNames = pageFunctionNames(functionNames, limit, startAfter); + if (pagedNames.isEmpty()) { + return FunctionStatusPage.builder() + .summaries(Collections.emptyList()) + .nextStartAfter(null) + .build(); + } + + int configuredParallelism = worker().getWorkerConfig() != null + ? worker().getWorkerConfig().getFunctionsStatusSummaryMaxParallelism() : 4; + int maxConcurrency = Math.max(1, Math.min(configuredParallelism, pagedNames.size())); + summaryExecutor = Executors.newFixedThreadPool(maxConcurrency); + List> futures = new ArrayList<>(pagedNames.size()); + for (String name : pagedNames) { + futures.add(CompletableFuture.supplyAsync( + () -> buildSummary(tenant, namespace, name, authParams), summaryExecutor)); + } + + List summaries = futures.stream() + .map(CompletableFuture::join) + .collect(Collectors.toList()); + + String nextStartAfter = null; + if (limit != null && !pagedNames.isEmpty()) { + String lastReturned = pagedNames.get(pagedNames.size() - 1); + int lastIndex = sorted.indexOf(lastReturned); + if (lastIndex >= 0 && lastIndex < sorted.size() - 1) { + nextStartAfter = lastReturned; + } + } + + return FunctionStatusPage.builder() + .summaries(summaries) + .nextStartAfter(nextStartAfter) + .build(); + } finally { + if (summaryExecutor != null) { + summaryExecutor.shutdown(); + try { + if (!summaryExecutor.awaitTermination(5, TimeUnit.SECONDS)) { + summaryExecutor.shutdownNow(); + } + } catch (InterruptedException e) { + summaryExecutor.shutdownNow(); + Thread.currentThread().interrupt(); + } + } + if (worker().getWorkerStatsManager() != null) { + worker().getWorkerStatsManager() + .observeFunctionsStatusSummaryQueryTime(((double) System.nanoTime() - startNs) / 1.0E6D); + } + } + } + + private static List pageFunctionNames(List functionNames, Integer limit, String startAfter) { + if (functionNames.isEmpty()) { + return functionNames; + } + List sorted = new ArrayList<>(functionNames); + sorted.sort(String::compareTo); + + int startIndex = 0; + if (isNotBlank(startAfter)) { + while (startIndex < sorted.size() && sorted.get(startIndex).compareTo(startAfter) <= 0) { + startIndex++; + } + } + if (startIndex >= sorted.size()) { + return Collections.emptyList(); + } + if (limit == null) { + return sorted.subList(startIndex, sorted.size()); + } + + int endIndex = Math.min(sorted.size(), startIndex + limit); + return sorted.subList(startIndex, endIndex); + } + + private FunctionStatusSummary buildSummary(String tenant, String namespace, + String name, AuthenticationParameters authParams) { + try { + FunctionStatus status = getFunctionStatusForSummary(tenant, namespace, name, authParams); + return FunctionStatusSummary.builder() + .name(name) + .state(deriveState(status.getNumInstances(), status.getNumRunning())) + .numInstances(status.getNumInstances()) + .numRunning(status.getNumRunning()) + .build(); + } catch (Exception e) { + log.warn("{}/{}/{} Failed to get status for summary", tenant, namespace, name, e); + return FunctionStatusSummary.builder() + .name(name) + .state(FunctionStatusSummary.SummaryState.UNKNOWN) + .error(e.getMessage()) + .errorType(classifyError(e)) + .build(); + } + } + + private FunctionStatus getFunctionStatusForSummary(String tenant, String namespace, String name, + AuthenticationParameters authParams) + throws Exception { + try { + // Fast path: local worker service path. + return getFunctionStatus(tenant, namespace, name, null, authParams); + } catch (RestException localRestError) { + // Preserve local semantic 4xx errors (authn/authz/not-found/validation). + // 5xx errors are treated as recoverable and can still use admin fallback. + int status = localRestError.getResponse() != null ? localRestError.getResponse().getStatus() : 0; + if (status >= 400 && status < 500) { + throw localRestError; + } + return getFunctionStatusFromAdminFallback(tenant, namespace, name, localRestError); + } catch (Exception localError) { + return getFunctionStatusFromAdminFallback(tenant, namespace, name, localError); + } + } + + private FunctionStatus getFunctionStatusFromAdminFallback(String tenant, String namespace, String name, + Exception localError) throws Exception { + // Fallback: query through internal admin client to avoid local redirect/null-uri edge cases. + PulsarAdmin functionAdmin = worker().getFunctionAdmin(); + if (functionAdmin == null || functionAdmin.functions() == null) { + throw localError; + } + try { + return functionAdmin.functions().getFunctionStatus(tenant, namespace, name); + } catch (PulsarAdminException remoteError) { + remoteError.addSuppressed(localError); + throw remoteError; + } catch (RuntimeException remoteRuntimeError) { + if (remoteRuntimeError != localError) { + localError.addSuppressed(remoteRuntimeError); + } + throw localError; + } + } + + private static FunctionStatusSummary.SummaryState deriveState(int numInstances, int numRunning) { + if (numInstances <= 0) { + return FunctionStatusSummary.SummaryState.UNKNOWN; + } + if (numRunning == numInstances) { + return FunctionStatusSummary.SummaryState.RUNNING; + } + if (numRunning == 0) { + return FunctionStatusSummary.SummaryState.STOPPED; + } + return FunctionStatusSummary.SummaryState.PARTIAL; + } + + private static FunctionStatusSummary.ErrorType classifyError(Throwable error) { + if (isAuthenticationError(error)) { + return FunctionStatusSummary.ErrorType.AUTHENTICATION_FAILED; + } + if (isFunctionNotFoundError(error)) { + return FunctionStatusSummary.ErrorType.FUNCTION_NOT_FOUND; + } + if (isNetworkError(error)) { + return FunctionStatusSummary.ErrorType.NETWORK_ERROR; + } + return FunctionStatusSummary.ErrorType.INTERNAL_ERROR; + } + + private static boolean isAuthenticationError(Throwable error) { + if (error instanceof RestException) { + int status = getStatusCode((RestException) error); + return status == 401 || status == 403; + } + if (error instanceof PulsarAdminException) { + int status = ((PulsarAdminException) error).getStatusCode(); + return status == 401 || status == 403; + } + return false; + } + + private static boolean isFunctionNotFoundError(Throwable error) { + if (error instanceof RestException) { + return getStatusCode((RestException) error) == 404; + } + if (error instanceof PulsarAdminException) { + return ((PulsarAdminException) error).getStatusCode() == 404; + } + return false; + } + + private static boolean isNetworkError(Throwable error) { + if (error instanceof PulsarAdminException.ConnectException + || error instanceof PulsarAdminException.TimeoutException) { + return true; + } + Throwable current = error; + while (current != null) { + if (current instanceof ConnectException + || current instanceof SocketTimeoutException + || current instanceof UnknownHostException + || current instanceof UnresolvedAddressException) { + return true; + } + current = current.getCause(); + } + return false; + } + + private static int getStatusCode(RestException error) { + return error.getResponse() != null ? error.getResponse().getStatus() : -1; + } } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java index 7bdc86d5fae3e..ab7c779f87fbc 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionsApiV3Resource.java @@ -45,6 +45,8 @@ import org.apache.pulsar.common.policies.data.FunctionInstanceStatsDataImpl; import org.apache.pulsar.common.policies.data.FunctionStatsImpl; import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; import org.apache.pulsar.functions.worker.WorkerService; import org.apache.pulsar.functions.worker.rest.FunctionApiResource; import org.apache.pulsar.functions.worker.service.api.Functions; @@ -118,6 +120,26 @@ public List listFunctions(final @PathParam("tenant") String tenant, return functions().listFunctions(tenant, namespace, authParams()); } + @GET + @ApiOperation( + value = "Displays a batch status summary for all Pulsar Functions in a namespace", + response = FunctionStatusSummary.class, + responseContainer = "List" + ) + @ApiResponses(value = { + @ApiResponse(code = 400, message = "Invalid request"), + @ApiResponse(code = 403, message = "The requester doesn't have admin permissions") + }) + @Produces(MediaType.APPLICATION_JSON) + @Path("/{tenant}/{namespace}/status/summary") + public FunctionStatusPage listFunctionsWithStatus( + final @PathParam("tenant") String tenant, + final @PathParam("namespace") String namespace, + final @QueryParam("limit") Integer limit, + final @QueryParam("startAfter") String startAfter) { + return functions().listFunctionsWithStatus(tenant, namespace, limit, startAfter, authParams()); + } + @GET @ApiOperation( value = "Displays the status of a Pulsar Function instance", diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java index 28bc73fd0a831..f167d96e35728 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/service/api/Functions.java @@ -28,6 +28,7 @@ import org.apache.pulsar.common.functions.UpdateOptionsImpl; import org.apache.pulsar.common.policies.data.FunctionStatus; import org.apache.pulsar.common.policies.data.FunctionStatus.FunctionInstanceStatus.FunctionInstanceStatusData; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; import org.apache.pulsar.functions.worker.WorkerService; import org.glassfish.jersey.media.multipart.FormDataContentDisposition; @@ -102,4 +103,15 @@ FunctionInstanceStatusData getFunctionInstanceStatus(String tenant, void reloadBuiltinFunctions(AuthenticationParameters authParams) throws IOException; List getBuiltinFunctions(AuthenticationParameters authParams); + + FunctionStatusPage listFunctionsWithStatus(String tenant, String namespace, + AuthenticationParameters authParams); + + default FunctionStatusPage listFunctionsWithStatus(String tenant, + String namespace, + Integer limit, + String startAfter, + AuthenticationParameters authParams) { + return listFunctionsWithStatus(tenant, namespace, authParams); + } } diff --git a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java index b4e3862d1bbbb..392b5a3b9feda 100644 --- a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java +++ b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/FunctionsImplTest.java @@ -21,15 +21,20 @@ import static org.mockito.Mockito.any; import static org.mockito.Mockito.anyInt; import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.eq; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertThrows; import static org.testng.Assert.assertTrue; import java.io.InputStream; +import java.net.ConnectException; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; @@ -39,6 +44,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.CompletableFuture; +import javax.ws.rs.core.Response; import org.apache.distributedlog.api.namespace.Namespace; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.broker.authentication.AuthenticationParameters; @@ -46,6 +52,7 @@ import org.apache.pulsar.broker.resources.NamespaceResources; import org.apache.pulsar.broker.resources.PulsarResources; import org.apache.pulsar.broker.resources.TenantResources; +import org.apache.pulsar.client.admin.Functions; import org.apache.pulsar.client.admin.Namespaces; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -57,6 +64,9 @@ import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.FunctionInstanceStatsImpl; import org.apache.pulsar.common.policies.data.FunctionStatsImpl; +import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.util.RestException; @@ -147,6 +157,7 @@ public void setup() throws Exception { when(mockedWorkerService.getDlogNamespace()).thenReturn(mockedNamespace); when(mockedWorkerService.isInitialized()).thenReturn(true); when(mockedWorkerService.getBrokerAdmin()).thenReturn(mockedPulsarAdmin); + when(mockedWorkerService.getFunctionAdmin()).thenReturn(mockedPulsarAdmin); when(mockedPulsarAdmin.tenants()).thenReturn(mockedTenants); when(mockedPulsarAdmin.namespaces()).thenReturn(mockedNamespaces); when(mockedTenants.getTenantInfo(any())).thenReturn(mockedTenantInfo); @@ -369,6 +380,272 @@ public void testIsSuperUser() throws PulsarAdminException { assertFalse(functionImpl.isSuperUser("non-superuser", nonSuperuserAuthData)); } + @Test + public void testListFunctionsWithStatus_allRunning() { + List functionNames = List.of("func-a", "func-b"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + FunctionStatus statusA = new FunctionStatus(); + statusA.setNumInstances(2); + statusA.numRunning = 2; + doReturn(statusA).when(resource).getFunctionStatus(eq(tenant), eq(namespace), eq("func-a"), any(), any()); + + FunctionStatus statusB = new FunctionStatus(); + statusB.setNumInstances(3); + statusB.numRunning = 3; + doReturn(statusB).when(resource).getFunctionStatus(eq(tenant), eq(namespace), eq("func-b"), any(), any()); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 2); + assertEquals(result.getSummaries().get(0).getName(), "func-a"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.RUNNING); + assertEquals(result.getSummaries().get(0).getNumRunning(), 2); + assertEquals(result.getSummaries().get(0).getNumInstances(), 2); + assertEquals(result.getSummaries().get(1).getName(), "func-b"); + assertEquals(result.getSummaries().get(1).getState(), FunctionStatusSummary.SummaryState.RUNNING); + } + + @Test + public void testListFunctionsWithStatus_mixedStates() { + List functionNames = List.of("running-fn", "stopped-fn", "partial-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + FunctionStatus runningStatus = new FunctionStatus(); + runningStatus.setNumInstances(2); + runningStatus.numRunning = 2; + doReturn(runningStatus).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("running-fn"), any(), any()); + + FunctionStatus stoppedStatus = new FunctionStatus(); + stoppedStatus.setNumInstances(3); + stoppedStatus.numRunning = 0; + doReturn(stoppedStatus).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("stopped-fn"), any(), any()); + + FunctionStatus partialStatus = new FunctionStatus(); + partialStatus.setNumInstances(4); + partialStatus.numRunning = 2; + doReturn(partialStatus).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("partial-fn"), any(), any()); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 3); + assertEquals(result.getSummaries().get(0).getName(), "partial-fn"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.PARTIAL); + assertEquals(result.getSummaries().get(1).getName(), "running-fn"); + assertEquals(result.getSummaries().get(1).getState(), FunctionStatusSummary.SummaryState.RUNNING); + assertEquals(result.getSummaries().get(2).getName(), "stopped-fn"); + assertEquals(result.getSummaries().get(2).getState(), FunctionStatusSummary.SummaryState.STOPPED); + } + + @Test + public void testListFunctionsWithStatus_partialFailureIsolation() { + List functionNames = List.of("good-fn", "bad-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + FunctionStatus goodStatus = new FunctionStatus(); + goodStatus.setNumInstances(1); + goodStatus.numRunning = 1; + doReturn(goodStatus).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("good-fn"), any(), any()); + + doThrow(new RuntimeException("connection refused")).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("bad-fn"), any(), any()); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 2); + assertEquals(result.getSummaries().get(0).getName(), "bad-fn"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.UNKNOWN); + assertEquals(result.getSummaries().get(0).getError(), "connection refused"); + + assertEquals(result.getSummaries().get(1).getName(), "good-fn"); + assertEquals(result.getSummaries().get(1).getState(), FunctionStatusSummary.SummaryState.RUNNING); + assertEquals(result.getSummaries().get(1).getError(), null); + } + + @Test + public void testListFunctionsWithStatus_fallbackToFunctionAdmin() throws Exception { + List functionNames = List.of("remote-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + doThrow(new RuntimeException("local path failed")).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("remote-fn"), any(), any()); + + FunctionStatus remoteStatus = new FunctionStatus(); + remoteStatus.setNumInstances(2); + remoteStatus.numRunning = 1; + Functions mockedFunctionsAdmin = + mock(Functions.class); + when(mockedPulsarAdmin.functions()).thenReturn(mockedFunctionsAdmin); + when(mockedFunctionsAdmin.getFunctionStatus(eq(tenant), eq(namespace), eq("remote-fn"))) + .thenReturn(remoteStatus); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 1); + assertEquals(result.getSummaries().get(0).getName(), "remote-fn"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.PARTIAL); + assertEquals(result.getSummaries().get(0).getNumRunning(), 1); + assertEquals(result.getSummaries().get(0).getNumInstances(), 2); + assertEquals(result.getSummaries().get(0).getError(), null); + } + + @Test + public void testListFunctionsWithStatus_noFallbackForRestException() throws Exception { + List functionNames = List.of("auth-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + doThrow(new RestException(Response.Status.UNAUTHORIZED, "not authorized")).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("auth-fn"), any(), any()); + + Functions mockedFunctionsAdmin = + mock(Functions.class); + when(mockedPulsarAdmin.functions()).thenReturn(mockedFunctionsAdmin); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 1); + assertEquals(result.getSummaries().get(0).getName(), "auth-fn"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.UNKNOWN); + assertEquals(result.getSummaries().get(0).getError(), "not authorized"); + assertEquals( + result.getSummaries().get(0).getErrorType(), + FunctionStatusSummary.ErrorType.AUTHENTICATION_FAILED); + verify(mockedFunctionsAdmin, never()).getFunctionStatus(any(), any(), any()); + } + + @Test + public void testListFunctionsWithStatus_fallbackForServerSideRestException() throws Exception { + List functionNames = List.of("remote-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + doThrow(new RestException(Response.Status.INTERNAL_SERVER_ERROR, "local status failed")).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("remote-fn"), any(), any()); + + FunctionStatus remoteStatus = new FunctionStatus(); + remoteStatus.setNumInstances(2); + remoteStatus.numRunning = 1; + Functions mockedFunctionsAdmin = + mock(Functions.class); + when(mockedPulsarAdmin.functions()).thenReturn(mockedFunctionsAdmin); + when(mockedFunctionsAdmin.getFunctionStatus(eq(tenant), eq(namespace), eq("remote-fn"))) + .thenReturn(remoteStatus); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 1); + assertEquals(result.getSummaries().get(0).getName(), "remote-fn"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.PARTIAL); + assertEquals(result.getSummaries().get(0).getNumRunning(), 1); + assertEquals(result.getSummaries().get(0).getNumInstances(), 2); + assertEquals(result.getSummaries().get(0).getError(), null); + verify(mockedFunctionsAdmin).getFunctionStatus(eq(tenant), eq(namespace), eq("remote-fn")); + } + + @Test + public void testListFunctionsWithStatus_emptyNamespace() { + doReturn(Collections.emptyList()).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 0); + } + + @Test + public void testListFunctionsWithStatus_notFoundErrorType() throws Exception { + List functionNames = List.of("missing-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + doThrow(new RestException(Response.Status.NOT_FOUND, "function not found")).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("missing-fn"), any(), any()); + + Functions mockedFunctionsAdmin = + mock(Functions.class); + when(mockedPulsarAdmin.functions()).thenReturn(mockedFunctionsAdmin); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 1); + assertEquals(result.getSummaries().get(0).getErrorType(), FunctionStatusSummary.ErrorType.FUNCTION_NOT_FOUND); + verify(mockedFunctionsAdmin, never()).getFunctionStatus(any(), any(), any()); + } + + @Test + public void testListFunctionsWithStatus_networkErrorType() throws Exception { + List functionNames = List.of("network-fn"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + doThrow(new RuntimeException(new ConnectException("refused"))).when(resource) + .getFunctionStatus(eq(tenant), eq(namespace), eq("network-fn"), any(), any()); + + Functions mockedFunctionsAdmin = + mock(Functions.class); + when(mockedPulsarAdmin.functions()).thenReturn(mockedFunctionsAdmin); + when(mockedFunctionsAdmin.getFunctionStatus(eq(tenant), eq(namespace), eq("network-fn"))) + .thenThrow(new RuntimeException(new ConnectException("still refused"))); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 1); + assertEquals(result.getSummaries().get(0).getErrorType(), FunctionStatusSummary.ErrorType.NETWORK_ERROR); + } + + @Test + public void testListFunctionsWithStatus_pagination() { + List functionNames = List.of("fn-c", "fn-a", "fn-b"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + FunctionStatus statusB = new FunctionStatus(); + statusB.setNumInstances(1); + statusB.setNumRunning(1); + doReturn(statusB).when(resource).getFunctionStatus(eq(tenant), eq(namespace), eq("fn-b"), any(), any()); + + FunctionStatusPage result = + resource.listFunctionsWithStatus(tenant, namespace, 1, "fn-a", null); + + assertEquals(result.getSummaries().size(), 1); + assertEquals(result.getSummaries().get(0).getName(), "fn-b"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.RUNNING); + verify(resource, never()).getFunctionStatus(eq(tenant), eq(namespace), eq("fn-a"), any(), any()); + verify(resource, never()).getFunctionStatus(eq(tenant), eq(namespace), eq("fn-c"), any(), any()); + } + + @Test + public void testListFunctionsWithStatus_invalidLimit() { + try { + resource.listFunctionsWithStatus(tenant, namespace, 0, null, null); + org.testng.Assert.fail("Expected RestException"); + } catch (RestException e) { + assertEquals(e.getResponse().getStatus(), Response.Status.BAD_REQUEST.getStatusCode()); + } + } + + @Test + public void testListFunctionsWithStatus_usesConfiguredParallelism() { + mockedWorkerService.getWorkerConfig().setFunctionsStatusSummaryMaxParallelism(1); + List functionNames = List.of("fn-b", "fn-a"); + doReturn(functionNames).when(resource).listFunctions(eq(tenant), eq(namespace), any()); + + FunctionStatus statusA = new FunctionStatus(); + statusA.setNumInstances(1); + statusA.setNumRunning(1); + doReturn(statusA).when(resource).getFunctionStatus(eq(tenant), eq(namespace), eq("fn-a"), any(), any()); + + FunctionStatus statusB = new FunctionStatus(); + statusB.setNumInstances(1); + statusB.setNumRunning(0); + doReturn(statusB).when(resource).getFunctionStatus(eq(tenant), eq(namespace), eq("fn-b"), any(), any()); + + FunctionStatusPage result = resource.listFunctionsWithStatus(tenant, namespace, null); + + assertEquals(result.getSummaries().size(), 2); + assertEquals(result.getSummaries().get(0).getName(), "fn-a"); + assertEquals(result.getSummaries().get(0).getState(), FunctionStatusSummary.SummaryState.RUNNING); + assertEquals(result.getSummaries().get(1).getName(), "fn-b"); + assertEquals(result.getSummaries().get(1).getState(), FunctionStatusSummary.SummaryState.STOPPED); + } + public static FunctionConfig createDefaultFunctionConfig() { FunctionConfig functionConfig = new FunctionConfig(); functionConfig.setTenant(tenant); @@ -388,3 +665,4 @@ public static Function.FunctionDetails createDefaultFunctionDetails() { return FunctionConfigUtils.convert(functionConfig); } } + diff --git a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionApiV3ResourceTest.java b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionApiV3ResourceTest.java index 35af162507278..a417c66b8a146 100644 --- a/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionApiV3ResourceTest.java +++ b/pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker/rest/api/v3/FunctionApiV3ResourceTest.java @@ -22,6 +22,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.when; import static org.testng.Assert.assertEquals; @@ -41,6 +42,9 @@ import org.apache.pulsar.broker.authentication.AuthenticationParameters; import org.apache.pulsar.common.functions.FunctionConfig; import org.apache.pulsar.common.functions.UpdateOptionsImpl; +import org.apache.pulsar.common.policies.data.FunctionStatus; +import org.apache.pulsar.common.policies.data.FunctionStatusPage; +import org.apache.pulsar.common.policies.data.FunctionStatusSummary; import org.apache.pulsar.common.util.RestException; import org.apache.pulsar.functions.proto.Function; import org.apache.pulsar.functions.utils.FunctionConfigUtils; @@ -439,4 +443,49 @@ public void testRegisterFunctionSuccessK8sWithUpload() throws Exception { } } + @Test + public void testListFunctionsWithStatusSuccess() { + mockInstanceUtils(); + List metaDataList = List.of( + Function.FunctionMetaData.newBuilder().setFunctionDetails( + Function.FunctionDetails.newBuilder().setName("fn-a").build()).build(), + Function.FunctionMetaData.newBuilder().setFunctionDetails( + Function.FunctionDetails.newBuilder().setName("fn-b").build()).build() + ); + when(mockedManager.listFunctions(eq(TENANT), eq(NAMESPACE))).thenReturn(metaDataList); + + FunctionStatusSummary.SummaryState running = FunctionStatusSummary.SummaryState.RUNNING; + FunctionStatusSummary.SummaryState stopped = FunctionStatusSummary.SummaryState.STOPPED; + + FunctionStatus statusA = new FunctionStatus(); + statusA.setNumInstances(2); + statusA.setNumRunning(2); + doReturn(statusA).when(resource) + .getFunctionStatus(eq(TENANT), eq(NAMESPACE), eq("fn-a"), any(), any()); + + FunctionStatus statusB = new FunctionStatus(); + statusB.setNumInstances(1); + statusB.setNumRunning(0); + doReturn(statusB).when(resource) + .getFunctionStatus(eq(TENANT), eq(NAMESPACE), eq("fn-b"), any(), any()); + + FunctionStatusPage result = resource.listFunctionsWithStatus(TENANT, NAMESPACE, null); + + assertEquals(result.getSummaries().size(), 2); + assertEquals(result.getSummaries().get(0).getName(), "fn-a"); + assertEquals(result.getSummaries().get(0).getState(), running); + assertEquals(result.getSummaries().get(1).getName(), "fn-b"); + assertEquals(result.getSummaries().get(1).getState(), stopped); + } + + @Test + public void testListFunctionsWithStatusMissingNamespace() { + try { + resource.listFunctionsWithStatus(TENANT, null, null); + Assert.fail("Expected RestException"); + } catch (RestException e) { + assertEquals(e.getResponse().getStatus(), Response.Status.BAD_REQUEST.getStatusCode()); + } + } + }