SCT-12 Creates the SCSt sink that can launch tasks.

resolves  spring-cloud/spring-cloud-task#12
This commit is contained in:
Glenn Renfro
2016-03-02 15:12:18 -05:00
committed by Michael Minella
parent f35f8ef52d
commit aad0a0b1ee
33 changed files with 2154 additions and 2 deletions

View File

@@ -0,0 +1,167 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher;
import java.io.Serializable;
import java.util.HashMap;
import java.util.Map;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
/**
* Request that contains the maven repository and property information required by the
* TaskLauncherSink to launch the task.
*
* @author Glenn Renfro
*/
public class TaskLaunchRequest implements Serializable{
private static final long serialVersionUID = 1L;
private String artifact;
private String taskGroupId;
private String taskVersion;
private String taskExtension;
private String taskClassifier;
private Map<String, String> properties;
/**
* Constructor for the TaskLaunchRequest;
* @param artifact is maven artifact coordinate for the task. Must not be empty nor null.
* @param taskGroupId is maven groupId coordinate for the task. Must not be empty nor null.
* @param taskVersion is maven version coordinate for the task. Must not be empty nor null.
* @param taskExtension is maven extension coordinate for the task.
* @param taskClassifier is maven classifier coordinate for the task.
* @param properties is the environment variables for this task.
*/
public TaskLaunchRequest(String artifact, String taskGroupId, String taskVersion,
String taskExtension, String taskClassifier,
Map<String, String> properties) {
Assert.hasText(artifact, "artifact must not be empty nor null.");
Assert.hasText(taskGroupId, "taskGroupID must not be empty nor null.");
Assert.hasText(taskVersion, "taskVersion must not be empty nor null.");
Assert.hasText(taskExtension, "taskExtension must not be empty nor null.");
this.artifact = artifact;
this.taskGroupId = taskGroupId;
this.taskVersion = taskVersion;
this.taskExtension = taskExtension;
this.taskClassifier = taskClassifier;
this.properties = properties == null ? new HashMap() : properties;
}
/**
* Retrieves the group maven coordinate for the task.
* @return group maven coordinate for the task.
*/
public String getTaskGroupId() {
return taskGroupId;
}
/**
* Retrieves the version maven coordinate for the task.
* @return version maven coordinate for the task.
*/
public String getTaskVersion() {
return taskVersion;
}
/**
* Retrieves the extension maven coordinate for the task.
* @return extension maven coordinate for the task.
*/
public String getTaskExtension() {
return taskExtension;
}
/**
* Retrieves the classifier maven coordinate for the task.
* @return classifier maven coordinate for the task.
*/
public String getTaskClassifier() {
return taskClassifier;
}
/**
* Retrieves the artifact maven coordinate for the task.
* @return artifact maven coordinate for the task.
*/
public String getArtifact() {
return artifact;
}
/**
* Retrieves the environment variables for the task.
* @return map containing the environment variables for the task.
*/
public Map<String, String> getProperties() {
return properties;
}
@Override
public String toString() {
String coordinates = taskGroupId + ":" + artifact + ":" + taskVersion ;
if(StringUtils.hasText(taskClassifier)){
coordinates = coordinates + ":" + taskClassifier;
}
coordinates = coordinates + ":" + taskExtension;
return coordinates;
}
@Override
public boolean equals(Object o) {
if (this == o){
return true;
}
if (o == null || getClass() != o.getClass()){
return false;
}
TaskLaunchRequest that = (TaskLaunchRequest) o;
if (!artifact.equals(that.artifact)){
return false;
}
if (!taskGroupId.equals(that.taskGroupId)){
return false;
}
if (!taskVersion.equals(that.taskVersion)){
return false;
}
if (!taskExtension.equals(that.taskExtension)){
return false;
}
if (taskClassifier != null ? !taskClassifier.equals(that.taskClassifier) : that.taskClassifier != null){
return false;
}
return properties != null ? properties.equals(that.properties) : that.properties == null;
}
@Override
public int hashCode() {
int result = artifact != null ? artifact.hashCode() : 0;
result = 31 * result + (taskGroupId != null ? taskGroupId.hashCode() : 0);
result = 31 * result + (taskVersion != null ? taskVersion.hashCode() : 0);
result = 31 * result + (taskExtension != null ? taskExtension.hashCode() : 0);
result = 31 * result + (taskClassifier != null ? taskClassifier.hashCode() : 0);
result = 31 * result + (properties != null ? properties.hashCode() : 0);
return result;
}
}

View File

@@ -0,0 +1,47 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.cloud.deployer.spi.local.LocalDeployerProperties;
import org.springframework.cloud.deployer.spi.local.LocalTaskLauncher;
import org.springframework.cloud.deployer.spi.task.TaskLauncher;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* Creates the appropriate Task Launcher Configuration based on the TaskLauncher
* that is available in the classpath.
* @author Glenn Renfro
*/
@Configuration
@ConditionalOnClass({TaskLauncher.class})
public class TaskLauncherConfiguration {
@Configuration
@ConditionalOnMissingBean(name = "taskLauncher")
@ConditionalOnClass({LocalTaskLauncher.class})
protected static class LocalTaskDeployerConfiguration {
@Bean
public TaskLauncher taskLauncher() {
return new LocalTaskLauncher(new LocalDeployerProperties());
}
}
}

View File

@@ -0,0 +1,70 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.deployer.resource.maven.MavenResource;
import org.springframework.cloud.deployer.spi.core.AppDefinition;
import org.springframework.cloud.deployer.spi.core.AppDeploymentRequest;
import org.springframework.cloud.deployer.spi.task.TaskLauncher;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.util.Assert;
/**
* A sink stream application that launches a tasks.
*
* @author Glenn Renfro
*/
@EnableBinding(Sink.class)
public class TaskLauncherSink {
private final static Logger logger = LoggerFactory.getLogger(TaskLauncherSink.class);
@Autowired
public TaskLauncher taskLauncher;
/**
* Launches a task upon the receipt of a valid TaskLaunchRequest.
* @param request is a TaskLaunchRequest containing the information required to launch
* a task.
*/
@ServiceActivator(inputChannel = Sink.INPUT)
public void taskLauncherSink(TaskLaunchRequest request) {
launchTask(request);
}
private void launchTask(TaskLaunchRequest taskLaunchRequest) {
Assert.notNull(taskLauncher, "TaskLauncher has not been initialized");
logger.info("Launching Task for the following resource " + taskLaunchRequest);
MavenResource resource = new MavenResource.Builder()
.artifactId(taskLaunchRequest.getArtifact())
.groupId(taskLaunchRequest.getTaskGroupId())
.version(taskLaunchRequest.getTaskVersion())
.extension(taskLaunchRequest.getTaskExtension())
.classifier(taskLaunchRequest.getTaskClassifier())
.build();
AppDefinition definition = new AppDefinition(taskLaunchRequest.getArtifact(), taskLaunchRequest.getProperties());
AppDeploymentRequest request = new AppDeploymentRequest(definition, resource);
taskLauncher.launch(request);
}
}

View File

@@ -0,0 +1,61 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher.annotation;
import java.lang.annotation.Documented;
import java.lang.annotation.ElementType;
import java.lang.annotation.Inherited;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
import org.springframework.cloud.deployer.spi.task.TaskLauncher;
import org.springframework.cloud.task.launcher.TaskLaunchRequest;
import org.springframework.cloud.task.launcher.TaskLauncherConfiguration;
import org.springframework.cloud.task.launcher.TaskLauncherSink;
import org.springframework.context.annotation.Import;
/**
* <p>
* Enable this boot app to be a sink to receive a {@link TaskLaunchRequest} and use the
* {@link TaskLauncher} to launch the task.
* </p>
*
* <pre class="code">
* &#064;Configuration
* &#064;EnableTaskLauncher
* public class AppConfig {
*
* &#064;Bean
* public MyCommandLineRunner myCommandLineRunner() {
* return new MyCommandLineRunner()
* }
* }
* </pre>
*
* Note that only one of your configuration classes needs to have the <code>&#064;EnableTaskLauncher</code>
* annotation.
*
* @author Glenn Renfro
*/
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
@Documented
@Inherited
@Import({TaskLauncherConfiguration.class, TaskLauncherSink.class})
public @interface EnableTaskLauncher {
}

View File

@@ -0,0 +1,83 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher;
import static org.junit.Assert.assertEquals;
import java.util.HashMap;
import java.util.Map;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.SpringApplicationConfiguration;
import org.springframework.cloud.deployer.spi.task.LaunchState;
import org.springframework.cloud.stream.annotation.Bindings;
import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.cloud.task.launcher.configuration.TaskConfiguration;
import org.springframework.cloud.task.launcher.util.TaskLauncherSinkApplication;
import org.springframework.context.ApplicationContext;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@RunWith(SpringJUnit4ClassRunner.class)
@SpringApplicationConfiguration(classes = {TaskLauncherSinkApplication.class, TaskConfiguration.class} )
public class TaskLauncherSinkTests {
private final static String DEFAULT_STATUS = "test_status";
@Autowired
private ApplicationContext context;
@Autowired
@Bindings(TaskLauncherSink.class)
private Sink sink;
@Test
public void testSuccess() {
TaskConfiguration.TestTaskLauncher testTaskLauncher =
context.getBean(TaskConfiguration.TestTaskLauncher.class);
Map<String, String> properties = new HashMap<>();
properties.put("server.port", "0");
TaskLaunchRequest request = new TaskLaunchRequest("timestamp-task",
"org.springframework.cloud.task.module","1.0.0.BUILD-SNAPSHOT", "jar",
"exec", properties);
GenericMessage<TaskLaunchRequest> message = new GenericMessage<>(request);
this.sink.input().send(message);
assertEquals(LaunchState.complete, testTaskLauncher.status(DEFAULT_STATUS).getState());
}
@Test
public void testNoRun() {
TaskConfiguration.TestTaskLauncher testTaskLauncher =
context.getBean(TaskConfiguration.TestTaskLauncher.class);
assertEquals(LaunchState.unknown, testTaskLauncher.status(DEFAULT_STATUS).getState());
}
@Test(expected = IllegalArgumentException.class)
public void testNoTaskLauncher() {
Map<String, String> properties = new HashMap<>();
properties.put("server.port", "0");
TaskLauncherSink sink = new TaskLauncherSink();
sink.taskLauncherSink(new TaskLaunchRequest("timestamp-task",
"org.springframework.cloud.task.module","1.0.0.BUILD-SNAPSHOT", "jar",
"exec", properties));
}
}

View File

@@ -0,0 +1,60 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher.configuration;
import org.springframework.cloud.deployer.spi.core.AppDeploymentRequest;
import org.springframework.cloud.deployer.spi.task.LaunchState;
import org.springframework.cloud.deployer.spi.task.TaskLauncher;
import org.springframework.cloud.deployer.spi.task.TaskStatus;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* @author Glenn Renfro
*/
@Configuration
public class TaskConfiguration {
@Bean
public TaskLauncher taskLauncher(){
return new TestTaskLauncher();
}
public static class TestTaskLauncher implements TaskLauncher{
public static final String LAUNCH_ID = "TEST_LAUNCH_ID";
private LaunchState state = LaunchState.unknown;
@Override
public String launch(AppDeploymentRequest request) {
state = LaunchState.complete;
return null;
}
@Override
public void cancel(String id) {
}
@Override
public TaskStatus status(String id) {
return new TaskStatus(LAUNCH_ID, state, null);
}
}
}

View File

@@ -0,0 +1,33 @@
/*
* Copyright 2016 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
*
* 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.springframework.cloud.task.launcher.util;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.task.launcher.annotation.EnableTaskLauncher;
/**
* @author Glenn Renfro
*/
@SpringBootApplication
@EnableTaskLauncher
public class TaskLauncherSinkApplication {
public static void main(String[] args) {
SpringApplication.run(TaskLauncherSinkApplication.class, args);
}
}