Skip to content

Commit 2122d66

Browse files
Apply Palantir Java Format
1 parent 6fdee02 commit 2122d66

5 files changed

Lines changed: 32 additions & 41 deletions

File tree

rqueue-core/src/main/java/com/github/sonus21/rqueue/core/impl/BaseMessageSender.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -220,8 +220,7 @@ protected Object deleteAllMessages(QueueDetail queueDetail) {
220220
MessageDeleteRequest.builder().queueDetail(queueDetail).build());
221221
}
222222

223-
protected void registerQueueInternal(String queueName, QueueType type,
224-
String... priorities) {
223+
protected void registerQueueInternal(String queueName, QueueType type, String... priorities) {
225224
validateQueue(queueName);
226225
notNull(priorities, "priorities cannot be null");
227226
Map<String, Integer> priorityMap = new HashMap<>();

rqueue-nats/src/main/java/com/github/sonus21/rqueue/nats/internal/NatsProvisioner.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -153,9 +153,8 @@ public void ensureStream(String streamName, List<String> subjects, QueueType que
153153
}
154154
try {
155155
StreamInfo existing = safeGetStreamInfo(streamName);
156-
RetentionPolicy desired = queueType == QueueType.STREAM
157-
? RetentionPolicy.Limits
158-
: RetentionPolicy.WorkQueue;
156+
RetentionPolicy desired =
157+
queueType == QueueType.STREAM ? RetentionPolicy.Limits : RetentionPolicy.WorkQueue;
159158
if (existing == null) {
160159
if (!config.isAutoCreateStreams()) {
161160
throw new RqueueNatsException(

rqueue-nats/src/main/java/com/github/sonus21/rqueue/nats/js/NatsStreamValidator.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -163,8 +163,7 @@ private void tryEnsureConsumer(
163163
}
164164
}
165165

166-
private int tryEnsure(
167-
List<String> failures, String streamName, String subject, QueueDetail q) {
166+
private int tryEnsure(List<String> failures, String streamName, String subject, QueueDetail q) {
168167
try {
169168
provisioner.ensureStream(streamName, List.of(subject), q.getType());
170169
return 1;

rqueue-nats/src/test/java/com/github/sonus21/rqueue/nats/JetStreamMessageBrokerProducerOnlyIT.java

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,8 @@ void enqueue_messagesAccumulateInStream() throws Exception {
3333
JetStreamMessageBroker.builder().connection(connection).build()) {
3434
int count = 10;
3535
for (int i = 0; i < count; i++) {
36-
broker.enqueue(q, RqueueMessage.builder().id("m-" + i).message("payload-" + i).build());
36+
broker.enqueue(
37+
q, RqueueMessage.builder().id("m-" + i).message("payload-" + i).build());
3738
}
3839
assertEquals(count, broker.size(q), "all enqueued messages should be visible in the stream");
3940
}
@@ -76,13 +77,18 @@ void enqueueReactive_messagesAccumulateInStream() {
7677
try (JetStreamMessageBroker broker =
7778
JetStreamMessageBroker.builder().connection(connection).build()) {
7879
int count = 8;
79-
Flux<Void> publishes = Flux.range(0, count).flatMap(i ->
80-
broker.enqueueReactive(
81-
q, RqueueMessage.builder().id("rm-" + i).message("reactive-payload-" + i).build()));
80+
Flux<Void> publishes = Flux.range(0, count)
81+
.flatMap(i -> broker.enqueueReactive(
82+
q,
83+
RqueueMessage.builder()
84+
.id("rm-" + i)
85+
.message("reactive-payload-" + i)
86+
.build()));
8287

8388
StepVerifier.create(publishes).verifyComplete();
8489

85-
assertEquals(count, broker.size(q), "all reactively enqueued messages should be in the stream");
90+
assertEquals(
91+
count, broker.size(q), "all reactively enqueued messages should be in the stream");
8692
}
8793
}
8894

@@ -97,7 +103,8 @@ void mixedEnqueue_allVariantsLandInCorrectStreams() {
97103

98104
// 3 plain messages on the main queue
99105
for (int i = 0; i < 3; i++) {
100-
broker.enqueue(mainQ, RqueueMessage.builder().id("plain-" + i).message("p" + i).build());
106+
broker.enqueue(
107+
mainQ, RqueueMessage.builder().id("plain-" + i).message("p" + i).build());
101108
}
102109
// 2 priority messages on the "high" sub-queue
103110
for (int i = 0; i < 2; i++) {
@@ -107,9 +114,8 @@ void mixedEnqueue_allVariantsLandInCorrectStreams() {
107114
RqueueMessage.builder().id("high-" + i).message("h" + i).build());
108115
}
109116
// 1 reactive message on the main queue
110-
StepVerifier.create(
111-
broker.enqueueReactive(
112-
mainQ, RqueueMessage.builder().id("react-0").message("r0").build()))
117+
StepVerifier.create(broker.enqueueReactive(
118+
mainQ, RqueueMessage.builder().id("react-0").message("r0").build()))
113119
.verifyComplete();
114120

115121
assertEquals(4L, broker.size(mainQ), "main stream: 3 plain + 1 reactive");

rqueue-nats/src/test/java/com/github/sonus21/rqueue/nats/JetStreamQueueModeIT.java

Lines changed: 13 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -55,35 +55,27 @@ void queueMode_queue_createsWorkQueueStream() throws Exception {
5555
String streamName = "rqueue-" + "qm-queue-" + System.nanoTime();
5656
String subject = "rqueue." + "qm-queue-" + System.nanoTime();
5757
JetStreamManagement jsm = connection.jetStreamManagement();
58-
NatsProvisioner provisioner =
59-
new NatsProvisioner(connection, jsm, RqueueNatsConfig.defaults());
58+
NatsProvisioner provisioner = new NatsProvisioner(connection, jsm, RqueueNatsConfig.defaults());
6059

6160
provisioner.ensureStream(streamName, List.of(subject), QueueType.QUEUE);
6261

63-
RetentionPolicy actual =
64-
jsm.getStreamInfo(streamName).getConfiguration().getRetentionPolicy();
62+
RetentionPolicy actual = jsm.getStreamInfo(streamName).getConfiguration().getRetentionPolicy();
6563
assertEquals(
66-
RetentionPolicy.WorkQueue,
67-
actual,
68-
"QUEUE mode must create a WorkQueue-retention stream");
64+
RetentionPolicy.WorkQueue, actual, "QUEUE mode must create a WorkQueue-retention stream");
6965
}
7066

7167
@Test
7268
void queueMode_stream_createsLimitsStream() throws Exception {
7369
String streamName = "rqueue-" + "qm-stream-" + System.nanoTime();
7470
String subject = "rqueue." + "qm-stream-" + System.nanoTime();
7571
JetStreamManagement jsm = connection.jetStreamManagement();
76-
NatsProvisioner provisioner =
77-
new NatsProvisioner(connection, jsm, RqueueNatsConfig.defaults());
72+
NatsProvisioner provisioner = new NatsProvisioner(connection, jsm, RqueueNatsConfig.defaults());
7873

7974
provisioner.ensureStream(streamName, List.of(subject), QueueType.STREAM);
8075

81-
RetentionPolicy actual =
82-
jsm.getStreamInfo(streamName).getConfiguration().getRetentionPolicy();
76+
RetentionPolicy actual = jsm.getStreamInfo(streamName).getConfiguration().getRetentionPolicy();
8377
assertEquals(
84-
RetentionPolicy.Limits,
85-
actual,
86-
"STREAM mode must create a Limits-retention stream");
78+
RetentionPolicy.Limits, actual, "STREAM mode must create a Limits-retention stream");
8779
}
8880

8981
// ---- Contract 2: consumer reuse preserves delivery position -----------
@@ -145,7 +137,8 @@ void queueMode_consumerReuse_preservesDeliveryPosition() throws Exception {
145137
cd.getMaxAckPending());
146138

147139
// Verify the consumer info still reflects the already-delivered messages.
148-
ConsumerInfo info = jsm.getConsumerInfo(RqueueNatsConfig.defaults().getStreamPrefix() + q.getName(), consumerName);
140+
ConsumerInfo info = jsm.getConsumerInfo(
141+
RqueueNatsConfig.defaults().getStreamPrefix() + q.getName(), consumerName);
149142
long numAcked = info.getNumAckPending() == 0
150143
? total - info.getNumPending()
151144
: total - info.getNumPending() - info.getNumAckPending();
@@ -154,8 +147,8 @@ void queueMode_consumerReuse_preservesDeliveryPosition() throws Exception {
154147
assertEquals(
155148
total - firstBatch,
156149
remaining,
157-
"consumer position must be preserved across ensureConsumer calls; "
158-
+ "remaining=" + remaining + " but expected " + (total - firstBatch));
150+
"consumer position must be preserved across ensureConsumer calls; " + "remaining="
151+
+ remaining + " but expected " + (total - firstBatch));
159152
}
160153
}
161154

@@ -179,8 +172,7 @@ void queueMode_queue_competingConsumers_eachMessageDeliveredOnce() throws Except
179172
for (int t = 0; t < 2; t++) {
180173
pool.submit(() -> {
181174
for (int round = 0; round < 100 && done.getCount() > 0; round++) {
182-
List<RqueueMessage> msgs =
183-
broker.pop(q, sharedConsumer, 5, Duration.ofMillis(300));
175+
List<RqueueMessage> msgs = broker.pop(q, sharedConsumer, 5, Duration.ofMillis(300));
184176
for (RqueueMessage m : msgs) {
185177
if (seen.add(m.getId())) {
186178
done.countDown();
@@ -194,9 +186,7 @@ void queueMode_queue_competingConsumers_eachMessageDeliveredOnce() throws Except
194186
pool.shutdownNow();
195187

196188
assertEquals(
197-
total,
198-
seen.size(),
199-
"QUEUE mode: each message must be delivered to exactly one worker");
189+
total, seen.size(), "QUEUE mode: each message must be delivered to exactly one worker");
200190
}
201191
}
202192

@@ -219,9 +209,7 @@ void queueMode_stream_fanOut_everyConsumerReceivesAllMessages() throws Exception
219209
Set<String> listenerTwoSeen = drain(broker, q, "listener-svc-2", total);
220210

221211
assertEquals(
222-
total,
223-
listenerOneSeen.size(),
224-
"STREAM mode: listener-svc-1 must receive all messages");
212+
total, listenerOneSeen.size(), "STREAM mode: listener-svc-1 must receive all messages");
225213
assertEquals(
226214
total,
227215
listenerTwoSeen.size(),

0 commit comments

Comments
 (0)