http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/main/java/org/apache/usergrid/persistence/graph/serialization/impl/shard/impl/ShardGroupCompactionImpl.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/main/java/org/apache/usergrid/persistence/graph/serialization/impl/shard/impl/ShardGroupCompactionImpl.java b/stack/corepersistence/graph/src/main/java/org/apache/usergrid/persistence/graph/serialization/impl/shard/impl/ShardGroupCompactionImpl.java index e663d5a..f398fa2 100644 --- a/stack/corepersistence/graph/src/main/java/org/apache/usergrid/persistence/graph/serialization/impl/shard/impl/ShardGroupCompactionImpl.java +++ b/stack/corepersistence/graph/src/main/java/org/apache/usergrid/persistence/graph/serialization/impl/shard/impl/ShardGroupCompactionImpl.java @@ -40,7 +40,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.usergrid.persistence.core.consistency.TimeService; -import org.apache.usergrid.persistence.core.executor.TaskExecutorFactory; import org.apache.usergrid.persistence.core.scope.ApplicationScope; import org.apache.usergrid.persistence.graph.GraphFig; import org.apache.usergrid.persistence.graph.MarkedEdge; @@ -66,7 +65,6 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; -import com.google.common.util.concurrent.MoreExecutors; import com.google.inject.Inject; import com.google.inject.Singleton; import com.netflix.astyanax.Keyspace; @@ -82,7 +80,7 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { private final AtomicLong countAudits; - private static final Logger LOG = LoggerFactory.getLogger( ShardGroupCompactionImpl.class ); + private static final Logger logger = LoggerFactory.getLogger( ShardGroupCompactionImpl.class ); private static final Charset CHARSET = Charset.forName( "UTF-8" ); @@ -149,8 +147,8 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { Preconditions .checkArgument( group.shouldCompact( startTime ), "Compaction cannot be run yet. Ignoring compaction." ); - if(LOG.isDebugEnabled()) { - LOG.debug("Compacting shard group. Audit count is {} ", countAudits.get()); + if(logger.isDebugEnabled()) { + logger.debug("Compacting shard group. Audit count is {} ", countAudits.get()); } final CompactionResult.CompactionBuilder resultBuilder = CompactionResult.builder(); @@ -215,7 +213,7 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { deleteRowBatch.execute(); } catch ( Throwable t ) { - LOG.error( "Unable to move edges from shard {} to shard {}", sourceShard, targetShard ); + logger.error( "Unable to move edges from shard {} to shard {}", sourceShard, targetShard ); } } } @@ -227,11 +225,11 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { deleteRowBatch.execute(); } catch ( Throwable t ) { - LOG.error( "Unable to move edges to target shard {}", targetShard ); + logger.error( "Unable to move edges to target shard {}", targetShard ); } - LOG.info( "Finished compacting {} shards and moved {} edges", sourceShards, edgeCount ); + logger.info( "Finished compacting {} shards and moved {} edges", sourceShards, edgeCount ); resultBuilder.withCopiedEdges( edgeCount ).withSourceShards( sourceShards ).withTargetShard( targetShard ); @@ -255,7 +253,7 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { continue; } - LOG.info( "Source shards have been fully drained. Removing shard {}", source ); + logger.info( "Source shards have been fully drained. Removing shard {}", source ); final MutationBatch shardRemoval = edgeShardSerialization.removeShardMeta( scope, source, edgeMeta ); shardRemovalRollup.mergeShallow( shardRemoval ); @@ -272,7 +270,7 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { } - LOG.info( "Shard has been fully compacted. Marking shard {} as compacted in Cassandra", targetShard ); + logger.info( "Shard has been fully compacted. Marking shard {} as compacted in Cassandra", targetShard ); //Overwrite our shard index with a newly created one that has been marked as compacted Shard compactedShard = new Shard( targetShard.getShardIndex(), timeService.getCurrentTime(), true ); @@ -306,8 +304,8 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { countAudits.getAndIncrement(); - if(LOG.isDebugEnabled()) { - LOG.debug("Auditing shard group {}. count is {} ", group, countAudits.get()); + if(logger.isDebugEnabled()) { + logger.debug("Auditing shard group {}. count is {} ", group, countAudits.get()); } /** @@ -322,7 +320,7 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { catch ( RejectedExecutionException ree ) { //ignore, if this happens we don't care, we're saturated, we can check later - LOG.error( "Rejected audit for shard of scope {} edge, meta {} and group {}", scope, edgeMeta, group ); + logger.error( "Rejected audit for shard of scope {} edge, meta {} and group {}", scope, edgeMeta, group ); return Futures.immediateFuture( AuditResult.NOT_CHECKED ); } @@ -333,13 +331,13 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { Futures.addCallback( future, new FutureCallback<AuditResult>() { @Override public void onSuccess( @Nullable final AuditResult result ) { - LOG.debug( "Successfully completed audit of task {}", result ); + logger.debug( "Successfully completed audit of task {}", result ); } @Override public void onFailure( final Throwable t ) { - LOG.error( "Unable to perform audit. Exception is ", t ); + logger.error( "Unable to perform audit. Exception is ", t ); } } ); @@ -414,7 +412,7 @@ public class ShardGroupCompactionImpl implements ShardGroupCompaction { */ try { CompactionResult result = compact( scope, edgeMeta, group ); - LOG.info( "Compaction result for compaction of scope {} with edge meta data of {} and shard group " + logger.info( "Compaction result for compaction of scope {} with edge meta data of {} and shard group " + "{} is {}", new Object[] { scope, edgeMeta, group, result } ); } finally {
http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerLoadTest.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerLoadTest.java b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerLoadTest.java index a0be6a6..864397e 100644 --- a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerLoadTest.java +++ b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerLoadTest.java @@ -31,7 +31,6 @@ import java.util.concurrent.Future; import org.apache.usergrid.StressTest; import org.junit.Before; -import org.junit.Ignore; import org.junit.Rule; import org.junit.Test; import org.junit.experimental.categories.Category; @@ -67,7 +66,7 @@ import static org.junit.Assert.fail; @UseModules( TestGraphModule.class ) @Category(StressTest.class) public class GraphManagerLoadTest { - private static final Logger log = LoggerFactory.getLogger( GraphManagerLoadTest.class ); + private static final Logger logger = LoggerFactory.getLogger( GraphManagerLoadTest.class ); @Inject private GraphManagerFactory factory; @@ -208,12 +207,12 @@ public class GraphManagerLoadTest { if ( i % 1000 == 0 ) { - log.info( " Wrote: " + i ); + logger.info( " Wrote: " + i ); } } timer.stop(); - log.info( "Total time to write {} entries {} ms", writeLimit, timer.getTime() ); + logger.info( "Total time to write {} entries {} ms", writeLimit, timer.getTime() ); timer.reset(); timer.start(); @@ -237,7 +236,7 @@ public class GraphManagerLoadTest { @Override public void onNext( final List<MarkedEdge> edges ) { - log.info("Read {} edges", edges.size()); + logger.info("Read {} edges", edges.size()); } } ); @@ -245,7 +244,7 @@ public class GraphManagerLoadTest { latch.await(); - log.info( "Total time to read {} entries {} ms", readCount, timer.getTime() ); + logger.info( "Total time to read {} entries {} ms", readCount, timer.getTime() ); return true; } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerShardConsistencyIT.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerShardConsistencyIT.java b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerShardConsistencyIT.java index 0cc0e88..3ae3ff1 100644 --- a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerShardConsistencyIT.java +++ b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerShardConsistencyIT.java @@ -42,7 +42,6 @@ import javax.annotation.Nullable; import org.apache.usergrid.StressTest; import org.junit.After; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.junit.experimental.categories.Category; import org.slf4j.Logger; @@ -83,7 +82,7 @@ import static org.junit.Assert.fail; public class GraphManagerShardConsistencyIT { - private static final Logger log = LoggerFactory.getLogger( GraphManagerShardConsistencyIT.class ); + private static final Logger logger = LoggerFactory.getLogger( GraphManagerShardConsistencyIT.class ); private static final MetricRegistry registry = new MetricRegistry(); @@ -135,7 +134,7 @@ public class GraphManagerShardConsistencyIT { reporter = - Slf4jReporter.forRegistry( registry ).outputTo( log ).convertRatesTo( TimeUnit.SECONDS ) + Slf4jReporter.forRegistry( registry ).outputTo(logger).convertRatesTo( TimeUnit.SECONDS ) .convertDurationsTo( TimeUnit.MILLISECONDS ).build(); @@ -228,7 +227,7 @@ public class GraphManagerShardConsistencyIT { final long minExecutionTime = graphFig.getShardMinDelta() + graphFig.getShardCacheTimeout(); - log.info( "Writing {} edges per worker on {} workers in {} injectors", workerWriteLimit, numWorkersPerInjector, + logger.info( "Writing {} edges per worker on {} workers in {} injectors", workerWriteLimit, numWorkersPerInjector, numInjectors ); @@ -284,7 +283,7 @@ public class GraphManagerShardConsistencyIT { @Override public void onSuccess( @Nullable final Long result ) { - log.info( "Successfully ran the read, re-running" ); + logger.info( "Successfully ran the read, re-running" ); executor.submit( new ReadWorker( gmf, generator, writeCount, readMeter ) ); } @@ -292,7 +291,7 @@ public class GraphManagerShardConsistencyIT { @Override public void onFailure( final Throwable t ) { failures.add( t ); - log.error( "Failed test!", t ); + logger.error( "Failed test!", t ); } } ); } @@ -337,7 +336,7 @@ public class GraphManagerShardConsistencyIT { final ShardEntryGroup group = groups.next(); shardEntryGroups.add( group ); - log.info( "Compaction pending status for group {} is {}", group, group.isCompactionPending() ); + logger.info( "Compaction pending status for group {} is {}", group, group.isCompactionPending() ); if ( !group.isCompactionPending() ) { compactedCount++; @@ -347,7 +346,7 @@ public class GraphManagerShardConsistencyIT { //we're done if ( compactedCount >= expectedShardCount ) { - log.info( "All compactions complete, sleeping" ); + logger.info( "All compactions complete, sleeping" ); // final Object mutex = new Object(); // @@ -467,7 +466,7 @@ public class GraphManagerShardConsistencyIT { final long minExecutionTime = graphFig.getShardMinDelta() + graphFig.getShardCacheTimeout(); - log.info( "Writing {} edges per worker on {} workers in {} injectors", workerWriteLimit, numWorkersPerInjector, + logger.info( "Writing {} edges per worker on {} workers in {} injectors", workerWriteLimit, numWorkersPerInjector, numInjectors ); @@ -518,11 +517,11 @@ public class GraphManagerShardConsistencyIT { shardCount++; - log.info( "Compaction pending status for group {} is {}", group, group.isCompactionPending() ); + logger.info( "Compaction pending status for group {} is {}", group, group.isCompactionPending() ); } - log.info( "found {} shard groups", shardCount ); + logger.info( "found {} shard groups", shardCount ); //now mark and delete all the edges @@ -559,7 +558,7 @@ public class GraphManagerShardConsistencyIT { @Override public void onSuccess( @Nullable final Long result ) { - log.info( "Successfully ran the read, re-running" ); + logger.info( "Successfully ran the read, re-running" ); executor.submit( new ReadWorker( gmf, generator, writeCount, readMeter ) ); } @@ -567,7 +566,7 @@ public class GraphManagerShardConsistencyIT { @Override public void onFailure( final Throwable t ) { failures.add( t ); - log.error( "Failed test!", t ); + logger.error( "Failed test!", t ); } } ); @@ -608,7 +607,7 @@ public class GraphManagerShardConsistencyIT { group = groups.next(); - log.info( "Shard size for group is {}", group.getReadShards() ); + logger.info( "Shard size for group is {}", group.getReadShards() ); shardCount += group.getReadShards().size(); } @@ -616,7 +615,7 @@ public class GraphManagerShardConsistencyIT { //we're done, 1 shard remains, we have a group, and it's our default shard if ( shardCount == 1 && group != null && group.getMinShard().getShardIndex() == Shard.MIN_SHARD.getShardIndex() ) { - log.info( "All compactions complete," ); + logger.info( "All compactions complete," ); break; } @@ -673,7 +672,7 @@ public class GraphManagerShardConsistencyIT { if ( i % 1000 == 0 ) { - log.info( " Wrote: " + i ); + logger.info( " Wrote: " + i ); } } @@ -716,10 +715,10 @@ public class GraphManagerShardConsistencyIT { .countLong().toBlocking().last(); - log.info( "Completed reading {} edges", returnedEdgeCount ); + logger.info( "Completed reading {} edges", returnedEdgeCount ); if ( writeCount != returnedEdgeCount ) { - log.warn( "Unexpected edge count returned!!! Expected {} but was {}", writeCount, + logger.warn( "Unexpected edge count returned!!! Expected {} but was {}", writeCount, returnedEdgeCount ); } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerStressTest.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerStressTest.java b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerStressTest.java index 98065ce..be1eee4 100644 --- a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerStressTest.java +++ b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/GraphManagerStressTest.java @@ -56,7 +56,7 @@ import static org.mockito.Mockito.when; @UseModules(TestGraphModule.class) @Category(StressTest.class) public class GraphManagerStressTest { - private static final Logger log = LoggerFactory.getLogger( GraphManagerStressTest.class ); + private static final Logger logger = LoggerFactory.getLogger( GraphManagerStressTest.class ); @Inject private GraphManagerFactory factory; @@ -120,7 +120,7 @@ public class GraphManagerStressTest { .toBlocking().toIterable(); for ( MarkedEdge edge : edges ) { - log.debug( "Firing on next for edge {}", edge ); + logger.debug( "Firing on next for edge {}", edge ); subscriber.onNext( edge ); } @@ -250,12 +250,12 @@ public class GraphManagerStressTest { ids.add( returned ); if ( i % 1000 == 0 ) { - log.info( " Wrote: " + i ); + logger.info( " Wrote: " + i ); } } timer.stop(); - log.info( "Total time to write {} entries {}ms", limit, timer.getTime() ); + logger.info( "Total time to write {} entries {}ms", limit, timer.getTime() ); timer.reset(); timer.start(); @@ -290,7 +290,7 @@ public class GraphManagerStressTest { assertEquals( 0, ids.size() ); - log.info( "Total time to read {} entries {}ms", limit, timer.getTime() ); + logger.info( "Total time to read {} entries {}ms", limit, timer.getTime() ); } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/impl/NodeDeleteListenerTest.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/impl/NodeDeleteListenerTest.java b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/impl/NodeDeleteListenerTest.java index 5a5df79..438a978 100644 --- a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/impl/NodeDeleteListenerTest.java +++ b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/impl/NodeDeleteListenerTest.java @@ -70,7 +70,7 @@ import static org.mockito.Mockito.when; @UseModules( { TestGraphModule.class } ) public class NodeDeleteListenerTest { - private static final Logger log = LoggerFactory.getLogger( NodeDeleteListenerTest.class ); + private static final Logger logger = LoggerFactory.getLogger( NodeDeleteListenerTest.class ); @Inject @@ -359,8 +359,8 @@ public class NodeDeleteListenerTest { assertEquals( edgeCount, countSaved ); - log.info( "Saved {} source edges", sourceCount ); - log.info( "Saved {} target edges", targetCount ); + logger.info( "Saved {} source edges", sourceCount ); + logger.info( "Saved {} target edges", targetCount ); long deleteVersion = Long.MAX_VALUE; http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/EdgeSerializationTest.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/EdgeSerializationTest.java b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/EdgeSerializationTest.java index d81413e..928d519 100644 --- a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/EdgeSerializationTest.java +++ b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/EdgeSerializationTest.java @@ -27,7 +27,6 @@ import java.util.UUID; import org.apache.usergrid.StressTest; import org.junit.Before; -import org.junit.Ignore; import org.junit.Rule; import org.junit.Test; import org.junit.experimental.categories.Category; @@ -76,7 +75,7 @@ import static org.mockito.Mockito.when; public abstract class EdgeSerializationTest { - private static final Logger log = LoggerFactory.getLogger( EdgeSerializationTest.class ); + private static final Logger logger = LoggerFactory.getLogger( EdgeSerializationTest.class ); @Inject @Rule @@ -719,7 +718,7 @@ public abstract class EdgeSerializationTest { timestamp++; } - log.info( "Flushing edges" ); + logger.info( "Flushing edges" ); batch.execute(); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/TestCount.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/TestCount.java b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/TestCount.java index 0c09e81..07eb9ab 100644 --- a/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/TestCount.java +++ b/stack/corepersistence/graph/src/test/java/org/apache/usergrid/persistence/graph/serialization/TestCount.java @@ -40,7 +40,7 @@ import static org.junit.Assert.assertEquals; */ public class TestCount { - private static final Logger log = LoggerFactory.getLogger( TestCount.class ); + private static final Logger logger = LoggerFactory.getLogger( TestCount.class ); @Test @@ -125,7 +125,7 @@ public class TestCount { final Integer value = values.get( i ); - log.info( "Emitting {}", value ); + logger.info( "Emitting {}", value ); subscriber.onNext( value ); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/corepersistence/queryindex/src/test/java/org/apache/usergrid/persistence/index/impl/EntityIndexTest.java ---------------------------------------------------------------------- diff --git a/stack/corepersistence/queryindex/src/test/java/org/apache/usergrid/persistence/index/impl/EntityIndexTest.java b/stack/corepersistence/queryindex/src/test/java/org/apache/usergrid/persistence/index/impl/EntityIndexTest.java index f26a018..8e790ff 100644 --- a/stack/corepersistence/queryindex/src/test/java/org/apache/usergrid/persistence/index/impl/EntityIndexTest.java +++ b/stack/corepersistence/queryindex/src/test/java/org/apache/usergrid/persistence/index/impl/EntityIndexTest.java @@ -69,7 +69,7 @@ import static org.junit.Assert.fail; @UseModules( { TestIndexModule.class } ) public class EntityIndexTest extends BaseIT { - private static final Logger log = LoggerFactory.getLogger(EntityIndexTest.class); + private static final Logger logger = LoggerFactory.getLogger(EntityIndexTest.class); @Inject public EntityIndexFactory eif; @@ -164,7 +164,7 @@ public class EntityIndexTest extends BaseIT { timer.stop(); assertEquals(2, candidateResults.size()); - log.debug("Query time {}ms", timer.getTime()); + logger.debug("Query time {}ms", timer.getTime()); final CandidateResult candidate1 = candidateResults.get(0); @@ -266,7 +266,7 @@ public class EntityIndexTest extends BaseIT { testQuery( searchEdge, searchTypes, "name = 'Lowe Kelley'", 1 ); - log.info("hi"); + logger.info("hi"); } @@ -309,7 +309,7 @@ public class EntityIndexTest extends BaseIT { indexProducer.put(batch.build()).subscribe();; EntityIndex.IndexRefreshCommandInfo info = entityIndex.refreshAsync().toBlocking().first(); long time = info.getExecutionTime(); - log.info("refresh took ms:" + time); + logger.info("refresh took ms:" + time); } @@ -342,7 +342,7 @@ public class EntityIndexTest extends BaseIT { } timer.stop(); - log.info( "Total time to index {} entries {}ms, average {}ms/entry", + logger.info( "Total time to index {} entries {}ms, average {}ms/entry", new Object[] { count, timer.getTime(), timer.getTime() / count } ); } @@ -435,7 +435,7 @@ public class EntityIndexTest extends BaseIT { timer.stop(); assertEquals(num, candidateResults.size()); - log.debug("Query time {}ms", timer.getTime()); + logger.debug("Query time {}ms", timer.getTime()); return candidateResults; } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/launcher/src/main/java/org/apache/usergrid/launcher/EmbeddedServerHelper.java ---------------------------------------------------------------------- diff --git a/stack/launcher/src/main/java/org/apache/usergrid/launcher/EmbeddedServerHelper.java b/stack/launcher/src/main/java/org/apache/usergrid/launcher/EmbeddedServerHelper.java index b3b72c2..441fe21 100644 --- a/stack/launcher/src/main/java/org/apache/usergrid/launcher/EmbeddedServerHelper.java +++ b/stack/launcher/src/main/java/org/apache/usergrid/launcher/EmbeddedServerHelper.java @@ -19,7 +19,7 @@ package org.apache.usergrid.launcher; /** @author Ran Tavory ([email protected]) */ public class EmbeddedServerHelper { - private static org.slf4j.Logger log = org.slf4j.LoggerFactory.getLogger( EmbeddedServerHelper.class ); + private static org.slf4j.Logger logger = org.slf4j.LoggerFactory.getLogger( EmbeddedServerHelper.class ); private static final String TMP = "tmp"; @@ -66,10 +66,10 @@ public class EmbeddedServerHelper { System.setProperty( "log4j.configuration", "file:" + TMP + "/log4j.properties" ); System.setProperty( "cassandra-foreground", "true" ); - log.info( "Starting executor" ); + logger.info( "Starting executor" ); executor.execute( new ITRunner() ); - log.info( "Started executor" ); + logger.info( "Started executor" ); } else { cassandraDaemon.start(); @@ -78,7 +78,7 @@ public class EmbeddedServerHelper { try { java.util.concurrent.TimeUnit.SECONDS.sleep( 3 ); - log.info( "Done sleeping" ); + logger.info( "Done sleeping" ); } catch ( InterruptedException e ) { throw new AssertionError( e ); @@ -93,7 +93,7 @@ public class EmbeddedServerHelper { } executor.shutdown(); executor.shutdownNow(); - log.info( "Teardown complete" ); + logger.info( "Teardown complete" ); } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/PartialUpdateTest.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/PartialUpdateTest.java b/stack/rest/src/test/java/org/apache/usergrid/rest/PartialUpdateTest.java index 5ce3c2d..1067365 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/PartialUpdateTest.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/PartialUpdateTest.java @@ -36,7 +36,7 @@ import static org.junit.Assert.*; * Partial update test. */ public class PartialUpdateTest extends AbstractRestIT { - private static final Logger log = LoggerFactory.getLogger(PartialUpdateTest.class); + private static final Logger logger = LoggerFactory.getLogger(PartialUpdateTest.class); double latitude = 37.772837; double longitude = -122.409895; @@ -85,7 +85,7 @@ public class PartialUpdateTest extends AbstractRestIT { // retrieve the user from the backend userNode = this.app().collection("users").entity(userNode).get(); - log.info(userNode.toString()); + logger.info(userNode.toString()); // verify that the user was returned assertNotNull(userNode); @@ -98,9 +98,9 @@ public class PartialUpdateTest extends AbstractRestIT { // Verify that the location was updated correctly AND that // it is not the same object reference from the original POST - log.info(geolocation.get("latitude") + " != " + logger.info(geolocation.get("latitude") + " != " + Double.parseDouble(((Map<String, Object>) userNode.get("location")).get("latitude").toString())); - log.info(geolocation.get("longitude") + " != " + logger.info(geolocation.get("longitude") + " != " + Double.parseDouble(((Map<String, Object>) userNode.get("location")).get("longitude").toString())); assertNotSame(geolocation.get("latitude"), Double.parseDouble(((Map<String, Object>) userNode.get("location")).get("latitude").toString())); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/TomcatRuntime.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/TomcatRuntime.java b/stack/rest/src/test/java/org/apache/usergrid/rest/TomcatRuntime.java index 4c09565..e1d2a09 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/TomcatRuntime.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/TomcatRuntime.java @@ -41,7 +41,7 @@ import javax.ws.rs.client.WebTarget; * Start and stop embedded Tomcat. */ public class TomcatRuntime extends ExternalResource { - private static final Logger log = LoggerFactory.getLogger( TomcatRuntime.class ); + private static final Logger logger = LoggerFactory.getLogger( TomcatRuntime.class ); private static final String WEBAPP_PATH = System.getProperty("webapp.directory"); @@ -137,9 +137,9 @@ public class TomcatRuntime extends ExternalResource { tomcat.addWebapp( "/", new File( webAppsPath ).getAbsolutePath() ); - log.info( "-----------------------------------------------------------------" ); - log.info( "Starting Tomcat embedded port {} dir {}", port, dataDir.getAbsolutePath() ); - log.info( "-----------------------------------------------------------------" ); + logger.info( "-----------------------------------------------------------------" ); + logger.info( "Starting Tomcat embedded port {} dir {}", port, dataDir.getAbsolutePath() ); + logger.info( "-----------------------------------------------------------------" ); tomcat.start(); waitForTomcat(); @@ -178,12 +178,12 @@ public class TomcatRuntime extends ExternalResource { Client c = ClientBuilder.newClient(); WebTarget wr = c.target( url ); wr.request().get( String.class ); - log.info( "Tomcat is started." ); + logger.info( "Tomcat is started." ); started = true; break; } catch ( Exception e ) { - log.info( "Waiting for Tomcat on url {}", url ); + logger.info( "Waiting for Tomcat on url {}", url ); } } if ( !started ) { http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AssetResourceIT.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AssetResourceIT.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AssetResourceIT.java index 3e591cb..8d647e7 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AssetResourceIT.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AssetResourceIT.java @@ -25,8 +25,6 @@ import org.apache.usergrid.rest.test.resource.model.Entity; import org.apache.usergrid.services.assets.data.AssetUtils; import org.glassfish.jersey.media.multipart.FormDataMultiPart; import org.glassfish.jersey.media.multipart.MultiPart; -import org.glassfish.jersey.media.multipart.MultiPartFeature; -import org.glassfish.jersey.media.multipart.file.FileDataBodyPart; import org.glassfish.jersey.media.multipart.file.StreamDataBodyPart; import org.junit.After; import org.junit.Assert; @@ -50,7 +48,7 @@ import static org.junit.Assert.*; public class AssetResourceIT extends AbstractRestIT { private String access_token; - private static final Logger LOG = LoggerFactory.getLogger( AssetResourceIT.class ); + private static final Logger logger = LoggerFactory.getLogger( AssetResourceIT.class ); private Map<String, Object> originalProperties; @@ -213,13 +211,13 @@ public class AssetResourceIT extends AbstractRestIT { FormDataMultiPart form = new FormDataMultiPart().field( "file", data, MediaType.MULTIPART_FORM_DATA_TYPE ); ApiResponse postResponse = pathResource( getOrgAppPath( "foos" ) ).post( form ); UUID assetId = postResponse.getEntities().get(0).getUuid(); - LOG.info( "Waiting for upload to finish..." ); + logger.info( "Waiting for upload to finish..." ); Thread.sleep( 2000 ); // check that entire file was uploaded ApiResponse getResponse = pathResource( getOrgAppPath( "foos/" +assetId ) ).get( ApiResponse.class ); - LOG.info( "Upload complete!" ); + logger.info( "Upload complete!" ); InputStream is = pathResource( getOrgAppPath( "foos/" + assetId ) ).getAssetAsStream(); byte[] foundData = IOUtils.toByteArray( is ); assertEquals( data.length, foundData.length ); @@ -252,7 +250,7 @@ public class AssetResourceIT extends AbstractRestIT { UUID assetId = postResponse.getEntities().get(0).getUuid(); String errorMessage = null; - LOG.info( "Waiting for upload to finish..." ); + logger.info( "Waiting for upload to finish..." ); Thread.sleep( 2000 ); // attempt to get asset entity, it should contain error http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AwsAssetResourceIT.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AwsAssetResourceIT.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AwsAssetResourceIT.java index b3f3008..c1e87a6 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AwsAssetResourceIT.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/assets/AwsAssetResourceIT.java @@ -48,7 +48,7 @@ public class AwsAssetResourceIT extends AbstractRestIT { private String access_token; private Map<String, Object> originalProperties; - private static final Logger LOG = LoggerFactory.getLogger( AwsAssetResourceIT.class ); + private static final Logger logger = LoggerFactory.getLogger( AwsAssetResourceIT.class ); /** * Mark tests as ignored if no AWS creds are present @@ -152,13 +152,13 @@ public class AwsAssetResourceIT extends AbstractRestIT { FormDataMultiPart form = new FormDataMultiPart().field( "file", data, MediaType.MULTIPART_FORM_DATA_TYPE ); ApiResponse postResponse = pathResource( getOrgAppPath( "foos" ) ).post( form ); UUID assetId = postResponse.getEntities().get(0).getUuid(); - LOG.info( "Waiting for upload to finish..." ); + logger.info( "Waiting for upload to finish..." ); Thread.sleep( 5000 ); // check that entire file was uploaded ApiResponse getResponse = pathResource( getOrgAppPath( "foos/" +assetId ) ).get( ApiResponse.class ); - LOG.info( "Upload complete!" ); + logger.info( "Upload complete!" ); InputStream is = pathResource( getOrgAppPath( "foos/" + assetId ) ).getAssetAsStream(); byte[] foundData = IOUtils.toByteArray( is ); assertEquals( data.length, foundData.length ); @@ -321,13 +321,13 @@ public class AwsAssetResourceIT extends AbstractRestIT { FormDataMultiPart form = new FormDataMultiPart().field( "file", data, MediaType.MULTIPART_FORM_DATA_TYPE ); ApiResponse postResponse = pathResource( getOrgAppPath( "foos" ) ).post( form ); UUID assetId = postResponse.getEntities().get(0).getUuid(); - LOG.info( "Waiting for upload to finish..." ); + logger.info( "Waiting for upload to finish..." ); Thread.sleep( 5000 ); // check that entire file was uploaded ApiResponse getResponse = pathResource( getOrgAppPath( "foos/" +assetId ) ).get( ApiResponse.class ); - LOG.info( "Upload complete!" ); + logger.info( "Upload complete!" ); InputStream is = pathResource( getOrgAppPath( "foos/" + assetId ) ).getAssetAsStream(); byte[] foundData = IOUtils.toByteArray( is ); assertEquals( data.length, foundData.length ); @@ -355,7 +355,7 @@ public class AwsAssetResourceIT extends AbstractRestIT { UUID assetId = postResponse.getEntities().get(0).getUuid(); String errorMessage = null; - LOG.info( "Waiting for upload to finish..." ); + logger.info( "Waiting for upload to finish..." ); Thread.sleep( 1000 ); // attempt to get asset entity, it should contain error http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/UserResourceIT.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/UserResourceIT.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/UserResourceIT.java index cc9eaa5..d36d0a1 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/UserResourceIT.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/UserResourceIT.java @@ -18,9 +18,6 @@ package org.apache.usergrid.rest.applications.collection.users; import com.fasterxml.jackson.databind.JsonNode; -import com.sun.jersey.api.client.UniformInterfaceException; -import com.sun.jersey.api.client.WebResource; -import com.sun.jersey.api.client.filter.HTTPBasicAuthFilter; import java.util.ArrayList; import java.util.HashMap; @@ -65,7 +62,7 @@ import static org.junit.Assert.fail; public class UserResourceIT extends AbstractRestIT { - private static final Logger log = LoggerFactory.getLogger(UserResourceIT.class); + private static final Logger logger = LoggerFactory.getLogger(UserResourceIT.class); UserRepo userRepo; CollectionEndpoint usersResource; CollectionEndpoint userResource; @@ -700,7 +697,7 @@ public class UserResourceIT extends AbstractRestIT { @Test public void test_POST_batch() throws IOException { - log.info("UserResourceIT.test_POST_batch"); + logger.info("UserResourceIT.test_POST_batch"); List<Entity> batch = new ArrayList<>(); @@ -1014,7 +1011,7 @@ public class UserResourceIT extends AbstractRestIT { assertTrue(false); } catch (ClientErrorException uie) { status = uie.getResponse().getStatus(); - log.info("Error Response Body: " + uie.getResponse().readEntity(String.class)); + logger.info("Error Response Body: " + uie.getResponse().readEntity(String.class)); } assertEquals(Response.Status.UNAUTHORIZED.getStatusCode(), status); @@ -1025,7 +1022,7 @@ public class UserResourceIT extends AbstractRestIT { assertTrue(false); } catch (ClientErrorException uie) { status = uie.getResponse().getStatus(); - log.info("Error Response Body: " + uie.getResponse().readEntity(String.class)); + logger.info("Error Response Body: " + uie.getResponse().readEntity(String.class)); } assertEquals(Response.Status.FORBIDDEN.getStatusCode(), status); @@ -1215,7 +1212,7 @@ public class UserResourceIT extends AbstractRestIT { fail("Should have thrown an exception"); } catch (BadRequestException uie) { status = uie.getResponse().getStatus(); - log.info("Error Response Body: {}" , uie.getResponse().getEntity()); + logger.info("Error Response Body: {}" , uie.getResponse().getEntity()); } assertEquals( Response.Status.BAD_REQUEST.getStatusCode(), status); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/extensions/TestResource.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/extensions/TestResource.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/extensions/TestResource.java index 7e6f278..81af791 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/extensions/TestResource.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/collection/users/extensions/TestResource.java @@ -21,7 +21,6 @@ import javax.ws.rs.GET; import javax.ws.rs.Produces; import javax.ws.rs.core.MediaType; -import org.junit.Ignore; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.context.annotation.Scope; @@ -34,11 +33,11 @@ import org.apache.usergrid.rest.applications.users.AbstractUserExtensionResource @Produces(MediaType.APPLICATION_JSON) public class TestResource extends AbstractUserExtensionResource { - private static final Logger log = LoggerFactory.getLogger( TestResource.class ); + private static final Logger logger = LoggerFactory.getLogger( TestResource.class ); public TestResource() { - log.info( "TestResource" ); + logger.info( "TestResource" ); } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/AndOrQueryTest.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/AndOrQueryTest.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/AndOrQueryTest.java index 16441fe..4bdd3fc 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/AndOrQueryTest.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/AndOrQueryTest.java @@ -36,7 +36,7 @@ import static org.junit.Assert.*; * @since 4.0 */ public class AndOrQueryTest extends QueryTestBase { - private static final Logger log = LoggerFactory.getLogger(AndOrQueryTest.class); + private static final Logger logger = LoggerFactory.getLogger(AndOrQueryTest.class); /** @@ -154,7 +154,7 @@ public class AndOrQueryTest extends QueryTestBase { for (int i = 0; i < returnSize; i++, index--) { count++; Entity activity = activities.getResponse().getEntities().get(i); - log.info(String.valueOf(activity.get("ordinal")) + " " + String.valueOf(activity.get("madeup"))); + logger.info(String.valueOf(activity.get("ordinal")) + " " + String.valueOf(activity.get("madeup"))); //if the entity is in the first half, the property "madeup" should be false if (index < numOfEntities / 2) { assertFalse(Boolean.parseBoolean(String.valueOf(activity.get("madeup")))); @@ -198,7 +198,7 @@ public class AndOrQueryTest extends QueryTestBase { count++; Entity activity = activities.getResponse().getEntities().get(i); long ordinal = Long.parseLong(String.valueOf(activity.get("ordinal"))); - log.info(ordinal + " " + String.valueOf(activity.get("verb"))); + logger.info(ordinal + " " + String.valueOf(activity.get("verb"))); //if the entity is in the first three quarters, the property "verb" should be "go" if (ordinal < (numOfEntities - numOfEntities / 4)) { assertEquals("go", String.valueOf(activity.get("verb"))); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/OrderByTest.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/OrderByTest.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/OrderByTest.java index a3d8391..e052264 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/OrderByTest.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/OrderByTest.java @@ -37,7 +37,7 @@ import static org.junit.Assert.assertEquals; * @since 4.0 */ public class OrderByTest extends QueryTestBase { - private static final Logger log = LoggerFactory.getLogger(OrderByTest.class); + private static final Logger logger = LoggerFactory.getLogger(OrderByTest.class); @Test public void orderByLongAsc() throws IOException { @@ -241,7 +241,7 @@ public class OrderByTest extends QueryTestBase { for (int i = 0; i < 20; i++) { props.put("ordinal", i); Entity activity = this.app().collection("activity").post(props); - log.info("Created", activity.get("created").toString()); + logger.info("Created", activity.get("created").toString()); if (i == 5) { created = Long.parseLong(activity.get("created").toString()); } @@ -331,8 +331,8 @@ public class OrderByTest extends QueryTestBase { props.put("ordinal", i); Entity e = this.app().collection("activity").post(props); activities[i] = e; - log.info(String.valueOf(e.get("uuid").toString())); - log.info(String.valueOf(Long.parseLong(activities[0].get("created").toString()))); + logger.info(String.valueOf(e.get("uuid").toString())); + logger.info(String.valueOf(Long.parseLong(activities[0].get("created").toString()))); } refreshIndex(); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/QueryTestBase.java ---------------------------------------------------------------------- diff --git a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/QueryTestBase.java b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/QueryTestBase.java index 48ddf51..f7eb3fe 100644 --- a/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/QueryTestBase.java +++ b/stack/rest/src/test/java/org/apache/usergrid/rest/applications/queries/QueryTestBase.java @@ -25,7 +25,7 @@ import org.slf4j.LoggerFactory; * A base class containing common methods used by query tests */ public class QueryTestBase extends AbstractRestIT { - private static final Logger log = LoggerFactory.getLogger(QueryTestBase.class); + private static final Logger logger = LoggerFactory.getLogger(QueryTestBase.class); /** * Create a number of entities in the specified collection * with properties to make them independently searchable @@ -62,7 +62,7 @@ public class QueryTestBase extends AbstractRestIT { } //create the entity in the desired collection and add it to the return array entities[i] = this.app().collection(collectionName).post(props); - log.info(entities[i].entrySet().toString()); + logger.info(entities[i].entrySet().toString()); } //refresh the index so that they are immediately searchable this.refreshIndex(); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/services/src/main/java/org/apache/usergrid/services/notifications/JobScheduler.java ---------------------------------------------------------------------- diff --git a/stack/services/src/main/java/org/apache/usergrid/services/notifications/JobScheduler.java b/stack/services/src/main/java/org/apache/usergrid/services/notifications/JobScheduler.java index f2a8976..f9fe8b8 100644 --- a/stack/services/src/main/java/org/apache/usergrid/services/notifications/JobScheduler.java +++ b/stack/services/src/main/java/org/apache/usergrid/services/notifications/JobScheduler.java @@ -24,15 +24,12 @@ import org.apache.usergrid.services.ServiceManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.util.HashMap; -import java.util.Map; - public class JobScheduler{ public static final long SCHEDULER_GRACE_PERIOD = 250; private final EntityManager em; protected ServiceManager sm; - private static final Logger LOG = LoggerFactory.getLogger(NotificationsService.class); + private static final Logger logger = LoggerFactory.getLogger(NotificationsService.class); public JobScheduler(ServiceManager sm,EntityManager em){ this.sm=sm; this.em = em; @@ -49,7 +46,7 @@ public class JobScheduler{ SchedulerService scheduler = getSchedulerService(); scheduler.createJob("notificationBatchJob", soonestPossible, jobData); - LOG.info("notification {} batch scheduled for delivery", notification.getUuid()); + logger.info("notification {} batch scheduled for delivery", notification.getUuid()); } public boolean scheduleQueueJob(Notification notification) throws Exception { return scheduleQueueJob(notification,false); @@ -73,7 +70,7 @@ public class JobScheduler{ jobData.setProperty("deliver", notification.getDeliver()); SchedulerService scheduler = getSchedulerService(); scheduler.createJob("queueJob", scheduleAt, jobData); - LOG.info("notification {} scheduled for queuing", notification.getUuid()); + logger.info("notification {} scheduled for queuing", notification.getUuid()); } return scheduled; } http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/services/src/main/java/org/apache/usergrid/services/notifications/NotificationsService.java ---------------------------------------------------------------------- diff --git a/stack/services/src/main/java/org/apache/usergrid/services/notifications/NotificationsService.java b/stack/services/src/main/java/org/apache/usergrid/services/notifications/NotificationsService.java index 202971d..36c6776 100644 --- a/stack/services/src/main/java/org/apache/usergrid/services/notifications/NotificationsService.java +++ b/stack/services/src/main/java/org/apache/usergrid/services/notifications/NotificationsService.java @@ -68,7 +68,7 @@ public class NotificationsService extends AbstractCollectionService { private Timer postTimer; private static final int PAGE = 100; - private static final Logger LOG = LoggerFactory.getLogger(NotificationsService.class); + private static final Logger logger = LoggerFactory.getLogger(NotificationsService.class); //need a mocking framework, this is to substitute for no mocking static final String MESSAGE_PROPERTY_DEVICE_UUID = "deviceUUID"; @@ -86,7 +86,7 @@ public class NotificationsService extends AbstractCollectionService { private QueueManagerFactory queueManagerFactory; public NotificationsService() { - LOG.info("/notifications"); + logger.info("/notifications"); } @Override @@ -129,7 +129,7 @@ public class NotificationsService extends AbstractCollectionService { @Override public ServiceResults postCollection(ServiceContext context) throws Exception { - LOG.info("NotificationService: start request."); + logger.info("NotificationService: start request."); Timer.Context timer = postTimer.time(); postMeter.mark(); try { @@ -149,16 +149,16 @@ public class NotificationsService extends AbstractCollectionService { properties.put("started", notification.getStarted()); properties.put("state", notification.getState()); notification.addProperties(properties); - LOG.info("ApplicationQueueMessage: notification {} properties updated in duration {} ms", notification.getUuid(), System.currentTimeMillis() - now); + logger.info("ApplicationQueueMessage: notification {} properties updated in duration {} ms", notification.getUuid(), System.currentTimeMillis() - now); } long now = System.currentTimeMillis(); notificationQueueManager.queueNotification(notification, null); - LOG.info("NotificationService: notification {} post queue duration {} ms ", notification.getUuid(), System.currentTimeMillis() - now); + logger.info("NotificationService: notification {} post queue duration {} ms ", notification.getUuid(), System.currentTimeMillis() - now); // future: somehow return 202? return results; }catch (Exception e){ - LOG.error("serialization failed",e); + logger.error("serialization failed",e); throw e; }finally { timer.stop(); http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/services/src/main/java/org/apache/usergrid/services/notifications/QueueListener.java ---------------------------------------------------------------------- diff --git a/stack/services/src/main/java/org/apache/usergrid/services/notifications/QueueListener.java b/stack/services/src/main/java/org/apache/usergrid/services/notifications/QueueListener.java index bcf7b49..e7d36e2 100644 --- a/stack/services/src/main/java/org/apache/usergrid/services/notifications/QueueListener.java +++ b/stack/services/src/main/java/org/apache/usergrid/services/notifications/QueueListener.java @@ -50,7 +50,7 @@ public class QueueListener { public long DEFAULT_SLEEP = 5000; - private static final Logger LOG = LoggerFactory.getLogger(QueueListener.class); + private static final Logger logger = LoggerFactory.getLogger(QueueListener.class); private MetricsFactory metricsService; @@ -93,7 +93,7 @@ public class QueueListener { // boolean shouldRun = new Boolean(properties.getProperty("usergrid.notifications.listener.run", "false")); - LOG.info("QueueListener: starting."); + logger.info("QueueListener: starting."); int threadCount = 0; try { @@ -113,23 +113,23 @@ public class QueueListener { pool = Executors.newFixedThreadPool(maxThreads); while (threadCount++ < maxThreads) { - LOG.info("QueueListener: Starting thread {}.", threadCount); + logger.info("QueueListener: Starting thread {}.", threadCount); Runnable task = new Runnable() { @Override public void run() { try { execute(); } catch (Exception e) { - LOG.error("failed to start push", e); + logger.error("failed to start push", e); } } }; futures.add( pool.submit(task)); } } catch (Exception e) { - LOG.error("QueueListener: failed to start:", e); + logger.error("QueueListener: failed to start:", e); } - LOG.info("QueueListener: done starting."); + logger.info("QueueListener: done starting."); } private void execute(){ @@ -139,11 +139,11 @@ public class QueueListener { Thread.currentThread().setName("Notifications_Processor"+UUID.randomUUID()); final AtomicInteger consecutiveExceptions = new AtomicInteger(); - LOG.info("QueueListener: Starting execute process."); + logger.info("QueueListener: Starting execute process."); Meter meter = metricsService.getMeter(QueueListener.class, "execute.commit"); com.codahale.metrics.Timer timer = metricsService.getTimer(QueueListener.class, "execute.dequeue"); svcMgr = smf.getServiceManager(smf.getManagementAppId()); - LOG.info("getting from queue {} ", queueName); + logger.info("getting from queue {} ", queueName); QueueScope queueScope = new QueueScopeImpl( queueName, QueueScope.RegionImplementation.LOCAL); QueueManager queueManager = TEST_QUEUE_MANAGER != null ? TEST_QUEUE_MANAGER : queueManagerFactory.getQueueManager(queueScope); // run until there are no more active jobs @@ -157,16 +157,19 @@ public class QueueListener { rx.Observable.from(queueManager.getMessages(getBatchSize(), MESSAGE_TRANSACTION_TIMEOUT, 10000, ApplicationQueueMessage.class)) .buffer(getBatchSize()) .doOnNext(messages -> { + try { - LOG.info("retrieved batch of {} messages from queue {} ", messages.size(),queueName); + logger.info("retrieved batch of {} messages from queue {} ", messages.size(),queueName); if (messages.size() > 0) { HashMap<UUID, List<QueueMessage>> messageMap = new HashMap<>(messages.size()); + //group messages into hash map by app id for (QueueMessage message : messages) { //TODO: stop copying around this area as it gets notification specific. ApplicationQueueMessage queueMessage = (ApplicationQueueMessage) message.getBody(); UUID applicationId = queueMessage.getApplicationId(); + //Groups queue messages by application Id, ( they are all probably going to the same place ) if (!messageMap.containsKey(applicationId)) { //For each app id it sends the set. @@ -177,61 +180,67 @@ public class QueueListener { messageMap.get(applicationId).add(message); } } + long now = System.currentTimeMillis(); Observable merge = null; + //send each set of app ids together for (Map.Entry<UUID, List<QueueMessage>> entry : messageMap.entrySet()) { UUID applicationId = entry.getKey(); ApplicationQueueManager manager = queueManagerMap.get(applicationId); - LOG.info("send batch for app {} of {} messages", entry.getKey(), entry.getValue().size()); + logger.info("send batch for app {} of {} messages", entry.getKey(), entry.getValue().size()); Observable current = manager.sendBatchToProviders(entry.getValue(),queueName); + if(merge == null) merge = current; else { merge = Observable.merge(merge,current); } } + if(merge!=null) { merge.toBlocking().lastOrDefault(null); } queueManager.commitMessages(messages); meter.mark(messages.size()); - LOG.info("sent batch {} messages duration {} ms", messages.size(),System.currentTimeMillis() - now); + logger.info("sent batch {} messages duration {} ms", messages.size(),System.currentTimeMillis() - now); if(sleepBetweenRuns > 0) { - LOG.info("sleep between rounds...sleep...{}", sleepBetweenRuns); + logger.info("sleep between rounds...sleep...{}", sleepBetweenRuns); Thread.sleep(sleepBetweenRuns); } + if(runCount.incrementAndGet() % consecutiveCallsToRemoveDevices == 0){ for(ApplicationQueueManager applicationQueueManager : queueManagerMap.asMap().values()){ try { applicationQueueManager.asyncCheckForInactiveDevices(); }catch (Exception inactiveDeviceException){ - LOG.error("Inactive Device Get failed",inactiveDeviceException); + logger.error("Inactive Device Get failed",inactiveDeviceException); } } //clear everything runCount.set(0); } } + else{ - LOG.info("no messages...sleep...{}", sleepWhenNoneFound); + logger.info("no messages...sleep...{}", sleepWhenNoneFound); Thread.sleep(sleepWhenNoneFound); } timerContext.stop(); //send to the providers consecutiveExceptions.set(0); }catch (Exception ex){ - LOG.error("failed to dequeue",ex); + logger.error("failed to dequeue",ex); try { long sleeptime = sleepWhenNoneFound*consecutiveExceptions.incrementAndGet(); long maxSleep = 15000; sleeptime = sleeptime > maxSleep ? maxSleep : sleeptime ; - LOG.info("sleeping due to failures {} ms", sleeptime); + logger.info("sleeping due to failures {} ms", sleeptime); Thread.sleep(sleeptime); }catch (InterruptedException ie){ - LOG.info("sleep interrupted"); + logger.info("sleep interrupted"); } } }) @@ -252,7 +261,7 @@ public class QueueListener { try { queueManagerNotifiication.getValue().stop(); } catch (Exception ie) { - LOG.error("Failed to shutdown from cache", ie); + logger.error("Failed to shutdown from cache", ie); } } }).build(new CacheLoader<UUID, ApplicationQueueManager>() { @@ -271,7 +280,7 @@ public class QueueListener { ); return manager; } catch (Exception e) { - LOG.error("Could not instantiate queue manager", e); + logger.error("Could not instantiate queue manager", e); return null; } } @@ -279,7 +288,7 @@ public class QueueListener { } public void stop(){ - LOG.info("stop processes"); + logger.info("stop processes"); if(futures == null){ return; http://git-wip-us.apache.org/repos/asf/usergrid/blob/bc33c88d/stack/services/src/main/java/org/apache/usergrid/services/notifications/TaskManager.java ---------------------------------------------------------------------- diff --git a/stack/services/src/main/java/org/apache/usergrid/services/notifications/TaskManager.java b/stack/services/src/main/java/org/apache/usergrid/services/notifications/TaskManager.java index a74cc14..6c7e53c 100644 --- a/stack/services/src/main/java/org/apache/usergrid/services/notifications/TaskManager.java +++ b/stack/services/src/main/java/org/apache/usergrid/services/notifications/TaskManager.java @@ -31,7 +31,7 @@ import java.util.concurrent.atomic.AtomicLong; public class TaskManager { - private static final Logger LOG = LoggerFactory.getLogger(TaskManager.class); + private static final Logger logger = LoggerFactory.getLogger(TaskManager.class); private Notification notification; private AtomicLong successes = new AtomicLong(); @@ -53,41 +53,41 @@ public class TaskManager { completed(notifier,null,deviceUUID,null); } public void completed(Notifier notifier, Receipt receipt, UUID deviceUUID, String newProviderId) throws Exception { - if (LOG.isDebugEnabled()) { - LOG.debug("REMOVED {}", deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("REMOVED {}", deviceUUID); } try { - if (LOG.isDebugEnabled()) { - LOG.debug("notification {} removing device {} from remaining", notification.getUuid(), deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("notification {} removing device {} from remaining", notification.getUuid(), deviceUUID); } EntityRef deviceRef = new SimpleEntityRef(Device.ENTITY_TYPE, deviceUUID); if (receipt != null) { - if (LOG.isDebugEnabled()) { - LOG.debug("notification {} sent to device {}. saving receipt.", notification.getUuid(), deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("notification {} sent to device {}. saving receipt.", notification.getUuid(), deviceUUID); } receipt.setSent(System.currentTimeMillis()); this.saveReceipt(notification, deviceRef, receipt,false); - if (LOG.isDebugEnabled()) { - LOG.debug("notification {} receipt saved for device {}", notification.getUuid(), deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("notification {} receipt saved for device {}", notification.getUuid(), deviceUUID); } successes.incrementAndGet(); } if (newProviderId != null) { - if (LOG.isDebugEnabled()) { - LOG.debug("notification {} replacing device {} notifierId", notification.getUuid(), deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("notification {} replacing device {} notifierId", notification.getUuid(), deviceUUID); } replaceProviderId(deviceRef, notifier, newProviderId); } - if (LOG.isDebugEnabled()) { - LOG.debug("notification {} completed device {}", notification.getUuid(), deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("notification {} completed device {}", notification.getUuid(), deviceUUID); } } finally { - if (LOG.isDebugEnabled()) { - LOG.debug("COUNT is: {}", successes.get()); + if (logger.isDebugEnabled()) { + logger.debug("COUNT is: {}", successes.get()); } if (hasFinished) { //process has finished but notifications are still coming in finishedBatch(); @@ -99,12 +99,12 @@ public class TaskManager { public void failed(Notifier notifier, Receipt receipt, UUID deviceUUID, Object code, String message) throws Exception { try { - if (LOG.isDebugEnabled()) { + if (logger.isDebugEnabled()) { StringBuilder sb = new StringBuilder(); sb.append("notification ").append(notification.getUuid()); sb.append(" for device ").append(deviceUUID); sb.append(" got error ").append(code); - LOG.debug(sb.toString()); + logger.debug(sb.toString()); } failures.incrementAndGet(); @@ -114,8 +114,8 @@ public class TaskManager { receipt.setErrorCode(code); receipt.setErrorMessage(message); this.saveReceipt(notification, new SimpleEntityRef(Device.ENTITY_TYPE, deviceUUID), receipt,true); - if (LOG.isDebugEnabled()) { - LOG.debug("notification {} receipt saved for device {}", notification.getUuid(), deviceUUID); + if (logger.isDebugEnabled()) { + logger.debug("notification {} receipt saved for device {}", notification.getUuid(), deviceUUID); } } finally { completed(notifier, deviceUUID); @@ -158,12 +158,15 @@ public class TaskManager { finishedBatch(true,false); } public void finishedBatch(boolean fetch, boolean force) throws Exception { + if (notification.getDebug() || getFailures() > 0 || force) { long successes = this.successes.get(); //reset counters long failures = this.failures.get(); //reset counters + for (int i = 0; i < successes; i++) { this.successes.decrementAndGet(); } + for (int i = 0; i < failures; i++) { this.failures.decrementAndGet(); } @@ -191,7 +194,7 @@ public class TaskManager { notification.addProperties(properties); long latency = notification.getFinished() - notification.getStarted(); - LOG.info("notification finished batch: {} of {} devices in " + latency + "ms", notification.getUuid(), totals); + logger.info("notification finished batch: {} of {} devices in " + latency + "ms", notification.getUuid(), totals); em.update(notification); // Set<Notifier> notifiers = new HashSet<>(proxy.getAdapterMap().values()); // remove dups
