Skip to content

Commit 45c9d1b

Browse files
authored
Event stream implementation for RPCv2 (#849)
1 parent 3b6b6ff commit 45c9d1b

20 files changed

Lines changed: 582 additions & 141 deletions

File tree

aws/aws-event-streams/src/main/java/software/amazon/smithy/java/aws/events/AwsEventDecoderFactory.java

Lines changed: 47 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
package software.amazon.smithy.java.aws.events;
77

8+
import java.util.Objects;
89
import java.util.function.Supplier;
910
import software.amazon.smithy.java.core.schema.InputEventStreamingApiOperation;
1011
import software.amazon.smithy.java.core.schema.OutputEventStreamingApiOperation;
@@ -17,43 +18,79 @@
1718
import software.amazon.smithy.java.core.serde.event.FrameDecoder;
1819
import software.amazon.smithy.java.core.serde.event.FrameTransformer;
1920

20-
public class AwsEventDecoderFactory<E extends SerializableStruct> implements EventDecoderFactory<AwsEventFrame> {
21+
/**
22+
* A {@link EventDecoderFactory} for AWS events.
23+
*
24+
* @param <E> The event shape type
25+
* @param <IR> The initial request shape type
26+
*/
27+
public final class AwsEventDecoderFactory<E extends SerializableStruct, IR extends SerializableStruct>
28+
implements EventDecoderFactory<AwsEventFrame> {
2129

22-
private final Schema schema;
30+
private final InitialEventType initialEventType;
31+
private final Supplier<ShapeBuilder<IR>> initialEventBuilder;
32+
private final Schema eventSchema;
2333
private final Codec codec;
2434
private final Supplier<ShapeBuilder<E>> eventBuilder;
2535
private final FrameTransformer<AwsEventFrame> transformer;
2636

2737
private AwsEventDecoderFactory(
28-
Schema schema,
38+
InitialEventType initialEventType,
39+
Supplier<ShapeBuilder<IR>> initialEventBuilder,
40+
Schema eventSchema,
2941
Codec codec,
3042
Supplier<ShapeBuilder<E>> eventBuilder,
3143
FrameTransformer<AwsEventFrame> transformer
3244
) {
33-
this.schema = schema.isMember() ? schema.memberTarget() : schema;
34-
this.codec = codec;
35-
this.eventBuilder = eventBuilder;
36-
this.transformer = transformer;
45+
this.initialEventType = Objects.requireNonNull(initialEventType, "initialEventType");
46+
this.initialEventBuilder = Objects.requireNonNull(initialEventBuilder, "initialEventBuilder");
47+
this.eventSchema = Objects.requireNonNull(eventSchema, "eventSchema").isMember() ? eventSchema.memberTarget()
48+
: eventSchema;
49+
this.codec = Objects.requireNonNull(codec, "codec");
50+
this.eventBuilder = Objects.requireNonNull(eventBuilder, "eventBuilder");
51+
this.transformer = Objects.requireNonNull(transformer, "transformer");
3752
}
3853

39-
public static <IE extends SerializableStruct> AwsEventDecoderFactory<IE> forInputStream(
54+
/**
55+
* Creates a new input stream decoder factory.
56+
*
57+
* @param operation The input operation for the factory
58+
* @param codec The protocol codec to decode the payload
59+
* @param transformer The frame transformer
60+
* @param <IE> The output event type
61+
* @return A new event decoder factory
62+
*/
63+
public static <IE extends SerializableStruct> AwsEventDecoderFactory<IE, ?> forInputStream(
4064
InputEventStreamingApiOperation<?, ?, IE> operation,
4165
Codec codec,
4266
FrameTransformer<AwsEventFrame> transformer
4367
) {
4468
return new AwsEventDecoderFactory<>(
69+
InitialEventType.INITIAL_REQUEST,
70+
operation::inputBuilder,
4571
operation.inputStreamMember(),
4672
codec,
4773
operation.inputEventBuilderSupplier(),
4874
transformer);
4975
}
5076

51-
public static <OE extends SerializableStruct> AwsEventDecoderFactory<OE> forOutputStream(
77+
/**
78+
* Creates a new output stream decoder factory.
79+
*
80+
* @param operation The output operation for the factory
81+
* @param codec The protocol codec to decode the payload
82+
* @param transformer The frame transformer
83+
* @param <OE> The output event type
84+
* @return A new event decoder factory
85+
*/
86+
public static <OE extends SerializableStruct> AwsEventDecoderFactory<OE, ?> forOutputStream(
5287
OutputEventStreamingApiOperation<?, ?, OE> operation,
5388
Codec codec,
5489
FrameTransformer<AwsEventFrame> transformer
5590
) {
5691
return new AwsEventDecoderFactory<>(
92+
InitialEventType.INITIAL_RESPONSE,
93+
operation::outputBuilder,
5794
operation.outputStreamMember(),
5895
codec,
5996
operation.outputEventBuilderSupplier(),
@@ -62,7 +99,7 @@ public static <OE extends SerializableStruct> AwsEventDecoderFactory<OE> forOutp
6299

63100
@Override
64101
public EventDecoder<AwsEventFrame> newEventDecoder() {
65-
return new AwsEventShapeDecoder<>(eventBuilder, schema, codec);
102+
return new AwsEventShapeDecoder<>(initialEventType, initialEventBuilder, eventBuilder, eventSchema, codec);
66103
}
67104

68105
@Override

aws/aws-event-streams/src/main/java/software/amazon/smithy/java/aws/events/AwsEventEncoderFactory.java

Lines changed: 41 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
package software.amazon.smithy.java.aws.events;
77

8+
import java.util.Objects;
89
import java.util.function.Function;
910
import software.amazon.smithy.java.core.schema.InputEventStreamingApiOperation;
1011
import software.amazon.smithy.java.core.schema.OutputEventStreamingApiOperation;
@@ -15,46 +16,77 @@
1516
import software.amazon.smithy.java.core.serde.event.EventStreamingException;
1617
import software.amazon.smithy.java.core.serde.event.FrameEncoder;
1718

18-
public class AwsEventEncoderFactory implements EventEncoderFactory<AwsEventFrame> {
19-
19+
/**
20+
* A {@link EventEncoderFactory} for AWS events.
21+
*/
22+
public final class AwsEventEncoderFactory implements EventEncoderFactory<AwsEventFrame> {
23+
private final InitialEventType initialEventType;
2024
private final Schema schema;
2125
private final Codec codec;
2226
private final String payloadMediaType;
2327
private final Function<Throwable, EventStreamingException> exceptionHandler;
2428

2529
private AwsEventEncoderFactory(
30+
InitialEventType initialEventType,
2631
Schema schema,
2732
Codec codec,
2833
String payloadMediaType,
2934
Function<Throwable, EventStreamingException> exceptionHandler
3035
) {
31-
this.schema = schema.isMember() ? schema.memberTarget() : schema;
32-
this.codec = codec;
33-
this.payloadMediaType = payloadMediaType;
34-
this.exceptionHandler = exceptionHandler;
36+
this.initialEventType = Objects.requireNonNull(initialEventType, "initialEventType");
37+
this.schema = Objects.requireNonNull(schema, "schema").isMember() ? schema.memberTarget() : schema;
38+
this.codec = Objects.requireNonNull(codec, "codec");
39+
this.payloadMediaType = Objects.requireNonNull(payloadMediaType, "payloadMediaType");
40+
this.exceptionHandler = Objects.requireNonNull(exceptionHandler, "exceptionHandler");
3541
}
3642

43+
/**
44+
* Creates a new input stream encoder factory.
45+
*
46+
* @param operation The input operation for the factory
47+
* @param codec The protocol codec to decode the payload
48+
* @param payloadMediaType The payload media type
49+
* @param exceptionHandler The handler to convert exceptions for event streaming
50+
* @return A new event encoder factory
51+
*/
3752
public static AwsEventEncoderFactory forInputStream(
3853
InputEventStreamingApiOperation<?, ?, ?> operation,
3954
Codec codec,
4055
String payloadMediaType,
4156
Function<Throwable, EventStreamingException> exceptionHandler
4257
) {
43-
return new AwsEventEncoderFactory(operation.inputStreamMember(), codec, payloadMediaType, exceptionHandler);
58+
return new AwsEventEncoderFactory(InitialEventType.INITIAL_REQUEST,
59+
operation.inputStreamMember(),
60+
codec,
61+
payloadMediaType,
62+
exceptionHandler);
4463
}
4564

65+
/**
66+
* Creates a new output stream encoder factory.
67+
*
68+
* @param operation The output operation for the factory
69+
* @param codec The protocol codec to decode the payload
70+
* @param payloadMediaType The payload media type
71+
* @param exceptionHandler The handler to convert exceptions for event streaming
72+
* @return A new event encoder factory
73+
*/
4674
public static AwsEventEncoderFactory forOutputStream(
4775
OutputEventStreamingApiOperation<?, ?, ?> operation,
4876
Codec codec,
4977
String payloadMediaType,
5078
Function<Throwable, EventStreamingException> exceptionHandler
5179
) {
52-
return new AwsEventEncoderFactory(operation.outputStreamMember(), codec, payloadMediaType, exceptionHandler);
80+
return new AwsEventEncoderFactory(InitialEventType.INITIAL_RESPONSE,
81+
operation.outputStreamMember(),
82+
codec,
83+
payloadMediaType,
84+
exceptionHandler);
5385
}
5486

5587
@Override
5688
public EventEncoder<AwsEventFrame> newEventEncoder() {
57-
return new AwsEventShapeEncoder(schema, codec, payloadMediaType, exceptionHandler);
89+
return new AwsEventShapeEncoder(initialEventType, schema, codec, payloadMediaType, exceptionHandler);
5890
}
5991

6092
@Override
@@ -66,5 +98,4 @@ public FrameEncoder<AwsEventFrame> newFrameEncoder() {
6698
public String contentType() {
6799
return "application/vnd.amazon.eventstream";
68100
}
69-
70101
}

aws/aws-event-streams/src/main/java/software/amazon/smithy/java/aws/events/AwsEventShapeDecoder.java

Lines changed: 82 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -5,56 +5,118 @@
55

66
package software.amazon.smithy.java.aws.events;
77

8+
import java.util.Objects;
9+
import java.util.concurrent.Flow;
810
import java.util.function.Supplier;
911
import software.amazon.eventstream.Message;
1012
import software.amazon.smithy.java.core.schema.Schema;
1113
import software.amazon.smithy.java.core.schema.SerializableStruct;
1214
import software.amazon.smithy.java.core.schema.ShapeBuilder;
15+
import software.amazon.smithy.java.core.schema.TraitKey;
1316
import software.amazon.smithy.java.core.serde.Codec;
17+
import software.amazon.smithy.java.core.serde.ShapeDeserializer;
18+
import software.amazon.smithy.java.core.serde.SpecificShapeDeserializer;
1419
import software.amazon.smithy.java.core.serde.event.EventDecoder;
1520

16-
public final class AwsEventShapeDecoder<E extends SerializableStruct> implements EventDecoder<AwsEventFrame> {
21+
/**
22+
* A decoder for AWS events
23+
*
24+
* @param <E> The type of the event
25+
* @param <IR> The type of the initial event
26+
*/
27+
public final class AwsEventShapeDecoder<E extends SerializableStruct, IR extends SerializableStruct>
28+
implements EventDecoder<AwsEventFrame> {
1729

30+
private final InitialEventType initialEventType;
31+
private final Supplier<ShapeBuilder<IR>> initialEventBuilder;
1832
private final Supplier<ShapeBuilder<E>> eventBuilder;
1933
private final Schema eventSchema;
2034
private final Codec codec;
35+
private volatile Flow.Publisher<SerializableStruct> publisher;
2136

22-
public AwsEventShapeDecoder(
37+
AwsEventShapeDecoder(
38+
InitialEventType initialEventType,
39+
Supplier<ShapeBuilder<IR>> initialEventBuilder,
2340
Supplier<ShapeBuilder<E>> eventBuilder,
2441
Schema eventSchema,
2542
Codec codec
2643
) {
27-
this.eventBuilder = eventBuilder;
28-
this.eventSchema = eventSchema;
29-
this.codec = codec;
44+
this.initialEventType = Objects.requireNonNull(initialEventType, "initialEventType");
45+
this.initialEventBuilder = Objects.requireNonNull(initialEventBuilder, "initialEventBuilder");
46+
this.eventBuilder = Objects.requireNonNull(eventBuilder, "eventBuilder");
47+
this.eventSchema = Objects.requireNonNull(eventSchema, "eventSchema");
48+
this.codec = Objects.requireNonNull(codec, "codec");
3049
}
3150

3251
@Override
33-
public E decode(AwsEventFrame frame) {
34-
Message message = frame.unwrap();
35-
String messageType = getMessageType(message);
36-
if (!messageType.equals("event")) {
37-
throw new UnsupportedOperationException("Unsupported frame type: " + messageType);
52+
public SerializableStruct decode(AwsEventFrame frame) {
53+
var message = frame.unwrap();
54+
var eventType = getEventType(message);
55+
if (initialEventType.getName().equals(eventType)) {
56+
return decodeInitialResponse(frame);
3857
}
39-
String eventType = getEventType(message);
40-
Schema memberSchema = eventSchema.member(eventType);
58+
return decodeEvent(frame);
59+
}
60+
61+
@Override
62+
public void onPrepare(Flow.Publisher<SerializableStruct> publisher) {
63+
this.publisher = publisher;
64+
}
65+
66+
private E decodeEvent(AwsEventFrame frame) {
67+
var message = frame.unwrap();
68+
var eventType = getEventType(message);
69+
var memberSchema = eventSchema.member(eventType);
4170
if (memberSchema == null) {
4271
throw new IllegalArgumentException("Unsupported event type: " + eventType);
4372
}
73+
var codecDeserializer = codec.createDeserializer(message.getPayload());
74+
var eventDeserializer = new AwsEventDeserializer(memberSchema, codecDeserializer);
75+
var builder = eventBuilder.get();
76+
return builder.deserialize(eventDeserializer).build();
77+
}
78+
79+
private IR decodeInitialResponse(AwsEventFrame frame) {
80+
var message = frame.unwrap();
81+
var codecDeserializer = codec.createDeserializer(message.getPayload());
82+
var builder = initialEventBuilder.get();
83+
builder.deserialize(codecDeserializer);
84+
var publisherMember = getPublisherMember(builder.schema());
85+
var responseDeserializer = new EventStreamDeserializer(publisherMember, publisher);
86+
builder.deserialize(responseDeserializer);
87+
return builder.build();
88+
}
4489

45-
return eventBuilder.get()
46-
.deserialize(
47-
new AwsEventDeserializer(
48-
memberSchema,
49-
codec.createDeserializer(message.getPayload())))
50-
.build();
90+
private Schema getPublisherMember(Schema schema) {
91+
for (var member : schema.members()) {
92+
if (member.memberTarget().hasTrait(TraitKey.STREAMING_TRAIT)) {
93+
return member;
94+
}
95+
}
96+
throw new IllegalArgumentException("cannot find streaming member");
5197
}
5298

5399
private String getEventType(Message message) {
54100
return message.getHeaders().get(":event-type").getString();
55101
}
56102

57-
private String getMessageType(Message message) {
58-
return message.getHeaders().get(":message-type").getString();
103+
static class EventStreamDeserializer extends SpecificShapeDeserializer {
104+
private final Schema publisherMember;
105+
private final Flow.Publisher<? extends SerializableStruct> publisher;
106+
107+
EventStreamDeserializer(Schema publisherMember, Flow.Publisher<? extends SerializableStruct> publisher) {
108+
this.publisherMember = publisherMember;
109+
this.publisher = publisher;
110+
}
111+
112+
@Override
113+
public Flow.Publisher<? extends SerializableStruct> readEventStream(Schema schema) {
114+
return publisher;
115+
}
116+
117+
@Override
118+
public <T> void readStruct(Schema schema, T state, ShapeDeserializer.StructMemberConsumer<T> consumer) {
119+
consumer.accept(state, publisherMember, this);
120+
}
59121
}
60122
}

0 commit comments

Comments
 (0)