PluginProbe
Media Cloud Sync / 1.4.2
Media Cloud Sync v1.4.2
1.4.2 1.4.1 1.4.0 1.3.12 1.3.11 1.3.10 trunk 1.0.0 1.0.1 1.0.2 1.0.3 1.1.0 1.1.1 1.2.0 1.2.10 1.2.11 1.2.12 1.2.13 1.2.2 1.2.3 1.2.4 1.2.5 1.2.6 1.2.7 1.2.8 All 36 releases
← All changes | includes/sdk/s3/Aws/Api/Parser/EventParsingIterator.php +94 -13 1.2.2 → 1.4.2 View file →
@@ -20,32 +20,63 @@
20 20 /** @var AbstractParser */
21 21 private $parser;
22 22 public function __construct(StreamInterface $stream, StructureShape $shape, AbstractParser $parser)
23 23 {
24 - $this->decodingIterator = new DecodingEventStreamIterator($stream);
24 + $this->decodingIterator = $this->chooseDecodingIterator($stream);
25 25 $this->shape = $shape;
26 26 $this->parser = $parser;
27 27 }
28 + /**
29 + * This method choose a decoding iterator implementation based on if the stream
30 + * is seekable or not.
31 + *
32 + * @param $stream
33 + *
34 + * @return Iterator
35 + */
36 + private function chooseDecodingIterator($stream)
37 + {
38 + if ($stream->isSeekable()) {
39 + return new DecodingEventStreamIterator($stream);
40 + } else {
41 + return new NonSeekableStreamDecodingEventStreamIterator($stream);
42 + }
43 + }
44 + /**
45 + * @return mixed
46 + */
28 47 #[\ReturnTypeWillChange]
29 48 public function current()
30 49 {
31 50 return $this->parseEvent($this->decodingIterator->current());
32 51 }
52 + /**
53 + * @return mixed
54 + */
33 55 #[\ReturnTypeWillChange]
34 56 public function key()
35 57 {
36 58 return $this->decodingIterator->key();
37 59 }
60 + /**
61 + * @return void
62 + */
38 63 #[\ReturnTypeWillChange]
39 64 public function next()
40 65 {
41 66 $this->decodingIterator->next();
42 67 }
68 + /**
69 + * @return void
70 + */
43 71 #[\ReturnTypeWillChange]
44 72 public function rewind()
45 73 {
46 74 $this->decodingIterator->rewind();
47 75 }
76 + /**
77 + * @return bool
78 + */
48 79 #[\ReturnTypeWillChange]
49 80 public function valid()
50 81 {
51 82 return $this->decodingIterator->valid();
@@ -55,32 +86,82 @@
55 86 if (!empty($event['headers'][':message-type'])) {
56 87 if ($event['headers'][':message-type'] === 'error') {
57 88 return $this->parseError($event);
58 89 }
90 + if ($event['headers'][':message-type'] === 'exception') {
91 + return $this->parseException($event);
92 + }
59 93 if ($event['headers'][':message-type'] !== 'event') {
60 94 throw new ParserException('Failed to parse unknown message type.');
61 95 }
62 96 }
63 - if (empty($event['headers'][':event-type'])) {
97 + $eventType = $event['headers'][':event-type'] ?? null;
98 + if (empty($eventType)) {
64 99 throw new ParserException('Failed to parse without event type.');
65 100 }
66 - $eventShape = $this->shape->getMember($event['headers'][':event-type']);
67 - $parsedEvent = [];
68 - foreach ($eventShape['members'] as $shape => $details) {
69 - if (!empty($details['eventpayload'])) {
70 - $payloadShape = $eventShape->getMember($shape);
71 - if ($payloadShape['type'] === 'blob') {
72 - $parsedEvent[$shape] = $event['payload'];
101 + $eventPayload = $event['payload'];
102 + if ($eventType === 'initial-response') {
103 + return $this->parseInitialResponseEvent($eventPayload);
104 + }
105 + $eventShape = $this->shape->getMember($eventType);
106 + return [$eventType => \array_merge($this->parseEventHeaders($event['headers'], $eventShape), $this->parseEventPayload($eventPayload, $eventShape))];
107 + }
108 + /**
109 + * @param $headers
110 + * @param $eventShape
111 + *
112 + * @return array
113 + */
114 + private function parseEventHeaders($headers, $eventShape) : array
115 + {
116 + $parsedHeaders = [];
117 + foreach ($eventShape->getMembers() as $memberName => $memberProps) {
118 + if (isset($memberProps['eventheader'])) {
119 + $parsedHeaders[$memberName] = $headers[$memberName];
120 + }
121 + }
122 + return $parsedHeaders;
123 + }
124 + /**
125 + * @param $payload
126 + * @param $eventShape
127 + *
128 + * @return array
129 + */
130 + private function parseEventPayload($payload, $eventShape) : array
131 + {
132 + $parsedPayload = [];
133 + foreach ($eventShape->getMembers() as $memberName => $memberProps) {
134 + $memberShape = $eventShape->getMember($memberName);
135 + if (isset($memberProps['eventpayload'])) {
136 + if ($memberShape->getType() === 'blob') {
137 + $parsedPayload[$memberName] = $payload;
73 138 } else {
74 - $parsedEvent[$shape] = $this->parser->parseMemberFromStream($event['payload'], $payloadShape, null);
139 + $parsedPayload[$memberName] = $this->parser->parseMemberFromStream($payload, $memberShape, null);
75 140 }
76 - } else {
77 - $parsedEvent[$shape] = $event['headers'][$shape];
141 + break;
78 142 }
79 143 }
80 - return [$event['headers'][':event-type'] => $parsedEvent];
144 + if (empty($parsedPayload) && !empty($payload->getContents())) {
145 + /**
146 + * If we did not find a member with an eventpayload trait, then we should deserialize the payload
147 + * using the event's shape.
148 + */
149 + $parsedPayload = $this->parser->parseMemberFromStream($payload, $eventShape, null);
150 + }
151 + return $parsedPayload;
81 152 }
82 153 private function parseError(array $event)
83 154 {
84 155 throw new EventStreamDataException($event['headers'][':error-code'], $event['headers'][':error-message']);
156 + }
157 + private function parseException(array $event)
158 + {
159 + $payload = $event['payload']?->getContents();
160 + $parsedPayload = \json_decode($payload, \true);
161 + throw new EventStreamDataException($event['headers'][':exception-type'] ?? 'Unknown', $parsedPayload['message'] ?? $payload);
162 + }
163 + private function parseInitialResponseEvent($payload) : array
164 + {
165 + return ['initial-response' => \json_decode($payload, \true)];
85 166 }
86 167 }