Ensure chunks released on cancel in StringDecoder
The current test were not catching the issue because they request 1 via StepVerifier, wait for it, and then cancel. In the case of StringDecoder it means all chunks are used up to produce that first String and so the cancel doesn't catch any cached chunks. Closes gh-30299
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2020 the original author or authors.
|
||||
* Copyright 2002-2023 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -18,12 +18,15 @@ package org.springframework.core.codec;
|
||||
|
||||
import java.nio.charset.Charset;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.time.Duration;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.reactivestreams.Subscription;
|
||||
import reactor.core.publisher.BaseSubscriber;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
@@ -75,16 +78,17 @@ class StringDecoderTests extends AbstractDecoderTests<StringDecoder> {
|
||||
String s = String.format("%s\n%s\n%s", u, e, o);
|
||||
Flux<DataBuffer> input = toDataBuffers(s, 1, UTF_8);
|
||||
|
||||
// TODO: temporarily replace testDecodeAll with explicit decode/cancel/empty
|
||||
// see https://github.com/reactor/reactor-core/issues/2041
|
||||
|
||||
// testDecode(input, TYPE, step -> step.expectNext(u, e, o).verifyComplete(), null, null);
|
||||
// testDecodeCancel(input, TYPE, null, null);
|
||||
// testDecodeEmpty(TYPE, null, null);
|
||||
|
||||
testDecodeAll(input, TYPE, step -> step.expectNext(u, e, o).verifyComplete(), null, null);
|
||||
}
|
||||
|
||||
@Test // gh-30299
|
||||
public void decodeAndCancelWithPendingChunks() {
|
||||
Flux<DataBuffer> input = toDataBuffers("abc", 1, UTF_8).concatWith(Flux.never());
|
||||
Flux<String> result = this.decoder.decode(input, TYPE, null, null);
|
||||
|
||||
StepVerifier.create(result).thenAwait(Duration.ofMillis(100)).thenCancel().verify();
|
||||
}
|
||||
|
||||
@Test
|
||||
void decodeMultibyteCharacterUtf16() {
|
||||
String u = "ü";
|
||||
@@ -264,4 +268,13 @@ class StringDecoderTests extends AbstractDecoderTests<StringDecoder> {
|
||||
return buffer;
|
||||
}
|
||||
|
||||
|
||||
private static class SingleRequestSubscriber extends BaseSubscriber<String> {
|
||||
|
||||
@Override
|
||||
protected void hookOnSubscribe(Subscription subscription) {
|
||||
subscription.request(1);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user