INT-1553, added more test for the Inbound side
This commit is contained in:
@@ -33,13 +33,14 @@ import org.springframework.integration.history.TrackableComponent;
|
||||
import org.springframework.integration.store.MetadataStore;
|
||||
import org.springframework.integration.store.SimpleMetadataStore;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.scheduling.TaskScheduler;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import twitter4j.DirectMessage;
|
||||
import twitter4j.Status;
|
||||
import twitter4j.Twitter;
|
||||
import twitter4j.http.OAuthAuthorization;
|
||||
|
||||
/**
|
||||
* Abstract class that defines common operations for receiving various types of
|
||||
@@ -94,6 +95,8 @@ public abstract class AbstractTwitterMessageSource<T> extends AbstractEndpoint
|
||||
|
||||
@Override
|
||||
protected void onInit() throws Exception{
|
||||
Assert.notNull(this.getTaskScheduler(),
|
||||
"Can not locate TaskScheduler. You must inject one explicitly or define a bean by the name 'taskScheduler'");
|
||||
super.onInit();
|
||||
|
||||
if (this.metadataStore == null) {
|
||||
@@ -119,13 +122,14 @@ public abstract class AbstractTwitterMessageSource<T> extends AbstractEndpoint
|
||||
else if (logger.isWarnEnabled()) {
|
||||
logger.warn(this.getClass().getSimpleName() + " has no name. MetadataStore key might not be unique.");
|
||||
}
|
||||
String accessToken = ((OAuthAuthorization)twitter.getAuthorization()).getOAuthAccessToken().getToken();
|
||||
String accessToken = twitter.getOAuthAccessToken().getToken();
|
||||
metadataKeyBuilder.append(accessToken);
|
||||
this.metadataKey = metadataKeyBuilder.toString();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
protected void forwardAll(List<T> tResponses) {
|
||||
Object o = tResponses.iterator();
|
||||
Collections.sort(tResponses, this.getComparator());
|
||||
for (T twitterResponse : tResponses) {
|
||||
forward(twitterResponse);
|
||||
|
||||
@@ -19,6 +19,7 @@ import java.util.Comparator;
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.integration.MessagingException;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
|
||||
import twitter4j.DirectMessage;
|
||||
import twitter4j.Paging;
|
||||
@@ -53,7 +54,9 @@ public class DirectMessageReceivingMessageSource extends AbstractTwitterMessageS
|
||||
? twitter.getDirectMessages()
|
||||
: twitter.getDirectMessages(new Paging(sinceId));
|
||||
|
||||
forwardAll(dms);
|
||||
if (!CollectionUtils.isEmpty(dms)){
|
||||
forwardAll(dms);
|
||||
}
|
||||
}
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
/*
|
||||
* Copyright 2002-2010 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.twitter.inbound;
|
||||
|
||||
import static junit.framework.Assert.assertEquals;
|
||||
import static junit.framework.Assert.assertTrue;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
import java.util.Queue;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
|
||||
import twitter4j.DirectMessage;
|
||||
import twitter4j.Paging;
|
||||
import twitter4j.RateLimitStatus;
|
||||
import twitter4j.ResponseList;
|
||||
import twitter4j.Twitter;
|
||||
import twitter4j.http.AccessToken;
|
||||
|
||||
/**
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class DirectMessageReceivingMessageSourceTests {
|
||||
|
||||
private DirectMessage firstMessage;
|
||||
|
||||
private DirectMessage secondMessage;
|
||||
|
||||
private Twitter twitter = mock(Twitter.class);
|
||||
|
||||
|
||||
@Before
|
||||
public void prepare() throws Exception{
|
||||
twitter = mock(Twitter.class);
|
||||
firstMessage = mock(DirectMessage.class);
|
||||
when(firstMessage.getCreatedAt()).thenReturn(new Date(5555555555L));
|
||||
when(firstMessage.getId()).thenReturn(200);
|
||||
secondMessage = mock(DirectMessage.class);
|
||||
when(secondMessage.getCreatedAt()).thenReturn(new Date(2222222222L));
|
||||
when(secondMessage.getId()).thenReturn(2000);
|
||||
|
||||
|
||||
when(twitter.getOAuthAccessToken()).thenReturn(new AccessToken("token123", "tokenSecret123"));
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testSuccessfullInitialization() throws Exception{
|
||||
DirectMessageReceivingMessageSource source = new DirectMessageReceivingMessageSource(twitter);
|
||||
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
|
||||
scheduler.afterPropertiesSet();
|
||||
source.setTaskScheduler(scheduler);
|
||||
source.setBeanName("twitterEndpoint");
|
||||
source.afterPropertiesSet();
|
||||
source.start();
|
||||
assertEquals("twitter:inbound-dm-channel-adapter.twitterEndpoint.token123", TestUtils.getPropertyValue(source, "metadataKey"));
|
||||
assertTrue(source.isRunning());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSuccessfullInitializationWithMessages() throws Exception{
|
||||
this.setUpMockScenarioForMessagePolling();
|
||||
|
||||
DirectMessageReceivingMessageSource source = new DirectMessageReceivingMessageSource(twitter);
|
||||
ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
|
||||
scheduler.afterPropertiesSet();
|
||||
source.setTaskScheduler(scheduler);
|
||||
source.setBeanName("twitterEndpoint");
|
||||
source.afterPropertiesSet();
|
||||
source.start();
|
||||
Thread.sleep(1000);
|
||||
System.out.println("Tweets: " + TestUtils.getPropertyValue(source, "tweets"));
|
||||
Queue msg = (Queue) TestUtils.getPropertyValue(source, "tweets");
|
||||
assertTrue(!CollectionUtils.isEmpty(msg));
|
||||
assertEquals(1, msg.size()); // because the other message has a older timestamp and is assumed to be read by
|
||||
DirectMessage message = (DirectMessage) msg.poll();
|
||||
assertEquals(secondMessage, message);
|
||||
|
||||
}
|
||||
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private void setUpMockScenarioForMessagePolling() throws Exception{
|
||||
RateLimitStatus rateLimitStatus = mock(RateLimitStatus.class);
|
||||
when(twitter.getRateLimitStatus()).thenReturn(rateLimitStatus);
|
||||
when(rateLimitStatus.getSecondsUntilReset()).thenReturn(2464);
|
||||
when(rateLimitStatus.getRemainingHits()).thenReturn(250);
|
||||
|
||||
//ResponseList<DirectMessage> responses = mock(ResponseList.class);
|
||||
SampleResoponceList testMessages = new SampleResoponceList();
|
||||
testMessages.add(firstMessage);
|
||||
testMessages.add(secondMessage);
|
||||
//when(responses.iterator()).thenReturn(testMessages.iterator());
|
||||
when(twitter.getDirectMessages()).thenReturn(testMessages);
|
||||
when(twitter.getDirectMessages(Mockito.any(Paging.class))).thenReturn(testMessages);
|
||||
}
|
||||
|
||||
public static class SampleResoponceList extends ArrayList implements ResponseList {
|
||||
|
||||
@Override
|
||||
public RateLimitStatus getRateLimitStatus() {
|
||||
return mock(RateLimitStatus.class);
|
||||
}
|
||||
|
||||
@Override
|
||||
public RateLimitStatus getFeatureSpecificRateLimitStatus() {
|
||||
return mock(RateLimitStatus.class);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -1,98 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2010 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.twitter.inbound;
|
||||
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import java.util.Date;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
|
||||
import twitter4j.DirectMessage;
|
||||
import twitter4j.Twitter;
|
||||
|
||||
/**
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class InboundDirectMessageStatusEndpointTests {
|
||||
|
||||
private DirectMessage firstMessage;
|
||||
|
||||
private DirectMessage secondMessage;
|
||||
|
||||
private Twitter twitter = mock(Twitter.class);
|
||||
|
||||
|
||||
@Before
|
||||
public void prepare() {
|
||||
twitter = mock(Twitter.class);
|
||||
firstMessage = mock(DirectMessage.class);
|
||||
when(firstMessage.getCreatedAt()).thenReturn(new Date(5555555555L));
|
||||
when(firstMessage.getId()).thenReturn(200);
|
||||
secondMessage = mock(DirectMessage.class);
|
||||
when(secondMessage.getCreatedAt()).thenReturn(new Date(2222222222L));
|
||||
when(secondMessage.getId()).thenReturn(2000);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testTwitterMockedUpdates() throws Exception{
|
||||
// QueueChannel channel = new QueueChannel();
|
||||
// DirectMessageReceivingMessageSource endpoint = new DirectMessageReceivingMessageSource(twitter);
|
||||
// //endpoint.setOutputChannel(channel);
|
||||
// ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler();
|
||||
// scheduler.afterPropertiesSet();
|
||||
// endpoint.setTaskScheduler(scheduler);
|
||||
// //endpoint.setConfiguration(this.getTestConfigurationForDirectMessages());
|
||||
// endpoint.setBeanName("twitterEndpoint");
|
||||
// endpoint.afterPropertiesSet();
|
||||
// endpoint.start();
|
||||
// Message<?> message1 = channel.receive(3000);
|
||||
// assertNotNull(message1);
|
||||
// // should be second message since its timestamp is newer
|
||||
// assertEquals(secondMessage.getId(), ((DirectMessage)message1.getPayload()).getId());
|
||||
// Message<?> message2 = channel.receive(100);
|
||||
// assertNull(message2); // should be null, since
|
||||
}
|
||||
|
||||
|
||||
// @SuppressWarnings("unchecked")
|
||||
// private OAuthConfiguration getTestConfigurationForDirectMessages() throws Exception{
|
||||
// OAuthConfiguration configuration = mock(OAuthConfiguration.class);
|
||||
// RateLimitStatus rateLimitStatus = mock(RateLimitStatus.class);
|
||||
// when(twitter.getRateLimitStatus()).thenReturn(rateLimitStatus);
|
||||
// when(configuration.getTwitter()).thenReturn(twitter);
|
||||
// when(rateLimitStatus.getSecondsUntilReset()).thenReturn(2464);
|
||||
// when(rateLimitStatus.getRemainingHits()).thenReturn(250);
|
||||
//
|
||||
// ResponseList<DirectMessage> responses = mock(ResponseList.class);
|
||||
// List<DirectMessage> testMessages = new ArrayList<DirectMessage>();
|
||||
// testMessages.add(firstMessage);
|
||||
// testMessages.add(secondMessage);
|
||||
//
|
||||
// when(responses.iterator()).thenReturn(testMessages.iterator());
|
||||
// when(twitter.getDirectMessages()).thenReturn(responses);
|
||||
// when(twitter.getDirectMessages(Mockito.any(Paging.class))).thenReturn(responses);
|
||||
// return configuration;
|
||||
// }
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user