From 9a87fa6d12ea4536e9bd72d813cac1859ab11759 Mon Sep 17 00:00:00 2001 From: Brian Clozel Date: Wed, 30 Sep 2026 15:58:46 +0200 Subject: [PATCH] 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 --- .../codec/AbstractCharSequenceDecoder.java | 5 +- .../core/io/buffer/DataBuffer.java | 2 +- .../core/io/buffer/DataBufferInputStream.java | 2 +- .../core/io/buffer/DataBufferMatcher.java | 55 +++ .../core/io/buffer/DataBufferUtils.java | 373 +---------------- .../core/io/buffer/DataBuffers.java | 386 ++++++++++++++++++ .../io/buffer/DefaultDataBufferFactory.java | 2 +- .../core/io/buffer/LimitedDataBufferList.java | 2 +- .../core/io/buffer/DataBufferUtilsTests.java | 12 +- .../http/codec/multipart/MultipartParser.java | 7 +- .../converter/multipart/MultipartParser.java | 37 +- .../converter/multipart/PartGenerator.java | 14 +- 12 files changed, 491 insertions(+), 406 deletions(-) create mode 100644 spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferMatcher.java create mode 100644 spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffers.java diff --git a/spring-core/src/main/java/org/springframework/core/codec/AbstractCharSequenceDecoder.java b/spring-core/src/main/java/org/springframework/core/codec/AbstractCharSequenceDecoder.java index 1918d7182db..99fc857c502 100644 --- a/spring-core/src/main/java/org/springframework/core/codec/AbstractCharSequenceDecoder.java +++ b/spring-core/src/main/java/org/springframework/core/codec/AbstractCharSequenceDecoder.java @@ -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 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 extend }); } - private Collection processDataBuffer(DataBuffer buffer, DataBufferUtils.Matcher matcher, + private Collection processDataBuffer(DataBuffer buffer, DataBufferMatcher matcher, LimitedDataBufferList chunks) { boolean release = true; diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffer.java b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffer.java index 0ac03091a48..4fa912495f9 100644 --- a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffer.java +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffer.java @@ -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)); } /** diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferInputStream.java b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferInputStream.java index 386e156851d..5d4b3c95fa2 100644 --- a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferInputStream.java +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferInputStream.java @@ -100,7 +100,7 @@ final class DataBufferInputStream extends InputStream { return; } if (this.releaseOnClose) { - DataBufferUtils.release(this.dataBuffer); + DataBuffers.release(this.dataBuffer); } this.closed = true; } diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferMatcher.java b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferMatcher.java new file mode 100644 index 00000000000..fd1e61c0732 --- /dev/null +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferMatcher.java @@ -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(); + +} diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferUtils.java b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferUtils.java index 6b9b1014ce9..0d6ff4f21cb 100644 --- a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferUtils.java +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBufferUtils.java @@ -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}. * + *

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 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 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 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 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 Knuth-Morris-Pratt string matching - */ - 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 { } diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffers.java b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffers.java new file mode 100644 index 00000000000..1381b6df0ef --- /dev/null +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/DataBuffers.java @@ -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. + * + *

{@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 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 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 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 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 Knuth-Morris-Pratt string matching + */ + 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); + } + } + +} diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/DefaultDataBufferFactory.java b/spring-core/src/main/java/org/springframework/core/io/buffer/DefaultDataBufferFactory.java index e2e1ce9b4df..8860a14955e 100644 --- a/spring-core/src/main/java/org/springframework/core/io/buffer/DefaultDataBufferFactory.java +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/DefaultDataBufferFactory.java @@ -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; } diff --git a/spring-core/src/main/java/org/springframework/core/io/buffer/LimitedDataBufferList.java b/spring-core/src/main/java/org/springframework/core/io/buffer/LimitedDataBufferList.java index 912a406a8f2..b37e4a07878 100644 --- a/spring-core/src/main/java/org/springframework/core/io/buffer/LimitedDataBufferList.java +++ b/spring-core/src/main/java/org/springframework/core/io/buffer/LimitedDataBufferList.java @@ -143,7 +143,7 @@ public class LimitedDataBufferList extends ArrayList { public void releaseAndClear() { forEach(buf -> { try { - DataBufferUtils.release(buf); + DataBuffers.release(buf); } catch (Throwable ex) { // Keep going.. diff --git a/spring-core/src/test/java/org/springframework/core/io/buffer/DataBufferUtilsTests.java b/spring-core/src/test/java/org/springframework/core/io/buffer/DataBufferUtilsTests.java index 8885adef9a0..abff1e95b51 100644 --- a/spring-core/src/test/java/org/springframework/core/io/buffer/DataBufferUtilsTests.java +++ b/spring-core/src/test/java/org/springframework/core/io/buffer/DataBufferUtilsTests.java @@ -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); diff --git a/spring-web/src/main/java/org/springframework/http/codec/multipart/MultipartParser.java b/spring-web/src/main/java/org/springframework/http/codec/multipart/MultipartParser.java index 879b05c41a5..cd339298384 100644 --- a/spring-web/src/main/java/org/springframework/http/codec/multipart/MultipartParser.java +++ b/spring-web/src/main/java/org/springframework/http/codec/multipart/MultipartParser.java @@ -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 { */ 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 { */ 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 { */ private final class BodyState implements State { - private final DataBufferUtils.Matcher boundary; + private final DataBufferMatcher boundary; private final int boundaryLength; diff --git a/spring-web/src/main/java/org/springframework/http/converter/multipart/MultipartParser.java b/spring-web/src/main/java/org/springframework/http/converter/multipart/MultipartParser.java index a76a7f482b0..dd942eb032d 100644 --- a/spring-web/src/main/java/org/springframework/http/converter/multipart/MultipartParser.java +++ b/spring-web/src/main/java/org/springframework/http/converter/multipart/MultipartParser.java @@ -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 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 diff --git a/spring-web/src/main/java/org/springframework/http/converter/multipart/PartGenerator.java b/spring-web/src/main/java/org/springframework/http/converter/multipart/PartGenerator.java index 479ba37e0b7..2da203a853d 100644 --- a/spring-web/src/main/java/org/springframework/http/converter/multipart/PartGenerator.java +++ b/spring-web/src/main/java/org/springframework/http/converter/multipart/PartGenerator.java @@ -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); } }