added debug statements.
This commit is contained in:
parent
966e847bbd
commit
e20365a8e6
2 changed files with 17 additions and 6 deletions
|
|
@ -140,6 +140,8 @@ class AsyncPool:
|
||||||
if self._log_every_n is not None and (self._total_queued % self._log_every_n) == 0:
|
if self._log_every_n is not None and (self._total_queued % self._log_every_n) == 0:
|
||||||
self._logger.info(f"pushed {self._total_queued}/{self._expected_total} items to {self._name} pool.")
|
self._logger.info(f"pushed {self._total_queued}/{self._expected_total} items to {self._name} pool.")
|
||||||
|
|
||||||
|
self._logger.debug(f"'{self._name}' pool has received a new job. {args} {kwargs}")
|
||||||
|
|
||||||
return future
|
return future
|
||||||
|
|
||||||
def start(self):
|
def start(self):
|
||||||
|
|
|
||||||
|
|
@ -337,30 +337,39 @@ class DownloadQueue:
|
||||||
name='download_pool',
|
name='download_pool',
|
||||||
logger=logging.getLogger('WorkerPool'),
|
logger=logging.getLogger('WorkerPool'),
|
||||||
) as executor:
|
) as executor:
|
||||||
|
lastLog = time.time()
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
while True:
|
while True:
|
||||||
if executor.has_open_workers() is False:
|
if executor.has_open_workers() is True:
|
||||||
LOG.debug(f'Waiting for workers to be available.')
|
|
||||||
await asyncio.sleep(1)
|
|
||||||
else:
|
|
||||||
break
|
break
|
||||||
|
if time.time() - lastLog > 120:
|
||||||
|
lastLog = time.time()
|
||||||
|
LOG.info(f'Waiting for workers to be free.')
|
||||||
|
await asyncio.sleep(1)
|
||||||
|
|
||||||
|
LOG.debug(f"Has '{executor.get_available_workers()}' free workers.")
|
||||||
|
|
||||||
while not self.queue.hasDownloads():
|
while not self.queue.hasDownloads():
|
||||||
LOG.info(f'Waiting for item to download. [{executor.get_available_workers()}] workers available.')
|
LOG.info(f'Waiting for item to download.')
|
||||||
await self.event.wait()
|
await self.event.wait()
|
||||||
self.event.clear()
|
self.event.clear()
|
||||||
|
LOG.debug(f"Cleared wait event.")
|
||||||
|
|
||||||
entry = self.queue.getNextDownload()
|
entry = self.queue.getNextDownload()
|
||||||
await asyncio.sleep(0.2)
|
await asyncio.sleep(0.2)
|
||||||
|
|
||||||
|
LOG.debug(f"Pushing {entry=} to executor.")
|
||||||
|
|
||||||
if entry.started() is False and entry.is_canceled() is False:
|
if entry.started() is False and entry.is_canceled() is False:
|
||||||
await executor.push(id=entry.info._id, entry=entry)
|
await executor.push(id=entry.info._id, entry=entry)
|
||||||
|
LOG.debug(f"Pushed {entry=} to executor.")
|
||||||
await asyncio.sleep(1)
|
await asyncio.sleep(1)
|
||||||
|
|
||||||
async def __download(self):
|
async def __download(self):
|
||||||
while True:
|
while True:
|
||||||
while self.queue.empty():
|
while self.queue.empty():
|
||||||
LOG.info('waiting for item to download.')
|
LOG.info('Waiting for item to download.')
|
||||||
await self.event.wait()
|
await self.event.wait()
|
||||||
self.event.clear()
|
self.event.clear()
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue