()
| 352 | self._safe_close(server_writer) |
| 353 | |
| 354 | async def _process_server_data(): |
| 355 | nonlocal server_buffer, should_sniff |
| 356 | try: |
| 357 | while True: |
| 358 | data = await server_reader.read(8192) |
| 359 | if not data: |
| 360 | break |
| 361 | |
| 362 | server_buffer.extend(data) |
| 363 | if b"\r\n\r\n" in server_buffer: |
| 364 | headers_end = server_buffer.find(b"\r\n\r\n") + 4 |
| 365 | headers_data = server_buffer[:headers_end] |
| 366 | body_data = server_buffer[headers_end:] |
| 367 | |
| 368 | lines = headers_data.split(b"\r\n") |
| 369 | status_code = 200 |
| 370 | status_message = "OK" |
| 371 | if lines and lines[0]: |
| 372 | try: |
| 373 | status_line = lines[0].decode("utf-8") |
| 374 | parts = status_line.split(" ", 2) |
| 375 | if len(parts) >= 2: |
| 376 | status_code = int(parts[1]) |
| 377 | status_message = parts[2] if len(parts) > 2 else "" |
| 378 | except (ValueError, UnicodeDecodeError): |
| 379 | pass |
| 380 | |
| 381 | headers: dict[str, str] = {} |
| 382 | for i in range(1, len(lines)): |
| 383 | if not lines[i]: |
| 384 | continue |
| 385 | try: |
| 386 | key, value = lines[i].decode("utf-8").split(":", 1) |
| 387 | headers[key.strip()] = value.strip() |
| 388 | except ValueError: |
| 389 | continue |
| 390 | |
| 391 | if should_sniff: |
| 392 | try: |
| 393 | if status_code >= 400: |
| 394 | self.logger.error( |
| 395 | f"[UPSTREAM ERROR] {status_code} {status_message}" |
| 396 | ) |
| 397 | if self.queue is not None: |
| 398 | error_payload = { |
| 399 | "error": True, |
| 400 | "status": status_code, |
| 401 | "message": f"{status_code} {status_message}", |
| 402 | "done": True, |
| 403 | } |
| 404 | self.queue.put(json.dumps(error_payload)) |
| 405 | else: |
| 406 | resp = await self.interceptor.process_response( |
| 407 | bytes(body_data), host, "", headers |
| 408 | ) |
| 409 | if self.queue is not None: |
| 410 | payload = { |
| 411 | "ts": request_context.get("request_ts", 0), |
nothing calls this directly
no test coverage detected