INT-1603 fixed SearchReceivingMessageSource and tests
This commit is contained in:
@@ -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<twitter4j.Tweet> t4jTweets = result.getTweets();
|
||||
List<Tweet> tweets = this.buildTweetsFromTwitterResponses(t4jTweets);
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -59,7 +59,7 @@ public abstract class AbstractTwitterMessageSource<T> extends AbstractEndpoint
|
||||
|
||||
private volatile String metadataKey;
|
||||
|
||||
protected final Queue<Object> tweets = new LinkedBlockingQueue<Object>();
|
||||
protected final Queue<Tweet> tweets = new LinkedBlockingQueue<Tweet>();
|
||||
|
||||
protected volatile int prefetchThreshold = 0;
|
||||
|
||||
@@ -129,9 +129,9 @@ public abstract class AbstractTwitterMessageSource<T> extends AbstractEndpoint
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
protected void forwardAll(List<T> tResponses) {
|
||||
protected void forwardAll(List<Tweet> tResponses) {
|
||||
Collections.sort(tResponses, this.getComparator());
|
||||
for (T twitterResponse : tResponses) {
|
||||
for (Tweet twitterResponse : tResponses) {
|
||||
forward(twitterResponse);
|
||||
}
|
||||
}
|
||||
@@ -164,23 +164,19 @@ public abstract class AbstractTwitterMessageSource<T> 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<T> extends AbstractEndpoint
|
||||
lastTweetId = Long.parseLong(lastId);
|
||||
}
|
||||
if (id > lastTweetId) {
|
||||
markLastStatusId(tweet.getId());
|
||||
tweets.add(tweet);
|
||||
markLastStatusId(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -29,13 +29,7 @@ import org.springframework.util.Assert;
|
||||
* @since 2.0
|
||||
*/
|
||||
public class SearchReceivingMessageSource extends AbstractTwitterMessageSource<Tweet> {
|
||||
/* 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<T
|
||||
public void run() {
|
||||
try {
|
||||
if (tweets.size() <= prefetchThreshold){
|
||||
if (currentPage == pageDepth){
|
||||
currentPage = 1;
|
||||
}
|
||||
SearchResults results = twitter.search(query, currentPage, 0);
|
||||
SearchResults results = twitter.search(query);
|
||||
|
||||
List<Tweet> 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){
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user