mirror of
				https://github.com/RicterZ/nhentai.git
				synced 2025-11-03 18:50:53 +01:00 
			
		
		
		
	fix: semaphore bound to different event loop
This commit is contained in:
		@@ -13,6 +13,7 @@ import httpx
 | 
			
		||||
 | 
			
		||||
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
 | 
			
		||||
 | 
			
		||||
 | 
			
		||||
class NHentaiImageNotExistException(Exception):
 | 
			
		||||
    pass
 | 
			
		||||
 | 
			
		||||
@@ -32,16 +33,24 @@ def download_callback(result):
 | 
			
		||||
        logger.log(16, f'{data} downloaded successfully')
 | 
			
		||||
 | 
			
		||||
 | 
			
		||||
async def fiber(tasks):
 | 
			
		||||
    for completed_task in asyncio.as_completed(tasks):
 | 
			
		||||
        try:
 | 
			
		||||
            result = await completed_task
 | 
			
		||||
            logger.info(f'{result[1]} download completed')
 | 
			
		||||
        except Exception as e:
 | 
			
		||||
            logger.error(f'An error occurred: {e}')
 | 
			
		||||
 | 
			
		||||
 | 
			
		||||
class Downloader(Singleton):
 | 
			
		||||
    def __init__(self, path='', threads=5, timeout=30, delay=0):
 | 
			
		||||
        self.semaphore = asyncio.Semaphore(threads)
 | 
			
		||||
        self.threads = threads
 | 
			
		||||
        self.path = str(path)
 | 
			
		||||
        self.timeout = timeout
 | 
			
		||||
        self.delay = delay
 | 
			
		||||
 | 
			
		||||
    async def _semaphore_download(self, *args, **kwargs):
 | 
			
		||||
        # This sets a concurrency limit for AsyncIO
 | 
			
		||||
        async with self.semaphore:
 | 
			
		||||
    async def _semaphore_download(self, semaphore, *args, **kwargs):
 | 
			
		||||
        async with semaphore:
 | 
			
		||||
            return await self.download(*args, **kwargs)
 | 
			
		||||
 | 
			
		||||
    async def download(self, url, folder='', filename='', retried=0, proxy=None):
 | 
			
		||||
@@ -143,19 +152,14 @@ class Downloader(Singleton):
 | 
			
		||||
            # Assuming we want to continue with rest of process.
 | 
			
		||||
            return True
 | 
			
		||||
 | 
			
		||||
        async def fiber(tasks):
 | 
			
		||||
            for completed_task in asyncio.as_completed(tasks):
 | 
			
		||||
                try:
 | 
			
		||||
                    result = await completed_task
 | 
			
		||||
                    logger.info(f'{result[1]} download completed')
 | 
			
		||||
                except Exception as e:
 | 
			
		||||
                    logger.error(f'An error occurred: {e}')
 | 
			
		||||
        semaphore = asyncio.Semaphore(self.threads)
 | 
			
		||||
 | 
			
		||||
        tasks = [
 | 
			
		||||
            self._semaphore_download(url, filename=os.path.basename(urlparse(url).path))
 | 
			
		||||
        coroutines = [
 | 
			
		||||
            self._semaphore_download(semaphore, url, filename=os.path.basename(urlparse(url).path))
 | 
			
		||||
            for url in queue
 | 
			
		||||
        ]
 | 
			
		||||
 | 
			
		||||
        # Prevent coroutines infection
 | 
			
		||||
        asyncio.run(fiber(tasks))
 | 
			
		||||
        asyncio.run(fiber(coroutines))
 | 
			
		||||
 | 
			
		||||
        return True
 | 
			
		||||
 
 | 
			
		||||
		Reference in New Issue
	
	Block a user