(self, event: Event)
| 503 | self.stream_queue = collections.defaultdict(list) |
| 504 | |
| 505 | def _handle_event(self, event: Event) -> CommandGenerator[None]: |
| 506 | # We can't reuse stream ids from the client because they may arrived reordered here |
| 507 | # and HTTP/2 forbids opening a stream on a lower id than what was previously sent (see test_stream_concurrency). |
| 508 | # To mitigate this, we transparently map the outside's stream id to our stream id. |
| 509 | if isinstance(event, HttpEvent): |
| 510 | ours = self.our_stream_id.get(event.stream_id, None) |
| 511 | if ours is None: |
| 512 | no_free_streams = self.h2_conn.open_outbound_streams >= ( |
| 513 | self.provisional_max_concurrency |
| 514 | or self.h2_conn.remote_settings.max_concurrent_streams |
| 515 | ) |
| 516 | if no_free_streams: |
| 517 | self.stream_queue[event.stream_id].append(event) |
| 518 | return |
| 519 | ours = self.h2_conn.get_next_available_stream_id() |
| 520 | self.our_stream_id[event.stream_id] = ours |
| 521 | self.their_stream_id[ours] = event.stream_id |
| 522 | event.stream_id = ours |
| 523 | |
| 524 | for cmd in self._handle_event2(event): |
| 525 | if isinstance(cmd, ReceiveHttp): |
| 526 | cmd.event.stream_id = self.their_stream_id[cmd.event.stream_id] |
| 527 | yield cmd |
| 528 | |
| 529 | can_resume_queue = self.stream_queue and self.h2_conn.open_outbound_streams < ( |
| 530 | self.provisional_max_concurrency |
| 531 | or self.h2_conn.remote_settings.max_concurrent_streams |
| 532 | ) |
| 533 | if can_resume_queue: |
| 534 | # popitem would be LIFO, but we want FIFO. |
| 535 | events = self.stream_queue.pop(next(iter(self.stream_queue))) |
| 536 | for event in events: |
| 537 | yield from self._handle_event(event) |
| 538 | |
| 539 | def _handle_event2(self, event: Event) -> CommandGenerator[None]: |
| 540 | if isinstance(event, Wakeup): |
nothing calls this directly
no test coverage detected