()
| 290 | request_context = {"request_ts": 0.0} |
| 291 | |
| 292 | async def _process_client_data(): |
| 293 | nonlocal client_buffer, should_sniff |
| 294 | try: |
| 295 | while True: |
| 296 | data = await client_reader.read(8192) |
| 297 | if not data: |
| 298 | break |
| 299 | client_buffer.extend(data) |
| 300 | |
| 301 | if b"\r\n\r\n" in client_buffer: |
| 302 | headers_end = client_buffer.find(b"\r\n\r\n") + 4 |
| 303 | headers_data = client_buffer[:headers_end] |
| 304 | body_data = client_buffer[headers_end:] |
| 305 | |
| 306 | lines = headers_data.split(b"\r\n") |
| 307 | request_line = lines[0].decode("utf-8") |
| 308 | |
| 309 | try: |
| 310 | _method, path, _ = request_line.split(" ") |
| 311 | except ValueError: |
| 312 | server_writer.write(client_buffer) |
| 313 | await server_writer.drain() |
| 314 | client_buffer.clear() |
| 315 | continue |
| 316 | |
| 317 | if "GenerateContent" in path or "generateContent" in path: |
| 318 | should_sniff = True |
| 319 | request_context["request_ts"] = time.time() |
| 320 | # Reset interceptor state for new request to prevent |
| 321 | # state leakage from previous requests |
| 322 | self.interceptor.reset_for_new_request() |
| 323 | self.logger.debug( |
| 324 | f"[Proxy] Detected GenerateContent request: {path[:60]}..." |
| 325 | ) |
| 326 | processed_body = await self.interceptor.process_request( |
| 327 | bytes(body_data), host, path |
| 328 | ) |
| 329 | server_writer.write(headers_data) |
| 330 | if isinstance(processed_body, bytes): |
| 331 | server_writer.write(processed_body) |
| 332 | else: |
| 333 | should_sniff = False |
| 334 | server_writer.write(client_buffer) |
| 335 | |
| 336 | await server_writer.drain() |
| 337 | client_buffer.clear() |
| 338 | else: |
| 339 | server_writer.write(data) |
| 340 | await server_writer.drain() |
| 341 | client_buffer.clear() |
| 342 | except ConnectionResetError: |
| 343 | self.logger.debug("Connection reset by peer processing client data.") |
| 344 | except Exception as e: |
| 345 | if "Broken pipe" in str(e) or "Connection reset" in str(e): |
| 346 | self.logger.debug(f"[Proxy] Client disconnected: {e}") |
| 347 | else: |
| 348 | self.logger.error( |
| 349 | f"Error processing client data: {e}", exc_info=True |
nothing calls this directly
no test coverage detected