Fix off-by-one error in PartEvent part count
This commit fixes an off-by-one error in the PartEventHttpMessageReader, so that it no longer counts empty windows. Closes gh-32122
This commit is contained in:
@@ -155,29 +155,22 @@ public class PartEventHttpMessageReader extends LoggingCodecSupport implements H
|
|||||||
AtomicInteger partCount = new AtomicInteger();
|
AtomicInteger partCount = new AtomicInteger();
|
||||||
return allPartsTokens
|
return allPartsTokens
|
||||||
.windowUntil(t -> t instanceof MultipartParser.HeadersToken, true)
|
.windowUntil(t -> t instanceof MultipartParser.HeadersToken, true)
|
||||||
.concatMap(partTokens -> {
|
.concatMap(partTokens -> partTokens
|
||||||
if (tooManyParts(partCount)) {
|
.switchOnFirst((signal, flux) -> {
|
||||||
return Mono.error(new DecodingException("Too many parts (" + partCount.get() + "/" +
|
if (!signal.hasValue()) {
|
||||||
this.maxParts + " allowed)"));
|
|
||||||
}
|
|
||||||
else {
|
|
||||||
return partTokens.switchOnFirst((signal, flux) -> {
|
|
||||||
if (signal.hasValue()) {
|
|
||||||
MultipartParser.HeadersToken headersToken = (MultipartParser.HeadersToken) signal.get();
|
|
||||||
Assert.state(headersToken != null, "Signal should be headers token");
|
|
||||||
|
|
||||||
HttpHeaders headers = headersToken.headers();
|
|
||||||
Flux<MultipartParser.BodyToken> bodyTokens = flux.ofType(
|
|
||||||
MultipartParser.BodyToken.class);
|
|
||||||
return createEvents(headers, bodyTokens);
|
|
||||||
}
|
|
||||||
else {
|
|
||||||
// complete or error signal
|
// complete or error signal
|
||||||
return flux.cast(PartEvent.class);
|
return flux.cast(PartEvent.class);
|
||||||
}
|
}
|
||||||
});
|
else if (tooManyParts(partCount)) {
|
||||||
}
|
return Mono.error(new DecodingException("Too many parts (" + partCount.get() +
|
||||||
});
|
"/" + this.maxParts + " allowed)"));
|
||||||
|
}
|
||||||
|
MultipartParser.HeadersToken headersToken = (MultipartParser.HeadersToken) signal.get();
|
||||||
|
Assert.state(headersToken != null, "Signal should be headers token");
|
||||||
|
|
||||||
|
HttpHeaders headers = headersToken.headers();
|
||||||
|
return createEvents(headers, flux.ofType(MultipartParser.BodyToken.class));
|
||||||
|
}));
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -238,6 +238,7 @@ class PartEventHttpMessageReaderTests {
|
|||||||
Flux<PartEvent> result = reader.read(forClass(PartEvent.class), request, emptyMap());
|
Flux<PartEvent> result = reader.read(forClass(PartEvent.class), request, emptyMap());
|
||||||
|
|
||||||
StepVerifier.create(result)
|
StepVerifier.create(result)
|
||||||
|
.assertNext(form(headers -> assertThat(headers).isEmpty(), "This is implicitly typed plain ASCII text.\r\nIt does NOT end with a linebreak."))
|
||||||
.expectError(DecodingException.class)
|
.expectError(DecodingException.class)
|
||||||
.verify();
|
.verify();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user