From 4bb80039eb9d69384adf24fbc5186e5d9bd9b082 Mon Sep 17 00:00:00 2001 From: Bikram Sharma Date: Thu, 24 Sep 2026 12:54:11 -0700 Subject: [PATCH 1/4] fix: don't signal onComplete when a chunk is fully consumed by the range skip AdjustedRangeSubscriber called onComplete (and fell through to an NPE) when numBytesToSkip exceeded a chunk, so a small first chunk ended ranged GETs with an empty stream. Consume the chunk toward the skip and wait for the next one instead. Fixes #517. --- .../internal/AdjustedRangeSubscriber.java | 14 +- .../internal/AdjustedRangeSubscriberTest.java | 124 ++++++++++++++++++ 2 files changed, 134 insertions(+), 4 deletions(-) create mode 100644 src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java diff --git a/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java b/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java index 73de3b214..011497d64 100644 --- a/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java +++ b/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java @@ -61,11 +61,17 @@ public void onNext(ByteBuffer byteBuffer) { if (numBytesToSkip != 0) { byte[] buf = byteBuffer.array(); - if (numBytesToSkip > buf.length) { - // If we need to skip past the available data, - // we are returning nothing, so signal completion + if (numBytesToSkip >= buf.length) { + // This chunk is entirely consumed by the leading-byte skip, so there is no + // in-range data to deliver yet. We must still signal the wrapped subscriber to + // keep the reactive-streams demand flowing: downstream drives the stream one + // element at a time (request(1)), and it only requests the next element after it + // receives an onNext. Returning without signaling would leave it waiting forever. + // Forward an empty buffer, mirroring CipherSubscriber's "avoid blocking" idiom, + // and wait for the next chunk instead of completing the stream. numBytesToSkip -= buf.length; - wrappedSubscriber.onComplete(); + wrappedSubscriber.onNext(ByteBuffer.wrap(new byte[0])); + return; } else { outputBuffer = Arrays.copyOfRange(buf, numBytesToSkip, buf.length); numBytesToSkip = 0; diff --git a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java new file mode 100644 index 000000000..1173a472a --- /dev/null +++ b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java @@ -0,0 +1,124 @@ +// Copyright Amazon.com Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.encryption.s3.legacy.internal; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; + +import java.io.ByteArrayOutputStream; +import java.nio.ByteBuffer; +import java.util.concurrent.atomic.AtomicInteger; + +import org.junit.jupiter.api.Test; +import org.reactivestreams.Subscriber; +import org.reactivestreams.Subscription; + +public class AdjustedRangeSubscriberTest { + + /** + * Records everything delivered downstream so tests can assert on bytes and + * terminal signals. + */ + private static class RecordingSubscriber implements Subscriber { + final ByteArrayOutputStream data = new ByteArrayOutputStream(); + final AtomicInteger completeCount = new AtomicInteger(); + final AtomicInteger onNextCount = new AtomicInteger(); + Throwable error; + + @Override + public void onSubscribe(Subscription s) { + } + + @Override + public void onNext(ByteBuffer byteBuffer) { + onNextCount.incrementAndGet(); + byte[] b = new byte[byteBuffer.remaining()]; + byteBuffer.get(b); + data.write(b, 0, b.length); + } + + @Override + public void onError(Throwable t) { + error = t; + } + + @Override + public void onComplete() { + completeCount.incrementAndGet(); + } + } + + private static ByteBuffer bytes(int start, int length) { + byte[] b = new byte[length]; + for (int i = 0; i < length; i++) { + b[i] = (byte) (start + i); + } + return ByteBuffer.wrap(b); + } + + @Test + public void testFirstChunkSmallerThanSkipDoesNotCompleteOrThrow() throws Exception { + // rangeBeginning=20 => numBytesToSkip=20, virtualAvailable=100 + RecordingSubscriber downstream = new RecordingSubscriber(); + AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); + + // First chunk (10 bytes) is smaller than the 20-byte skip. + subscriber.onNext(bytes(0, 10)); + assertEquals(0, downstream.completeCount.get()); + assertNull(downstream.error); + assertEquals(0, downstream.data.size()); + + // Second chunk (10 bytes) exactly finishes the skip; still no data delivered. + subscriber.onNext(bytes(10, 10)); + assertEquals(0, downstream.completeCount.get()); + assertEquals(0, downstream.data.size()); + + // Third chunk carries the actual payload, which must now be delivered. + subscriber.onNext(bytes(100, 100)); + assertArrayEquals(bytes(100, 100).array(), downstream.data.toByteArray()); + } + + @Test + public void testEmptyFirstChunkIsSkippedNotTreatedAsCompletion() throws Exception { + RecordingSubscriber downstream = new RecordingSubscriber(); + AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); + + // CipherSubscriber can emit an empty buffer; it must not complete the stream. + subscriber.onNext(ByteBuffer.allocate(0)); + assertEquals(0, downstream.completeCount.get()); + assertNull(downstream.error); + + subscriber.onNext(bytes(0, 20)); // finish the skip + subscriber.onNext(bytes(50, 100)); + assertArrayEquals(bytes(50, 100).array(), downstream.data.toByteArray()); + } + + @Test + public void testSkipOnlyChunkSignalsEmptyOnNextToKeepDemandFlowing() throws Exception { + // Under one-at-a-time (request(1)) demand, a chunk fully consumed by the skip must still + // signal the wrapped subscriber, or downstream waits forever for the next element. The + // subscriber forwards an empty buffer (no in-range data, but the demand chain advances). + RecordingSubscriber downstream = new RecordingSubscriber(); + AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); + + subscriber.onNext(bytes(0, 10)); // 10 bytes < 20-byte skip: entirely skipped + + assertEquals(1, downstream.onNextCount.get(), "skip-only chunk must forward exactly one onNext"); + assertEquals(0, downstream.data.size(), "the forwarded onNext must carry no in-range bytes"); + assertEquals(0, downstream.completeCount.get(), "skip-only chunk must not complete the stream"); + assertNull(downstream.error); + } + + @Test + public void testChunkLargerThanSkipDeliversRemainder() throws Exception { + RecordingSubscriber downstream = new RecordingSubscriber(); + AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); + + // A single 120-byte chunk: 20 skipped, 100 delivered. + subscriber.onNext(bytes(0, 120)); + assertArrayEquals(bytes(20, 100).array(), downstream.data.toByteArray()); + assertFalse(downstream.completeCount.get() == 0); + } +} From 3fd5b0b9ef78f7a996d9bb3d4b172cb3089f0b6a Mon Sep 17 00:00:00 2001 From: Bikram Sharma Date: Wed, 30 Sep 2026 13:28:28 -0700 Subject: [PATCH 2/4] test: cover AdjustedRangeSubscriber demand across sub-skip chunks The existing coverage exercises ranged GETs over real transport, which only intermittently produces a first chunk smaller than the skip, so it cannot deterministically catch the premature-onComplete bug (#517). RangedGetUtilsTest covers only the range string math and never drives the subscriber. Add tests that wire the real production topology in-JVM -- a backpressure- honoring publisher feeding AdjustedRangeSubscriber wrapping the real InputStreamSubscriber (as toBlockingInputStream uses) -- and assert the full in-range payload is delivered when the first chunk is smaller than, empty, or exactly equal to the skip. Each read is bounded by a timeout so a demand stall fails fast instead of hanging. These fail on the pre-fix code (empty stream) and pass with the fix. No AWS credentials required; runs in the normal test phase. --- .../AdjustedRangeSubscriberDemandTest.java | 212 ++++++++++++++++++ 1 file changed, 212 insertions(+) create mode 100644 src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java diff --git a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java new file mode 100644 index 000000000..2a30133a9 --- /dev/null +++ b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java @@ -0,0 +1,212 @@ +// Copyright Amazon.com Inc. or its affiliates. All Rights Reserved. +// SPDX-License-Identifier: Apache-2.0 +package software.amazon.encryption.s3.legacy.internal; + +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; + +import java.io.ByteArrayOutputStream; +import java.io.InputStream; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; +import org.reactivestreams.Publisher; +import org.reactivestreams.Subscriber; +import org.reactivestreams.Subscription; + +import software.amazon.awssdk.utils.async.InputStreamSubscriber; + +/** + * Drives the REAL production topology for issue #517: + * + * BackpressurePublisher -> AdjustedRangeSubscriber -> InputStreamSubscriber (== toBlockingInputStream) + * + * AdjustedRangeSubscriber forwards the upstream Subscription straight to the InputStreamSubscriber, + * so demand is driven by the blocking InputStream reader on the shared subscription -- exactly the + * synchronous S3EncryptionClient.getObject path. The publisher honors request(n) and delivers on a + * separate thread, mimicking async S3 delivery. Each read is bounded by a timeout so that a + * demand-stall (the suspected failure mode of a wrong fix) surfaces as a test failure, not a hang. + */ +public class AdjustedRangeSubscriberDemandTest { + + /** Reactive-streams publisher that respects request(n) and delivers chunks on its own thread. */ + private static class BackpressurePublisher implements Publisher { + private final ConcurrentLinkedQueue chunks; + private final ExecutorService exec = Executors.newSingleThreadExecutor(r -> { + Thread t = new Thread(r, "delivery"); + t.setDaemon(true); + return t; + }); + + BackpressurePublisher(List chunks) { + this.chunks = new ConcurrentLinkedQueue<>(chunks); + } + + @Override + public void subscribe(Subscriber s) { + AtomicBoolean terminated = new AtomicBoolean(false); + s.onSubscribe(new Subscription() { + @Override + public void request(long n) { + exec.submit(() -> { + for (long i = 0; i < n; i++) { + if (terminated.get()) { + return; + } + ByteBuffer b = chunks.poll(); + if (b == null) { + if (terminated.compareAndSet(false, true)) { + s.onComplete(); + } + return; + } + s.onNext(b); + } + }); + } + + @Override + public void cancel() { + terminated.set(true); + exec.shutdownNow(); + } + }); + } + } + + private static ByteBuffer bytes(int start, int length) { + byte[] b = new byte[length]; + for (int i = 0; i < length; i++) { + b[i] = (byte) (start + i); + } + return ByteBuffer.wrap(b); + } + + private static byte[] readAllWithTimeout(InputStream in, long timeoutSeconds) throws Exception { + ExecutorService reader = Executors.newSingleThreadExecutor(r -> { + Thread t = new Thread(r, "reader"); + t.setDaemon(true); + return t; + }); + try { + Future f = reader.submit(() -> { + ByteArrayOutputStream out = new ByteArrayOutputStream(); + byte[] tmp = new byte[64]; + int n; + while ((n = in.read(tmp)) != -1) { + out.write(tmp, 0, n); + } + return out.toByteArray(); + }); + return f.get(timeoutSeconds, TimeUnit.SECONDS); + } finally { + reader.shutdownNow(); + } + } + + /** + * CTR case: a tiny non-empty first chunk (< skip), then a second small chunk finishing the skip, + * then the payload. Must deliver the full in-range payload, not an empty stream. + */ + @Test + public void firstChunkSmallerThanSkip_deliversFullPayload() throws Exception { + // rangeBeginning=20 => numBytesToSkip=20 ; rangeEnd=119 => virtualAvailable=100 + List chunks = new ArrayList<>(); + chunks.add(bytes(0, 10)); // 10 bytes -> entirely skipped + chunks.add(bytes(10, 10)); // 10 bytes -> finishes the 20-byte skip + chunks.add(bytes(100, 100)); // 100 bytes -> the in-range payload + + InputStreamSubscriber iss = new InputStreamSubscriber(); + AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); + new BackpressurePublisher(chunks).subscribe(ars); + + byte[] out = readAllWithTimeout(iss, 10); + assertEquals(100, out.length, "must receive the full in-range payload, not an empty stream"); + assertArrayEquals(bytes(100, 100).array(), out); + } + + /** + * AES/CBC (v1) case from the issue: CipherSubscriber emits an empty ByteBuffer first. The empty + * buffer satisfies "chunk <= skip"; the fix must not treat it as completion and must keep demand + * flowing so the payload still arrives. + */ + @Test + public void emptyFirstChunk_thenPayload_deliversFullPayload() throws Exception { + List chunks = new ArrayList<>(); + chunks.add(ByteBuffer.allocate(0)); // empty, as CipherSubscriber emits for CBC + chunks.add(bytes(0, 20)); // finishes the 20-byte skip + chunks.add(bytes(50, 100)); // in-range payload + + InputStreamSubscriber iss = new InputStreamSubscriber(); + AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); + new BackpressurePublisher(chunks).subscribe(ars); + + byte[] out = readAllWithTimeout(iss, 10); + assertEquals(100, out.length, "empty first chunk must not truncate the stream"); + assertArrayEquals(bytes(50, 100).array(), out); + } + + /** + * Many tiny sub-skip chunks in a row (aggressive TLS/TCP fragmentation), then payload split + * across several chunks. Exercises repeated empty-onNext forwarding. + */ + @Test + public void manyTinyChunksBeforeSkipCompletes_deliversFullPayload() throws Exception { + List chunks = new ArrayList<>(); + for (int i = 0; i < 20; i++) { + chunks.add(bytes(i, 1)); // 20 x 1-byte chunks == the whole 20-byte skip + } + // payload 100 bytes, split + chunks.add(bytes(100, 40)); + chunks.add(bytes(140, 60)); + + InputStreamSubscriber iss = new InputStreamSubscriber(); + AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); + new BackpressurePublisher(chunks).subscribe(ars); + + byte[] out = readAllWithTimeout(iss, 10); + assertEquals(100, out.length); + assertArrayEquals(bytes(100, 100).array(), out); + } + + /** Single chunk larger than the skip: 20 skipped, remainder delivered (regression guard). */ + @Test + public void singleChunkLargerThanSkip_deliversRemainder() throws Exception { + List chunks = new ArrayList<>(); + chunks.add(bytes(0, 120)); // 20 skipped, 100 in range + + InputStreamSubscriber iss = new InputStreamSubscriber(); + AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); + new BackpressurePublisher(chunks).subscribe(ars); + + byte[] out = readAllWithTimeout(iss, 10); + assertArrayEquals(bytes(20, 100).array(), out); + } + + /** + * A chunk exactly equal to the remaining skip must be fully consumed (the {@code >=} boundary), + * with the payload in the following chunk still delivered in full. + */ + @Test + public void chunkEqualToSkip_thenPayload_deliversFullPayload() throws Exception { + List chunks = new ArrayList<>(); + chunks.add(bytes(0, 20)); // exactly the 20-byte skip + chunks.add(bytes(100, 100)); // in-range payload + + InputStreamSubscriber iss = new InputStreamSubscriber(); + AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); + new BackpressurePublisher(chunks).subscribe(ars); + + byte[] out = readAllWithTimeout(iss, 10); + assertEquals(100, out.length, "chunk equal to skip must not truncate the stream"); + assertArrayEquals(bytes(100, 100).array(), out); + } +} From 44b8858c7faaefa2045cba7ec5417b588c13bc08 Mon Sep 17 00:00:00 2001 From: Bikram Sharma Date: Wed, 30 Sep 2026 13:50:18 -0700 Subject: [PATCH 3/4] test: reword AdjustedRangeSubscriberDemandTest comments Drop editorializing ("REAL", issue references) from the javadoc in favor of plain descriptions of the topology and each case. --- .../AdjustedRangeSubscriberDemandTest.java | 21 ++++++++----------- 1 file changed, 9 insertions(+), 12 deletions(-) diff --git a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java index 2a30133a9..e37863a28 100644 --- a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java +++ b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java @@ -25,15 +25,12 @@ import software.amazon.awssdk.utils.async.InputStreamSubscriber; /** - * Drives the REAL production topology for issue #517: - * - * BackpressurePublisher -> AdjustedRangeSubscriber -> InputStreamSubscriber (== toBlockingInputStream) - * - * AdjustedRangeSubscriber forwards the upstream Subscription straight to the InputStreamSubscriber, - * so demand is driven by the blocking InputStream reader on the shared subscription -- exactly the - * synchronous S3EncryptionClient.getObject path. The publisher honors request(n) and delivers on a - * separate thread, mimicking async S3 delivery. Each read is bounded by a timeout so that a - * demand-stall (the suspected failure mode of a wrong fix) surfaces as a test failure, not a hang. + * Exercises {@link AdjustedRangeSubscriber} through the subscriber chain used by a synchronous + * ranged GET: a backpressure-honoring publisher feeds the subscriber, which wraps an + * {@link InputStreamSubscriber} (the subscriber behind {@code toBlockingInputStream()}). Because + * the upstream subscription is forwarded straight to the inner subscriber, demand is driven by the + * blocking reader on the shared subscription. The publisher delivers on a separate thread to mimic + * async delivery, and each read is bounded by a timeout so a demand stall fails rather than hangs. */ public class AdjustedRangeSubscriberDemandTest { @@ -134,9 +131,9 @@ public void firstChunkSmallerThanSkip_deliversFullPayload() throws Exception { } /** - * AES/CBC (v1) case from the issue: CipherSubscriber emits an empty ByteBuffer first. The empty - * buffer satisfies "chunk <= skip"; the fix must not treat it as completion and must keep demand - * flowing so the payload still arrives. + * AES/CBC (v1) case: CipherSubscriber can emit an empty ByteBuffer first. The empty buffer + * satisfies "chunk <= skip"; it must not be treated as completion, and demand must keep flowing + * so the payload still arrives. */ @Test public void emptyFirstChunk_thenPayload_deliversFullPayload() throws Exception { From 3afe291579723fe2428d0d916c15ffbd5966404a Mon Sep 17 00:00:00 2001 From: Bikram Sharma Date: Wed, 30 Sep 2026 13:53:10 -0700 Subject: [PATCH 4/4] chore: condense comments in AdjustedRangeSubscriber and its tests Shorten the skip-branch comment and the test comments to one-liners, remove issue references, and drop narration that just restates the adjacent code. --- .../internal/AdjustedRangeSubscriber.java | 11 ++-- .../AdjustedRangeSubscriberDemandTest.java | 52 +++++++------------ .../internal/AdjustedRangeSubscriberTest.java | 21 +++----- 3 files changed, 30 insertions(+), 54 deletions(-) diff --git a/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java b/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java index 011497d64..7607fd2d3 100644 --- a/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java +++ b/src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java @@ -62,13 +62,10 @@ public void onNext(ByteBuffer byteBuffer) { if (numBytesToSkip != 0) { byte[] buf = byteBuffer.array(); if (numBytesToSkip >= buf.length) { - // This chunk is entirely consumed by the leading-byte skip, so there is no - // in-range data to deliver yet. We must still signal the wrapped subscriber to - // keep the reactive-streams demand flowing: downstream drives the stream one - // element at a time (request(1)), and it only requests the next element after it - // receives an onNext. Returning without signaling would leave it waiting forever. - // Forward an empty buffer, mirroring CipherSubscriber's "avoid blocking" idiom, - // and wait for the next chunk instead of completing the stream. + // This chunk is entirely consumed by the skip, so there is no in-range data yet. + // Downstream only requests the next element after receiving an onNext, so forward + // an empty buffer (mirroring CipherSubscriber) to keep demand flowing, rather than + // completing the stream or staying silent and stalling it. numBytesToSkip -= buf.length; wrappedSubscriber.onNext(ByteBuffer.wrap(new byte[0])); return; diff --git a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java index e37863a28..2640db37b 100644 --- a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java +++ b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberDemandTest.java @@ -26,11 +26,9 @@ /** * Exercises {@link AdjustedRangeSubscriber} through the subscriber chain used by a synchronous - * ranged GET: a backpressure-honoring publisher feeds the subscriber, which wraps an - * {@link InputStreamSubscriber} (the subscriber behind {@code toBlockingInputStream()}). Because - * the upstream subscription is forwarded straight to the inner subscriber, demand is driven by the - * blocking reader on the shared subscription. The publisher delivers on a separate thread to mimic - * async delivery, and each read is bounded by a timeout so a demand stall fails rather than hangs. + * ranged GET: a backpressure publisher feeds the subscriber, which wraps an + * {@link InputStreamSubscriber} (the one behind {@code toBlockingInputStream()}). Reads are + * timeout-bounded so a demand stall fails instead of hanging. */ public class AdjustedRangeSubscriberDemandTest { @@ -109,17 +107,14 @@ private static byte[] readAllWithTimeout(InputStream in, long timeoutSeconds) th } } - /** - * CTR case: a tiny non-empty first chunk (< skip), then a second small chunk finishing the skip, - * then the payload. Must deliver the full in-range payload, not an empty stream. - */ + /** CTR case: first chunk smaller than the skip, then the payload. */ @Test public void firstChunkSmallerThanSkip_deliversFullPayload() throws Exception { - // rangeBeginning=20 => numBytesToSkip=20 ; rangeEnd=119 => virtualAvailable=100 + // rangeBeginning=20 => skip=20, rangeEnd=119 => virtualAvailable=100 List chunks = new ArrayList<>(); - chunks.add(bytes(0, 10)); // 10 bytes -> entirely skipped - chunks.add(bytes(10, 10)); // 10 bytes -> finishes the 20-byte skip - chunks.add(bytes(100, 100)); // 100 bytes -> the in-range payload + chunks.add(bytes(0, 10)); // skipped + chunks.add(bytes(10, 10)); // finishes the skip + chunks.add(bytes(100, 100)); // payload InputStreamSubscriber iss = new InputStreamSubscriber(); AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); @@ -130,17 +125,13 @@ public void firstChunkSmallerThanSkip_deliversFullPayload() throws Exception { assertArrayEquals(bytes(100, 100).array(), out); } - /** - * AES/CBC (v1) case: CipherSubscriber can emit an empty ByteBuffer first. The empty buffer - * satisfies "chunk <= skip"; it must not be treated as completion, and demand must keep flowing - * so the payload still arrives. - */ + /** AES/CBC case: an empty first buffer (as CipherSubscriber emits) must not complete the stream. */ @Test public void emptyFirstChunk_thenPayload_deliversFullPayload() throws Exception { List chunks = new ArrayList<>(); - chunks.add(ByteBuffer.allocate(0)); // empty, as CipherSubscriber emits for CBC - chunks.add(bytes(0, 20)); // finishes the 20-byte skip - chunks.add(bytes(50, 100)); // in-range payload + chunks.add(ByteBuffer.allocate(0)); // empty + chunks.add(bytes(0, 20)); // finishes the skip + chunks.add(bytes(50, 100)); // payload InputStreamSubscriber iss = new InputStreamSubscriber(); AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); @@ -151,17 +142,13 @@ public void emptyFirstChunk_thenPayload_deliversFullPayload() throws Exception { assertArrayEquals(bytes(50, 100).array(), out); } - /** - * Many tiny sub-skip chunks in a row (aggressive TLS/TCP fragmentation), then payload split - * across several chunks. Exercises repeated empty-onNext forwarding. - */ + /** Many 1-byte chunks span the skip (heavy fragmentation), then a split payload. */ @Test public void manyTinyChunksBeforeSkipCompletes_deliversFullPayload() throws Exception { List chunks = new ArrayList<>(); for (int i = 0; i < 20; i++) { - chunks.add(bytes(i, 1)); // 20 x 1-byte chunks == the whole 20-byte skip + chunks.add(bytes(i, 1)); // 20 x 1 byte == the skip } - // payload 100 bytes, split chunks.add(bytes(100, 40)); chunks.add(bytes(140, 60)); @@ -174,7 +161,7 @@ public void manyTinyChunksBeforeSkipCompletes_deliversFullPayload() throws Excep assertArrayEquals(bytes(100, 100).array(), out); } - /** Single chunk larger than the skip: 20 skipped, remainder delivered (regression guard). */ + /** Single chunk larger than the skip: remainder delivered. */ @Test public void singleChunkLargerThanSkip_deliversRemainder() throws Exception { List chunks = new ArrayList<>(); @@ -188,15 +175,12 @@ public void singleChunkLargerThanSkip_deliversRemainder() throws Exception { assertArrayEquals(bytes(20, 100).array(), out); } - /** - * A chunk exactly equal to the remaining skip must be fully consumed (the {@code >=} boundary), - * with the payload in the following chunk still delivered in full. - */ + /** Chunk exactly equal to the skip (the {@code >=} boundary), then the payload. */ @Test public void chunkEqualToSkip_thenPayload_deliversFullPayload() throws Exception { List chunks = new ArrayList<>(); - chunks.add(bytes(0, 20)); // exactly the 20-byte skip - chunks.add(bytes(100, 100)); // in-range payload + chunks.add(bytes(0, 20)); // exactly the skip + chunks.add(bytes(100, 100)); // payload InputStreamSubscriber iss = new InputStreamSubscriber(); AdjustedRangeSubscriber ars = new AdjustedRangeSubscriber(iss, 20L, 119L); diff --git a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java index 1173a472a..0b3a91171 100644 --- a/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java +++ b/src/test/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriberTest.java @@ -60,23 +60,20 @@ private static ByteBuffer bytes(int start, int length) { @Test public void testFirstChunkSmallerThanSkipDoesNotCompleteOrThrow() throws Exception { - // rangeBeginning=20 => numBytesToSkip=20, virtualAvailable=100 + // rangeBeginning=20 => skip=20, virtualAvailable=100 RecordingSubscriber downstream = new RecordingSubscriber(); AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); - // First chunk (10 bytes) is smaller than the 20-byte skip. - subscriber.onNext(bytes(0, 10)); + subscriber.onNext(bytes(0, 10)); // smaller than the skip assertEquals(0, downstream.completeCount.get()); assertNull(downstream.error); assertEquals(0, downstream.data.size()); - // Second chunk (10 bytes) exactly finishes the skip; still no data delivered. - subscriber.onNext(bytes(10, 10)); + subscriber.onNext(bytes(10, 10)); // finishes the skip, no data yet assertEquals(0, downstream.completeCount.get()); assertEquals(0, downstream.data.size()); - // Third chunk carries the actual payload, which must now be delivered. - subscriber.onNext(bytes(100, 100)); + subscriber.onNext(bytes(100, 100)); // payload assertArrayEquals(bytes(100, 100).array(), downstream.data.toByteArray()); } @@ -97,13 +94,12 @@ public void testEmptyFirstChunkIsSkippedNotTreatedAsCompletion() throws Exceptio @Test public void testSkipOnlyChunkSignalsEmptyOnNextToKeepDemandFlowing() throws Exception { - // Under one-at-a-time (request(1)) demand, a chunk fully consumed by the skip must still - // signal the wrapped subscriber, or downstream waits forever for the next element. The - // subscriber forwards an empty buffer (no in-range data, but the demand chain advances). + // A chunk fully consumed by the skip must still forward an (empty) onNext so downstream + // requests the next element instead of waiting forever. RecordingSubscriber downstream = new RecordingSubscriber(); AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); - subscriber.onNext(bytes(0, 10)); // 10 bytes < 20-byte skip: entirely skipped + subscriber.onNext(bytes(0, 10)); // smaller than the skip assertEquals(1, downstream.onNextCount.get(), "skip-only chunk must forward exactly one onNext"); assertEquals(0, downstream.data.size(), "the forwarded onNext must carry no in-range bytes"); @@ -116,8 +112,7 @@ public void testChunkLargerThanSkipDeliversRemainder() throws Exception { RecordingSubscriber downstream = new RecordingSubscriber(); AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L); - // A single 120-byte chunk: 20 skipped, 100 delivered. - subscriber.onNext(bytes(0, 120)); + subscriber.onNext(bytes(0, 120)); // 20 skipped, 100 delivered assertArrayEquals(bytes(20, 100).array(), downstream.data.toByteArray()); assertFalse(downstream.completeCount.get() == 0); }