MCPcopy Create free account
hub / github.com/CJackHwang/AIstudioProxyAPI / _forward_data

Method _forward_data

stream/proxy_server.py:234–277  ·  view source on GitHub ↗
(
        self,
        client_reader: asyncio.StreamReader,
        client_writer: asyncio.StreamWriter,
        server_reader: asyncio.StreamReader,
        server_writer: asyncio.StreamWriter,
    )

Source from the content-addressed store, hash-verified

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,

Calls

no outgoing calls