Skip to content

Commit 3df32ba

Browse files
authored
fix(grpc-gcp): correct channel lifecycle bookkeeping (#14196)
Use a copy-on-write test metric recorder because synchronized-list iteration raced with asynchronous exports.
1 parent f07e1a4 commit 3df32ba

5 files changed

Lines changed: 163 additions & 48 deletions

File tree

grpc-gcp-java/src/main/java/com/google/cloud/grpc/GcpClientCall.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -125,7 +125,10 @@ public void sendMessage(ReqT message) {
125125
delegateChannelRef.activeStreamsCountIncr();
126126

127127
// Create the client call and do the previous operations.
128-
delegateCall = delegateChannelRef.getChannel().newCall(methodDescriptor, callOptions);
128+
CallOptions callOptionsWithChannelId =
129+
callOptions.withOption(GcpManagedChannel.CHANNEL_ID_KEY, delegateChannelRef.getId());
130+
delegateCall =
131+
delegateChannelRef.getChannel().newCall(methodDescriptor, callOptionsWithChannelId);
129132
for (Runnable call : calls) {
130133
call.run();
131134
}

grpc-gcp-java/src/main/java/com/google/cloud/grpc/GcpManagedChannel.java

Lines changed: 88 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -269,6 +269,16 @@ private static int stateFromChannelId(int channelId) {
269269
// Clock supplier for nanoTime, injectable for testing.
270270
private Supplier<Long> nanoClock = System::nanoTime;
271271

272+
@VisibleForTesting
273+
Map<Integer, Map<String, Integer>> fallbackMapForTest() {
274+
return fallbackMap;
275+
}
276+
277+
@VisibleForTesting
278+
int readyChannelCountForTest() {
279+
return readyChannels.get();
280+
}
281+
272282
@VisibleForTesting
273283
void setNanoClock(Supplier<Long> nanoClock) {
274284
this.nanoClock = nanoClock;
@@ -408,10 +418,7 @@ private void removeOldestChannels(int num) {
408418

409419
for (ChannelRef channelRef : channelsToRemove) {
410420
channelRef.resetAffinityCount();
411-
channelRef.deactivate();
412-
if (channelRef.getState() == ConnectivityState.READY) {
413-
decReadyChannels(false);
414-
}
421+
channelRef.deactivateAndAccountReadiness();
415422
}
416423

417424
// Remove affinity keys mapping for the channels.
@@ -1482,6 +1489,9 @@ private class ChannelStateMonitor implements Runnable {
14821489
private long connectingStartNanos;
14831490
private long connectedSinceNanos;
14841491

1492+
@GuardedBy("channelRef")
1493+
private boolean readyAccounted;
1494+
14851495
private ChannelStateMonitor(ManagedChannel channel, ChannelRef channelRef) {
14861496
this.channelRef = channelRef;
14871497
this.channel = channel;
@@ -1496,15 +1506,26 @@ public ConnectivityState getCurrentState() {
14961506
return currentState;
14971507
}
14981508

1509+
private void accountReadyIfNeeded() {
1510+
if (currentState == ConnectivityState.READY && !readyAccounted) {
1511+
readyAccounted = true;
1512+
incReadyChannels(false);
1513+
}
1514+
}
1515+
1516+
private void unaccountReadyIfNeeded() {
1517+
if (readyAccounted) {
1518+
readyAccounted = false;
1519+
decReadyChannels(false);
1520+
}
1521+
}
1522+
14991523
@Override
15001524
public void run() {
15011525
if (channel == null) {
15021526
return;
15031527
}
15041528

1505-
// Is the channel in the pool?
1506-
boolean isActive = channelRefs.contains(this.channelRef);
1507-
15081529
// Keep minSize channels always connected.
15091530
boolean requestConnection =
15101531
channelRefs.size() < minSize
@@ -1515,35 +1536,39 @@ public void run() {
15151536
.anyMatch(id -> (id == channelRef.getId()));
15161537

15171538
ConnectivityState newState = channel.getState(requestConnection);
1518-
if (logger.isLoggable(Level.FINER)) {
1519-
logger.finer(
1520-
log(
1521-
"Channel %d state change detected: %s -> %s",
1522-
channelRef.getId(), currentState, newState));
1523-
}
1524-
if (newState == ConnectivityState.READY && currentState != ConnectivityState.READY) {
1525-
connectedSinceNanos = System.nanoTime();
1526-
if (isActive) {
1527-
incReadyChannels(true);
1528-
if (connectingStartNanos > 0) {
1529-
saveReadinessTime(System.nanoTime() - connectingStartNanos);
1539+
boolean isActive;
1540+
synchronized (channelRef) {
1541+
isActive = channelRef.isActive() && channelRefs.contains(channelRef);
1542+
if (logger.isLoggable(Level.FINER)) {
1543+
logger.finer(
1544+
log(
1545+
"Channel %d state change detected: %s -> %s",
1546+
channelRef.getId(), currentState, newState));
1547+
}
1548+
if (newState == ConnectivityState.READY && currentState != ConnectivityState.READY) {
1549+
connectedSinceNanos = System.nanoTime();
1550+
if (isActive && !readyAccounted) {
1551+
readyAccounted = true;
1552+
incReadyChannels(true);
1553+
if (connectingStartNanos > 0) {
1554+
saveReadinessTime(System.nanoTime() - connectingStartNanos);
1555+
}
15301556
}
1557+
connectingStartNanos = 0;
15311558
}
1532-
connectingStartNanos = 0;
1533-
}
1534-
if (isActive
1535-
&& newState != ConnectivityState.READY
1536-
&& currentState == ConnectivityState.READY) {
1537-
decReadyChannels(true);
1538-
}
1539-
if (newState == ConnectivityState.CONNECTING
1540-
&& currentState != ConnectivityState.CONNECTING) {
1541-
connectingStartNanos = System.nanoTime();
1542-
}
1543-
if (newState != ConnectivityState.READY) {
1544-
connectedSinceNanos = 0;
1559+
if (newState != ConnectivityState.READY && readyAccounted) {
1560+
readyAccounted = false;
1561+
decReadyChannels(true);
1562+
}
1563+
if (newState == ConnectivityState.CONNECTING
1564+
&& currentState != ConnectivityState.CONNECTING) {
1565+
connectingStartNanos = System.nanoTime();
1566+
}
1567+
if (newState != ConnectivityState.READY) {
1568+
connectedSinceNanos = 0;
1569+
}
1570+
currentState = newState;
15451571
}
1546-
currentState = newState;
15471572

15481573
processChannelStateChange(channelRef.getId(), newState);
15491574
if (isActive) {
@@ -1686,8 +1711,12 @@ protected ChannelRef getChannelRef(@Nullable String key) {
16861711
if (logger.isLoggable(Level.FINEST)) {
16871712
logger.finest(log("Using fallback channel: %d -> %d", mappedChannel.getId(), channelId));
16881713
}
1689-
fallbacksSucceeded.incrementAndGet();
1690-
return channelRefs.get(channelId);
1714+
ChannelRef fallbackChannel = channelIdToChannelRef.get(channelId);
1715+
if (fallbackChannel != null && fallbackChannel.isActive()) {
1716+
fallbacksSucceeded.incrementAndGet();
1717+
return fallbackChannel;
1718+
}
1719+
tempMap.remove(key, channelId);
16911720
}
16921721
// No temp mapping for this key or fallback channel is also broken.
16931722
ChannelRef channelRef = pickLeastBusyChannel(/* forFallback= */ true);
@@ -1710,7 +1739,10 @@ protected ChannelRef getChannelRef(@Nullable String key) {
17101739
fallbacksFailed.incrementAndGet();
17111740
if (channelId != null) {
17121741
// Stick with previous mapping if fallback has failed.
1713-
return channelRefs.get(channelId);
1742+
ChannelRef fallbackChannel = channelIdToChannelRef.get(channelId);
1743+
if (fallbackChannel != null && fallbackChannel.isActive()) {
1744+
return fallbackChannel;
1745+
}
17141746
}
17151747
return mappedChannel;
17161748
}
@@ -1779,16 +1811,16 @@ ChannelRef createNewChannel() {
17791811
channelRefs.add(chRef);
17801812
removedChannelRefs.remove(chRef);
17811813
channelIdToChannelRef.put(chRef.getId(), chRef);
1782-
chRef.activate();
1814+
chRef.activateAndAccountReadiness();
17831815
logger.finer(log("Channel %d reused.", chRef.getId()));
1784-
incReadyChannels(false);
17851816
maxChannels.accumulateAndGet(getNumberOfChannels(), Math::max);
17861817
return chRef;
17871818
}
17881819

17891820
ChannelRef channelRef = new ChannelRef(delegateChannelBuilder.build());
17901821
channelRefs.add(channelRef);
17911822
channelIdToChannelRef.put(channelRef.getId(), channelRef);
1823+
channelRef.activateAndAccountReadiness();
17921824
logger.finer(log("Channel %d created.", channelRef.getId()));
17931825
maxChannels.accumulateAndGet(getNumberOfChannels(), Math::max);
17941826
return channelRef;
@@ -2447,12 +2479,27 @@ protected boolean isActive() {
24472479
return active;
24482480
}
24492481

2450-
private void activate() {
2451-
active = true;
2482+
private void activateAndAccountReadiness() {
2483+
synchronized (this) {
2484+
active = true;
2485+
channelStateMonitor.accountReadyIfNeeded();
2486+
}
2487+
}
2488+
2489+
private void deactivateAndAccountReadiness() {
2490+
synchronized (this) {
2491+
channelStateMonitor.unaccountReadyIfNeeded();
2492+
active = false;
2493+
}
24522494
}
24532495

24542496
private void deactivate() {
2455-
active = false;
2497+
deactivateAndAccountReadiness();
2498+
}
2499+
2500+
@VisibleForTesting
2501+
void deactivateForTest() {
2502+
deactivateAndAccountReadiness();
24562503
}
24572504

24582505
protected void affinityCountIncr() {

grpc-gcp-java/src/test/java/com/google/cloud/grpc/GcpManagedChannelTest.java

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -364,6 +364,72 @@ public void testChannelAffinityRefInvalidChannelIdPicksAvailableChannel() throws
364364
}
365365
}
366366

367+
@Test
368+
public void fallbackUsesChannelIdMapAfterPoolHasIndexGap() {
369+
resetGcpChannel();
370+
ExecutorService executorService = Executors.newSingleThreadExecutor();
371+
try {
372+
List<FakeManagedChannel> channels = new ArrayList<>();
373+
for (int i = 0; i < 3; i++) {
374+
FakeManagedChannel channel = new FakeManagedChannel(executorService);
375+
channel.setState(ConnectivityState.READY);
376+
channels.add(channel);
377+
}
378+
gcpChannel =
379+
(GcpManagedChannel)
380+
GcpManagedChannelBuilder.forDelegateBuilder(new FakeManagedChannelBuilder(channels))
381+
.withOptions(
382+
GcpManagedChannelOptions.newBuilder()
383+
.withChannelPoolOptions(
384+
GcpChannelPoolOptions.newBuilder()
385+
.setMinSize(3)
386+
.setMaxSize(3)
387+
.build())
388+
.withResiliencyOptions(
389+
GcpResiliencyOptions.newBuilder().setNotReadyFallback(true).build())
390+
.build())
391+
.build();
392+
ChannelRef removed = gcpChannel.channelRefs.get(0);
393+
ChannelRef mapped = gcpChannel.channelRefs.get(1);
394+
ChannelRef fallback = gcpChannel.channelRefs.get(2);
395+
String key = "session";
396+
gcpChannel.bind(mapped, Collections.singletonList(key));
397+
gcpChannel.processChannelStateChange(mapped.getId(), ConnectivityState.TRANSIENT_FAILURE);
398+
gcpChannel.fallbackMapForTest().get(mapped.getId()).put(key, fallback.getId());
399+
gcpChannel.channelRefs.remove(removed);
400+
401+
assertThat(gcpChannel.getChannelRef(key)).isSameInstanceAs(fallback);
402+
} finally {
403+
gcpChannel.shutdownNow();
404+
executorService.shutdownNow();
405+
}
406+
}
407+
408+
@Test
409+
public void readyAccountingRemainsExactWhenReadyChannelIsReused() {
410+
resetGcpChannel();
411+
ExecutorService executorService = Executors.newSingleThreadExecutor();
412+
try {
413+
gcpChannel = createPoolWithFakeReadyChannels(executorService, 2);
414+
assertThat(gcpChannel.readyChannelCountForTest()).isEqualTo(2);
415+
416+
ChannelRef reused = gcpChannel.channelRefs.get(0);
417+
gcpChannel.channelRefs.remove(reused);
418+
reused.deactivateForTest();
419+
gcpChannel.removedChannelRefs.add(reused);
420+
assertThat(gcpChannel.readyChannelCountForTest()).isEqualTo(1);
421+
422+
assertThat(gcpChannel.createNewChannel()).isSameInstanceAs(reused);
423+
assertThat(gcpChannel.readyChannelCountForTest()).isEqualTo(2);
424+
reused.deactivateForTest();
425+
reused.deactivateForTest();
426+
assertThat(gcpChannel.readyChannelCountForTest()).isEqualTo(1);
427+
} finally {
428+
gcpChannel.shutdownNow();
429+
executorService.shutdownNow();
430+
}
431+
}
432+
367433
@Test
368434
public void testChannelAffinityRefRemovedChannelPicksAvailableChannel() throws Exception {
369435
resetGcpChannel();

grpc-gcp-java/src/test/java/com/google/cloud/grpc/fallback/GcpFallbackChannelTest.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -75,11 +75,10 @@
7575
import java.io.InputStream;
7676
import java.nio.charset.StandardCharsets;
7777
import java.time.Duration;
78-
import java.util.ArrayList;
7978
import java.util.Collection;
80-
import java.util.Collections;
8179
import java.util.EnumSet;
8280
import java.util.List;
81+
import java.util.concurrent.CopyOnWriteArrayList;
8382
import java.util.concurrent.ScheduledExecutorService;
8483
import java.util.concurrent.TimeUnit;
8584
import java.util.concurrent.atomic.AtomicLong;
@@ -352,7 +351,7 @@ private void simulateCanceledCall(boolean expectFallbackRouting) {
352351
}
353352

354353
private class TestMetricExporter implements MetricExporter {
355-
public final List<MetricData> exportedMetrics = Collections.synchronizedList(new ArrayList<>());
354+
public final List<MetricData> exportedMetrics = new CopyOnWriteArrayList<>();
356355

357356
@Override
358357
public CompletableResultCode export(@Nonnull Collection<MetricData> metrics) {

java-spanner/google-cloud-spanner/src/test/java/com/google/cloud/spanner/RequestIdMockServerTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -680,9 +680,9 @@ public void testOtherClientId() {
680680
ImmutableList.of(
681681
XGoogSpannerRequestId.of(getClientId(), -1, 1, 1),
682682
// The CreateSession RPC from the initialization of the second client is included in
683-
// the requests that we see. This request does not include a channel hint, hence the
684-
// zero value for the channel number in the request ID.
685-
XGoogSpannerRequestId.of(otherClientId, 0, 1, 1),
683+
// the requests that we see. Its request ID carries the channel that grpc-gcp selected
684+
// for the affinity-key path.
685+
XGoogSpannerRequestId.of(otherClientId, -1, 1, 1),
686686
XGoogSpannerRequestId.of(otherClientId, -1, 2, 1),
687687
XGoogSpannerRequestId.of(getClientId(), -1, 2, 1)),
688688
actual);

0 commit comments

Comments
 (0)