Skip to content

Commit 736963a

Browse files
committed
Add pendingResponse to ServerMetrics
1 parent 19068f6 commit 736963a

4 files changed

Lines changed: 64 additions & 23 deletions

File tree

core/src/main/java/com/linecorp/armeria/server/GracefulShutdownSupport.java

Lines changed: 33 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919
import java.time.Duration;
2020
import java.util.concurrent.Executor;
2121
import java.util.concurrent.ThreadPoolExecutor;
22-
import java.util.concurrent.atomic.LongAdder;
2322

2423
import com.google.common.base.Ticker;
2524

@@ -29,39 +28,48 @@
2928
*/
3029
abstract class GracefulShutdownSupport {
3130

32-
static GracefulShutdownSupport create(Duration quietPeriod, Executor blockingTaskExecutor) {
33-
return create(quietPeriod, blockingTaskExecutor, Ticker.systemTicker());
31+
private final ServerMetrics serverMetrics;
32+
33+
/**
34+
* Creates a new instance.
35+
*/
36+
GracefulShutdownSupport(ServerMetrics serverMetrics) {
37+
this.serverMetrics = serverMetrics;
3438
}
3539

36-
static GracefulShutdownSupport create(Duration quietPeriod, Executor blockingTaskExecutor, Ticker ticker) {
37-
return new DefaultGracefulShutdownSupport(quietPeriod, blockingTaskExecutor, ticker);
40+
static GracefulShutdownSupport create(Duration quietPeriod, Executor blockingTaskExecutor,
41+
ServerMetrics serverMetrics) {
42+
return create(quietPeriod, blockingTaskExecutor, Ticker.systemTicker(), serverMetrics);
3843
}
3944

40-
static GracefulShutdownSupport createDisabled() {
41-
return new DisabledGracefulShutdownSupport();
45+
static GracefulShutdownSupport create(Duration quietPeriod, Executor blockingTaskExecutor, Ticker ticker,
46+
ServerMetrics serverMetrics) {
47+
return new DefaultGracefulShutdownSupport(quietPeriod, blockingTaskExecutor, ticker, serverMetrics);
4248
}
4349

44-
private final LongAdder pendingResponses = new LongAdder();
50+
static GracefulShutdownSupport createDisabled(ServerMetrics serverMetrics) {
51+
return new DisabledGracefulShutdownSupport(serverMetrics);
52+
}
4553

4654
/**
4755
* Increases the number of pending responses.
4856
*/
4957
final void inc() {
50-
pendingResponses.increment();
58+
serverMetrics.increaseNonTransientRequests();
5159
}
5260

5361
/**
5462
* Decreases the number of pending responses.
5563
*/
5664
void dec() {
57-
pendingResponses.decrement();
65+
serverMetrics.decreaseNonTransientRequests();
5866
}
5967

6068
/**
6169
* Returns the number of pending responses.
6270
*/
63-
final long pendingResponses() {
64-
return pendingResponses.sum();
71+
final long activeNonTransientResponses() {
72+
return serverMetrics.activeNonTransientRequests();
6573
}
6674

6775
/**
@@ -78,6 +86,14 @@ private static final class DisabledGracefulShutdownSupport extends GracefulShutd
7886

7987
private volatile boolean shuttingDown;
8088

89+
/**
90+
* Creates a new instance.
91+
*
92+
*/
93+
DisabledGracefulShutdownSupport(ServerMetrics serverMetrics) {
94+
super(serverMetrics);
95+
}
96+
8197
@Override
8298
boolean isShuttingDown() {
8399
return shuttingDown;
@@ -97,12 +113,14 @@ private static final class DefaultGracefulShutdownSupport extends GracefulShutdo
97113
private final Executor blockingTaskExecutor;
98114

99115
/**
100-
* Declared as non-volatile because using {@link #pendingResponses} as a memory barrier.
116+
* Declared as non-volatile because using {@link #activeNonTransientResponses} as a memory barrier.
101117
*/
102118
private long lastResTimeNanos;
103119
private volatile long shutdownStartTimeNanos;
104120

105-
DefaultGracefulShutdownSupport(Duration quietPeriod, Executor blockingTaskExecutor, Ticker ticker) {
121+
DefaultGracefulShutdownSupport(Duration quietPeriod, Executor blockingTaskExecutor, Ticker ticker,
122+
ServerMetrics serverMetrics) {
123+
super(serverMetrics);
106124
quietPeriodNanos = quietPeriod.toNanos();
107125
this.blockingTaskExecutor = blockingTaskExecutor;
108126
this.ticker = ticker;
@@ -125,7 +143,7 @@ boolean completedQuietPeriod() {
125143
shutdownStartTimeNanos = readTicker();
126144
}
127145

128-
if (pendingResponses() != 0 || !completedBlockingTasks()) {
146+
if (activeNonTransientResponses() != 0 || !completedBlockingTasks()) {
129147
return false;
130148
}
131149

core/src/main/java/com/linecorp/armeria/server/Server.java

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -502,11 +502,11 @@ private final class ServerStartStopSupport extends StartStopSupport<Void, Void,
502502
@Override
503503
protected CompletionStage<Void> doStart(@Nullable Void arg) {
504504
if (config().gracefulShutdownQuietPeriod().isZero()) {
505-
gracefulShutdownSupport = GracefulShutdownSupport.createDisabled();
505+
gracefulShutdownSupport = GracefulShutdownSupport.createDisabled(config.serverMetrics());
506506
} else {
507507
gracefulShutdownSupport =
508508
GracefulShutdownSupport.create(config().gracefulShutdownQuietPeriod(),
509-
config().blockingTaskExecutor());
509+
config().blockingTaskExecutor(), config.serverMetrics());
510510
}
511511

512512
// Initialize the server sockets asynchronously.
@@ -585,8 +585,6 @@ private void setupServerMetrics() {
585585
final GracefulShutdownSupport gracefulShutdownSupport = this.gracefulShutdownSupport;
586586
assert gracefulShutdownSupport != null;
587587

588-
meterRegistry.gauge("armeria.server.pending.responses", gracefulShutdownSupport,
589-
GracefulShutdownSupport::pendingResponses);
590588
config.serverMetrics().bindTo(meterRegistry);
591589
}
592590

core/src/main/java/com/linecorp/armeria/server/ServerMetrics.java

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ public final class ServerMetrics implements MeterBinder {
3939
private final LongAdder activeHttp1WebSocketRequests = new LongAdder();
4040
private final LongAdder activeHttp1Requests = new LongAdder();
4141
private final LongAdder activeHttp2Requests = new LongAdder();
42+
private final LongAdder activeNonTransientRequests = new LongAdder();
4243

4344
/**
4445
* AtomicInteger is used to read the number of active connections frequently.
@@ -98,6 +99,13 @@ public long activeHttp2Requests() {
9899
return activeHttp2Requests.longValue();
99100
}
100101

102+
/**
103+
* Returns the number of pending http responses.
104+
*/
105+
public long activeNonTransientRequests() {
106+
return activeNonTransientRequests.longValue();
107+
}
108+
101109
/**
102110
* Returns the number of open connections.
103111
*/
@@ -153,6 +161,14 @@ void decreaseActiveConnections() {
153161
activeConnections.decrementAndGet();
154162
}
155163

164+
void increaseNonTransientRequests() {
165+
activeNonTransientRequests.increment();
166+
}
167+
168+
void decreaseNonTransientRequests() {
169+
activeNonTransientRequests.decrement();
170+
}
171+
156172
@Override
157173
public void bindTo(MeterRegistry meterRegistry) {
158174
meterRegistry.gauge("armeria.server.connections", activeConnections);
@@ -174,6 +190,10 @@ public void bindTo(MeterRegistry meterRegistry) {
174190
meterRegistry.gauge(allRequestsMeterName,
175191
ImmutableList.of(Tag.of("protocol", "http1.websocket"), Tag.of("state", "active")),
176192
activeHttp1WebSocketRequests);
193+
// pending non-transient responses
194+
meterRegistry.gauge(allRequestsMeterName,
195+
ImmutableList.of(Tag.of("protocol", "all"), Tag.of("state", "active")),
196+
activeNonTransientRequests);
177197
}
178198

179199
@Override
@@ -185,6 +205,7 @@ public String toString() {
185205
.add("pendingHttp2Requests", pendingHttp2Requests)
186206
.add("activeHttp2Requests", activeHttp2Requests)
187207
.add("activeConnections", activeConnections)
208+
.add("activeNonTransientRequests", activeNonTransientRequests)
188209
.toString();
189210
}
190211
}

core/src/test/java/com/linecorp/armeria/server/GracefulShutdownSupportTest.java

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,8 @@ class GracefulShutdownSupportTest {
4545
@Mock
4646
private Ticker ticker;
4747

48+
private ServerMetrics serverMetrics;
49+
4850
private GracefulShutdownSupport support;
4951
private ThreadPoolExecutor executor;
5052

@@ -54,7 +56,9 @@ void setUp() {
5456
0, 1, 1, TimeUnit.SECONDS, new LinkedTransferQueue<>(),
5557
ThreadFactories.newThreadFactory("graceful-shutdown-test", true));
5658

57-
support = GracefulShutdownSupport.create(Duration.ofNanos(QUIET_PERIOD_NANOS), executor, ticker);
59+
serverMetrics = new ServerMetrics();
60+
support = GracefulShutdownSupport.create(Duration.ofNanos(QUIET_PERIOD_NANOS), executor, ticker,
61+
serverMetrics);
5862
}
5963

6064
@AfterEach
@@ -64,17 +68,17 @@ void tearDown() {
6468

6569
@Test
6670
void testDisabled() {
67-
final GracefulShutdownSupport support = GracefulShutdownSupport.createDisabled();
71+
final GracefulShutdownSupport support = GracefulShutdownSupport.createDisabled(serverMetrics);
6872
assertThat(support.isShuttingDown()).isFalse();
6973
assertThat(support.completedQuietPeriod()).isTrue();
7074
assertThat(support.isShuttingDown()).isTrue();
7175
support.inc();
72-
assertThat(support.pendingResponses()).isOne();
76+
assertThat(support.activeNonTransientResponses()).isOne();
7377
assertThat(support.completedQuietPeriod()).isTrue();
7478

7579
// pendingResponses must be updated even if disabled, because it's part of metrics.
7680
support.dec();
77-
assertThat(support.pendingResponses()).isZero();
81+
assertThat(support.activeNonTransientResponses()).isZero();
7882
}
7983

8084
@Test

0 commit comments

Comments
 (0)