(
self,
client_reader: asyncio.StreamReader,
client_writer: asyncio.StreamWriter,
server_reader: asyncio.StreamReader,
server_writer: asyncio.StreamWriter,
)
| 232 | pass |
| 233 | |
| 234 | async def _forward_data( |
| 235 | self, |
| 236 | client_reader: asyncio.StreamReader, |
| 237 | client_writer: asyncio.StreamWriter, |
| 238 | server_reader: asyncio.StreamReader, |
| 239 | server_writer: asyncio.StreamWriter, |
| 240 | ) -> None: |
| 241 | async def _forward( |
| 242 | reader: asyncio.StreamReader, writer: asyncio.StreamWriter |
| 243 | ) -> None: |
| 244 | try: |
| 245 | while True: |
| 246 | data = await reader.read(8192) |
| 247 | if not data: |
| 248 | break |
| 249 | writer.write(data) |
| 250 | await writer.drain() |
| 251 | except ConnectionResetError: |
| 252 | self.logger.debug("Connection reset by peer.") |
| 253 | except Exception as e: |
| 254 | self.logger.error(f"Error forwarding data: {e}", exc_info=True) |
| 255 | finally: |
| 256 | self._safe_close(writer) |
| 257 | |
| 258 | client_to_server = asyncio.create_task(_forward(client_reader, server_writer)) |
| 259 | server_to_client = asyncio.create_task(_forward(server_reader, client_writer)) |
| 260 | |
| 261 | tasks = [client_to_server, server_to_client] |
| 262 | try: |
| 263 | _done, pending = await asyncio.wait( |
| 264 | tasks, return_when=asyncio.FIRST_COMPLETED |
| 265 | ) |
| 266 | except asyncio.CancelledError: |
| 267 | for task in tasks: |
| 268 | task.cancel() |
| 269 | await asyncio.gather(*tasks, return_exceptions=True) |
| 270 | raise |
| 271 | |
| 272 | for task in pending: |
| 273 | task.cancel() |
| 274 | try: |
| 275 | await task |
| 276 | except asyncio.CancelledError: |
| 277 | pass |
| 278 | |
| 279 | async def _forward_data_with_interception( |
| 280 | self, |
no outgoing calls