From ae2413e68af2bf24bd736230faf64e0e91aef882 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 5 Nov 2010 17:17:49 -0400 Subject: [PATCH] INT-1553, added more test for the Inbound side --- .../inbound/AbstractTwitterMessageSource.java | 8 +- .../DirectMessageReceivingMessageSource.java | 5 +- ...ectMessageReceivingMessageSourceTests.java | 134 ++++++++++++++++++ ...boundDirectMessageStatusEndpointTests.java | 98 ------------- 4 files changed, 144 insertions(+), 101 deletions(-) create mode 100644 spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSourceTests.java delete mode 100644 spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/InboundDirectMessageStatusEndpointTests.java diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java index cf4d6870f5..73ddc759df 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/AbstractTwitterMessageSource.java @@ -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 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 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 tResponses) { + Object o = tResponses.iterator(); Collections.sort(tResponses, this.getComparator()); for (T twitterResponse : tResponses) { forward(twitterResponse); diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSource.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSource.java index 6b9d1ff94d..a245540ec6 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSource.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSource.java @@ -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(); diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSourceTests.java b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSourceTests.java new file mode 100644 index 0000000000..12c0acd924 --- /dev/null +++ b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/DirectMessageReceivingMessageSourceTests.java @@ -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 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); + } + + } +} diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/InboundDirectMessageStatusEndpointTests.java b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/InboundDirectMessageStatusEndpointTests.java deleted file mode 100644 index e51cfd2bfc..0000000000 --- a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/InboundDirectMessageStatusEndpointTests.java +++ /dev/null @@ -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 responses = mock(ResponseList.class); -// List testMessages = new ArrayList(); -// 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; -// } - -}