Skip to content

Commit f9fb658

Browse files
committed
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.
1 parent af4253f commit f9fb658

2 files changed

Lines changed: 110 additions & 4 deletions

File tree

‎src/main/java/software/amazon/encryption/s3/legacy/internal/AdjustedRangeSubscriber.java‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -61,11 +61,11 @@ public void onNext(ByteBuffer byteBuffer) {
6161

6262
if (numBytesToSkip != 0) {
6363
byte[] buf = byteBuffer.array();
64-
if (numBytesToSkip > buf.length) {
65-
// If we need to skip past the available data,
66-
// we are returning nothing, so signal completion
64+
if (numBytesToSkip >= buf.length) {
65+
// Chunk is fully consumed by the skip; wait for the next chunk
66+
// rather than signaling completion.
6767
numBytesToSkip -= buf.length;
68-
wrappedSubscriber.onComplete();
68+
return;
6969
} else {
7070
outputBuffer = Arrays.copyOfRange(buf, numBytesToSkip, buf.length);
7171
numBytesToSkip = 0;
Lines changed: 106 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,106 @@
1+
// Copyright Amazon.com Inc. or its affiliates. All Rights Reserved.
2+
// SPDX-License-Identifier: Apache-2.0
3+
package software.amazon.encryption.s3.legacy.internal;
4+
5+
import static org.junit.jupiter.api.Assertions.assertArrayEquals;
6+
import static org.junit.jupiter.api.Assertions.assertEquals;
7+
import static org.junit.jupiter.api.Assertions.assertFalse;
8+
import static org.junit.jupiter.api.Assertions.assertNull;
9+
10+
import java.io.ByteArrayOutputStream;
11+
import java.nio.ByteBuffer;
12+
import java.util.concurrent.atomic.AtomicInteger;
13+
14+
import org.junit.jupiter.api.Test;
15+
import org.reactivestreams.Subscriber;
16+
import org.reactivestreams.Subscription;
17+
18+
public class AdjustedRangeSubscriberTest {
19+
20+
/**
21+
* Records everything delivered downstream so tests can assert on bytes and
22+
* terminal signals.
23+
*/
24+
private static class RecordingSubscriber implements Subscriber<ByteBuffer> {
25+
final ByteArrayOutputStream data = new ByteArrayOutputStream();
26+
final AtomicInteger completeCount = new AtomicInteger();
27+
Throwable error;
28+
29+
@Override
30+
public void onSubscribe(Subscription s) {
31+
}
32+
33+
@Override
34+
public void onNext(ByteBuffer byteBuffer) {
35+
byte[] b = new byte[byteBuffer.remaining()];
36+
byteBuffer.get(b);
37+
data.write(b, 0, b.length);
38+
}
39+
40+
@Override
41+
public void onError(Throwable t) {
42+
error = t;
43+
}
44+
45+
@Override
46+
public void onComplete() {
47+
completeCount.incrementAndGet();
48+
}
49+
}
50+
51+
private static ByteBuffer bytes(int start, int length) {
52+
byte[] b = new byte[length];
53+
for (int i = 0; i < length; i++) {
54+
b[i] = (byte) (start + i);
55+
}
56+
return ByteBuffer.wrap(b);
57+
}
58+
59+
@Test
60+
public void testFirstChunkSmallerThanSkipDoesNotCompleteOrThrow() throws Exception {
61+
// rangeBeginning=20 => numBytesToSkip=20, virtualAvailable=100
62+
RecordingSubscriber downstream = new RecordingSubscriber();
63+
AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L);
64+
65+
// First chunk (10 bytes) is smaller than the 20-byte skip.
66+
subscriber.onNext(bytes(0, 10));
67+
assertEquals(0, downstream.completeCount.get());
68+
assertNull(downstream.error);
69+
assertEquals(0, downstream.data.size());
70+
71+
// Second chunk (10 bytes) exactly finishes the skip; still no data delivered.
72+
subscriber.onNext(bytes(10, 10));
73+
assertEquals(0, downstream.completeCount.get());
74+
assertEquals(0, downstream.data.size());
75+
76+
// Third chunk carries the actual payload, which must now be delivered.
77+
subscriber.onNext(bytes(100, 100));
78+
assertArrayEquals(bytes(100, 100).array(), downstream.data.toByteArray());
79+
}
80+
81+
@Test
82+
public void testEmptyFirstChunkIsSkippedNotTreatedAsCompletion() throws Exception {
83+
RecordingSubscriber downstream = new RecordingSubscriber();
84+
AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L);
85+
86+
// CipherSubscriber can emit an empty buffer; it must not complete the stream.
87+
subscriber.onNext(ByteBuffer.allocate(0));
88+
assertEquals(0, downstream.completeCount.get());
89+
assertNull(downstream.error);
90+
91+
subscriber.onNext(bytes(0, 20)); // finish the skip
92+
subscriber.onNext(bytes(50, 100));
93+
assertArrayEquals(bytes(50, 100).array(), downstream.data.toByteArray());
94+
}
95+
96+
@Test
97+
public void testChunkLargerThanSkipDeliversRemainder() throws Exception {
98+
RecordingSubscriber downstream = new RecordingSubscriber();
99+
AdjustedRangeSubscriber subscriber = new AdjustedRangeSubscriber(downstream, 20L, 119L);
100+
101+
// A single 120-byte chunk: 20 skipped, 100 delivered.
102+
subscriber.onNext(bytes(0, 120));
103+
assertArrayEquals(bytes(20, 100).array(), downstream.data.toByteArray());
104+
assertFalse(downstream.completeCount.get() == 0);
105+
}
106+
}

0 commit comments

Comments
 (0)