mirror of
https://github.com/LmeSzinc/StarRailCopilot.git
synced 2024-11-22 00:35:34 +00:00
Fix: [ALAS] Limit pool size when using run_in_executor()
(cherry picked from commit e1bb483aaafa293a14e2faec7c9baadb236e1a01)
This commit is contained in:
parent
5a9a1beaa9
commit
ce7d971485
@ -721,7 +721,12 @@ class Connection(ConnectionAttr):
|
|||||||
serial_list (list[str]):
|
serial_list (list[str]):
|
||||||
"""
|
"""
|
||||||
import asyncio
|
import asyncio
|
||||||
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
ev = asyncio.new_event_loop()
|
ev = asyncio.new_event_loop()
|
||||||
|
pool = ThreadPoolExecutor(
|
||||||
|
max_workers=len(serial_list),
|
||||||
|
thread_name_prefix='adb_brute_force_connect',
|
||||||
|
)
|
||||||
|
|
||||||
def _connect(serial):
|
def _connect(serial):
|
||||||
msg = self.adb_client.connect(serial)
|
msg = self.adb_client.connect(serial)
|
||||||
@ -729,10 +734,12 @@ class Connection(ConnectionAttr):
|
|||||||
return msg
|
return msg
|
||||||
|
|
||||||
async def connect():
|
async def connect():
|
||||||
tasks = [ev.run_in_executor(None, _connect, serial) for serial in serial_list]
|
tasks = [ev.run_in_executor(pool, _connect, serial) for serial in serial_list]
|
||||||
await asyncio.gather(*tasks)
|
await asyncio.gather(*tasks)
|
||||||
|
|
||||||
ev.run_until_complete(connect())
|
ev.run_until_complete(connect())
|
||||||
|
pool.shutdown(wait=False)
|
||||||
|
ev.close()
|
||||||
|
|
||||||
@Config.when(DEVICE_OVER_HTTP=True)
|
@Config.when(DEVICE_OVER_HTTP=True)
|
||||||
def adb_connect(self):
|
def adb_connect(self):
|
||||||
|
@ -275,11 +275,25 @@ class NemuIpcImpl:
|
|||||||
|
|
||||||
def __exit__(self, exc_type, exc_val, exc_tb):
|
def __exit__(self, exc_type, exc_val, exc_tb):
|
||||||
self.disconnect()
|
self.disconnect()
|
||||||
|
if has_cached_property(self, '_ev'):
|
||||||
|
self._ev.close()
|
||||||
|
del_cached_property(self, '_ev')
|
||||||
|
if has_cached_property(self, '_pool'):
|
||||||
|
self._pool.shutdown(wait=False)
|
||||||
|
del_cached_property(self, '_pool')
|
||||||
|
|
||||||
@cached_property
|
@cached_property
|
||||||
def _ev(self):
|
def _ev(self):
|
||||||
return asyncio.new_event_loop()
|
return asyncio.new_event_loop()
|
||||||
|
|
||||||
|
@cached_property
|
||||||
|
def _pool(self):
|
||||||
|
from concurrent.futures import ThreadPoolExecutor
|
||||||
|
return ThreadPoolExecutor(
|
||||||
|
max_workers=1,
|
||||||
|
thread_name_prefix='NemuIpc',
|
||||||
|
)
|
||||||
|
|
||||||
async def ev_run_async(self, func, *args, timeout=0.15, **kwargs):
|
async def ev_run_async(self, func, *args, timeout=0.15, **kwargs):
|
||||||
"""
|
"""
|
||||||
Args:
|
Args:
|
||||||
@ -294,7 +308,7 @@ class NemuIpcImpl:
|
|||||||
func_wrapped = partial(func, *args, **kwargs)
|
func_wrapped = partial(func, *args, **kwargs)
|
||||||
# Increased timeout for slow PCs
|
# Increased timeout for slow PCs
|
||||||
# Default screenshot interval is 0.2s, so a 0.15s timeout would have a fast retry without extra time costs
|
# Default screenshot interval is 0.2s, so a 0.15s timeout would have a fast retry without extra time costs
|
||||||
result = await asyncio.wait_for(self._ev.run_in_executor(None, func_wrapped), timeout=timeout)
|
result = await asyncio.wait_for(self._ev.run_in_executor(self._pool, func_wrapped), timeout=timeout)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
def ev_run_sync(self, func, *args, **kwargs):
|
def ev_run_sync(self, func, *args, **kwargs):
|
||||||
|
Loading…
Reference in New Issue
Block a user