Skip to content

Commit 4d5d2e8

Browse files
committed
Add pendingResponse to ServerMetrics
1 parent 19068f6 commit 4d5d2e8

4 files changed

Lines changed: 73 additions & 31 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,47 @@
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+
public 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, ServerMetrics serverMetrics) {
46+
return new DefaultGracefulShutdownSupport(quietPeriod, blockingTaskExecutor, ticker, serverMetrics);
4247
}
4348

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

4653
/**
4754
* Increases the number of pending responses.
4855
*/
4956
final void inc() {
50-
pendingResponses.increment();
57+
serverMetrics.increaseNonTransientRequests();
5158
}
5259

5360
/**
5461
* Decreases the number of pending responses.
5562
*/
5663
void dec() {
57-
pendingResponses.decrement();
64+
serverMetrics.decreaseNonTransientRequests();
5865
}
5966

6067
/**
6168
* Returns the number of pending responses.
6269
*/
63-
final long pendingResponses() {
64-
return pendingResponses.sum();
70+
final long activeNonTransientResponses() {
71+
return serverMetrics.activeNonTransientRequests();
6572
}
6673

6774
/**
@@ -78,6 +85,15 @@ private static final class DisabledGracefulShutdownSupport extends GracefulShutd
7885

7986
private volatile boolean shuttingDown;
8087

88+
/**
89+
* Creates a new instance.
90+
*
91+
* @param serverMetrics
92+
*/
93+
public 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: 12 additions & 12 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.
@@ -520,7 +520,7 @@ protected CompletionStage<Void> doStart(@Nullable Void arg) {
520520
try {
521521
doStart(primary).addListener(new ServerPortStartListener(primary))
522522
.addListener(new NextServerPortStartListener(this, it, future));
523-
setupServerMetrics();
523+
// setupServerMetrics();
524524
} catch (Throwable cause) {
525525
future.completeExceptionally(cause);
526526
}
@@ -580,15 +580,15 @@ private ChannelFuture doStart(ServerPort port) {
580580
return b.bind(localAddress);
581581
}
582582

583-
private void setupServerMetrics() {
584-
final MeterRegistry meterRegistry = config.meterRegistry();
585-
final GracefulShutdownSupport gracefulShutdownSupport = this.gracefulShutdownSupport;
586-
assert gracefulShutdownSupport != null;
587-
588-
meterRegistry.gauge("armeria.server.pending.responses", gracefulShutdownSupport,
589-
GracefulShutdownSupport::pendingResponses);
590-
config.serverMetrics().bindTo(meterRegistry);
591-
}
583+
// private void setupServerMetrics() {
584+
// final MeterRegistry meterRegistry = config.meterRegistry();
585+
// final GracefulShutdownSupport gracefulShutdownSupport = this.gracefulShutdownSupport;
586+
// assert gracefulShutdownSupport != null;
587+
//
588+
// meterRegistry.gauge("armeria.server.pending.responses", gracefulShutdownSupport,
589+
//// GracefulShutdownSupport::pendingResponses);
590+
// config.serverMetrics().bindTo(meterRegistry);
591+
// }
592592

593593
@Override
594594
protected CompletionStage<Void> doStop(@Nullable Void arg) {

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: 7 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,8 @@ 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, serverMetrics);
5861
}
5962

6063
@AfterEach
@@ -64,17 +67,17 @@ void tearDown() {
6467

6568
@Test
6669
void testDisabled() {
67-
final GracefulShutdownSupport support = GracefulShutdownSupport.createDisabled();
70+
final GracefulShutdownSupport support = GracefulShutdownSupport.createDisabled(serverMetrics);
6871
assertThat(support.isShuttingDown()).isFalse();
6972
assertThat(support.completedQuietPeriod()).isTrue();
7073
assertThat(support.isShuttingDown()).isTrue();
7174
support.inc();
72-
assertThat(support.pendingResponses()).isOne();
75+
assertThat(support.activeNonTransientResponses()).isOne();
7376
assertThat(support.completedQuietPeriod()).isTrue();
7477

7578
// pendingResponses must be updated even if disabled, because it's part of metrics.
7679
support.dec();
77-
assertThat(support.pendingResponses()).isZero();
80+
assertThat(support.activeNonTransientResponses()).isZero();
7881
}
7982

8083
@Test

0 commit comments

Comments
 (0)