Implement Tasklauncher function
Clean up test
This commit is contained in:
committed by
Soby Chacko
parent
e08ea45ebe
commit
8bdd18a34a
39
function/tasklauncher-function/README.adoc
Normal file
39
function/tasklauncher-function/README.adoc
Normal file
@@ -0,0 +1,39 @@
|
||||
# TaskLauncher Function
|
||||
|
||||
This module provides a `Function` that uses the Data Flow REST client to launch a registered task on a configured https://docs.spring.io/spring-cloud-dataflow/docs/current/reference/htmlsingle/#configuration-local-tasks[task platform].
|
||||
The client must be configured to connect to a remote Data Flow Server, including any required authentication (see Configuration Options below).
|
||||
|
||||
## Beans for injection
|
||||
|
||||
You can import the `TaskLauncherFunctionConfiguration` configuration in a Spring Boot application and then inject the following bean.
|
||||
|
||||
`taskLauncherFunction`
|
||||
|
||||
You may inject this as `TaskLauncherFunction` which implements `Function<LaunchRequest,<Optional<Long>>`.
|
||||
|
||||
You can use `taskLauncherFunction` as a qualifier when injecting.
|
||||
|
||||
Once injected, you can use the `apply` method of the `Function` to launch a task.
|
||||
The function takes a link:src/main/java/org/springframework/cloud/fn/tasklauncher/LaunchRequest.java[LaunchRequest] as input.
|
||||
This is a simple value object that specifies, at a minimum, the name of the task registered in Data Flow.
|
||||
Optionally, you can pass command line arguments as a `List<String>` and deployment properties as a `Map<String, String>`.
|
||||
The return value is an `Optional` containing the unique task ID of the launched instance if the launch request is successful.
|
||||
|
||||
NOTE: This version of the tasklauncher requires Spring Cloud Data Flow version 2.4.x or higher
|
||||
|
||||
## Configuration Options
|
||||
|
||||
Specific properties, including the task platform name are prefixed with `tasklauncher`.
|
||||
|
||||
For more information on the various options available, please see link:src/main/java/org/springframework/cloud/fn/tasklauncher/TaskLauncherFunctionProperties.java[TaskLauncherFunctionProperties.java]
|
||||
|
||||
Data Flow client configuration properties are prefixed with `spring.cloud.dataflow.client`.
|
||||
Please see https://github.com/spring-cloud/spring-cloud-dataflow/blob/master/spring-cloud-dataflow-rest-client/src/main/java/org/springframework/cloud/dataflow/rest/client/config/DataFlowClientProperties.java[DataFlowClientProperties.java] for more details.
|
||||
|
||||
## Examples
|
||||
|
||||
See this link:src/test/java/org/springframework/cloud/fn/tasklauncher/TaskLauncherFunctionApplicationTests.java[test suite] for examples of how this function is used.
|
||||
|
||||
## Other usage
|
||||
|
||||
See this link:../../../applications/sink/tasklauncher-sink/README.adoc[README] where this function is used to create a Spring Cloud Stream application to submit task launch requests.
|
||||
55
function/tasklauncher-function/pom.xml
Normal file
55
function/tasklauncher-function/pom.xml
Normal file
@@ -0,0 +1,55 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns="http://maven.apache.org/POM/4.0.0"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
<artifactId>tasklauncher-function</artifactId>
|
||||
<version>1.0.0-SNAPSHOT</version>
|
||||
<name>tasklauncher-function</name>
|
||||
<description>Spring Native Function for applying filter SpEL expressions</description>
|
||||
|
||||
<parent>
|
||||
<groupId>org.springframework.cloud.fn</groupId>
|
||||
<artifactId>spring-functions-parent</artifactId>
|
||||
<version>1.0.0-SNAPSHOT</version>
|
||||
<relativePath>../../spring-functions-parent</relativePath>
|
||||
</parent>
|
||||
|
||||
<properties>
|
||||
<spring.cloud.dataflow.version>2.5.1.RELEASE</spring.cloud.dataflow.version>
|
||||
</properties>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-configuration-processor</artifactId>
|
||||
<scope>provided</scope>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-dataflow-rest-client</artifactId>
|
||||
<version>${spring.cloud.dataflow.version}</version>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-scheduler-spi</artifactId>
|
||||
</exclusion>
|
||||
<exclusion>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<artifactId>spring-cloud-skipper</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-test</artifactId>
|
||||
<scope>test</scope>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>org.junit.vintage</groupId>
|
||||
<artifactId>junit-vintage-engine</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
||||
</project>
|
||||
@@ -0,0 +1,66 @@
|
||||
/*
|
||||
* Copyright 2019-2020 the original author or authors.
|
||||
*
|
||||
* Licensed 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
|
||||
*
|
||||
* https://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.springframework.cloud.fn.tasklauncher;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* @author David Turanski
|
||||
**/
|
||||
|
||||
public class LaunchRequest {
|
||||
@JsonProperty("args")
|
||||
private List<String> commandlineArguments = new ArrayList<>();
|
||||
@JsonProperty("deploymentProps")
|
||||
private Map<String, String> deploymentProperties = new HashMap<>();
|
||||
@JsonProperty("name")
|
||||
private String taskName;
|
||||
|
||||
public List<String> getCommandlineArguments() {
|
||||
return commandlineArguments;
|
||||
}
|
||||
|
||||
public void setCommandlineArguments(List<String> commandlineArguments) {
|
||||
Assert.notNull(commandlineArguments, "'commandLineArguments' cannot be null.");
|
||||
this.commandlineArguments = commandlineArguments;
|
||||
}
|
||||
|
||||
public Map<String, String> getDeploymentProperties() {
|
||||
return deploymentProperties;
|
||||
}
|
||||
|
||||
public void setDeploymentProperties(Map<String, String> deploymentProperties) {
|
||||
Assert.notNull(commandlineArguments, "'deploymentProperties' cannot be null.");
|
||||
this.deploymentProperties = deploymentProperties;
|
||||
}
|
||||
|
||||
public String getTaskName() {
|
||||
return taskName;
|
||||
}
|
||||
|
||||
public void setTaskName(String taskName) {
|
||||
Assert.hasText(taskName, "'taskName' cannot be blank.");
|
||||
this.taskName = taskName;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,167 @@
|
||||
/*
|
||||
* Copyright 2019-2020 the original author or authors.
|
||||
*
|
||||
* Licensed 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
|
||||
*
|
||||
* https://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.springframework.cloud.fn.tasklauncher;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.cloud.dataflow.rest.client.TaskOperations;
|
||||
import org.springframework.cloud.dataflow.rest.resource.CurrentTaskExecutionsResource;
|
||||
import org.springframework.cloud.dataflow.rest.resource.LauncherResource;
|
||||
import org.springframework.hateoas.PagedModel;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
*
|
||||
* A {@link Function} that submits a task {@link LaunchRequest} to a Data Flow server.
|
||||
* This will check if the Data Flow task platform is at capacity. If not, will submit the
|
||||
* task launch request, otherwise it will return and log a warning message.
|
||||
*
|
||||
* @author David Turanski
|
||||
**/
|
||||
public class TaskLauncherFunction implements Function<LaunchRequest, Optional<Long>>, InitializingBean {
|
||||
private static final Log log = LogFactory.getLog(TaskLauncherFunction.class);
|
||||
|
||||
static final String TASK_PLATFORM_NAME = "spring.cloud.dataflow.task.platformName";
|
||||
|
||||
private final TaskOperations taskOperations;
|
||||
|
||||
private String platformName = "default";
|
||||
|
||||
public TaskLauncherFunction(TaskOperations taskOperations) {
|
||||
Assert.notNull(taskOperations, "`taskOperations` cannot be null.");
|
||||
this.taskOperations = taskOperations;
|
||||
}
|
||||
|
||||
/**
|
||||
*
|
||||
* @param launchRequest the task launch request for the Data Flow server.
|
||||
* @return an {@code Optional<Long>} containing the task Id if the request is accepted or
|
||||
* empty otherwise.
|
||||
*/
|
||||
@Override
|
||||
public Optional<Long> apply(LaunchRequest launchRequest) {
|
||||
if (platformIsAcceptingNewTasks()) {
|
||||
return Optional.of(launchTask(launchRequest));
|
||||
}
|
||||
log.warn(String.format("Platform is at capacity. Did not submit task launch request for task %s.",
|
||||
launchRequest.getTaskName()));
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
public boolean platformIsAcceptingNewTasks() {
|
||||
|
||||
boolean availableForNewTasks;
|
||||
int maximumTaskExecutions = 0;
|
||||
int runningExecutionCount = 0;
|
||||
|
||||
List<String> currentPlatforms = new ArrayList<>();
|
||||
|
||||
boolean validPlatform = false;
|
||||
for (CurrentTaskExecutionsResource currentTaskExecutionsResource : taskOperations.currentTaskExecutions()) {
|
||||
if (currentTaskExecutionsResource.getName().equals(platformName)) {
|
||||
maximumTaskExecutions = currentTaskExecutionsResource.getMaximumTaskExecutions();
|
||||
runningExecutionCount = currentTaskExecutionsResource.getRunningExecutionCount();
|
||||
validPlatform = true;
|
||||
}
|
||||
currentPlatforms.add(currentTaskExecutionsResource.getName());
|
||||
}
|
||||
|
||||
// Verify for each request as configuration may have changed on server.
|
||||
assertValidPlatform(validPlatform, currentPlatforms);
|
||||
|
||||
availableForNewTasks = runningExecutionCount < maximumTaskExecutions;
|
||||
if (!availableForNewTasks) {
|
||||
log.warn(String.format(
|
||||
"The data Flow task platform %s has reached its concurrent task execution limit: (%d)",
|
||||
platformName,
|
||||
maximumTaskExecutions));
|
||||
}
|
||||
|
||||
return availableForNewTasks;
|
||||
|
||||
}
|
||||
|
||||
private long launchTask(LaunchRequest request) {
|
||||
String requestPlatformName = request.getDeploymentProperties().get(TASK_PLATFORM_NAME);
|
||||
if (StringUtils.hasText(requestPlatformName) && !platformName.equals(requestPlatformName)) {
|
||||
throw new IllegalStateException(
|
||||
String.format(
|
||||
"Task Launch request for Task %s contains deployment property '%s=%s' which does not " +
|
||||
"match the platform configured for the Task Launcher: '%s'",
|
||||
request.getTaskName(),
|
||||
TASK_PLATFORM_NAME,
|
||||
request.getDeploymentProperties().get(TASK_PLATFORM_NAME),
|
||||
platformName));
|
||||
}
|
||||
log.info(String.format("Launching Task %s on platform %s", request.getTaskName(), platformName));
|
||||
long taskId = taskOperations.launch(request.getTaskName(),
|
||||
enrichDeploymentProperties(request.getDeploymentProperties()),
|
||||
request.getCommandlineArguments(), null);
|
||||
log.info(String.format("Launched Task %s - task ID is %d", request.getTaskName(), taskId));
|
||||
return taskId;
|
||||
}
|
||||
|
||||
private Map<String, String> enrichDeploymentProperties(Map<String, String> deploymentProperties) {
|
||||
if (!deploymentProperties.containsKey(TASK_PLATFORM_NAME)) {
|
||||
Map<String, String> enrichedProperties = new HashMap<>();
|
||||
enrichedProperties.putAll(deploymentProperties);
|
||||
enrichedProperties.put(TASK_PLATFORM_NAME, platformName);
|
||||
return enrichedProperties;
|
||||
}
|
||||
return deploymentProperties;
|
||||
}
|
||||
|
||||
public void setPlatformName(String platformName) {
|
||||
this.platformName = platformName;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterPropertiesSet() {
|
||||
PagedModel<LauncherResource> launchers = taskOperations.listPlatforms();
|
||||
|
||||
boolean validPlatform = false;
|
||||
List<String> currentPlatforms = new ArrayList<>();
|
||||
|
||||
for (LauncherResource launcherResource : launchers) {
|
||||
currentPlatforms.add(launcherResource.getName());
|
||||
if (launcherResource.getName().equals(platformName)) {
|
||||
validPlatform = true;
|
||||
}
|
||||
}
|
||||
|
||||
assertValidPlatform(validPlatform, currentPlatforms);
|
||||
}
|
||||
|
||||
private void assertValidPlatform(boolean validPlatform, List<String> currentPlatforms) {
|
||||
Assert.notEmpty(currentPlatforms, "The Data Flow Server has no task platforms configured");
|
||||
|
||||
Assert.isTrue(validPlatform, String.format(
|
||||
"The task launcher's platform name '%s' does not match one of the Data Flow server's configured task "
|
||||
+ "platforms: [%s].",
|
||||
platformName, StringUtils.collectionToCommaDelimitedString(currentPlatforms)));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
/*
|
||||
* Copyright 2019-2020 the original author or authors.
|
||||
*
|
||||
* Licensed 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
|
||||
*
|
||||
* https://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.springframework.cloud.fn.tasklauncher;
|
||||
|
||||
import org.springframework.boot.context.properties.EnableConfigurationProperties;
|
||||
import org.springframework.cloud.dataflow.rest.client.DataFlowOperations;
|
||||
import org.springframework.cloud.dataflow.rest.client.config.DataFlowClientProperties;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
/**
|
||||
* Configuration class for the TaskLauncher Data Flow Sink.
|
||||
*
|
||||
* @author David Turanski
|
||||
* @author Gunnar Hillert
|
||||
*/
|
||||
@Configuration
|
||||
@EnableConfigurationProperties({TaskLauncherFunctionProperties.class, DataFlowClientProperties.class})
|
||||
public class TaskLauncherFunctionConfiguration {
|
||||
|
||||
@Bean
|
||||
public TaskLauncherFunction taskLauncherFunction(
|
||||
DataFlowOperations dataFlowOperations, TaskLauncherFunctionProperties functionProperties) {
|
||||
|
||||
if (dataFlowOperations.taskOperations() == null) {
|
||||
throw new IllegalArgumentException("The SCDF server does not support task operations");
|
||||
}
|
||||
TaskLauncherFunction function = new TaskLauncherFunction(dataFlowOperations.taskOperations());
|
||||
function.setPlatformName(functionProperties.getPlatformName());
|
||||
return function;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
/*
|
||||
* Copyright 2019-2020 the original author or authors.
|
||||
*
|
||||
* Licensed 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
|
||||
*
|
||||
* https://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.springframework.cloud.fn.tasklauncher;
|
||||
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
|
||||
/**
|
||||
* @author David Turanski
|
||||
**/
|
||||
@ConfigurationProperties("tasklauncher")
|
||||
public class TaskLauncherFunctionProperties {
|
||||
/**
|
||||
* The Spring Cloud Data Flow platform to use for launching tasks.
|
||||
*/
|
||||
private String platformName = "default";
|
||||
|
||||
public String getPlatformName() {
|
||||
return platformName;
|
||||
}
|
||||
|
||||
public void setPlatformName(String platformName) {
|
||||
this.platformName = platformName;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,151 @@
|
||||
/*
|
||||
* Copyright 2019-2020 the original author or authors.
|
||||
*
|
||||
* Licensed 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
|
||||
*
|
||||
* https://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.springframework.cloud.fn.tasklauncher;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.Optional;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import org.springframework.beans.factory.BeanCreationException;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.boot.test.context.SpringBootTest;
|
||||
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
|
||||
import org.springframework.cloud.dataflow.rest.client.DataFlowOperations;
|
||||
import org.springframework.cloud.dataflow.rest.client.TaskOperations;
|
||||
import org.springframework.cloud.dataflow.rest.resource.CurrentTaskExecutionsResource;
|
||||
import org.springframework.cloud.dataflow.rest.resource.LauncherResource;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.context.annotation.Import;
|
||||
import org.springframework.context.annotation.Profile;
|
||||
import org.springframework.hateoas.PagedModel;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
|
||||
import static org.assertj.core.api.Assertions.assertThatIllegalStateException;
|
||||
import static org.mockito.ArgumentMatchers.anyList;
|
||||
import static org.mockito.ArgumentMatchers.anyMap;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.isNull;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.verify;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
@SpringBootTest
|
||||
public class TaskLauncherFunctionApplicationTests {
|
||||
|
||||
@Autowired
|
||||
private TaskLauncherFunction taskLauncherFunction;
|
||||
|
||||
@Autowired
|
||||
private TaskOperations taskOperations;
|
||||
|
||||
@Test
|
||||
public void successfulLaunch() {
|
||||
LaunchRequest launchRequest = new LaunchRequest();
|
||||
launchRequest.setTaskName("someTask");
|
||||
setCurrentExecutionState(0);
|
||||
Optional<Long> taskId = taskLauncherFunction.apply(launchRequest);
|
||||
assertThat(taskId.isPresent()).isTrue();
|
||||
assertThat(taskId.get()).isEqualTo(1L);
|
||||
|
||||
verify(taskOperations).launch("someTask",
|
||||
Collections.singletonMap(TaskLauncherFunction.TASK_PLATFORM_NAME, "default"),
|
||||
Collections.emptyList(),
|
||||
null);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void taskPlatformAtCapacity() {
|
||||
LaunchRequest launchRequest = new LaunchRequest();
|
||||
launchRequest.setTaskName("someTask");
|
||||
setCurrentExecutionState(3);
|
||||
Optional<Long> taskId = taskLauncherFunction.apply(launchRequest);
|
||||
assertThat(taskId.isPresent()).isFalse();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void platformMismatch() {
|
||||
LaunchRequest launchRequest = new LaunchRequest();
|
||||
launchRequest.setTaskName("someTask");
|
||||
launchRequest
|
||||
.setDeploymentProperties(Collections.singletonMap(TaskLauncherFunction.TASK_PLATFORM_NAME, "other"));
|
||||
setCurrentExecutionState(0);
|
||||
assertThatIllegalStateException().isThrownBy(() -> taskLauncherFunction.apply(launchRequest))
|
||||
.withStackTraceContaining("does not match the platform configured for the Task Launcher");
|
||||
}
|
||||
|
||||
private void setCurrentExecutionState(int runningExecutions) {
|
||||
CurrentTaskExecutionsResource currentTaskExecutionsResource = new CurrentTaskExecutionsResource();
|
||||
currentTaskExecutionsResource.setMaximumTaskExecutions(3);
|
||||
currentTaskExecutionsResource.setRunningExecutionCount(runningExecutions);
|
||||
currentTaskExecutionsResource.setName("default");
|
||||
when(taskOperations.currentTaskExecutions())
|
||||
.thenReturn(Collections.singletonList(currentTaskExecutionsResource));
|
||||
when(taskOperations.launch(anyString(), anyMap(), anyList(), isNull())).thenReturn(1L);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void noLaunchersConfigured() {
|
||||
ApplicationContextRunner contextRunner = new ApplicationContextRunner().withUserConfiguration(TestConfig.class);
|
||||
assertThatExceptionOfType(IllegalStateException.class).isThrownBy(() -> contextRunner
|
||||
.withPropertyValues("spring.profiles.active=nolaunchers")
|
||||
.run(context -> context.start()))
|
||||
.withCauseInstanceOf(BeanCreationException.class)
|
||||
.withRootCauseInstanceOf(IllegalArgumentException.class)
|
||||
.withStackTraceContaining("The Data Flow Server has no task platforms configured");
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@Import(TaskLauncherFunctionConfiguration.class)
|
||||
static class TestConfig {
|
||||
|
||||
@Bean
|
||||
@Profile("default")
|
||||
TaskOperations taskOperations() {
|
||||
TaskOperations taskOperations = mock(TaskOperations.class);
|
||||
LauncherResource launcherResource = mock(LauncherResource.class);
|
||||
when(launcherResource.getName()).thenReturn("default");
|
||||
|
||||
when(taskOperations.listPlatforms()).thenReturn(PagedModel.of(
|
||||
Collections.singletonList(launcherResource), (PagedModel.PageMetadata) null));
|
||||
return taskOperations;
|
||||
}
|
||||
|
||||
@Bean
|
||||
@Profile("nolaunchers")
|
||||
TaskOperations taskOperationsNoLaunchers() {
|
||||
TaskOperations taskOperations = mock(TaskOperations.class);
|
||||
when(taskOperations.listPlatforms()).thenReturn(PagedModel.of(
|
||||
Collections.emptyList(), (PagedModel.PageMetadata) null));
|
||||
return taskOperations;
|
||||
}
|
||||
|
||||
@Bean
|
||||
DataFlowOperations dataFlowOperations(TaskOperations taskOperations) {
|
||||
DataFlowOperations dataFlowOperations = mock(DataFlowOperations.class);
|
||||
when(dataFlowOperations.taskOperations()).thenReturn(taskOperations);
|
||||
return dataFlowOperations;
|
||||
}
|
||||
}
|
||||
|
||||
@SpringBootApplication
|
||||
static class TestApplication {
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user