Skip to content

Commit 6229e7a

Browse files
Apply Palantir Java Format
1 parent 9b1939d commit 6229e7a

6 files changed

Lines changed: 48 additions & 41 deletions

File tree

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

Lines changed: 19 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,8 @@ void enqueue_publishesToPrefixedSubject() throws Exception {
7373
Fixture f = newFixture(RqueueNatsConfig.defaults());
7474
when(f.js.publish(any(String.class), any(Headers.class), any(byte[].class)))
7575
.thenReturn(mock(PublishAck.class));
76-
f.broker.enqueue(queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build());
76+
f.broker.enqueue(
77+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build());
7778
verify(f.js, times(1)).publish(eq("rqueue.orders"), any(Headers.class), any(byte[].class));
7879
}
7980

@@ -83,7 +84,9 @@ void enqueueWithPriority_appendsPrioritySuffixToSubject() throws Exception {
8384
when(f.js.publish(any(String.class), any(Headers.class), any(byte[].class)))
8485
.thenReturn(mock(PublishAck.class));
8586
f.broker.enqueue(
86-
queueNamed("orders"), "high", RqueueMessage.builder().id("m1").message("hi").build());
87+
queueNamed("orders"),
88+
"high",
89+
RqueueMessage.builder().id("m1").message("hi").build());
8790
verify(f.js, times(1)).publish(eq("rqueue.orders.high"), any(Headers.class), any(byte[].class));
8891
}
8992

@@ -103,7 +106,8 @@ void enqueue_honorsCustomSubjectPrefix() throws Exception {
103106
Fixture f = newFixture(cfg);
104107
when(f.js.publish(any(String.class), any(Headers.class), any(byte[].class)))
105108
.thenReturn(mock(PublishAck.class));
106-
f.broker.enqueue(queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build());
109+
f.broker.enqueue(
110+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build());
107111
verify(f.js, times(1)).publish(eq("custom.orders"), any(Headers.class), any(byte[].class));
108112
}
109113

@@ -112,12 +116,10 @@ void enqueue_wrapsIoExceptionInRqueueNatsException() throws Exception {
112116
Fixture f = newFixture(RqueueNatsConfig.defaults());
113117
when(f.js.publish(any(String.class), any(Headers.class), any(byte[].class)))
114118
.thenThrow(new IOException("boom"));
115-
RqueueNatsException ex =
116-
assertThrows(
117-
RqueueNatsException.class,
118-
() ->
119-
f.broker.enqueue(
120-
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()));
119+
RqueueNatsException ex = assertThrows(
120+
RqueueNatsException.class,
121+
() -> f.broker.enqueue(
122+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()));
121123
assertNotNull(ex.getCause());
122124
}
123125

@@ -128,9 +130,8 @@ void enqueue_wrapsJetStreamApiExceptionInRqueueNatsException() throws Exception
128130
.thenThrow(mock(JetStreamApiException.class));
129131
assertThrows(
130132
RqueueNatsException.class,
131-
() ->
132-
f.broker.enqueue(
133-
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()));
133+
() -> f.broker.enqueue(
134+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()));
134135
}
135136

136137
@Test
@@ -164,9 +165,8 @@ void enqueueReactive_completesWhenPublishFutureCompletes() {
164165
when(f.js.publishAsync(any(String.class), any(Headers.class), any(byte[].class)))
165166
.thenReturn(done);
166167

167-
StepVerifier.create(
168-
f.broker.enqueueReactive(
169-
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()))
168+
StepVerifier.create(f.broker.enqueueReactive(
169+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()))
170170
.verifyComplete();
171171
verify(f.js, times(1)).publishAsync(eq("rqueue.orders"), any(Headers.class), any(byte[].class));
172172
}
@@ -179,21 +179,17 @@ void enqueueReactive_wrapsAsyncFailureInRqueueNatsException() {
179179
when(f.js.publishAsync(any(String.class), any(Headers.class), any(byte[].class)))
180180
.thenReturn(failed);
181181

182-
StepVerifier.create(
183-
f.broker.enqueueReactive(
184-
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()))
182+
StepVerifier.create(f.broker.enqueueReactive(
183+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build()))
185184
.expectError(RqueueNatsException.class)
186185
.verify();
187186
}
188187

189188
@Test
190189
void enqueueWithDelayReactive_returnsErrorMonoOfUOE() {
191190
Fixture f = newFixture(RqueueNatsConfig.defaults());
192-
StepVerifier.create(
193-
f.broker.enqueueWithDelayReactive(
194-
queueNamed("orders"),
195-
RqueueMessage.builder().id("m1").message("hi").build(),
196-
100))
191+
StepVerifier.create(f.broker.enqueueWithDelayReactive(
192+
queueNamed("orders"), RqueueMessage.builder().id("m1").message("hi").build(), 100))
197193
.expectError(UnsupportedOperationException.class)
198194
.verify();
199195
}

rqueue-spring-boot-starter/src/test/java/com/github/sonus21/rqueue/spring/boot/integration/NatsConcurrencyE2EIT.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,9 +44,11 @@
4444
@Tag("nats")
4545
class NatsConcurrencyE2EIT extends AbstractNatsBootIT {
4646

47-
@Autowired RqueueMessageEnqueuer enqueuer;
47+
@Autowired
48+
RqueueMessageEnqueuer enqueuer;
4849

49-
@Autowired ConcurrencyListener listener;
50+
@Autowired
51+
ConcurrencyListener listener;
5052

5153
@Test
5254
void parallelInvocationsAreObserved() throws Exception {

rqueue-spring-boot-starter/src/test/java/com/github/sonus21/rqueue/spring/boot/integration/NatsConsumerNameOverrideE2EIT.java

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -45,11 +45,14 @@
4545
@Tag("nats")
4646
class NatsConsumerNameOverrideE2EIT extends AbstractNatsBootIT {
4747

48-
@Autowired RqueueMessageEnqueuer enqueuer;
48+
@Autowired
49+
RqueueMessageEnqueuer enqueuer;
4950

50-
@Autowired CustomConsumerListener listener;
51+
@Autowired
52+
CustomConsumerListener listener;
5153

52-
@Autowired JetStreamManagement jsm;
54+
@Autowired
55+
JetStreamManagement jsm;
5356

5457
@Test
5558
void overriddenConsumerNameIsRegisteredOnTheStream() throws Exception {

rqueue-spring-boot-starter/src/test/java/com/github/sonus21/rqueue/spring/boot/integration/NatsMultipleListenersOnSameQueueE2EIT.java

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -47,16 +47,18 @@
4747
classes = NatsMultipleListenersOnSameQueueE2EIT.TestApp.class,
4848
properties = {"rqueue.backend=nats"})
4949
@Tag("nats")
50-
@Disabled(
51-
"Default JetStream retention=WorkQueue prevents true fan-out across multiple consumers; "
52-
+ "enable once retention is configurable per queue or defaulted to Limits/Interest.")
50+
@Disabled("Default JetStream retention=WorkQueue prevents true fan-out across multiple consumers; "
51+
+ "enable once retention is configurable per queue or defaulted to Limits/Interest.")
5352
class NatsMultipleListenersOnSameQueueE2EIT extends AbstractNatsBootIT {
5453

55-
@Autowired RqueueMessageEnqueuer enqueuer;
54+
@Autowired
55+
RqueueMessageEnqueuer enqueuer;
5656

57-
@Autowired ListenerOne one;
57+
@Autowired
58+
ListenerOne one;
5859

59-
@Autowired ListenerTwo two;
60+
@Autowired
61+
ListenerTwo two;
6062

6163
@Test
6264
void bothListenersReceiveAllMessages() throws Exception {

rqueue-spring-boot-starter/src/test/java/com/github/sonus21/rqueue/spring/boot/integration/NatsPriorityQueuesE2EIT.java

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -46,9 +46,11 @@
4646
@Tag("nats")
4747
class NatsPriorityQueuesE2EIT extends AbstractNatsBootIT {
4848

49-
@Autowired RqueueMessageEnqueuer enqueuer;
49+
@Autowired
50+
RqueueMessageEnqueuer enqueuer;
5051

51-
@Autowired PriorityListener listener;
52+
@Autowired
53+
PriorityListener listener;
5254

5355
@Test
5456
void messagesEnqueuedAtBothPrioritiesAreReceived() throws Exception {
@@ -58,7 +60,8 @@ void messagesEnqueuedAtBothPrioritiesAreReceived() throws Exception {
5860
}
5961
assertThat(listener.latch.await(30, TimeUnit.SECONDS)).isTrue();
6062

61-
long highCount = listener.received.stream().filter(s -> s.startsWith("high-")).count();
63+
long highCount =
64+
listener.received.stream().filter(s -> s.startsWith("high-")).count();
6265
long lowCount = listener.received.stream().filter(s -> s.startsWith("low-")).count();
6366
assertThat(highCount).isEqualTo(5);
6467
assertThat(lowCount).isEqualTo(5);

rqueue-spring-boot-starter/src/test/java/com/github/sonus21/rqueue/spring/boot/integration/NatsReactiveEnqueueE2EIT.java

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -46,9 +46,11 @@
4646
@Tag("nats")
4747
class NatsReactiveEnqueueE2EIT extends AbstractNatsBootIT {
4848

49-
@Autowired ReactiveRqueueMessageEnqueuer reactiveEnqueuer;
49+
@Autowired
50+
ReactiveRqueueMessageEnqueuer reactiveEnqueuer;
5051

51-
@Autowired ReactiveListener listener;
52+
@Autowired
53+
ReactiveListener listener;
5254

5355
@Test
5456
void reactivelyEnqueuedMessagesAreReceivedByListener() throws Exception {
@@ -60,8 +62,7 @@ void reactivelyEnqueuedMessagesAreReceivedByListener() throws Exception {
6062
assertThat(ids).hasSize(5);
6163

6264
assertThat(listener.latch.await(20, TimeUnit.SECONDS)).isTrue();
63-
assertThat(listener.received)
64-
.containsExactlyInAnyOrder("rx-0", "rx-1", "rx-2", "rx-3", "rx-4");
65+
assertThat(listener.received).containsExactlyInAnyOrder("rx-0", "rx-1", "rx-2", "rx-3", "rx-4");
6566
}
6667

6768
@SpringBootApplication(

0 commit comments

Comments
 (0)