GH-1210 Add S3EventSerializer logic to ensure S3Event deserialization
Resolves #1210
This commit is contained in:
@@ -19,7 +19,7 @@
|
|||||||
<properties>
|
<properties>
|
||||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||||
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
|
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
|
||||||
<aws-lambda-events.version>3.11.4</aws-lambda-events.version>
|
<aws-lambda-events.version>3.14.0</aws-lambda-events.version>
|
||||||
<aws-java-sdk.version>1.12.29</aws-java-sdk.version>
|
<aws-java-sdk.version>1.12.29</aws-java-sdk.version>
|
||||||
<aws-lambda-java-log4j.version>1.0.1</aws-lambda-java-log4j.version>
|
<aws-lambda-java-log4j.version>1.0.1</aws-lambda-java-log4j.version>
|
||||||
<aws-lambda-java-serialization.version>1.1.5</aws-lambda-java-serialization.version>
|
<aws-lambda-java-serialization.version>1.1.5</aws-lambda-java-serialization.version>
|
||||||
|
|||||||
@@ -17,11 +17,14 @@
|
|||||||
package org.springframework.cloud.function.adapter.aws;
|
package org.springframework.cloud.function.adapter.aws;
|
||||||
|
|
||||||
import java.io.ByteArrayInputStream;
|
import java.io.ByteArrayInputStream;
|
||||||
|
import java.io.ByteArrayOutputStream;
|
||||||
import java.nio.charset.StandardCharsets;
|
import java.nio.charset.StandardCharsets;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
|
|
||||||
import com.amazonaws.services.lambda.runtime.serialization.PojoSerializer;
|
import com.amazonaws.services.lambda.runtime.serialization.PojoSerializer;
|
||||||
import com.amazonaws.services.lambda.runtime.serialization.events.LambdaEventSerializers;
|
import com.amazonaws.services.lambda.runtime.serialization.events.LambdaEventSerializers;
|
||||||
|
import com.amazonaws.services.lambda.runtime.serialization.events.serializers.S3EventSerializer;
|
||||||
|
|
||||||
import org.springframework.cloud.function.cloudevent.CloudEventMessageUtils;
|
import org.springframework.cloud.function.cloudevent.CloudEventMessageUtils;
|
||||||
import org.springframework.cloud.function.context.config.JsonMessageConverter;
|
import org.springframework.cloud.function.context.config.JsonMessageConverter;
|
||||||
@@ -30,6 +33,7 @@ import org.springframework.lang.Nullable;
|
|||||||
import org.springframework.messaging.Message;
|
import org.springframework.messaging.Message;
|
||||||
import org.springframework.messaging.MessageHeaders;
|
import org.springframework.messaging.MessageHeaders;
|
||||||
import org.springframework.messaging.converter.MessageConverter;
|
import org.springframework.messaging.converter.MessageConverter;
|
||||||
|
import org.springframework.util.ClassUtils;
|
||||||
import org.springframework.util.MimeType;
|
import org.springframework.util.MimeType;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -44,6 +48,9 @@ class AWSTypesMessageConverter extends JsonMessageConverter {
|
|||||||
|
|
||||||
private final JsonMapper jsonMapper;
|
private final JsonMapper jsonMapper;
|
||||||
|
|
||||||
|
@SuppressWarnings("rawtypes")
|
||||||
|
private final AtomicReference<S3EventSerializer> s3EventSerializer = new AtomicReference<>();
|
||||||
|
|
||||||
AWSTypesMessageConverter(JsonMapper jsonMapper) {
|
AWSTypesMessageConverter(JsonMapper jsonMapper) {
|
||||||
this(jsonMapper, new MimeType("application", "json"), new MimeType(CloudEventMessageUtils.APPLICATION_CLOUDEVENTS.getType(),
|
this(jsonMapper, new MimeType("application", "json"), new MimeType(CloudEventMessageUtils.APPLICATION_CLOUDEVENTS.getType(),
|
||||||
CloudEventMessageUtils.APPLICATION_CLOUDEVENTS.getSubtype() + "+json"));
|
CloudEventMessageUtils.APPLICATION_CLOUDEVENTS.getSubtype() + "+json"));
|
||||||
@@ -75,7 +82,6 @@ class AWSTypesMessageConverter extends JsonMessageConverter {
|
|||||||
if (message.getPayload().getClass().isAssignableFrom(targetClass)) {
|
if (message.getPayload().getClass().isAssignableFrom(targetClass)) {
|
||||||
return message.getPayload();
|
return message.getPayload();
|
||||||
}
|
}
|
||||||
|
|
||||||
if (targetClass.getPackage() != null &&
|
if (targetClass.getPackage() != null &&
|
||||||
targetClass.getPackage().getName().startsWith("com.amazonaws.services.lambda.runtime.events")) {
|
targetClass.getPackage().getName().startsWith("com.amazonaws.services.lambda.runtime.events")) {
|
||||||
PojoSerializer<?> serializer = LambdaEventSerializers.serializerFor(targetClass, Thread.currentThread().getContextClassLoader());
|
PojoSerializer<?> serializer = LambdaEventSerializers.serializerFor(targetClass, Thread.currentThread().getContextClassLoader());
|
||||||
@@ -110,12 +116,23 @@ class AWSTypesMessageConverter extends JsonMessageConverter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
@Override
|
@Override
|
||||||
protected Object convertToInternal(Object payload, @Nullable MessageHeaders headers,
|
protected Object convertToInternal(Object payload, @Nullable MessageHeaders headers,
|
||||||
@Nullable Object conversionHint) {
|
@Nullable Object conversionHint) {
|
||||||
if (payload instanceof String && headers.containsKey(AWSLambdaUtils.IS_BASE64_ENCODED) && (boolean) headers.get(AWSLambdaUtils.IS_BASE64_ENCODED)) {
|
if (payload instanceof String && headers.containsKey(AWSLambdaUtils.IS_BASE64_ENCODED) && (boolean) headers.get(AWSLambdaUtils.IS_BASE64_ENCODED)) {
|
||||||
return ((String) payload).getBytes(StandardCharsets.UTF_8);
|
return ((String) payload).getBytes(StandardCharsets.UTF_8);
|
||||||
}
|
}
|
||||||
|
if (payload.getClass().getName().equals("com.amazonaws.services.lambda.runtime.events.S3Event")) {
|
||||||
|
if (this.s3EventSerializer.get() == null) {
|
||||||
|
this.s3EventSerializer.set(new S3EventSerializer<>().withClassLoader(ClassUtils.getDefaultClassLoader()));
|
||||||
|
}
|
||||||
|
ByteArrayOutputStream stream = new ByteArrayOutputStream();
|
||||||
|
this.s3EventSerializer.get().toJson(payload, stream);
|
||||||
|
return stream.toByteArray();
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
return jsonMapper.toJson(payload);
|
return jsonMapper.toJson(payload);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -999,6 +999,18 @@ public class FunctionInvokerTests {
|
|||||||
assertThat(result).contains("s3SchemaVersion");
|
assertThat(result).contains("s3SchemaVersion");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void testS3EventAsOutput() throws Exception {
|
||||||
|
System.setProperty("MAIN_CLASS", S3Configuration.class.getName());
|
||||||
|
System.setProperty("spring.cloud.function.definition", "outputS3Event");
|
||||||
|
FunctionInvoker invoker = new FunctionInvoker();
|
||||||
|
|
||||||
|
InputStream targetStream = new ByteArrayInputStream(this.s3Event.getBytes());
|
||||||
|
ByteArrayOutputStream output = new ByteArrayOutputStream();
|
||||||
|
invoker.handleRequest(targetStream, output, null);
|
||||||
|
assertThat(output.toByteArray()).isNotNull();
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void testS3Event() throws Exception {
|
public void testS3Event() throws Exception {
|
||||||
System.setProperty("MAIN_CLASS", S3Configuration.class.getName());
|
System.setProperty("MAIN_CLASS", S3Configuration.class.getName());
|
||||||
@@ -1679,6 +1691,13 @@ public class FunctionInvokerTests {
|
|||||||
@EnableAutoConfiguration
|
@EnableAutoConfiguration
|
||||||
@Configuration
|
@Configuration
|
||||||
public static class S3Configuration {
|
public static class S3Configuration {
|
||||||
|
|
||||||
|
@Bean
|
||||||
|
public Function<S3Event, S3Event> outputS3Event() {
|
||||||
|
return v -> {
|
||||||
|
return v;
|
||||||
|
};
|
||||||
|
}
|
||||||
@Bean
|
@Bean
|
||||||
public Function<String, String> echoString() {
|
public Function<String, String> echoString() {
|
||||||
return v -> v;
|
return v -> v;
|
||||||
|
|||||||
Reference in New Issue
Block a user