Removed the 'org.springframework.integration.quartz' module
This commit is contained in:
@@ -1,19 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<classpath>
|
||||
<classpathentry kind="src" path="src/main/java"/>
|
||||
<classpathentry kind="src" output="target/test-classes" path="src/test/java"/>
|
||||
<classpathentry kind="src" path="src/main/resources"/>
|
||||
<classpathentry kind="con" path="org.eclipse.jdt.launching.JRE_CONTAINER"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.apache.commons/com.springsource.org.apache.commons.logging/1.1.1/com.springsource.org.apache.commons.logging-1.1.1.jar" sourcepath="/IVY_CACHE/org.apache.commons/com.springsource.org.apache.commons.logging/1.1.1/com.springsource.org.apache.commons.logging-sources-1.1.1.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.apache.commons/com.springsource.org.apache.commons.collections/3.2.0/com.springsource.org.apache.commons.collections-3.2.0.jar" sourcepath="/IVY_CACHE/org.apache.commons/com.springsource.org.apache.commons.collections/3.2.0/com.springsource.org.apache.commons.collections-sources-3.2.0.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/com.opensymphony.quartz/com.springsource.org.quartz/1.6.0/com.springsource.org.quartz-1.6.0.jar" sourcepath="/IVY_CACHE/com.opensymphony.quartz/com.springsource.org.quartz/1.6.0/com.springsource.org.quartz-sources-1.6.0.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/javax.transaction/com.springsource.javax.transaction/1.1.0/com.springsource.javax.transaction-1.1.0.jar" sourcepath="/IVY_CACHE/javax.transaction/com.springsource.javax.transaction/1.1.0/com.springsource.javax.transaction-sources-1.1.0.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.junit/com.springsource.org.junit/4.4.0/com.springsource.org.junit-4.4.0.jar" sourcepath="/IVY_CACHE/org.junit/com.springsource.org.junit/4.4.0/com.springsource.org.junit-sources-4.4.0.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.springframework/org.springframework.beans/2.5.5.A/org.springframework.beans-2.5.5.A.jar" sourcepath="/IVY_CACHE/org.springframework/org.springframework.beans/2.5.5.A/org.springframework.beans-sources-2.5.5.A.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.springframework/org.springframework.context/2.5.5.A/org.springframework.context-2.5.5.A.jar" sourcepath="/IVY_CACHE/org.springframework/org.springframework.context/2.5.5.A/org.springframework.context-sources-2.5.5.A.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.springframework/org.springframework.context.support/2.5.5.A/org.springframework.context.support-2.5.5.A.jar" sourcepath="/IVY_CACHE/org.springframework/org.springframework.context.support/2.5.5.A/org.springframework.context.support-sources-2.5.5.A.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.springframework/org.springframework.core/2.5.5.A/org.springframework.core-2.5.5.A.jar" sourcepath="/IVY_CACHE/org.springframework/org.springframework.core/2.5.5.A/org.springframework.core-sources-2.5.5.A.jar"/>
|
||||
<classpathentry kind="var" path="IVY_CACHE/org.springframework/org.springframework.test/2.5.5.A/org.springframework.test-2.5.5.A.jar" sourcepath="/IVY_CACHE/org.springframework/org.springframework.test/2.5.5.A/org.springframework.test-sources-2.5.5.A.jar"/>
|
||||
<classpathentry combineaccessrules="false" kind="src" path="/org.springframework.integration"/>
|
||||
<classpathentry kind="output" path="target/classes"/>
|
||||
</classpath>
|
||||
@@ -1,17 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<projectDescription>
|
||||
<name>org.springframework.integration.quartz</name>
|
||||
<comment></comment>
|
||||
<projects>
|
||||
</projects>
|
||||
<buildSpec>
|
||||
<buildCommand>
|
||||
<name>org.eclipse.jdt.core.javabuilder</name>
|
||||
<arguments>
|
||||
</arguments>
|
||||
</buildCommand>
|
||||
</buildSpec>
|
||||
<natures>
|
||||
<nature>org.eclipse.jdt.core.javanature</nature>
|
||||
</natures>
|
||||
</projectDescription>
|
||||
@@ -1,8 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project name="org.springframework.integration.quartz">
|
||||
|
||||
<property file="${basedir}/../build.properties"/>
|
||||
<import file="${basedir}/../build-spring-integration/package-bundle.xml"/>
|
||||
<import file="${basedir}/../spring-build/standard/default.xml"/>
|
||||
|
||||
</project>
|
||||
@@ -1,33 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<?xml-stylesheet type="text/xsl" href="http://ivyrep.jayasoft.org/ivy-doc.xsl"?>
|
||||
<ivy-module
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:noNamespaceSchemaLocation="http://incubator.apache.org/ivy/schemas/ivy.xsd"
|
||||
version="1.3">
|
||||
|
||||
<info organisation="org.springframework.integration" module="${ant.project.name}">
|
||||
<license name="Apache 2.0" url="http://www.apache.org/licenses/LICENSE-2.0"/>
|
||||
<ivyauthor name="Mark Fisher"/>
|
||||
</info>
|
||||
|
||||
<configurations>
|
||||
<include file="${spring.build.dir}/common/default-ivy-configurations.xml"/>
|
||||
</configurations>
|
||||
|
||||
<publications>
|
||||
<artifact name="${ant.project.name}"/>
|
||||
<artifact name="${ant.project.name}-sources" type="src" ext="jar"/>
|
||||
</publications>
|
||||
|
||||
<dependencies>
|
||||
<dependency org="org.springframework" name="org.springframework.context" rev="2.5.5.A" conf="compile->runtime"/>
|
||||
<dependency org="org.springframework" name="org.springframework.context.support" rev="2.5.5.A" conf="compile->runtime"/>
|
||||
<dependency org="org.springframework" name="org.springframework.test" rev="2.5.5.A" conf="compile->runtime"/>
|
||||
<dependency org="com.opensymphony.quartz" name="com.springsource.org.quartz" rev="1.6.0" conf="compile->runtime"/>
|
||||
<dependency org="javax.transaction" name="com.springsource.javax.transaction" rev="1.1.0" conf="test->runtime"/>
|
||||
<dependency org="org.apache.commons" name="com.springsource.org.apache.commons.collections" rev="3.2.0" conf="test->runtime"/>
|
||||
<dependency org="org.junit" name="com.springsource.org.junit" rev="4.4.0" conf="test->runtime"/>
|
||||
<dependency org="org.springframework.integration" name="org.springframework.integration" rev="latest.integration" conf="compile->compile"/>
|
||||
</dependencies>
|
||||
|
||||
</ivy-module>
|
||||
@@ -1,308 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2008 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.integration.quartz;
|
||||
|
||||
import java.util.Date;
|
||||
import java.util.concurrent.CancellationException;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Delayed;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.ScheduledFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import org.quartz.CronTrigger;
|
||||
import org.quartz.InterruptableJob;
|
||||
import org.quartz.Job;
|
||||
import org.quartz.JobDetail;
|
||||
import org.quartz.JobExecutionContext;
|
||||
import org.quartz.JobExecutionException;
|
||||
import org.quartz.JobListener;
|
||||
import org.quartz.Scheduler;
|
||||
import org.quartz.SchedulerException;
|
||||
import org.quartz.SimpleTrigger;
|
||||
import org.quartz.Trigger;
|
||||
import org.quartz.TriggerUtils;
|
||||
import org.quartz.UnableToInterruptJobException;
|
||||
import org.quartz.listeners.JobListenerSupport;
|
||||
|
||||
import org.springframework.integration.scheduling.spi.ScheduleServiceProvider;
|
||||
import org.springframework.scheduling.quartz.MethodInvokingJobDetailFactoryBean;
|
||||
|
||||
/**
|
||||
* A Quartz-based implementation of the {@link org.springframework.integration.scheduling.spi.ScheduleServiceProvider}.
|
||||
*
|
||||
* @author Marius Bogoevici
|
||||
*/
|
||||
public class QuartzScheduleServiceProvider implements ScheduleServiceProvider {
|
||||
|
||||
private static final String RUN_METHOD_NAME = "run";
|
||||
|
||||
private static final String RUNNABLE_INSTANCE = "runnable.instance";
|
||||
|
||||
private static final String FIXED_DELAY_PARAMETER = "fixedDelay";
|
||||
|
||||
private static final String TIME_UNIT_PARAMETER = "timeUnit";
|
||||
|
||||
private static final String TRIGGER_NAME_PARAMETER = "triggerName";
|
||||
|
||||
private static final AtomicLong sequenceNumber = new AtomicLong(Long.MIN_VALUE);
|
||||
|
||||
private final Scheduler scheduler;
|
||||
|
||||
private final JobListener fixedDelayJobListener;
|
||||
|
||||
|
||||
public QuartzScheduleServiceProvider(Scheduler scheduler) {
|
||||
this.scheduler = scheduler;
|
||||
this.fixedDelayJobListener = new FixedDelayJobListener();
|
||||
try {
|
||||
this.scheduler.addJobListener(this.fixedDelayJobListener);
|
||||
}
|
||||
catch (SchedulerException e) {
|
||||
throw new QuartzSchedulingException(e);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
public void execute(Runnable runnable) {
|
||||
try {
|
||||
scheduler.scheduleJob(wrapAsJobDetail(runnable),
|
||||
TriggerUtils.makeImmediateTrigger(generateNameForInstance(runnable), 0, 0l));
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new QuartzSchedulingException(e);
|
||||
}
|
||||
}
|
||||
|
||||
public void shutdown(boolean waitForTasksToCompleteOnShutdown) {
|
||||
try {
|
||||
this.scheduler.shutdown(waitForTasksToCompleteOnShutdown);
|
||||
}
|
||||
catch (SchedulerException e) {
|
||||
throw new QuartzSchedulingException(e);
|
||||
}
|
||||
}
|
||||
|
||||
public ScheduledFuture<?> scheduleWithInitialDelay(Runnable runnable, long initialDelay, TimeUnit timeUnit)
|
||||
throws Exception {
|
||||
getFutureDate(initialDelay, timeUnit);
|
||||
Trigger initialDelayTrigger = new SimpleTrigger(generateNameForInstance(runnable), Scheduler.DEFAULT_GROUP,
|
||||
getFutureDate(initialDelay, timeUnit));
|
||||
return scheduleWithTrigger(runnable, initialDelayTrigger, null);
|
||||
}
|
||||
|
||||
public ScheduledFuture<?> scheduleAtFixedRate(Runnable runnable, long initialDelay, long period, TimeUnit timeUnit)
|
||||
throws Exception {
|
||||
Trigger fixedRateTrigger = new SimpleTrigger(generateNameForInstance(runnable), Scheduler.DEFAULT_GROUP,
|
||||
getFutureDate(initialDelay, timeUnit), null, SimpleTrigger.REPEAT_INDEFINITELY,
|
||||
TimeUnit.MILLISECONDS.convert(period, timeUnit));
|
||||
return scheduleWithTrigger(runnable, fixedRateTrigger, null);
|
||||
}
|
||||
|
||||
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable runnable, long initialDelay, long delay,
|
||||
TimeUnit timeUnit) throws Exception {
|
||||
getFutureDate(initialDelay, timeUnit);
|
||||
Trigger fixedDelayTrigger = createFixedDelayTrigger(runnable, initialDelay, delay, timeUnit);
|
||||
return scheduleWithTrigger(runnable, fixedDelayTrigger, new String[]{this.fixedDelayJobListener.getName()});
|
||||
}
|
||||
|
||||
public ScheduledFuture<?> scheduleWithCronExpression(Runnable runnable, String cronExpression)
|
||||
throws Exception {
|
||||
return scheduleWithTrigger(runnable, new CronTrigger(generateNameForInstance(runnable), Scheduler.DEFAULT_GROUP,
|
||||
cronExpression), null);
|
||||
}
|
||||
|
||||
private ScheduledFuture<?> scheduleWithTrigger(Runnable runnable, Trigger fixedRateTrigger,
|
||||
String[] jobListenerNames) throws Exception {
|
||||
JobDetail jobDetail = wrapAsJobDetail(runnable);
|
||||
if (null != jobListenerNames) {
|
||||
for (String jobListenerName : jobListenerNames) {
|
||||
jobDetail.addJobListener(jobListenerName);
|
||||
}
|
||||
}
|
||||
this.scheduler.scheduleJob(jobDetail, fixedRateTrigger);
|
||||
return new ScheduledFutureJobWrapper(scheduler, jobDetail, fixedRateTrigger);
|
||||
}
|
||||
|
||||
private static JobDetail wrapAsJobDetail(Runnable runnable) throws Exception {
|
||||
MethodInvokingJobDetailFactoryBean factoryBean = new MethodInvokingJobDetailFactoryBean();
|
||||
factoryBean.setTargetObject(runnable);
|
||||
factoryBean.setTargetMethod(RUN_METHOD_NAME);
|
||||
factoryBean.setBeanName(generateNameForInstance(runnable));
|
||||
factoryBean.afterPropertiesSet();
|
||||
JobDetail jobDetail = (JobDetail) factoryBean.getObject();
|
||||
jobDetail.setJobClass(InterruptableMethodInvokingJob.class);
|
||||
jobDetail.getJobDataMap().put(RUNNABLE_INSTANCE, runnable);
|
||||
return jobDetail;
|
||||
}
|
||||
|
||||
private static String generateNameForInstance(Object instance) {
|
||||
return instance.getClass().getName() + "#" + sequenceNumber.incrementAndGet();
|
||||
}
|
||||
|
||||
private static Date getFutureDate(long delay, TimeUnit timeUnit) {
|
||||
return new Date(new Date().getTime() + TimeUnit.MILLISECONDS.convert(delay, timeUnit));
|
||||
}
|
||||
|
||||
private static Trigger createFixedDelayTrigger(Runnable runnable, long initialDelay, long taskDelay, TimeUnit timeUnit) {
|
||||
Trigger fixedDelayTrigger = new SimpleTrigger(generateNameForInstance(runnable), Scheduler.DEFAULT_GROUP,
|
||||
getFutureDate(initialDelay, timeUnit));
|
||||
fixedDelayTrigger.getJobDataMap().put(FIXED_DELAY_PARAMETER, taskDelay);
|
||||
fixedDelayTrigger.getJobDataMap().put(TIME_UNIT_PARAMETER, timeUnit);
|
||||
fixedDelayTrigger.getJobDataMap().put(TRIGGER_NAME_PARAMETER, fixedDelayTrigger.getName());
|
||||
return fixedDelayTrigger;
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrapper class for a Quartz {@link Job}, allowing running Quartz jobs top be manipulated via
|
||||
* the {@link java.util.concurrent.ScheduledFuture} interface.
|
||||
* It is designed to be used for periodic tasks, therefore get() either block or throw
|
||||
* {@link java.util.concurrent.CancellationException}.
|
||||
*/
|
||||
public class ScheduledFutureJobWrapper implements ScheduledFuture<Object> {
|
||||
|
||||
private final Scheduler scheduler;
|
||||
|
||||
private final JobDetail jobDetail;
|
||||
|
||||
private final Trigger trigger;
|
||||
|
||||
private final CountDownLatch cancellationLatch = new CountDownLatch(1);
|
||||
|
||||
|
||||
private ScheduledFutureJobWrapper(Scheduler scheduler, JobDetail jobDetail, Trigger trigger) {
|
||||
this.scheduler = scheduler;
|
||||
this.jobDetail = jobDetail;
|
||||
this.trigger = trigger;
|
||||
}
|
||||
|
||||
|
||||
public long getDelay(TimeUnit unit) {
|
||||
long nextFireTime = this.trigger.getNextFireTime().getTime();
|
||||
long timeNow = new Date().getTime();
|
||||
return nextFireTime > timeNow ? unit.convert(nextFireTime - timeNow, TimeUnit.MILLISECONDS) : 0;
|
||||
}
|
||||
|
||||
public int compareTo(Delayed o) {
|
||||
return new Long(getDelay(TimeUnit.MILLISECONDS)).compareTo(o.getDelay(TimeUnit.MILLISECONDS));
|
||||
}
|
||||
|
||||
/*
|
||||
* Synchronized for thread-safety. This is not supposed to be a highly concurrent operation, therefore
|
||||
* contention should be minimal.
|
||||
*/
|
||||
public synchronized boolean cancel(boolean mayInterruptIfRunning) {
|
||||
if (this.cancellationLatch.getCount() == 0) {
|
||||
return true;
|
||||
}
|
||||
try {
|
||||
if (mayInterruptIfRunning) {
|
||||
this.scheduler.interrupt(this.jobDetail.getName(), Scheduler.DEFAULT_GROUP);
|
||||
}
|
||||
if (this.scheduler.deleteJob(this.jobDetail.getName(), Scheduler.DEFAULT_GROUP)) {
|
||||
this.cancellationLatch.countDown();
|
||||
return true;
|
||||
}
|
||||
else {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
catch (SchedulerException e) {
|
||||
throw new QuartzSchedulingException(e);
|
||||
}
|
||||
}
|
||||
|
||||
public boolean isCancelled() {
|
||||
return this.cancellationLatch.getCount() == 0;
|
||||
}
|
||||
|
||||
public boolean isDone() {
|
||||
return isCancelled();
|
||||
}
|
||||
|
||||
public Object get(long timeout, TimeUnit unit)
|
||||
throws InterruptedException, ExecutionException, TimeoutException {
|
||||
this.cancellationLatch.await(timeout, unit);
|
||||
throw new CancellationException();
|
||||
}
|
||||
|
||||
public Object get() throws InterruptedException, ExecutionException {
|
||||
this.cancellationLatch.await();
|
||||
throw new CancellationException();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrapper class allowing for Quartz jobs to be interrupted.
|
||||
*/
|
||||
public static class InterruptableMethodInvokingJob extends MethodInvokingJobDetailFactoryBean.MethodInvokingJob
|
||||
implements InterruptableJob {
|
||||
|
||||
private Thread executionThread;
|
||||
|
||||
protected void executeInternal(JobExecutionContext context) throws JobExecutionException {
|
||||
this.executionThread = Thread.currentThread();
|
||||
super.executeInternal(context);
|
||||
}
|
||||
|
||||
public void interrupt() throws UnableToInterruptJobException {
|
||||
this.executionThread.interrupt();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
/**
|
||||
* {@link org.quartz.JobListener} used for re-scheduling a fixed delay task, as Quartz does
|
||||
* not support fixed delay tasks out of the box.
|
||||
*/
|
||||
public class FixedDelayJobListener extends JobListenerSupport {
|
||||
|
||||
private final String name;
|
||||
|
||||
|
||||
private FixedDelayJobListener() {
|
||||
this.name = generateNameForInstance(this);
|
||||
}
|
||||
|
||||
|
||||
public String getName() {
|
||||
return name;
|
||||
}
|
||||
|
||||
public void jobWasExecuted(JobExecutionContext context, JobExecutionException jobException) {
|
||||
try {
|
||||
Trigger trigger = createFixedDelayTrigger((Runnable) context.getMergedJobDataMap().get(
|
||||
RUNNABLE_INSTANCE),
|
||||
context.getMergedJobDataMap().getLong(FIXED_DELAY_PARAMETER),
|
||||
context.getMergedJobDataMap().getLong(FIXED_DELAY_PARAMETER),
|
||||
(TimeUnit) context.getMergedJobDataMap().get(TIME_UNIT_PARAMETER));
|
||||
trigger.setJobGroup(context.getJobDetail().getGroup());
|
||||
trigger.setJobName(context.getJobDetail().getName());
|
||||
scheduler.rescheduleJob(context.getMergedJobDataMap().getString(TRIGGER_NAME_PARAMETER),
|
||||
Scheduler.DEFAULT_GROUP, trigger);
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new QuartzSchedulingException(e);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
@@ -1,35 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2008 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.integration.quartz;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
*/
|
||||
public class QuartzSchedulingException extends RuntimeException {
|
||||
|
||||
public QuartzSchedulingException() {
|
||||
super();
|
||||
}
|
||||
|
||||
public QuartzSchedulingException(Throwable cause) {
|
||||
super(cause);
|
||||
}
|
||||
|
||||
public QuartzSchedulingException(String message) {
|
||||
super(message);
|
||||
}
|
||||
}
|
||||
@@ -1,209 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2008 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.integration.quartz;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ScheduledFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import org.junit.Test;
|
||||
import org.quartz.Scheduler;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.integration.scheduling.spi.ScheduleServiceProvider;
|
||||
import org.springframework.test.annotation.DirtiesContext;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.AbstractJUnit4SpringContextTests;
|
||||
|
||||
/**
|
||||
* An integration test for the Quartz-based implementation of the Scheduling SPI.
|
||||
*
|
||||
* @author Marius Bogoevici
|
||||
*/
|
||||
@ContextConfiguration(locations = "quartz-scheduler-context.xml")
|
||||
public class TestQuartzScheduleServiceProvider extends AbstractJUnit4SpringContextTests {
|
||||
|
||||
@Autowired
|
||||
private ScheduleServiceProvider scheduleServiceProvider;
|
||||
|
||||
@Autowired
|
||||
private Scheduler scheduler;
|
||||
|
||||
@Test
|
||||
@DirtiesContext
|
||||
public void testExecute() throws InterruptedException {
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
AtomicBoolean hasRun = new AtomicBoolean(false);
|
||||
scheduleServiceProvider.execute(new SimpleRunnable(latch, hasRun));
|
||||
latch.await(1000, TimeUnit.MILLISECONDS);
|
||||
assertTrue(hasRun.get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DirtiesContext
|
||||
public void testReturnedScheduledFutureCancelsJobWithoutInterruption() throws Exception {
|
||||
CountDownLatch latch = new CountDownLatch(2);
|
||||
AtomicBoolean hasBeenInterrupted = new AtomicBoolean(false);
|
||||
ScheduledFuture<?> scheduledFuture = scheduleServiceProvider.scheduleAtFixedRate(new SimpleRunnable(latch,
|
||||
hasBeenInterrupted), 500, 500,
|
||||
TimeUnit.MILLISECONDS);
|
||||
// the job is scheduled
|
||||
assertEquals(1, scheduler.getJobNames(Scheduler.DEFAULT_GROUP).length);
|
||||
scheduledFuture.cancel(false);
|
||||
// the job is not scheduled anymore
|
||||
assertEquals(0, scheduler.getJobNames(Scheduler.DEFAULT_GROUP).length);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DirtiesContext
|
||||
public void testReturnedScheduledFutureCancelsJobWithInterruption() throws Exception {
|
||||
CountDownLatch startLatch = new CountDownLatch(1);
|
||||
CountDownLatch endLatch = new CountDownLatch(1);
|
||||
AtomicBoolean hasBeenInterrupted = new AtomicBoolean(false);
|
||||
ScheduledFuture<?> scheduledFuture = scheduleServiceProvider.scheduleAtFixedRate(new LongRunningRunnable(startLatch,
|
||||
endLatch, hasBeenInterrupted), 500, 10000, TimeUnit.MILLISECONDS);
|
||||
// the job is scheduled
|
||||
assertEquals(1, scheduler.getJobNames(Scheduler.DEFAULT_GROUP).length);
|
||||
startLatch.await(1000, TimeUnit.MILLISECONDS);
|
||||
scheduledFuture.cancel(true);
|
||||
endLatch.await(1000, TimeUnit.MILLISECONDS);
|
||||
// the job is not scheduled anymore
|
||||
assertEquals(0, scheduler.getJobNames(Scheduler.DEFAULT_GROUP).length);
|
||||
assertTrue(hasBeenInterrupted.get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DirtiesContext
|
||||
public void testExecuteWithFixedRate() throws Exception {
|
||||
CountDownLatch latch = new CountDownLatch(2);
|
||||
AtomicBoolean hasRun = new AtomicBoolean(false);
|
||||
SimpleRunnable runnable = new SimpleRunnable(latch, hasRun);
|
||||
scheduleServiceProvider.scheduleAtFixedRate(runnable, 500, 500, TimeUnit.MILLISECONDS);
|
||||
latch.await(3000, TimeUnit.MILLISECONDS);
|
||||
assertTrue(hasRun.get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DirtiesContext
|
||||
public void testExecuteWithCron() throws Exception {
|
||||
CountDownLatch latch = new CountDownLatch(2);
|
||||
AtomicBoolean hasRun = new AtomicBoolean(false);
|
||||
SimpleRunnable runnable = new SimpleRunnable(latch, hasRun);
|
||||
scheduleServiceProvider.scheduleWithCronExpression(runnable, "0/2 * * * * ?");
|
||||
latch.await(6000, TimeUnit.MILLISECONDS);
|
||||
assertTrue(hasRun.get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DirtiesContext
|
||||
public void testExecuteWithFixedDelay() throws Exception {
|
||||
CountDownLatch latch = new CountDownLatch(2);
|
||||
AtomicBoolean hasRun = new AtomicBoolean(false);
|
||||
SimpleRunnable runnable = new SimpleRunnable(latch, hasRun);
|
||||
scheduleServiceProvider.scheduleWithFixedDelay(runnable, 500, 500, TimeUnit.MILLISECONDS);
|
||||
latch.await(3000, TimeUnit.MILLISECONDS);
|
||||
assertTrue(hasRun.get());
|
||||
}
|
||||
|
||||
|
||||
private class SimpleRunnable implements Runnable {
|
||||
|
||||
private final CountDownLatch latch;
|
||||
|
||||
private final AtomicBoolean hasRun;
|
||||
|
||||
private final List<Long> executionTimes = new ArrayList<Long>();
|
||||
|
||||
private final long startTime;
|
||||
|
||||
|
||||
public SimpleRunnable(CountDownLatch latch, AtomicBoolean hasRun) {
|
||||
this.latch = latch;
|
||||
this.hasRun = hasRun;
|
||||
startTime = 0;
|
||||
}
|
||||
|
||||
|
||||
public List<Long> getExecutionTimes() {
|
||||
return executionTimes;
|
||||
}
|
||||
|
||||
public long getStartTime() {
|
||||
return startTime;
|
||||
}
|
||||
|
||||
public void run() {
|
||||
this.executionTimes.add(new Date().getTime());
|
||||
try {
|
||||
Thread.sleep(500);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
latch.countDown();
|
||||
if (latch.getCount() == 0) {
|
||||
hasRun.set(true);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private class LongRunningRunnable implements Runnable {
|
||||
|
||||
private final CountDownLatch startLatch;
|
||||
|
||||
private final CountDownLatch endLatch;
|
||||
|
||||
private final AtomicBoolean hasBeenInterrupted;
|
||||
|
||||
|
||||
public LongRunningRunnable(CountDownLatch startLatch, CountDownLatch endLatch, AtomicBoolean hasBeenInterupted) {
|
||||
this.startLatch = startLatch;
|
||||
this.endLatch = endLatch;
|
||||
this.hasBeenInterrupted = hasBeenInterupted;
|
||||
}
|
||||
|
||||
|
||||
public void run() {
|
||||
// This method differentiates between the interruptions caused by the test and the external ones
|
||||
try {
|
||||
startLatch.countDown();
|
||||
// Sleep for a long time, waiting for an interruption
|
||||
try {
|
||||
Thread.sleep(2000);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
hasBeenInterrupted.set(true);
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
finally {
|
||||
endLatch.countDown();
|
||||
}
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,11 +0,0 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-2.0.xsd">
|
||||
|
||||
<bean id="scheduleServiceProvider" class="org.springframework.integration.quartz.QuartzScheduleServiceProvider">
|
||||
<constructor-arg ref="scheduler"/>
|
||||
</bean>
|
||||
|
||||
<bean id="scheduler" class="org.springframework.scheduling.quartz.SchedulerFactoryBean"/>
|
||||
</beans>
|
||||
@@ -1,9 +0,0 @@
|
||||
Bundle-SymbolicName: org.springframework.integration.quartz
|
||||
Bundle-Name: Spring Integration Quartz Support
|
||||
Bundle-Vendor: SpringSource
|
||||
Bundle-ManifestVersion: 2
|
||||
Import-Template:
|
||||
org.springframework.integration.*;version="[1.0.0, 1.0.1)",
|
||||
org.springframework.scheduling.quartz;version="[2.5.5.A, 3.0.0)",
|
||||
org.quartz.*;version="[1.6.0, 2.0.0)"
|
||||
|
||||
Reference in New Issue
Block a user