returns true if further processing should be stopped.
(self, event: h2.events.Event)
| 223 | raise AssertionError(f"Unexpected event: {event!r}") |
| 224 | |
| 225 | def handle_h2_event(self, event: h2.events.Event) -> CommandGenerator[bool]: |
| 226 | """returns true if further processing should be stopped.""" |
| 227 | if isinstance(event, h2.events.DataReceived): |
| 228 | state = self.streams.get(event.stream_id, None) |
| 229 | if state is StreamState.HEADERS_RECEIVED: |
| 230 | is_empty_eos_data_frame = event.stream_ended and not event.data |
| 231 | if not is_empty_eos_data_frame: |
| 232 | yield ReceiveHttp(self.ReceiveData(event.stream_id, event.data)) |
| 233 | elif state is StreamState.EXPECTING_HEADERS: |
| 234 | yield from self.protocol_error( |
| 235 | f"Received HTTP/2 data frame, expected headers." |
| 236 | ) |
| 237 | return True |
| 238 | self.h2_conn.acknowledge_received_data( |
| 239 | event.flow_controlled_length, event.stream_id |
| 240 | ) |
| 241 | elif isinstance(event, h2.events.TrailersReceived): |
| 242 | trailers = http.Headers(event.headers) |
| 243 | yield ReceiveHttp(self.ReceiveTrailers(event.stream_id, trailers)) |
| 244 | elif isinstance(event, h2.events.StreamEnded): |
| 245 | state = self.streams.get(event.stream_id, None) |
| 246 | if state is StreamState.HEADERS_RECEIVED: |
| 247 | yield ReceiveHttp(self.ReceiveEndOfMessage(event.stream_id)) |
| 248 | elif state is StreamState.EXPECTING_HEADERS: |
| 249 | raise AssertionError("unreachable") |
| 250 | if self.is_closed(event.stream_id): |
| 251 | self.streams.pop(event.stream_id, None) |
| 252 | elif isinstance(event, h2.events.StreamReset): |
| 253 | if event.stream_id in self.streams: |
| 254 | try: |
| 255 | err_str = h2.errors.ErrorCodes(event.error_code).name |
| 256 | except ValueError: |
| 257 | err_str = str(event.error_code) |
| 258 | match event.error_code: |
| 259 | case h2.errors.ErrorCodes.CANCEL: |
| 260 | err_code = ErrorCode.CANCEL |
| 261 | case h2.errors.ErrorCodes.HTTP_1_1_REQUIRED: |
| 262 | err_code = ErrorCode.HTTP_1_1_REQUIRED |
| 263 | case _: |
| 264 | err_code = self.ReceiveProtocolError.code |
| 265 | yield ReceiveHttp( |
| 266 | self.ReceiveProtocolError( |
| 267 | event.stream_id, |
| 268 | f"stream reset by client ({err_str})", |
| 269 | code=err_code, |
| 270 | ) |
| 271 | ) |
| 272 | self.streams.pop(event.stream_id) |
| 273 | else: |
| 274 | pass # We don't track priority frames which could be followed by a stream reset here. |
| 275 | elif isinstance(event, h2.exceptions.ProtocolError): |
| 276 | yield from self.protocol_error(f"HTTP/2 protocol error: {event}") |
| 277 | return True |
| 278 | elif isinstance(event, h2.events.ConnectionTerminated): |
| 279 | yield from self.close_connection(f"HTTP/2 connection closed: {event!r}") |
| 280 | return True |
| 281 | # The implementation above isn't really ideal, we should probably only terminate streams > last_stream_id? |
| 282 | # We currently lack a mechanism to signal that connections are still active but cannot be reused. |
no test coverage detected