(self, event: events.Event)
| 135 | |
| 136 | @expect(events.DataReceived, events.ConnectionClosed, WebSocketMessageInjected) |
| 137 | def relay_messages(self, event: events.Event) -> layer.CommandGenerator[None]: |
| 138 | assert self.flow.websocket # satisfy type checker |
| 139 | |
| 140 | if isinstance(event, events.ConnectionEvent): |
| 141 | from_client = event.connection == self.context.client |
| 142 | injected = False |
| 143 | elif isinstance(event, WebSocketMessageInjected): |
| 144 | from_client = event.message.from_client |
| 145 | injected = True |
| 146 | else: |
| 147 | raise AssertionError(f"Unexpected event: {event}") |
| 148 | |
| 149 | from_str = "client" if from_client else "server" |
| 150 | if from_client: |
| 151 | src_ws = self.client_ws |
| 152 | dst_ws = self.server_ws |
| 153 | else: |
| 154 | src_ws = self.server_ws |
| 155 | dst_ws = self.client_ws |
| 156 | |
| 157 | if isinstance(event, events.DataReceived): |
| 158 | src_ws.receive_data(event.data) |
| 159 | elif isinstance(event, events.ConnectionClosed): |
| 160 | src_ws.receive_data(None) |
| 161 | elif isinstance(event, WebSocketMessageInjected): |
| 162 | fragmentizer = Fragmentizer([], event.message.type == Opcode.TEXT) |
| 163 | src_ws._events.extend(fragmentizer(event.message.content)) |
| 164 | else: # pragma: no cover |
| 165 | raise AssertionError(f"Unexpected event: {event}") |
| 166 | |
| 167 | for ws_event in src_ws.events(): |
| 168 | if isinstance(ws_event, wsproto.events.Message): |
| 169 | is_text = isinstance(ws_event.data, str) |
| 170 | if is_text: |
| 171 | typ = Opcode.TEXT |
| 172 | src_ws.frame_buf[-1] += ws_event.data.encode() |
| 173 | else: |
| 174 | typ = Opcode.BINARY |
| 175 | src_ws.frame_buf[-1] += ws_event.data |
| 176 | |
| 177 | if ws_event.message_finished: |
| 178 | content = b"".join(src_ws.frame_buf) |
| 179 | |
| 180 | fragmentizer = Fragmentizer(src_ws.frame_buf, is_text) |
| 181 | src_ws.frame_buf = [b""] |
| 182 | |
| 183 | message = websocket.WebSocketMessage( |
| 184 | typ, from_client, content, injected=injected |
| 185 | ) |
| 186 | self.flow.websocket.messages.append(message) |
| 187 | yield WebsocketMessageHook(self.flow) |
| 188 | |
| 189 | if not message.dropped: |
| 190 | for msg in fragmentizer(message.content): |
| 191 | yield dst_ws.send2(msg) |
| 192 | |
| 193 | elif ws_event.frame_finished: |
| 194 | src_ws.frame_buf.append(b"") |
nothing calls this directly
no test coverage detected