Merge branch '7.0.x'

This commit is contained in:
Brian Clozel
2026-06-04 12:27:18 +02:00
2 changed files with 55 additions and 5 deletions
@@ -1138,9 +1138,11 @@ public abstract class DataBufferUtils {
protected void hookOnNext(DataBuffer dataBuffer) {
try {
try (DataBuffer.ByteBufferIterator iterator = dataBuffer.readableByteBuffers()) {
ByteBuffer byteBuffer = iterator.next();
while (byteBuffer.hasRemaining()) {
this.channel.write(byteBuffer);
while (iterator.hasNext()) {
ByteBuffer byteBuffer = iterator.next();
while (byteBuffer.hasRemaining()) {
this.channel.write(byteBuffer);
}
}
}
this.sink.next(dataBuffer);
@@ -1213,6 +1215,11 @@ public abstract class DataBufferUtils {
failed(ex, attachment);
}
}
else {
iterator.close();
this.sink.next(dataBuffer);
request(1);
}
}
@Override
@@ -1236,7 +1243,6 @@ public abstract class DataBufferUtils {
@Override
public void completed(Integer written, Attachment attachment) {
DataBuffer.ByteBufferIterator iterator = attachment.iterator();
iterator.close();
long pos = this.position.addAndGet(written);
ByteBuffer byteBuffer = attachment.byteBuffer();
@@ -1246,9 +1252,11 @@ public abstract class DataBufferUtils {
}
else if (iterator.hasNext()) {
ByteBuffer next = iterator.next();
this.channel.write(next, pos, attachment, this);
Attachment nextAttachment = new Attachment(next, attachment.dataBuffer(), iterator);
this.channel.write(next, pos, nextAttachment, this);
}
else {
iterator.close();
this.sink.next(attachment.dataBuffer());
this.writing.set(false);
@@ -338,6 +338,27 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
channel.close();
}
@ParameterizedDataBufferAllocatingTest
void writeWritableByteChannelWithJoinedBuffer(DataBufferFactory bufferFactory) throws Exception {
super.bufferFactory = bufferFactory;
DataBuffer foo = stringBuffer("foo");
DataBuffer bar = stringBuffer("bar");
DataBuffer joined = bufferFactory.join(List.of(foo, bar));
WritableByteChannel channel = Files.newByteChannel(tempFile, StandardOpenOption.WRITE);
Flux<DataBuffer> writeResult = DataBufferUtils.write(Flux.just(joined), channel);
StepVerifier.create(writeResult)
.consumeNextWith(stringConsumer("foobar"))
.verifyComplete();
String result = String.join("", Files.readAllLines(tempFile));
assertThat(result).isEqualTo("foobar");
channel.close();
}
@ParameterizedDataBufferAllocatingTest
void writeWritableByteChannelErrorInFlux(DataBufferFactory bufferFactory) throws Exception {
super.bufferFactory = bufferFactory;
@@ -445,6 +466,27 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
assertThat(result).isEqualTo("foobarbazqux");
}
@ParameterizedDataBufferAllocatingTest
void writeAsynchronousFileChannelWithJoinedBuffer(DataBufferFactory bufferFactory) throws Exception {
super.bufferFactory = bufferFactory;
DataBuffer foo = stringBuffer("foo");
DataBuffer bar = stringBuffer("bar");
DataBuffer joined = bufferFactory.join(List.of(foo, bar));
AsynchronousFileChannel channel = AsynchronousFileChannel.open(tempFile, StandardOpenOption.WRITE);
Flux<DataBuffer> writeResult = DataBufferUtils.write(Flux.just(joined), channel);
StepVerifier.create(writeResult)
.consumeNextWith(stringConsumer("foobar"))
.verifyComplete();
String result = String.join("", Files.readAllLines(tempFile));
assertThat(result).isEqualTo("foobar");
channel.close();
}
@ParameterizedDataBufferAllocatingTest
void writeAsynchronousFileChannelErrorInFlux(DataBufferFactory bufferFactory) throws Exception {
super.bufferFactory = bufferFactory;