Skip to content

Commit fdaaa9a

Browse files
committed
Fix triggerCloseHandshake for observability mode.
Initialize closeNow to observabilityMode to ensure the call is closed immediately when closing in observability mode, preventing hangs and avoiding duplicate proceedWithClose calls. Update givenObservabilityMode_whenDataPlaneClosed_thenSidecarCloseIsDeferred to assert that onClose is called exactly once, preventing double onClose notifications from passing silently. TAG=agy CONV=2c1e4760-c239-4698-810a-162bf10fccc4
1 parent 2d3c77b commit fdaaa9a

2 files changed

Lines changed: 47 additions & 48 deletions

File tree

xds/src/main/java/io/grpc/xds/ExternalProcessorClientInterceptor.java

Lines changed: 44 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -1431,67 +1431,63 @@ private void triggerCloseHandshake() {
14311431
return;
14321432
}
14331433

1434-
if (dataPlaneClientCall.responseDrainComplete.get()) {
1435-
proceedWithClose();
1436-
if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1437-
dataPlaneClientCall.closeExtProcStream();
1438-
}
1439-
return;
1440-
}
1441-
1442-
boolean sendResponseHeaders =
1443-
dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1444-
== ProcessingMode.HeaderSendMode.SEND
1445-
|| dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1446-
== ProcessingMode.HeaderSendMode.DEFAULT;
1447-
1448-
boolean sendResponseTrailers =
1449-
dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1450-
== ProcessingMode.HeaderSendMode.SEND;
1434+
boolean closeNow = dataPlaneClientCall.getConfig().getObservabilityMode();
14511435

1452-
if (trailersOnly.get()) {
1453-
if (sendResponseHeaders) {
1436+
if (dataPlaneClientCall.responseDrainComplete.get()) {
1437+
closeNow = true;
1438+
} else {
1439+
boolean sendResponseHeaders =
1440+
dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1441+
== ProcessingMode.HeaderSendMode.SEND
1442+
|| dataPlaneClientCall.getCurrentProcessingMode().getResponseHeaderMode()
1443+
== ProcessingMode.HeaderSendMode.DEFAULT;
1444+
1445+
boolean sendResponseTrailers =
1446+
dataPlaneClientCall.getCurrentProcessingMode().getResponseTrailerMode()
1447+
== ProcessingMode.HeaderSendMode.SEND;
1448+
1449+
if (trailersOnly.get()) {
1450+
if (sendResponseHeaders) {
1451+
dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1452+
.setResponseHeaders(HttpHeaders.newBuilder()
1453+
.setHeaders(
1454+
toHeaderMap(
1455+
savedTrailers,
1456+
dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1457+
.setEndOfStream(true)
1458+
.build())
1459+
.build());
1460+
dataPlaneClientCall.responseTrailersSent.set(true);
1461+
} else {
1462+
closeNow = true;
1463+
}
1464+
} else if (sendResponseTrailers) {
1465+
dataPlaneClientCall.getIsProcessingTrailers().set(true);
14541466
dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1455-
.setResponseHeaders(HttpHeaders.newBuilder()
1456-
.setHeaders(
1467+
.setResponseTrailers(HttpTrailers.newBuilder()
1468+
.setTrailers(
14571469
toHeaderMap(
14581470
savedTrailers,
14591471
dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1460-
.setEndOfStream(true)
14611472
.build())
14621473
.build());
14631474
dataPlaneClientCall.responseTrailersSent.set(true);
14641475
} else {
1465-
proceedWithClose();
1466-
if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1467-
dataPlaneClientCall.closeExtProcStream();
1468-
}
1469-
}
1470-
} else if (sendResponseTrailers) {
1471-
dataPlaneClientCall.getIsProcessingTrailers().set(true);
1472-
dataPlaneClientCall.sendToExtProc(ProcessingRequest.newBuilder()
1473-
.setResponseTrailers(HttpTrailers.newBuilder()
1474-
.setTrailers(
1475-
toHeaderMap(
1476-
savedTrailers,
1477-
dataPlaneClientCall.getConfig().getForwardRulesConfig()))
1478-
.build())
1479-
.build());
1480-
dataPlaneClientCall.responseTrailersSent.set(true);
1481-
} else {
1482-
proceedWithClose();
1483-
if (!dataPlaneClientCall.getConfig().getObservabilityMode()) {
1484-
dataPlaneClientCall.closeExtProcStream();
1476+
closeNow = true;
14851477
}
14861478
}
14871479

1488-
if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1480+
if (closeNow) {
14891481
proceedWithClose();
1490-
@SuppressWarnings("unused")
1491-
ScheduledFuture<?> unused = dataPlaneClientCall.getScheduler().schedule(
1492-
dataPlaneClientCall::closeExtProcStream,
1493-
dataPlaneClientCall.getConfig().getDeferredCloseTimeoutNanos(),
1494-
TimeUnit.NANOSECONDS);
1482+
if (dataPlaneClientCall.getConfig().getObservabilityMode()) {
1483+
@SuppressWarnings("unused")
1484+
ScheduledFuture<?> unused = dataPlaneClientCall.getScheduler().schedule(
1485+
dataPlaneClientCall::closeExtProcStream,
1486+
dataPlaneClientCall.getConfig().getDeferredCloseTimeoutNanos(),
1487+
TimeUnit.NANOSECONDS);
1488+
} else {
1489+
dataPlaneClientCall.closeExtProcStream();
1490+
}
14951491
}
14961492
}
14971493

xds/src/test/java/io/grpc/xds/ExternalProcessorClientInterceptorTest.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10236,9 +10236,11 @@ public void onCompleted() {
1023610236
.build());
1023710237

1023810238
try {
10239+
final AtomicInteger onCloseCallCount = new AtomicInteger(0);
1023910240
final CountDownLatch appCloseLatch = new CountDownLatch(1);
1024010241
ClientCall.Listener<String> appListener = new ClientCall.Listener<String>() {
1024110242
@Override public void onClose(Status status, Metadata trailers) {
10243+
onCloseCallCount.incrementAndGet();
1024210244
appCloseLatch.countDown();
1024310245
}
1024410246
};
@@ -10265,6 +10267,7 @@ public void onCompleted() {
1026510267
fakeClock.forwardTime(1, TimeUnit.SECONDS);
1026610268
}
1026710269
assertThat(appCloseLatch.await(5, TimeUnit.SECONDS)).isTrue();
10270+
assertThat(onCloseCallCount.get()).isEqualTo(1);
1026810271

1026910272
// At this point, app received onClose, but sidecar should NOT be completed yet
1027010273
assertThat(sidecarCompletedLatch.getCount()).isEqualTo(1);

0 commit comments

Comments
 (0)