Skip to content

Commit 491c330

Browse files
committed
xds: set end_of_stream on client half-close and handle end_of_stream_without_message per ext_proc spec
TAG=agy CONV=86c74ebe-9fd7-4876-b27c-a4c1b230d346
1 parent b7b76d6 commit 491c330

4 files changed

Lines changed: 1326 additions & 220 deletions

File tree

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

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -890,6 +890,7 @@ public void halfClose() {
890890
// Mode is GRPC
891891
sendToExtProc(ProcessingRequest.newBuilder()
892892
.setRequestBody(HttpBody.newBuilder()
893+
.setEndOfStream(true)
893894
.setEndOfStreamWithoutMessage(true)
894895
.build())
895896
.build());
@@ -914,10 +915,13 @@ private void handleRequestBodyResponse(BodyResponse bodyResponse) {
914915
BodyMutation mutation = bodyResponse.getResponse().getBodyMutation();
915916
if (mutation.hasStreamedResponse()) {
916917
StreamedBodyResponse streamed = mutation.getStreamedResponse();
917-
if (!streamed.getEndOfStreamWithoutMessage()) {
918+
boolean isEndOfStream = streamed.getEndOfStream();
919+
boolean isEndOfStreamWithoutMessage =
920+
isEndOfStream && streamed.getEndOfStreamWithoutMessage();
921+
if (!isEndOfStreamWithoutMessage) {
918922
super.sendMessage(new KnownLengthInputStream(streamed.getBody()));
919923
}
920-
if (streamed.getEndOfStream() || streamed.getEndOfStreamWithoutMessage()) {
924+
if (isEndOfStream) {
921925
if (requestSideClosed.compareAndSet(false, true)) {
922926
proceedWithHalfClose();
923927
}

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

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1335,7 +1335,9 @@ int drainPendingMutatedRequestBodies() {
13351335
break;
13361336
}
13371337
StreamedBodyResponse peeked = pendingMutatedRequestBodies.peek();
1338-
if (peeked.getEndOfStreamWithoutMessage()) {
1338+
boolean isEosWithoutMessage =
1339+
peeked.getEndOfStream() && peeked.getEndOfStreamWithoutMessage();
1340+
if (isEosWithoutMessage) {
13391341
toDeliver.add(pendingMutatedRequestBodies.poll());
13401342
} else if (pendingAppRequests.get() > 0) {
13411343
pendingAppRequests.decrementAndGet();
@@ -1350,10 +1352,13 @@ int drainPendingMutatedRequestBodies() {
13501352
final int bodySize = streamed.getBody().size();
13511353
callContext.run(() -> {
13521354
try {
1353-
if (!finalStreamed.getEndOfStreamWithoutMessage()) {
1355+
boolean isEndOfStream = finalStreamed.getEndOfStream();
1356+
boolean isEndOfStreamWithoutMessage =
1357+
isEndOfStream && finalStreamed.getEndOfStreamWithoutMessage();
1358+
if (!isEndOfStreamWithoutMessage) {
13541359
wrappedListener.onExternalBody(finalStreamed.getBody());
13551360
}
1356-
if (finalStreamed.getEndOfStream() || finalStreamed.getEndOfStreamWithoutMessage()) {
1361+
if (isEndOfStream) {
13571362
wrappedListener.proceedWithHalfClose();
13581363
}
13591364
} finally {
@@ -1484,10 +1489,13 @@ private void drainRequestMessagesFailOpen() {
14841489
for (StreamedBodyResponse streamed : mutatedToDeliver) {
14851490
final StreamedBodyResponse finalStreamed = streamed;
14861491
callContext.run(() -> {
1487-
if (!finalStreamed.getEndOfStreamWithoutMessage()) {
1492+
boolean isEndOfStream = finalStreamed.getEndOfStream();
1493+
boolean isEndOfStreamWithoutMessage =
1494+
isEndOfStream && finalStreamed.getEndOfStreamWithoutMessage();
1495+
if (!isEndOfStreamWithoutMessage) {
14881496
wrappedListener.onExternalBody(finalStreamed.getBody());
14891497
}
1490-
if (finalStreamed.getEndOfStream() || finalStreamed.getEndOfStreamWithoutMessage()) {
1498+
if (isEndOfStream) {
14911499
wrappedListener.proceedWithHalfClose();
14921500
}
14931501
});
@@ -1908,6 +1916,7 @@ private void sendHalfCloseToExtProc() {
19081916
}
19091917

19101918
HttpBody.Builder bodyBuilder = HttpBody.newBuilder()
1919+
.setEndOfStream(true)
19111920
.setEndOfStreamWithoutMessage(true);
19121921

19131922
ProcessingRequest.Builder builder = ProcessingRequest.newBuilder()

0 commit comments

Comments
 (0)