From 6335dd44d73a75f01ce19918b86bc0d90635908c Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 11 Nov 2010 14:29:31 -0500 Subject: [PATCH] INT-1603 fixed SearchReceivingMessageSource and tests --- .../twitter/core/Twitter4jTemplate.java | 7 ++--- .../twitter/core/TwitterOperations.java | 3 +- .../inbound/AbstractTwitterMessageSource.java | 22 ++++++-------- .../DirectMessageReceivingMessageSource.java | 2 +- .../inbound/SearchReceivingMessageSource.java | 29 +++---------------- .../TimelineUpdateReceivingMessageSource.java | 2 +- .../twitter/core/Twitter4jTemplateTests.java | 1 + ...ectMessageReceivingMessageSourceTests.java | 10 +++---- ...lineUpdateReceivingMessageSourceTests.java | 4 +-- 9 files changed, 28 insertions(+), 52 deletions(-) diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/Twitter4jTemplate.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/Twitter4jTemplate.java index a01f8c047c..fdfcffcaa5 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/Twitter4jTemplate.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/Twitter4jTemplate.java @@ -15,7 +15,6 @@ */ package org.springframework.integration.twitter.core; -import java.util.ArrayList; import java.util.LinkedList; import java.util.List; @@ -24,7 +23,6 @@ import org.apache.commons.lang.NotImplementedException; import org.springframework.util.Assert; import twitter4j.DirectMessage; -import twitter4j.IDs; import twitter4j.Paging; import twitter4j.Query; import twitter4j.QueryResult; @@ -203,10 +201,11 @@ public class Twitter4jTemplate implements TwitterOperations{ } @Override - public SearchResults search(String query, int page, int pageSize) { + public SearchResults search(String query, int page, int sinceId) { Assert.hasText(query, "'query' must not be null"); Query q = new Query(query); q.setPage(page); + q.setSinceId(sinceId); return this.search(q); } @@ -218,6 +217,7 @@ public class Twitter4jTemplate implements TwitterOperations{ q.setPage(page); q.setSinceId(sinceId); q.setMaxId(maxId); + q.setRpp(resultsPerPage); return this.search(q); } @@ -228,7 +228,6 @@ public class Twitter4jTemplate implements TwitterOperations{ private SearchResults search(Query query){ try { QueryResult result = twitter.search(query); - if (result != null){ List t4jTweets = result.getTweets(); List tweets = this.buildTweetsFromTwitterResponses(t4jTweets); diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/TwitterOperations.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/TwitterOperations.java index a6bc676cf6..884cc26188 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/TwitterOperations.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/core/TwitterOperations.java @@ -87,7 +87,8 @@ public interface TwitterOperations { * @return a {@link SearchResults} containing {@link Tweet}s * */ - SearchResults search(String query, int page, int pageSize); + SearchResults search(String query, int page, int sinceId); + /** * Searches Twitter, returning a specific page out of the complete set of 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 f9f29b3cc2..c49b093199 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 @@ -59,7 +59,7 @@ public abstract class AbstractTwitterMessageSource extends AbstractEndpoint private volatile String metadataKey; - protected final Queue tweets = new LinkedBlockingQueue(); + protected final Queue tweets = new LinkedBlockingQueue(); protected volatile int prefetchThreshold = 0; @@ -129,9 +129,9 @@ public abstract class AbstractTwitterMessageSource extends AbstractEndpoint } @SuppressWarnings("unchecked") - protected void forwardAll(List tResponses) { + protected void forwardAll(List tResponses) { Collections.sort(tResponses, this.getComparator()); - for (T twitterResponse : tResponses) { + for (Tweet twitterResponse : tResponses) { forward(twitterResponse); } } @@ -164,23 +164,19 @@ public abstract class AbstractTwitterMessageSource extends AbstractEndpoint } public Message receive() { - Object tweet = tweets.poll(); + Tweet tweet = tweets.poll(); if (tweet != null){ + return MessageBuilder.withPayload(tweet).build(); } return null; } - protected void forward(T tweet) { + protected void forward(Tweet tweet) { synchronized (this.markerGuard) { - long id = 0; - if (tweet instanceof Tweet) { - id = ((Tweet) tweet).getId(); - } - else { - throw new IllegalArgumentException("Unsupported type of Twitter message: " + tweet.getClass()); - } + long id = tweet.getId(); + String lastId = this.metadataStore.get(this.metadataKey); long lastTweetId = 0; @@ -188,8 +184,8 @@ public abstract class AbstractTwitterMessageSource extends AbstractEndpoint lastTweetId = Long.parseLong(lastId); } if (id > lastTweetId) { + markLastStatusId(tweet.getId()); tweets.add(tweet); - markLastStatusId(id); } } } 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 4c283f3b30..0c6d1c8454 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 @@ -38,7 +38,7 @@ public class DirectMessageReceivingMessageSource extends AbstractTwitterMessageS @Override public String getComponentType() { - return "twitter:inbound-dm-channel-adapter"; + return "twitter:dm-inbound-channel-adapter"; } @Override diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSource.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSource.java index 23e263c54b..8ded71b0ac 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSource.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/SearchReceivingMessageSource.java @@ -29,13 +29,7 @@ import org.springframework.util.Assert; * @since 2.0 */ public class SearchReceivingMessageSource extends AbstractTwitterMessageSource { - /* since Twitter return 15 entries per page we need to be able to manage - * how many pages deep are we willing to go. Not sure yet about exposing this attribute via namespace - * but setting default to 10. - */ - private volatile int pageDepth = 10; - - private volatile int currentPage = 1; + private volatile String query; public SearchReceivingMessageSource(TwitterOperations twitter){ @@ -58,25 +52,10 @@ public class SearchReceivingMessageSource extends AbstractTwitterMessageSource twetList = results.getTweets(); - if (currentPage == 1){ - forwardAll(twetList); - } - else { - for (Tweet tweet : twetList) { - tweets.add(tweet); - } - } - if (twetList != null && twetList.size() > 0){ - currentPage++; - } - else { - currentPage = 1; - } + forwardAll(twetList); } } catch (Exception e) { if (e instanceof RuntimeException){ diff --git a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSource.java b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSource.java index 29b0b0d4d2..fa93daab73 100644 --- a/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSource.java +++ b/spring-integration-twitter/src/main/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSource.java @@ -37,7 +37,7 @@ public class TimelineUpdateReceivingMessageSource extends AbstractTwitterMessage } @Override public String getComponentType() { - return "twitter:inbound-update-channel-adapter"; + return "twitter:inbound-channel-adapter"; } @Override diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/core/Twitter4jTemplateTests.java b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/core/Twitter4jTemplateTests.java index bd4236cf62..0cc3990aa2 100644 --- a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/core/Twitter4jTemplateTests.java +++ b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/core/Twitter4jTemplateTests.java @@ -73,6 +73,7 @@ public class Twitter4jTemplateTests { @Test public void testProfileId() throws Exception{ when(twitter.getScreenName()).thenReturn("kermit"); + when(twitter.isOAuthEnabled()).thenReturn(true); assertEquals("kermit", template.getProfileId()); } 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 index b7d72c9e79..c101509578 100644 --- 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 @@ -99,7 +99,7 @@ public class DirectMessageReceivingMessageSourceTests { @Test public void testSuccessfullInitialization() throws Exception{ - + when(tw.isOAuthEnabled()).thenReturn(true); DirectMessageReceivingMessageSource source = new DirectMessageReceivingMessageSource(twitter); ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); scheduler.afterPropertiesSet(); @@ -107,7 +107,7 @@ public class DirectMessageReceivingMessageSourceTests { source.setBeanName("twitterEndpoint"); source.afterPropertiesSet(); source.start(); - assertEquals("twitter:inbound-dm-channel-adapter.twitterEndpoint.kermit", TestUtils.getPropertyValue(source, "metadataKey")); + assertEquals("twitter:dm-inbound-channel-adapter.twitterEndpoint.kermit", TestUtils.getPropertyValue(source, "metadataKey")); assertTrue(source.isRunning()); source.stop(); } @@ -127,7 +127,7 @@ public class DirectMessageReceivingMessageSourceTests { Thread.sleep(1000); 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 + assertEquals(1, msg.size()); Tweet message = (Tweet) msg.poll(); assertEquals(2000, message.getId()); Thread.sleep(1000); @@ -164,7 +164,7 @@ public class DirectMessageReceivingMessageSourceTests { Thread.sleep(1000); 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 + assertEquals(1, msg.size()); Tweet message = (Tweet) msg.poll(); assertEquals(2000, message.getId()); source.stop(); @@ -201,7 +201,7 @@ public class DirectMessageReceivingMessageSourceTests { @SuppressWarnings("unchecked") private void setUpMockScenarioForMessagePolling() throws Exception{ RateLimitStatus rateLimitStatus = mock(RateLimitStatus.class); - + when(tw.isOAuthEnabled()).thenReturn(true); when(tw.getRateLimitStatus()).thenReturn(rateLimitStatus); when(rateLimitStatus.getSecondsUntilReset()).thenReturn(1000); when(rateLimitStatus.getRemainingHits()).thenReturn(1000); diff --git a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSourceTests.java b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSourceTests.java index 62045b0f4d..25fec54fba 100644 --- a/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSourceTests.java +++ b/spring-integration-twitter/src/test/java/org/springframework/integration/twitter/inbound/TimelineUpdateReceivingMessageSourceTests.java @@ -99,7 +99,7 @@ public class TimelineUpdateReceivingMessageSourceTests { @Test public void testSuccessfullInitialization() throws Exception{ - + when(tw.isOAuthEnabled()).thenReturn(true); TimelineUpdateReceivingMessageSource source = new TimelineUpdateReceivingMessageSource(twitter); ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); scheduler.afterPropertiesSet(); @@ -107,7 +107,7 @@ public class TimelineUpdateReceivingMessageSourceTests { source.setBeanName("twitterEndpoint"); source.afterPropertiesSet(); source.start(); - assertEquals("twitter:inbound-update-channel-adapter.twitterEndpoint.kermit", TestUtils.getPropertyValue(source, "metadataKey")); + assertEquals("twitter:inbound-channel-adapter.twitterEndpoint.kermit", TestUtils.getPropertyValue(source, "metadataKey")); assertTrue(source.isRunning()); source.stop(); }