Add DataBuffers utility class

Prior to this commit, `DataBufferUtils` would implement many utility
static methods for managing and processing `DataBuffer` instances.
`DataBuffer` is tightly linked to the reactive space, but its usage
shoudn't be limited to applications that use `Publisher` extensively.

This commit gathers methods that do not depend on `Publisher` into a new
`DataBuffers` type that `DataBufferUtils` now extend. This allows to use
`DataBuffer` instances in a non-reactive case.

Fixes gh-37353
This commit is contained in:
Brian Clozel
2026-09-30 15:58:46 +02:00
parent 634d187130
commit 9a87fa6d12
12 changed files with 491 additions and 406 deletions
@@ -33,6 +33,7 @@ import reactor.core.publisher.Mono;
import org.springframework.core.ResolvableType;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferMatcher;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.LimitedDataBufferList;
import org.springframework.core.log.LogFormatUtils;
@@ -100,7 +101,7 @@ public abstract class AbstractCharSequenceDecoder<T extends CharSequence> extend
byte[][] delimiterBytes = getDelimiterBytes(mimeType);
LimitedDataBufferList chunks = new LimitedDataBufferList(getMaxInMemorySize());
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delimiterBytes);
DataBufferMatcher matcher = DataBufferUtils.matcher(delimiterBytes);
return Flux.from(input)
.concatMapIterable(buffer -> processDataBuffer(buffer, matcher, chunks))
@@ -127,7 +128,7 @@ public abstract class AbstractCharSequenceDecoder<T extends CharSequence> extend
});
}
private Collection<DataBuffer> processDataBuffer(DataBuffer buffer, DataBufferUtils.Matcher matcher,
private Collection<DataBuffer> processDataBuffer(DataBuffer buffer, DataBufferMatcher matcher,
LimitedDataBufferList chunks) {
boolean release = true;
@@ -357,7 +357,7 @@ public interface DataBuffer {
*/
@Deprecated(since = "6.0")
default DataBuffer retainedSlice(int index, int length) {
return DataBufferUtils.retain(slice(index, length));
return DataBuffers.retain(slice(index, length));
}
/**
@@ -100,7 +100,7 @@ final class DataBufferInputStream extends InputStream {
return;
}
if (this.releaseOnClose) {
DataBufferUtils.release(this.dataBuffer);
DataBuffers.release(this.dataBuffer);
}
this.closed = true;
}
@@ -0,0 +1,55 @@
/*
* Copyright 2002-present 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.core.io.buffer;
/**
* Contract to find delimiter(s) against one or more data buffers that can
* be passed one at a time to the {@link #match(DataBuffer)} method.
*
* @author Arjen Poutsma
* @since 7.1
* @see #of(byte[]...)
* @see #match(DataBuffer)
*/
public interface DataBufferMatcher {
/**
* Return a {@link DataBufferMatcher} for the given delimiters.
* @param delimiters the delimiter bytes to find
* @return the matcher
*/
static DataBufferMatcher of(byte[]... delimiters) {
return DataBuffers.matcher(delimiters);
}
/**
* Find the first matching delimiter and return the index of the last
* byte of the delimiter, or {@code -1} if not found.
*/
int match(DataBuffer dataBuffer);
/**
* Return the delimiter from the last invocation of {@link #match(DataBuffer)}.
*/
byte[] delimiter();
/**
* Reset the state of this matcher.
*/
void reset();
}
@@ -38,8 +38,6 @@ import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.jspecify.annotations.Nullable;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
@@ -59,16 +57,14 @@ import org.springframework.util.CollectionUtils;
/**
* Utility class for working with {@link DataBuffer DataBuffers}.
*
* <p>This class adds the methods based on Reactive Streams to the ones inherited
* from {@link DataBuffers}, which can be used without Reactor on the classpath.
*
* @author Arjen Poutsma
* @author Brian Clozel
* @since 5.0
*/
public abstract class DataBufferUtils {
private static final Log logger = LogFactory.getLog(DataBufferUtils.class);
private static final Consumer<DataBuffer> RELEASE_CONSUMER = DataBufferUtils::release;
public abstract class DataBufferUtils extends DataBuffers {
//---------------------------------------------------------------------
// Reading
@@ -555,86 +551,6 @@ public abstract class DataBufferUtils {
}).doOnDiscard(DataBuffer.class, DataBufferUtils::release);
}
/**
* Retain the given data buffer, if it is a {@link PooledDataBuffer}.
* @param dataBuffer the data buffer to retain
* @return the retained buffer
*/
@SuppressWarnings("unchecked")
public static <T extends DataBuffer> T retain(T dataBuffer) {
if (dataBuffer instanceof PooledDataBuffer pooledDataBuffer) {
return (T) pooledDataBuffer.retain();
}
else {
return dataBuffer;
}
}
/**
* Associate the given hint with the data buffer if it is a pooled buffer
* and supports leak tracking.
* @param dataBuffer the data buffer to attach the hint to
* @param hint the hint to attach
* @return the input buffer
* @since 5.3.2
*/
@SuppressWarnings("unchecked")
public static <T extends DataBuffer> T touch(T dataBuffer, Object hint) {
if (dataBuffer instanceof TouchableDataBuffer touchableDataBuffer) {
return (T) touchableDataBuffer.touch(hint);
}
else {
return dataBuffer;
}
}
/**
* Release the given data buffer. If it is a {@link PooledDataBuffer} and
* has been {@linkplain PooledDataBuffer#isAllocated() allocated}, this
* method will call {@link PooledDataBuffer#release()}. If it is a
* {@link CloseableDataBuffer}, this method will call
* {@link CloseableDataBuffer#close()}.
* @param dataBuffer the data buffer to release
* @return {@code true} if the buffer was released; {@code false} otherwise.
*/
public static boolean release(@Nullable DataBuffer dataBuffer) {
if (dataBuffer instanceof PooledDataBuffer pooledDataBuffer) {
if (pooledDataBuffer.isAllocated()) {
try {
return pooledDataBuffer.release();
}
catch (IllegalStateException ex) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to release PooledDataBuffer: " + dataBuffer, ex);
}
return false;
}
}
}
else if (dataBuffer instanceof CloseableDataBuffer closeableDataBuffer) {
try {
closeableDataBuffer.close();
return true;
}
catch (IllegalStateException ex) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to release CloseableDataBuffer " + dataBuffer, ex);
}
return false;
}
}
return false;
}
/**
* Return a consumer that calls {@link #release(DataBuffer)} on all
* passed data buffers.
*/
public static Consumer<DataBuffer> releaseConsumer() {
return RELEASE_CONSUMER;
}
/**
* Return a new {@code DataBuffer} composed of joining together the given
* {@code dataBuffers} elements. Depending on the {@link DataBuffer} type,
@@ -680,291 +596,16 @@ public abstract class DataBufferUtils {
.doOnDiscard(DataBuffer.class, DataBufferUtils::release);
}
/**
* Return a {@link Matcher} for the given delimiter.
* The matcher can be used to find the delimiters in a stream of data buffers.
* @param delimiter the delimiter bytes to find
* @return the matcher
* @since 5.2
*/
public static Matcher matcher(byte[] delimiter) {
return createMatcher(delimiter);
}
/**
* Return a {@link Matcher} for the given delimiters.
* The matcher can be used to find the delimiters in a stream of data buffers.
* @param delimiters the delimiters bytes to find
* @return the matcher
* @since 5.2
*/
public static Matcher matcher(byte[]... delimiters) {
Assert.isTrue(delimiters.length > 0, "Delimiters must not be empty");
return (delimiters.length == 1 ? createMatcher(delimiters[0]) : new CompositeMatcher(delimiters));
}
private static NestedMatcher createMatcher(byte[] delimiter) {
// extract length due to Eclipse IDE compiler error in switch expression
int length = delimiter.length;
Assert.isTrue(length > 0, "Delimiter must not be empty");
return switch (length) {
case 1 -> (delimiter[0] == 10 ? SingleByteMatcher.NEWLINE_MATCHER : new SingleByteMatcher(delimiter));
case 2 -> new TwoByteMatcher(delimiter);
default -> new KnuthMorrisPrattMatcher(delimiter);
};
}
/**
* Contract to find delimiter(s) against one or more data buffers that can
* be passed one at a time to the {@link #match(DataBuffer)} method.
*
* @since 5.2
* @see #match(DataBuffer)
* @deprecated as of 7.1 in favor of {@link DataBufferMatcher}
*/
public interface Matcher {
/**
* Find the first matching delimiter and return the index of the last
* byte of the delimiter, or {@code -1} if not found.
*/
int match(DataBuffer dataBuffer);
/**
* Return the delimiter from the last invocation of {@link #match(DataBuffer)}.
*/
byte[] delimiter();
/**
* Reset the state of this matcher.
*/
void reset();
}
/**
* Matcher that supports searching for multiple delimiters.
*/
private static class CompositeMatcher implements Matcher {
private static final byte[] NO_DELIMITER = new byte[0];
private final NestedMatcher[] matchers;
byte[] longestDelimiter = NO_DELIMITER;
CompositeMatcher(byte[][] delimiters) {
this.matchers = initMatchers(delimiters);
}
private static NestedMatcher[] initMatchers(byte[][] delimiters) {
NestedMatcher[] matchers = new NestedMatcher[delimiters.length];
for (int i = 0; i < delimiters.length; i++) {
matchers[i] = createMatcher(delimiters[i]);
}
return matchers;
}
@Override
public int match(DataBuffer dataBuffer) {
this.longestDelimiter = NO_DELIMITER;
for (int pos = dataBuffer.readPosition(); pos < dataBuffer.writePosition(); pos++) {
byte b = dataBuffer.getByte(pos);
for (NestedMatcher matcher : this.matchers) {
if (matcher.match(b) && matcher.delimiter().length > this.longestDelimiter.length) {
this.longestDelimiter = matcher.delimiter();
}
}
if (this.longestDelimiter != NO_DELIMITER) {
reset();
return pos;
}
}
return -1;
}
@Override
public byte[] delimiter() {
Assert.state(this.longestDelimiter != NO_DELIMITER, "'delimiter' not set");
return this.longestDelimiter;
}
@Override
public void reset() {
for (NestedMatcher matcher : this.matchers) {
matcher.reset();
}
}
}
/**
* Matcher that can be nested within {@link CompositeMatcher} where multiple
* matchers advance together using the same index, one byte at a time.
*/
private interface NestedMatcher extends Matcher {
/**
* Perform a match against the next byte of the stream and return true
* if the delimiter is fully matched.
*/
boolean match(byte b);
}
/**
* Matcher for a single byte delimiter.
*/
private static class SingleByteMatcher implements NestedMatcher {
static final SingleByteMatcher NEWLINE_MATCHER = new SingleByteMatcher(new byte[] {10});
private final byte[] delimiter;
SingleByteMatcher(byte[] delimiter) {
Assert.isTrue(delimiter.length == 1, "Expected a 1 byte delimiter");
this.delimiter = delimiter;
}
@Override
public int match(DataBuffer dataBuffer) {
int start = dataBuffer.readPosition();
int end = dataBuffer.writePosition();
return dataBuffer.forEachByte(start, end - start, b -> !this.match(b));
}
@Override
public boolean match(byte b) {
return this.delimiter[0] == b;
}
@Override
public byte[] delimiter() {
return this.delimiter;
}
@Override
public void reset() {
}
}
/**
* Base class for a {@link NestedMatcher}.
*/
private abstract static class AbstractNestedMatcher implements NestedMatcher {
private final byte[] delimiter;
private int matches = 0;
protected AbstractNestedMatcher(byte[] delimiter) {
this.delimiter = delimiter;
}
protected void setMatches(int index) {
this.matches = index;
}
protected int getMatches() {
return this.matches;
}
@Override
public int match(DataBuffer dataBuffer) {
int start = dataBuffer.readPosition();
int end = dataBuffer.writePosition();
int matchPosition = dataBuffer.forEachByte(start, end - start, b -> !this.match(b));
if (matchPosition != -1) {
reset();
}
return matchPosition;
}
@Override
public boolean match(byte b) {
if (b == this.delimiter[this.matches]) {
this.matches++;
return (this.matches == delimiter().length);
}
return false;
}
@Override
public byte[] delimiter() {
return this.delimiter;
}
@Override
public void reset() {
this.matches = 0;
}
}
/**
* Matcher with a 2 byte delimiter that does not benefit from a
* Knuth-Morris-Pratt suffix-prefix table.
*/
private static class TwoByteMatcher extends AbstractNestedMatcher {
protected TwoByteMatcher(byte[] delimiter) {
super(delimiter);
Assert.isTrue(delimiter.length == 2, "Expected a 2-byte delimiter");
}
@Override
public boolean match(byte b) {
if (getMatches() > 0 && b != delimiter()[getMatches()]) {
setMatches(0);
}
return super.match(b);
}
}
/**
* Implementation of {@link Matcher} that uses the Knuth-Morris-Pratt algorithm.
* @see <a href="https://www.nayuki.io/page/knuth-morris-pratt-string-matching">Knuth-Morris-Pratt string matching</a>
*/
private static class KnuthMorrisPrattMatcher extends AbstractNestedMatcher {
private final int[] table;
public KnuthMorrisPrattMatcher(byte[] delimiter) {
super(delimiter);
this.table = longestSuffixPrefixTable(delimiter);
}
private static int[] longestSuffixPrefixTable(byte[] delimiter) {
int[] result = new int[delimiter.length];
result[0] = 0;
for (int i = 1; i < delimiter.length; i++) {
int j = result[i - 1];
while (j > 0 && delimiter[i] != delimiter[j]) {
j = result[j - 1];
}
if (delimiter[i] == delimiter[j]) {
j++;
}
result[i] = j;
}
return result;
}
@Override
public boolean match(byte b) {
while (getMatches() > 0 && b != delimiter()[getMatches()]) {
setMatches(this.table[getMatches() - 1]);
}
return super.match(b);
}
@Deprecated(since = "7.1", forRemoval = true)
public interface Matcher extends DataBufferMatcher {
}
@@ -0,0 +1,386 @@
/*
* Copyright 2002-present 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.
* You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.core.io.buffer;
import java.util.function.Consumer;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.jspecify.annotations.Nullable;
import org.springframework.util.Assert;
/**
* Utility methods for working with {@link DataBuffer DataBuffers} that do not
* require Reactor or Reactive Streams on the classpath, and can therefore be
* used in any application.
*
* <p>{@link DataBufferUtils} extends this class and adds the methods that
* are based on Reactive Streams. This class must not depend on it.
*
* @author Arjen Poutsma
* @author Brian Clozel
* @since 7.1
*/
public abstract class DataBuffers {
// Same category as DataBufferUtils, without loading that class (and Reactive Streams)
private static final Log logger = LogFactory.getLog("org.springframework.core.io.buffer.DataBufferUtils");
private static final Consumer<DataBuffer> RELEASE_CONSUMER = DataBuffers::release;
/**
* Retain the given data buffer, if it is a {@link PooledDataBuffer}.
* @param dataBuffer the data buffer to retain
* @return the retained buffer
*/
@SuppressWarnings("unchecked")
public static <T extends DataBuffer> T retain(T dataBuffer) {
if (dataBuffer instanceof PooledDataBuffer pooledDataBuffer) {
return (T) pooledDataBuffer.retain();
}
else {
return dataBuffer;
}
}
/**
* Release the given data buffer. If it is a {@link PooledDataBuffer} and
* has been {@linkplain PooledDataBuffer#isAllocated() allocated}, this
* method will call {@link PooledDataBuffer#release()}. If it is a
* {@link CloseableDataBuffer}, this method will call
* {@link CloseableDataBuffer#close()}.
* @param dataBuffer the data buffer to release
* @return {@code true} if the buffer was released; {@code false} otherwise.
*/
public static boolean release(@Nullable DataBuffer dataBuffer) {
if (dataBuffer instanceof PooledDataBuffer pooledDataBuffer) {
if (pooledDataBuffer.isAllocated()) {
try {
return pooledDataBuffer.release();
}
catch (IllegalStateException ex) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to release PooledDataBuffer: " + dataBuffer, ex);
}
return false;
}
}
}
else if (dataBuffer instanceof CloseableDataBuffer closeableDataBuffer) {
try {
closeableDataBuffer.close();
return true;
}
catch (IllegalStateException ex) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to release CloseableDataBuffer " + dataBuffer, ex);
}
return false;
}
}
return false;
}
/**
* Associate the given hint with the data buffer if it is a pooled buffer
* and supports leak tracking.
* @param dataBuffer the data buffer to attach the hint to
* @param hint the hint to attach
* @return the input buffer
* @since 5.3.2
*/
@SuppressWarnings("unchecked")
public static <T extends DataBuffer> T touch(T dataBuffer, Object hint) {
if (dataBuffer instanceof TouchableDataBuffer touchableDataBuffer) {
return (T) touchableDataBuffer.touch(hint);
}
else {
return dataBuffer;
}
}
/**
* Return a consumer that calls {@link #release(DataBuffer)} on all
* passed data buffers.
*/
public static Consumer<DataBuffer> releaseConsumer() {
return RELEASE_CONSUMER;
}
/**
* Return a {@link DataBufferMatcher} for the given delimiter.
* The matcher can be used to find the delimiters in a stream of data buffers.
* @param delimiter the delimiter bytes to find
* @return the matcher
* @since 5.2
*/
public static DataBufferMatcher matcher(byte[] delimiter) {
return matcher(new byte[][] {delimiter});
}
/**
* Return a {@link DataBufferMatcher} for the given delimiters.
* The matcher can be used to find the delimiters in a stream of data buffers.
* @param delimiters the delimiters bytes to find
* @return the matcher
* @since 5.2
*/
public static DataBufferMatcher matcher(byte[]... delimiters) {
Assert.isTrue(delimiters.length > 0, "Delimiters must not be empty");
return (delimiters.length == 1 ? createMatcher(delimiters[0]) : new CompositeMatcher(delimiters));
}
private static NestedMatcher createMatcher(byte[] delimiter) {
// extract length due to Eclipse IDE compiler error in switch expression
int length = delimiter.length;
Assert.isTrue(length > 0, "Delimiter must not be empty");
return switch (length) {
case 1 -> (delimiter[0] == 10 ? SingleByteMatcher.NEWLINE_MATCHER : new SingleByteMatcher(delimiter));
case 2 -> new TwoByteMatcher(delimiter);
default -> new KnuthMorrisPrattMatcher(delimiter);
};
}
/**
* Matcher that supports searching for multiple delimiters.
*/
private static class CompositeMatcher implements DataBufferMatcher {
private static final byte[] NO_DELIMITER = new byte[0];
private final NestedMatcher[] matchers;
byte[] longestDelimiter = NO_DELIMITER;
CompositeMatcher(byte[][] delimiters) {
this.matchers = initMatchers(delimiters);
}
private static NestedMatcher[] initMatchers(byte[][] delimiters) {
NestedMatcher[] matchers = new NestedMatcher[delimiters.length];
for (int i = 0; i < delimiters.length; i++) {
matchers[i] = createMatcher(delimiters[i]);
}
return matchers;
}
@Override
public int match(DataBuffer dataBuffer) {
this.longestDelimiter = NO_DELIMITER;
for (int pos = dataBuffer.readPosition(); pos < dataBuffer.writePosition(); pos++) {
byte b = dataBuffer.getByte(pos);
for (NestedMatcher matcher : this.matchers) {
if (matcher.match(b) && matcher.delimiter().length > this.longestDelimiter.length) {
this.longestDelimiter = matcher.delimiter();
}
}
if (this.longestDelimiter != NO_DELIMITER) {
reset();
return pos;
}
}
return -1;
}
@Override
public byte[] delimiter() {
Assert.state(this.longestDelimiter != NO_DELIMITER, "'delimiter' not set");
return this.longestDelimiter;
}
@Override
public void reset() {
for (NestedMatcher matcher : this.matchers) {
matcher.reset();
}
}
}
/**
* Matcher that can be nested within {@link CompositeMatcher} where multiple
* matchers advance together using the same index, one byte at a time.
*/
private interface NestedMatcher extends DataBufferMatcher {
/**
* Perform a match against the next byte of the stream and return true
* if the delimiter is fully matched.
*/
boolean match(byte b);
}
/**
* Matcher for a single byte delimiter.
*/
private static class SingleByteMatcher implements NestedMatcher {
static final SingleByteMatcher NEWLINE_MATCHER = new SingleByteMatcher(new byte[] {10});
private final byte[] delimiter;
SingleByteMatcher(byte[] delimiter) {
Assert.isTrue(delimiter.length == 1, "Expected a 1 byte delimiter");
this.delimiter = delimiter;
}
@Override
public int match(DataBuffer dataBuffer) {
int start = dataBuffer.readPosition();
int end = dataBuffer.writePosition();
return dataBuffer.forEachByte(start, end - start, b -> !this.match(b));
}
@Override
public boolean match(byte b) {
return this.delimiter[0] == b;
}
@Override
public byte[] delimiter() {
return this.delimiter;
}
@Override
public void reset() {
}
}
/**
* Base class for a {@link NestedMatcher}.
*/
private abstract static class AbstractNestedMatcher implements NestedMatcher {
private final byte[] delimiter;
private int matches = 0;
protected AbstractNestedMatcher(byte[] delimiter) {
this.delimiter = delimiter;
}
protected void setMatches(int index) {
this.matches = index;
}
protected int getMatches() {
return this.matches;
}
@Override
public int match(DataBuffer dataBuffer) {
int start = dataBuffer.readPosition();
int end = dataBuffer.writePosition();
int matchPosition = dataBuffer.forEachByte(start, end - start, b -> !this.match(b));
if (matchPosition != -1) {
reset();
}
return matchPosition;
}
@Override
public boolean match(byte b) {
if (b == this.delimiter[this.matches]) {
this.matches++;
return (this.matches == delimiter().length);
}
return false;
}
@Override
public byte[] delimiter() {
return this.delimiter;
}
@Override
public void reset() {
this.matches = 0;
}
}
/**
* Matcher with a 2 byte delimiter that does not benefit from a
* Knuth-Morris-Pratt suffix-prefix table.
*/
private static class TwoByteMatcher extends AbstractNestedMatcher {
protected TwoByteMatcher(byte[] delimiter) {
super(delimiter);
Assert.isTrue(delimiter.length == 2, "Expected a 2-byte delimiter");
}
@Override
public boolean match(byte b) {
if (getMatches() > 0 && b != delimiter()[getMatches()]) {
setMatches(0);
}
return super.match(b);
}
}
/**
* Implementation of {@link DataBufferMatcher} that uses the Knuth-Morris-Pratt algorithm.
* @see <a href="https://www.nayuki.io/page/knuth-morris-pratt-string-matching">Knuth-Morris-Pratt string matching</a>
*/
private static class KnuthMorrisPrattMatcher extends AbstractNestedMatcher {
private final int[] table;
public KnuthMorrisPrattMatcher(byte[] delimiter) {
super(delimiter);
this.table = longestSuffixPrefixTable(delimiter);
}
private static int[] longestSuffixPrefixTable(byte[] delimiter) {
int[] result = new int[delimiter.length];
result[0] = 0;
for (int i = 1; i < delimiter.length; i++) {
int j = result[i - 1];
while (j > 0 && delimiter[i] != delimiter[j]) {
j = result[j - 1];
}
if (delimiter[i] == delimiter[j]) {
j++;
}
result[i] = j;
}
return result;
}
@Override
public boolean match(byte b) {
while (getMatches() > 0 && b != delimiter()[getMatches()]) {
setMatches(this.table[getMatches() - 1]);
}
return super.match(b);
}
}
}
@@ -119,7 +119,7 @@ public class DefaultDataBufferFactory implements DataBufferFactory {
int capacity = dataBuffers.stream().mapToInt(DataBuffer::readableByteCount).sum();
DefaultDataBuffer result = allocateBuffer(capacity);
dataBuffers.forEach(result::write);
dataBuffers.forEach(DataBufferUtils::release);
dataBuffers.forEach(DataBuffers::release);
return result;
}
@@ -143,7 +143,7 @@ public class LimitedDataBufferList extends ArrayList<DataBuffer> {
public void releaseAndClear() {
forEach(buf -> {
try {
DataBufferUtils.release(buf);
DataBuffers.release(buf);
}
catch (Throwable ex) {
// Keep going..
@@ -1393,7 +1393,7 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
DataBuffer bar = stringBuffer("bar");
byte[] delims = "ooba".getBytes(StandardCharsets.UTF_8);
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delims);
DataBufferMatcher matcher = DataBufferUtils.matcher(delims);
int result = matcher.match(foo);
assertThat(result).isEqualTo(-1);
result = matcher.match(bar);
@@ -1410,7 +1410,7 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
DataBuffer foo = stringBuffer("foooobar");
byte[] delims = "oo".getBytes(StandardCharsets.UTF_8);
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delims);
DataBufferMatcher matcher = DataBufferUtils.matcher(delims);
int endIndex = matcher.match(foo);
assertThat(endIndex).isEqualTo(2);
foo.readPosition(endIndex + 1);
@@ -1430,7 +1430,7 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
DataBuffer foo = stringBuffer("foooobar");
byte[] delims = "oo".getBytes(StandardCharsets.UTF_8);
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delims);
DataBufferMatcher matcher = DataBufferUtils.matcher(delims);
int endIndex = matcher.match(foo);
assertThat(endIndex).isEqualTo(2);
foo.readPosition(endIndex + 1);
@@ -1450,7 +1450,7 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
DataBuffer buffer = stringBuffer("a\rXY\nb");
byte[] delims = "\r\n".getBytes(StandardCharsets.UTF_8);
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delims);
DataBufferMatcher matcher = DataBufferUtils.matcher(delims);
int result = matcher.match(buffer);
assertThat(result).isEqualTo(-1);
@@ -1464,7 +1464,7 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
DataBuffer buffer = stringBuffer("a\r\nb");
byte[] delims = "\r\n".getBytes(StandardCharsets.UTF_8);
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delims);
DataBufferMatcher matcher = DataBufferUtils.matcher(delims);
int result = matcher.match(buffer);
assertThat(result).isEqualTo(2);
@@ -1478,7 +1478,7 @@ class DataBufferUtilsTests extends AbstractDataBufferAllocatingTests {
DataBuffer buffer = stringBuffer("a\r\rX\nb");
byte[] delims = "\r\n".getBytes(StandardCharsets.UTF_8);
DataBufferUtils.Matcher matcher = DataBufferUtils.matcher(delims);
DataBufferMatcher matcher = DataBufferUtils.matcher(delims);
int result = matcher.match(buffer);
assertThat(result).isEqualTo(-1);
@@ -40,6 +40,7 @@ import reactor.util.context.Context;
import org.springframework.core.codec.DecodingException;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferLimitException;
import org.springframework.core.io.buffer.DataBufferMatcher;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.http.HttpHeaders;
@@ -306,7 +307,7 @@ final class MultipartParser extends BaseSubscriber<DataBuffer> {
*/
private final class PreambleState implements State {
private final DataBufferUtils.Matcher firstBoundary;
private final DataBufferMatcher firstBoundary;
public PreambleState() {
@@ -359,7 +360,7 @@ final class MultipartParser extends BaseSubscriber<DataBuffer> {
*/
private final class HeadersState implements State {
private final DataBufferUtils.Matcher endHeaders = DataBufferUtils.matcher(MultipartUtils.concat(CR_LF, CR_LF));
private final DataBufferMatcher endHeaders = DataBufferUtils.matcher(MultipartUtils.concat(CR_LF, CR_LF));
private final AtomicInteger byteCount = new AtomicInteger();
@@ -504,7 +505,7 @@ final class MultipartParser extends BaseSubscriber<DataBuffer> {
*/
private final class BodyState implements State {
private final DataBufferUtils.Matcher boundary;
private final DataBufferMatcher boundary;
private final int boundaryLength;
@@ -31,7 +31,8 @@ import org.jspecify.annotations.Nullable;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferLimitException;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.DataBufferMatcher;
import org.springframework.core.io.buffer.DataBuffers;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.http.HttpHeaders;
import org.springframework.http.converter.HttpMessageConversionException;
@@ -113,7 +114,7 @@ final class MultipartParser {
newState.data(remainder);
}
else {
DataBufferUtils.release(remainder);
DataBuffers.release(remainder);
}
}
}
@@ -255,11 +256,11 @@ final class MultipartParser {
*/
private final class PreambleState implements State {
private final DataBufferUtils.Matcher firstBoundary;
private final DataBufferMatcher firstBoundary;
PreambleState() {
this.firstBoundary = DataBufferUtils.matcher(concat(TWO_HYPHENS, MultipartParser.this.boundary));
this.firstBoundary = DataBufferMatcher.of(concat(TWO_HYPHENS, MultipartParser.this.boundary));
}
/**
@@ -275,11 +276,11 @@ final class MultipartParser {
logger.trace("First boundary found @" + endIdx + " in " + buf);
}
DataBuffer preambleBuffer = buf.split(endIdx + 1);
DataBufferUtils.release(preambleBuffer);
DataBuffers.release(preambleBuffer);
changeState(new HeadersState(), buf);
}
else {
DataBufferUtils.release(buf);
DataBuffers.release(buf);
}
}
@@ -302,7 +303,7 @@ final class MultipartParser {
*/
private final class HeadersState implements State {
private final DataBufferUtils.Matcher endHeaders = DataBufferUtils.matcher(concat(CR_LF, CR_LF));
private final DataBufferMatcher endHeaders = DataBufferMatcher.of(concat(CR_LF, CR_LF));
private final List<DataBuffer> buffers = new ArrayList<>();
@@ -405,7 +406,7 @@ final class MultipartParser {
DataBuffer joined = this.buffers.get(0).factory().join(this.buffers);
this.buffers.clear();
String string = joined.toString(MultipartParser.this.headersCharset);
DataBufferUtils.release(joined);
DataBuffers.release(joined);
String[] lines = string.split(HEADER_ENTRY_SEPARATOR);
HttpHeaders result = new HttpHeaders();
for (String line : lines) {
@@ -429,7 +430,7 @@ final class MultipartParser {
@Override
public void dispose() {
this.buffers.forEach(DataBufferUtils::release);
this.buffers.forEach(DataBuffers::release);
}
@Override
@@ -446,7 +447,7 @@ final class MultipartParser {
*/
private final class BodyState implements State {
private final DataBufferUtils.Matcher boundaryMatcher;
private final DataBufferMatcher boundaryMatcher;
private final int boundaryLength;
@@ -454,7 +455,7 @@ final class MultipartParser {
public BodyState() {
byte[] delimiter = concat(CR_LF, TWO_HYPHENS, MultipartParser.this.boundary);
this.boundaryMatcher = DataBufferUtils.matcher(delimiter);
this.boundaryMatcher = DataBufferMatcher.of(delimiter);
this.boundaryLength = delimiter.length;
}
@@ -479,14 +480,14 @@ final class MultipartParser {
// whole boundary in buffer.
// slice off the body part, and flush
DataBuffer body = boundaryBuffer.split(len);
DataBufferUtils.release(boundaryBuffer);
DataBuffers.release(boundaryBuffer);
enqueue(body);
flush();
}
else if (len < 0) {
// boundary spans multiple buffers, and we've just found the end
// iterate over buffers in reverse order
DataBufferUtils.release(boundaryBuffer);
DataBuffers.release(boundaryBuffer);
DataBuffer prev;
boolean found = false;
while ((prev = this.queue.pollLast()) != null) {
@@ -495,7 +496,7 @@ final class MultipartParser {
if (prevLen >= 0) {
// slice body part of previous buffer, and flush it
DataBuffer body = prev.split(prevLen + prev.readPosition());
DataBufferUtils.release(prev);
DataBuffers.release(prev);
enqueue(body);
flush();
found = true;
@@ -503,7 +504,7 @@ final class MultipartParser {
}
else {
// previous buffer only contains boundary bytes
DataBufferUtils.release(prev);
DataBuffers.release(prev);
len += prevByteCount;
}
}
@@ -514,7 +515,7 @@ final class MultipartParser {
}
else /* if (len == 0) */ {
// buffer starts with complete delimiter, flush out the previous buffers
DataBufferUtils.release(boundaryBuffer);
DataBuffers.release(boundaryBuffer);
if (this.queue.isEmpty()) {
// nothing was ever buffered for this part: the part had an empty body
invokeListener(buffer.factory().allocateBuffer(0), true);
@@ -581,7 +582,7 @@ final class MultipartParser {
@Override
public void dispose() {
this.queue.forEach(DataBufferUtils::release);
this.queue.forEach(DataBuffers::release);
this.queue.clear();
}
@@ -605,7 +606,7 @@ final class MultipartParser {
@Override
public void data(DataBuffer buf) {
DataBufferUtils.release(buf);
DataBuffers.release(buf);
}
@Override
@@ -31,7 +31,7 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.DataBuffers;
import org.springframework.core.io.buffer.DefaultDataBuffer;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.http.HttpHeaders;
@@ -201,7 +201,7 @@ final class PartGenerator implements MultipartParser.PartListener {
@Override
public void onBody(DataBuffer dataBuffer, boolean last) {
DataBufferUtils.release(dataBuffer);
DataBuffers.release(dataBuffer);
throw new HttpMessageConversionException("Body token not expected");
}
@@ -235,7 +235,7 @@ final class PartGenerator implements MultipartParser.PartListener {
store(dataBuffer);
}
else {
DataBufferUtils.release(dataBuffer);
DataBuffers.release(dataBuffer);
throw new HttpMessageConversionException("Form field value exceeded the memory usage limit of " +
PartGenerator.this.maxInMemorySize + " bytes");
}
@@ -265,7 +265,7 @@ final class PartGenerator implements MultipartParser.PartListener {
throw new HttpMessageConversionException("Cannot store multipart body", ex);
}
finally {
DataBufferUtils.release(dataBuffer);
DataBuffers.release(dataBuffer);
}
}
@@ -335,7 +335,7 @@ final class PartGenerator implements MultipartParser.PartListener {
int len = buffer.readableByteCount();
buffer.read(bytes, idx, len);
idx += len;
DataBufferUtils.release(buffer);
DataBuffers.release(buffer);
}
this.content.clear();
DefaultDataBuffer content = DefaultDataBufferFactory.sharedInstance.wrap(bytes);
@@ -345,7 +345,7 @@ final class PartGenerator implements MultipartParser.PartListener {
@Override
public void dispose() {
this.content.forEach(DataBufferUtils::release);
this.content.forEach(DataBuffers::release);
}
@Override
@@ -446,7 +446,7 @@ final class PartGenerator implements MultipartParser.PartListener {
throw new UncheckedIOException("Could not write to temp file ", exc);
}
finally {
DataBufferUtils.release(dataBuffer);
DataBuffers.release(dataBuffer);
}
}