Skip to content

Commit b56e59d

Browse files
committed
Orient DataStream around InputStream
DataStream now is written around the idea of using an underyling InputStream more often than a Publisher, so default implementations were updated (now default provided for asPublisher, and asInputStream no longer has a default). Also closes InputStreams in a couple places that we could have leaked. Addresses smithy-lang#913
1 parent 636af59 commit b56e59d

30 files changed

Lines changed: 150 additions & 186 deletions

File tree

aws/aws-sigv4/src/main/java/software/amazon/smithy/java/aws/client/auth/scheme/sigv4/SigV4Signer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,7 @@ public HttpRequest sign(HttpRequest request, AwsCredentialsIdentity identity, Co
130130
}
131131

132132
private String getPayloadHash(DataStream dataStream) {
133-
return hexHash(dataStream.waitForByteBuffer());
133+
return hexHash(dataStream.asByteBuffer());
134134
}
135135

136136
private String hexHash(ByteBuffer bytes) {

aws/client/aws-client-awsjson/src/it/java/software/amazon/smithy/java/client/aws/jsonprotocols/AwsJson1ProtocolTests.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ public void requestTest(DataStream expected, DataStream actual) {
3939
if (expected.contentLength() != 0) {
4040
// Use the node parser to strip out white space.
4141
expectedJson = Node.printJson(
42-
Node.parse(new String(ByteBufferUtils.getBytes(expected.waitForByteBuffer()),
42+
Node.parse(new String(ByteBufferUtils.getBytes(expected.asByteBuffer()),
4343
StandardCharsets.UTF_8)));
4444
}
4545
assertEquals(expectedJson, new StringBuildingSubscriber(actual).getResult());

aws/client/aws-client-awsjson/src/main/java/software/amazon/smithy/java/aws/client/awsjson/AwsJsonProtocol.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -101,7 +101,7 @@ public <I extends SerializableStruct, O extends SerializableStruct> O deserializ
101101
return codec.deserializeShape(EMPTY_PAYLOAD, builder);
102102
}
103103

104-
var bytes = content.waitForByteBuffer();
104+
var bytes = content.asByteBuffer();
105105
return codec.deserializeShape(bytes, builder);
106106
}
107107
}

aws/client/aws-client-restjson/src/it/java/software/amazon/smithy/java/client/aws/restjson/RestJson1ProtocolTests.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,7 @@ public void requestTest(DataStream expected, DataStream actual) {
4141
var actualStr = new StringBuildingSubscriber(actual).getResult();
4242
if (expected.contentLength() != 0) {
4343
var expectedStr = new String(
44-
ByteBufferUtils.getBytes(expected.waitForByteBuffer()),
44+
ByteBufferUtils.getBytes(expected.asByteBuffer()),
4545
StandardCharsets.UTF_8);
4646
if ("application/json".equals(expected.contentType())) {
4747
var expectedNode = Node.parse(expectedStr);

aws/client/aws-client-restxml/src/it/java/software/amazon/smithy/java/aws/client/restxml/RestXmlProtocolTests.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,8 @@ public class RestXmlProtocolTests {
4141
})
4242
public void requestTest(DataStream expected, DataStream actual) {
4343
if (expected.contentLength() != 0) {
44-
var a = new String(ByteBufferUtils.getBytes(actual.waitForByteBuffer()), StandardCharsets.UTF_8);
45-
var b = new String(ByteBufferUtils.getBytes(expected.waitForByteBuffer()), StandardCharsets.UTF_8);
44+
var a = new String(ByteBufferUtils.getBytes(actual.asByteBuffer()), StandardCharsets.UTF_8);
45+
var b = new String(ByteBufferUtils.getBytes(expected.asByteBuffer()), StandardCharsets.UTF_8);
4646
if ("application/xml".equals(expected.contentType())) {
4747
if (!XMLComparator.compareXMLStrings(a, b)) {
4848
// Do this comparison to see what is different.

aws/integrations/aws-lambda-endpoint/src/main/java/software/amazon/smithy/java/aws/integrations/lambda/LambdaEndpoint.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,7 @@ private static ProxyResponse getResponse(HttpResponse httpResponse, boolean shou
157157

158158
DataStream val = httpResponse.getSerializedValue();
159159
if (val != null) {
160-
ByteBuffer buf = val.waitForByteBuffer();
160+
ByteBuffer buf = val.asByteBuffer();
161161
String body;
162162
// TODO: handle base64 encoding better
163163
if (shouldBase64Encode) {

aws/server/aws-server-restjson/src/it/java/software/amazon/smithy/java/aws/server/restjson/AwsRestJson1ProtocolTests.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -85,9 +85,9 @@ public void responseTest(DataStream expected, DataStream actual) {
8585
.isTrue()
8686
.isSameAs(actual.hasKnownLength());
8787

88-
String actualJson = new String(ByteBufferUtils.getBytes(actual.waitForByteBuffer()), StandardCharsets.UTF_8);
88+
String actualJson = new String(ByteBufferUtils.getBytes(actual.asByteBuffer()), StandardCharsets.UTF_8);
8989
String expectedJson = new String(
90-
ByteBufferUtils.getBytes(expected.waitForByteBuffer()),
90+
ByteBufferUtils.getBytes(expected.asByteBuffer()),
9191
StandardCharsets.UTF_8);
9292
if (expected.contentLength() == 0) {
9393
assertThat(actualJson).isIn("", "{}");

client/client-http/src/main/java/software/amazon/smithy/java/client/http/HttpErrorDeserializer.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -138,7 +138,7 @@ public ModeledException createError(
138138
ShapeBuilder<ModeledException> builder
139139
) {
140140
try {
141-
ByteBuffer bytes = createDataStream(response).waitForByteBuffer();
141+
ByteBuffer bytes = createDataStream(response).asByteBuffer();
142142
return codec.deserializeShape(bytes, builder);
143143
} catch (Exception e) {
144144
throw new RuntimeException("Failed to deserialize error", e);
@@ -248,7 +248,7 @@ private static CallException makeErrorFromPayload(
248248
try {
249249
// Read the payload into a JSON document so we can efficiently find __type and then directly
250250
// deserialize the document into the identified builder.
251-
ByteBuffer buffer = content.waitForByteBuffer();
251+
ByteBuffer buffer = content.asByteBuffer();
252252

253253
if (buffer.remaining() > 0) {
254254
var document = codec.createDeserializer(buffer).readDocument();

client/client-http/src/main/java/software/amazon/smithy/java/client/http/JavaHttpClientTransport.java

Lines changed: 19 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,9 @@
55

66
package software.amazon.smithy.java.client.http;
77

8+
import static java.net.http.HttpRequest.BodyPublishers;
9+
import static java.net.http.HttpResponse.BodyHandlers;
10+
811
import java.io.IOException;
912
import java.io.InputStream;
1013
import java.net.URI;
@@ -110,16 +113,15 @@ private java.net.http.HttpRequest createJavaRequest(Context context, HttpRequest
110113
java.net.http.HttpRequest.BodyPublisher bodyPublisher;
111114
if (request.body().hasKnownLength()) {
112115
if (request.body().contentLength() == 0) {
113-
bodyPublisher = java.net.http.HttpRequest.BodyPublishers.noBody();
116+
bodyPublisher = BodyPublishers.noBody();
114117
} else {
115-
bodyPublisher = java.net.http.HttpRequest.BodyPublishers.ofByteArray(
116-
ByteBufferUtils.getBytes(request.body().waitForByteBuffer()));
118+
bodyPublisher = BodyPublishers.ofByteArray(ByteBufferUtils.getBytes(request.body().asByteBuffer()));
117119
}
118120
} else {
119-
bodyPublisher = java.net.http.HttpRequest.BodyPublishers.fromPublisher(request.body());
121+
bodyPublisher = BodyPublishers.fromPublisher(request.body());
120122
}
121123

122-
java.net.http.HttpRequest.Builder httpRequestBuilder = java.net.http.HttpRequest.newBuilder()
124+
var httpRequestBuilder = java.net.http.HttpRequest.newBuilder()
123125
.version(smithyToHttpVersion(request.httpVersion()))
124126
.method(request.method(), bodyPublisher)
125127
.uri(request.uri());
@@ -141,13 +143,24 @@ private java.net.http.HttpRequest createJavaRequest(Context context, HttpRequest
141143
}
142144

143145
private HttpResponse sendRequest(java.net.http.HttpRequest request) {
146+
java.net.http.HttpResponse<InputStream> res = null;
144147
try {
145-
var res = client.send(request, java.net.http.HttpResponse.BodyHandlers.ofInputStream());
148+
res = client.send(request, BodyHandlers.ofInputStream());
146149
return createSmithyResponse(res);
147150
} catch (IOException | InterruptedException | RuntimeException e) {
151+
// Close the response body stream if we got a response but failed to process it
152+
if (res != null) {
153+
try {
154+
res.body().close();
155+
} catch (IOException closeException) {
156+
LOGGER.trace("Failed to close response body after error", closeException);
157+
}
158+
}
159+
148160
if (e instanceof HttpConnectTimeoutException) {
149161
throw new ConnectTimeoutException(e);
150162
}
163+
151164
// The client pipeline also does this remapping, but to adhere to the required contract of
152165
// ClientTransport, we remap here too if needed.
153166
throw ClientTransport.remapExceptions(e);

client/client-http/src/main/java/software/amazon/smithy/java/client/http/plugins/HttpChecksumPlugin.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ private static HttpRequest processRequest(RequestHook<?, ?, HttpRequest> hook) {
5050
static HttpRequest addContentMd5Header(HttpRequest request) {
5151
var body = request.body();
5252
if (body != null) {
53-
var buffer = body.waitForByteBuffer();
53+
var buffer = body.asByteBuffer();
5454
var bytes = ByteBufferUtils.getBytes(buffer);
5555
try {
5656
byte[] hash = MessageDigest.getInstance("MD5").digest(bytes);

0 commit comments

Comments
 (0)