This commit is contained in:
ArabCoders 2024-12-16 22:06:54 +03:00
parent 4787579d67
commit f614dbc217
6 changed files with 152 additions and 149 deletions

76
Pipfile.lock generated
View file

@ -117,11 +117,11 @@
}, },
"aiosignal": { "aiosignal": {
"hashes": [ "hashes": [
"sha256:54cd96e15e1649b75d6c87526a6ff0b6c1b0dd3459f43d9ca11d48c339b68cfc", "sha256:45cde58e409a301715980c2b01d0c28bdde3770d8290b5eb2173759d9acb31a5",
"sha256:f8376fb07dd1e86a584e4fcdec80b36b7f81aac666ebc724e2c090300dd83b17" "sha256:a8c255c66fafb1e499c9351d0bf32ff2d8a0321595ebac3b93713656d2436f54"
], ],
"markers": "python_version >= '3.7'", "markers": "python_version >= '3.9'",
"version": "==1.3.1" "version": "==1.3.2"
}, },
"anyio": { "anyio": {
"hashes": [ "hashes": [
@ -142,11 +142,11 @@
}, },
"attrs": { "attrs": {
"hashes": [ "hashes": [
"sha256:5cfb1b9148b5b086569baec03f20d7b6bf3bcacc9a42bebf87ffaaca362f6346", "sha256:8f5c07333d543103541ba7be0e2ce16eeee8130cb0b3f9238ab904ce1e85baff",
"sha256:81921eb96de3191c8258c199618104dd27ac608d9366f5e35d011eae1867ede2" "sha256:ac96cd038792094f438ad1f6ff80837353805ac950cd2aa0e0625ef19850c308"
], ],
"markers": "python_version >= '3.7'", "markers": "python_version >= '3.8'",
"version": "==24.2.0" "version": "==24.3.0"
}, },
"bidict": { "bidict": {
"hashes": [ "hashes": [
@ -166,11 +166,11 @@
}, },
"certifi": { "certifi": {
"hashes": [ "hashes": [
"sha256:922820b53db7a7257ffbda3f597266d435245903d80737e34f8a45ff3e3230d8", "sha256:1275f7a45be9464efc1173084eaa30f866fe2e47d389406136d332ed4967ec56",
"sha256:bec941d2aa8195e248a60b31ff9f0558284cf01a52591ceda73ea9afffd69fd9" "sha256:b650d30f370c2b724812bee08008be0c4163b163ddaec3f2546c1caf65f191db"
], ],
"markers": "python_version >= '3.6'", "markers": "python_version >= '3.6'",
"version": "==2024.8.30" "version": "==2024.12.14"
}, },
"cffi": { "cffi": {
"hashes": [ "hashes": [
@ -281,36 +281,36 @@
}, },
"debugpy": { "debugpy": {
"hashes": [ "hashes": [
"sha256:1339e14c7d980407248f09824d1b25ff5c5616651689f1e0f0e51bdead3ea13e", "sha256:0e22f846f4211383e6a416d04b4c13ed174d24cc5d43f5fd52e7821d0ebc8920",
"sha256:17c5e0297678442511cf00a745c9709e928ea4ca263d764e90d233208889a19e", "sha256:116bf8342062246ca749013df4f6ea106f23bc159305843491f64672a55af2e5",
"sha256:1efbb3ff61487e2c16b3e033bc8595aea578222c08aaf3c4bf0f93fadbd662ee", "sha256:189058d03a40103a57144752652b3ab08ff02b7595d0ce1f651b9acc3a3a35a0",
"sha256:365e556a4772d7d0d151d7eb0e77ec4db03bcd95f26b67b15742b88cacff88e9", "sha256:23dc34c5e03b0212fa3c49a874df2b8b1b8fda95160bd79c01eb3ab51ea8d851",
"sha256:3d9755e77a2d680ce3d2c5394a444cf42be4a592caaf246dbfbdd100ffcf7ae5", "sha256:28e45b3f827d3bf2592f3cf7ae63282e859f3259db44ed2b129093ca0ac7940b",
"sha256:3e59842d6c4569c65ceb3751075ff8d7e6a6ada209ceca6308c9bde932bcef11", "sha256:2b26fefc4e31ff85593d68b9022e35e8925714a10ab4858fb1b577a8a48cb8cd",
"sha256:472a3994999fe6c0756945ffa359e9e7e2d690fb55d251639d07208dbc37caea", "sha256:32db46ba45849daed7ccf3f2e26f7a386867b077f39b2a974bb5c4c2c3b0a280",
"sha256:54a7e6d3014c408eb37b0b06021366ee985f1539e12fe49ca2ee0d392d9ceca5", "sha256:40499a9979c55f72f4eb2fc38695419546b62594f8af194b879d2a18439c97a9",
"sha256:5e565fc54b680292b418bb809f1386f17081d1346dca9a871bf69a8ac4071afe", "sha256:44b1b8e6253bceada11f714acf4309ffb98bfa9ac55e4fce14f9e5d4484287a1",
"sha256:62d22dacdb0e296966d7d74a7141aaab4bec123fa43d1a35ddcb39bf9fd29d70", "sha256:52c3cf9ecda273a19cc092961ee34eb9ba8687d67ba34cc7b79a521c1c64c4c0",
"sha256:66eeae42f3137eb428ea3a86d4a55f28da9bd5a4a3d369ba95ecc3a92c1bba53", "sha256:52d8a3166c9f2815bfae05f386114b0b2d274456980d41f320299a8d9a5615a7",
"sha256:6953b335b804a41f16a192fa2e7851bdcfd92173cbb2f9f777bb934f49baab65", "sha256:61bc8b3b265e6949855300e84dc93d02d7a3a637f2aec6d382afd4ceb9120c9f",
"sha256:7c4d65d03bee875bcb211c76c1d8f10f600c305dbd734beaed4077e902606fee", "sha256:654130ca6ad5de73d978057eaf9e582244ff72d4574b3e106fb8d3d2a0d32458",
"sha256:7e646e62d4602bb8956db88b1e72fe63172148c1e25c041e03b103a25f36673c", "sha256:6ad2688b69235c43b020e04fecccdf6a96c8943ca9c2fb340b8adc103c655e57",
"sha256:7e8b079323a56f719977fde9d8115590cb5e7a1cba2fcee0986ef8817116e7c1", "sha256:6c1f6a173d1140e557347419767d2b14ac1c9cd847e0b4c5444c7f3144697e4e",
"sha256:8138efff315cd09b8dcd14226a21afda4ca582284bf4215126d87342bba1cc66", "sha256:84e511a7545d11683d32cdb8f809ef63fc17ea2a00455cc62d0a4dbb4ed1c308",
"sha256:8e99c0b1cc7bf86d83fb95d5ccdc4ad0586d4432d489d1f54e4055bcc795f693", "sha256:85de8474ad53ad546ff1c7c7c89230db215b9b8a02754d41cb5a76f70d0be296",
"sha256:957363d9a7a6612a37458d9a15e72d03a635047f946e5fceee74b50d52a9c8e2", "sha256:8988f7163e4381b0da7696f37eec7aca19deb02e500245df68a7159739bbd0d3",
"sha256:957ecffff80d47cafa9b6545de9e016ae8c9547c98a538ee96ab5947115fb3dd", "sha256:8da1db4ca4f22583e834dcabdc7832e56fe16275253ee53ba66627b86e304da1",
"sha256:ada7fb65102a4d2c9ab62e8908e9e9f12aed9d76ef44880367bc9308ebe49a0f", "sha256:8ffc382e4afa4aee367bf413f55ed17bd91b191dcaf979890af239dda435f2a1",
"sha256:b74a49753e21e33e7cf030883a92fa607bddc4ede1aa4145172debc637780040", "sha256:987bce16e86efa86f747d5151c54e91b3c1e36acc03ce1ddb50f9d09d16ded0e",
"sha256:c36856343cbaa448171cba62a721531e10e7ffb0abff838004701454149bc037", "sha256:ad7efe588c8f5cf940f40c3de0cd683cc5b76819446abaa50dc0829a30c094db",
"sha256:cc37a6c9987ad743d9c3a14fa1b1a14b7e4e6041f9dd0c8abf8895fe7a97b899", "sha256:bb3b15e25891f38da3ca0740271e63ab9db61f41d4d8541745cfc1824252cb28",
"sha256:cfe1e6c6ad7178265f74981edf1154ffce97b69005212fbc90ca22ddfe3d017e", "sha256:c928bbf47f65288574b78518449edaa46c82572d340e2750889bbf8cd92f3737",
"sha256:e46b420dc1bea64e5bbedd678148be512442bc589b0111bd799367cde051e71a", "sha256:ce291a5aca4985d82875d6779f61375e959208cdf09fcec40001e65fb0a54768",
"sha256:ff54ef77ad9f5c425398efb150239f6fe8e20c53ae2f68367eba7ece1e96226d" "sha256:d8768edcbeb34da9e11bcb8b5c2e0958d25218df7a6e56adf415ef262cd7b6d1"
], ],
"index": "pypi", "index": "pypi",
"markers": "python_version >= '3.8'", "markers": "python_version >= '3.8'",
"version": "==1.8.9" "version": "==1.8.11"
}, },
"frozenlist": { "frozenlist": {
"hashes": [ "hashes": [

View file

@ -45,14 +45,20 @@ class Emitter:
tasks = [] tasks = []
for emitter in self.emitters: for emitter in self.emitters:
_ret = emitter(event, data, **kwargs) try:
if _ret: _ret = emitter(event, data, **kwargs)
if isinstance(_ret, list): if _ret:
tasks.extend(_ret) if isinstance(_ret, list):
else: tasks.extend(_ret)
tasks.append(_ret) else:
tasks.append(_ret)
except Exception as e:
LOG.error(f"Emitter '{emitter}' failed with error '{e}'.")
if len(tasks) < 1:
return
try: try:
await asyncio.wait_for(asyncio.gather(*tasks), timeout=60) await asyncio.wait_for(asyncio.gather(*tasks), timeout=60)
except asyncio.TimeoutError: except asyncio.TimeoutError:
LOG.error(f"Timed out sending event {event}.") LOG.error(f"Timed out sending event '{event}'.")

View file

@ -109,7 +109,7 @@ class HttpAPI(common):
contentType = self.extToMime.get(os.path.splitext(file)[1], MIME.from_file(file)) contentType = self.extToMime.get(os.path.splitext(file)[1], MIME.from_file(file))
self.staticHolder[urlPath] = {'content': content, 'content_type': contentType} self.staticHolder[urlPath] = {'content': content, 'content_type': contentType}
LOG.debug(f'Preloading: [{urlPath}].') LOG.debug(f"Preloading '{urlPath}'.")
app.router.add_get(urlPath, self.staticFile) app.router.add_get(urlPath, self.staticFile)
if urlPath.endswith('/index.html') and urlPath != '/index.html': if urlPath.endswith('/index.html') and urlPath != '/index.html':
@ -148,7 +148,7 @@ class HttpAPI(common):
self.routes.get(self.config.url_prefix[:-1])(lambda _: web.HTTPFound(self.config.url_prefix)) self.routes.get(self.config.url_prefix[:-1])(lambda _: web.HTTPFound(self.config.url_prefix))
# add static files. # add static files.
self.routes.static(f'{self.config.url_prefix}download/', self.config.download_path) self.routes.static(f"{self.config.url_prefix}download/", self.config.download_path)
self.preloadStatic(app) self.preloadStatic(app)
try: try:
@ -304,7 +304,7 @@ class HttpAPI(common):
updated = True updated = True
setattr(item.info, k, v) setattr(item.info, k, v)
LOG.info(f'Updated [{k}] to [{v}] for [{item.info.id}]') LOG.info(f"Updated '{k}' to '{v}' for '{item.info.id}'")
status = 200 if updated else 304 status = 200 if updated else 304
if updated: if updated:
@ -478,7 +478,7 @@ class HttpAPI(common):
'X-Accel-Buffering': 'no', 'X-Accel-Buffering': 'no',
'Access-Control-Allow-Origin': '*', 'Access-Control-Allow-Origin': '*',
'Pragma': 'public', 'Pragma': 'public',
'Cache-Control': f'public, max-age={time.time() + 31536000}', 'Cache-Control': f"public, max-age={time.time() + 31536000}",
'Last-Modified': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(os.path.getmtime(file_path)).timetuple()), 'Last-Modified': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(os.path.getmtime(file_path)).timetuple()),
'Expires': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(time.time() + 31536000).timetuple()), 'Expires': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(time.time() + 31536000).timetuple()),
}) })
@ -504,7 +504,7 @@ class HttpAPI(common):
'X-Accel-Buffering': 'no', 'X-Accel-Buffering': 'no',
'Access-Control-Allow-Origin': '*', 'Access-Control-Allow-Origin': '*',
'Pragma': 'public', 'Pragma': 'public',
'Cache-Control': f'public, max-age={time.time() + 31536000}', 'Cache-Control': f"public, max-age={time.time() + 31536000}",
'Last-Modified': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(os.path.getmtime(file_path)).timetuple()), 'Last-Modified': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(os.path.getmtime(file_path)).timetuple()),
'Expires': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(time.time() + 31536000).timetuple()), 'Expires': time.strftime('%a, %d %b %Y %H:%M:%S GMT', datetime.fromtimestamp(time.time() + 31536000).timetuple()),
}) })

View file

@ -13,7 +13,7 @@ from .config import Config
from .DownloadQueue import DownloadQueue from .DownloadQueue import DownloadQueue
from .common import common from .common import common
from .encoder import Encoder from .encoder import Encoder
from .Utils import isDownloaded, load_file from .Utils import isDownloaded
from .Emitter import Emitter from .Emitter import Emitter
LOG = logging.getLogger('socket') LOG = logging.getLogger('socket')
@ -65,7 +65,7 @@ class HttpSocket(common):
return return
try: try:
LOG.info(f'Cli command from {sid=}: {data}') LOG.info(f"Cli command from client '{sid}'. '{data}'")
args = ['yt-dlp'] + shlex.split(data) args = ['yt-dlp'] + shlex.split(data)
_env = os.environ.copy() _env = os.environ.copy()
@ -93,7 +93,7 @@ class HttpSocket(common):
try: try:
os.close(slave_fd) os.close(slave_fd)
except Exception as e: except Exception as e:
LOG.error(f'Error closing PTY: {str(e)}') LOG.error(f"Error closing PTY. '{str(e)}'.")
async def read_pty(): async def read_pty():
loop = asyncio.get_running_loop() loop = asyncio.get_running_loop()
@ -126,7 +126,7 @@ class HttpSocket(common):
try: try:
os.close(master_fd) os.close(master_fd)
except Exception as e: except Exception as e:
LOG.error(f'Error closing PTY: {str(e)}') LOG.error(f"Error closing PTY. '{str(e)}'.")
# Start reading output from PTY # Start reading output from PTY
read_task = asyncio.create_task(read_pty()) read_task = asyncio.create_task(read_pty())
@ -139,7 +139,7 @@ class HttpSocket(common):
await self.emitter.emit('cli_close', {'exitcode': returncode}, to=sid) await self.emitter.emit('cli_close', {'exitcode': returncode}, to=sid)
except Exception as e: except Exception as e:
LOG.error(f'CLI Error for {sid=}: {str(e)}') LOG.error(f"CLI execute exception was thrown for client '{sid}'. {str(e)}")
await self.emitter.emit('cli_out1put', { await self.emitter.emit('cli_out1put', {
'type': 'stderr', 'type': 'stderr',
'line': str(e), 'line': str(e),
@ -152,7 +152,7 @@ class HttpSocket(common):
quality: str = data.get('quality') quality: str = data.get('quality')
if not url: if not url:
self.emitter.warning('No URL provided.') self.emitter.warning('No URL provided.', to=sid)
return return
format: str = data.get('format') format: str = data.get('format')
@ -178,7 +178,7 @@ class HttpSocket(common):
@ws_event @ws_event
async def item_cancel(self, sid: str, id: str): async def item_cancel(self, sid: str, id: str):
if not id: if not id:
await self.emitter.warning('Invalid request.') await self.emitter.warning('Invalid request.', to=sid)
return return
status: dict[str, str] = {} status: dict[str, str] = {}
@ -190,7 +190,7 @@ class HttpSocket(common):
@ws_event @ws_event
async def item_delete(self, sid: str, id: str): async def item_delete(self, sid: str, id: str):
if not id: if not id:
await self.emitter.warning('Invalid request.') await self.emitter.warning('Invalid request.', to=sid)
return return
status: dict[str, str] = {} status: dict[str, str] = {}
@ -230,30 +230,22 @@ class HttpSocket(common):
with open(manual_archive, 'a') as f: with open(manual_archive, 'a') as f:
f.write(f"{idDict['archive_id']} - at: {datetime.now().isoformat()}\n") f.write(f"{idDict['archive_id']} - at: {datetime.now().isoformat()}\n")
LOG.info(f'Archiving item: {data["url"]=}') LOG.info(f"Archiving url '{data['url']}' with id '{idDict['archive_id']}'.")
@ws_event @ws_event
async def connect(self, sid: str, _=None): async def connect(self, sid: str, _=None):
data: dict = { data: dict = {
**self.queue.get(), **self.queue.get(),
"config": self.config.frontend(), "config": self.config.frontend(),
"tasks": [], "tasks": self.config.tasks,
"presets": [], "presets": [],
"directories": [],
} }
if os.path.exists(os.path.join(self.config.config_path, 'tasks.json')):
try:
(tasks, status, error) = load_file(os.path.join(self.config.config_path, 'tasks.json'), list)
if status is False:
LOG.error(f"Could not load tasks file. Error message '{error}'.")
else:
data['tasks'] = tasks
except Exception as e:
pass
# get directory listing # get directory listing
dlDir: str = self.config.download_path dlDir: str = self.config.download_path
data['directories'] = [name for name in os.listdir(dlDir) if os.path.isdir(os.path.join(dlDir, name))] data['directories'] = [
name for name in os.listdir(dlDir) if os.path.isdir(os.path.join(dlDir, name))
]
await self.emitter.emit('initial_data', data, to=sid) await self.emitter.emit('initial_data', data, to=sid)

View file

@ -64,9 +64,13 @@ class Config:
auth_password: str = None auth_password: str = None
ytdlp_version: str = YTDLP_VERSION ytdlp_version: str = YTDLP_VERSION
tasks: list = []
_manual_vars: tuple = ('temp_path', 'config_path', 'download_path',) _manual_vars: tuple = ('temp_path', 'config_path', 'download_path',)
_immutable: tuple = ('version', '__instance', 'ytdl_options', 'new_version_available', 'ytdlp_version', 'started') _immutable: tuple = (
'version', '__instance', 'ytdl_options', 'tasks',
'new_version_available', 'ytdlp_version', 'started',
)
_int_vars: tuple = ('port', 'max_workers', 'socket_timeout', 'extract_info_timeout', 'debugpy_port',) _int_vars: tuple = ('port', 'max_workers', 'socket_timeout', 'extract_info_timeout', 'debugpy_port',)
_boolean_vars: tuple = ('keep_archive', 'ytdl_debug', 'debug', 'temp_keep', 'allow_manifestless',) _boolean_vars: tuple = ('keep_archive', 'ytdl_debug', 'debug', 'temp_keep', 'allow_manifestless',)
@ -122,7 +126,7 @@ class Config:
for key in re.findall(r'\{.*?\}', v): for key in re.findall(r'\{.*?\}', v):
localKey: str = key[1:-1] localKey: str = key[1:-1]
if localKey not in self.__dict__: if localKey not in self.__dict__:
logging.error(f'Config variable "{k}" had non existing config reference "{key}"') logging.error(f"Config variable '{k}' had non-existing config reference '{key}'.")
sys.exit(1) sys.exit(1)
v = v.replace(key, getattr(self, localKey)) v = v.replace(key, getattr(self, localKey))
@ -131,7 +135,7 @@ class Config:
if k in self._boolean_vars: if k in self._boolean_vars:
if str(v).lower() not in (True, False, 'true', 'false', 'on', 'off', '1', '0'): if str(v).lower() not in (True, False, 'true', 'false', 'on', 'off', '1', '0'):
raise ValueError(f'Config variable "{k}" is set to a non-boolean value "{v}".') raise ValueError(f"Config variable '{k}' is set to a non-boolean value '{v}'.")
setattr(self, k, str(v).lower() in (True, 'true', 'on', '1')) setattr(self, k, str(v).lower() in (True, 'true', 'on', '1'))
@ -143,7 +147,7 @@ class Config:
numeric_level = getattr(logging, self.log_level.upper(), None) numeric_level = getattr(logging, self.log_level.upper(), None)
if not isinstance(numeric_level, int): if not isinstance(numeric_level, int):
raise ValueError(f"Invalid log level: {self.log_level}") raise ValueError(f"Invalid log level '{self.log_level}' specified.")
coloredlogs.install( coloredlogs.install(
level=numeric_level, level=numeric_level,
@ -157,36 +161,48 @@ class Config:
try: try:
import debugpy import debugpy
debugpy.listen(("0.0.0.0", self.debugpy_port)) debugpy.listen(("0.0.0.0", self.debugpy_port))
LOG.info(f"starting debugpy server on [0.0.0.0:{self.debugpy_port}]") LOG.info(f"starting debugpy server on '0.0.0.0:{self.debugpy_port}'.")
except ImportError: except ImportError:
LOG.error("debugpy not found, please install it with 'pip install debugpy'") LOG.error("debugpy package not found, please install it with 'pip install debugpy'.")
except Exception as e: except Exception as e:
LOG.error(f"Error starting debugpy server at [0.0.0.0:{self.debugpy_port}]: {e}") LOG.error(f"Error starting debugpy server at '0.0.0.0:{self.debugpy_port}'. {e}")
optsFile: str = os.path.join(self.config_path, 'ytdlp.json') optsFile: str = os.path.join(self.config_path, 'ytdlp.json')
if os.path.exists(optsFile) and os.path.getsize(optsFile) > 0: if os.path.exists(optsFile) and os.path.getsize(optsFile) > 0:
LOG.info(f'Loading yt-dlp custom options from "{optsFile}"') LOG.info(f"Loading yt-dlp custom options from '{optsFile}'.")
(opts, status, error) = load_file(optsFile, dict) (opts, status, error) = load_file(optsFile, dict)
if not status: if not status:
LOG.error(f'Error loading "{optsFile}": {error}') LOG.error(f"Could not load yt-dlp custom options from '{optsFile}'. {error}")
sys.exit(1) sys.exit(1)
self.ytdl_options.update(opts) self.ytdl_options.update(opts)
else: else:
LOG.info(f'No custom yt-dlp options found in "{self.config_path}"') LOG.info(f"No yt-dlp custom options found at '{optsFile}'.")
tasksFile = os.path.join(self.config_path, 'tasks.json')
if os.path.exists(tasksFile) and os.path.getsize(tasksFile) > 0:
LOG.info(f"Loading tasks from '{tasksFile}'.")
try:
(tasks, status, error) = load_file(tasksFile, list)
if not status:
LOG.error(f"Could not load tasks file from '{tasksFile}'. '{error}'.")
sys.exit(1)
self.tasks.extend(tasks)
except Exception as e:
pass
self.ytdl_options['socket_timeout'] = self.socket_timeout self.ytdl_options['socket_timeout'] = self.socket_timeout
if self.keep_archive: if self.keep_archive:
LOG.info(f'keep archive: {self.keep_archive}') LOG.info(f"keep archive option is enabled.")
self.ytdl_options['download_archive'] = os.path.join( self.ytdl_options['download_archive'] = os.path.join(self.config_path, 'archive.log')
self.config_path, 'archive.log')
LOG.info(f'Keep temp: {self.temp_keep}') if self.temp_keep:
LOG.info(f'Keep temp files option is enabled.')
if self.auth_password and self.auth_username: if self.auth_password and self.auth_username:
LOG.warn(f"Basic authentication enabled with username: '{self.auth_username}'.") LOG.warning(f"Basic authentication enabled with username '{self.auth_username}'.")
self.started = time.time() self.started = time.time()

View file

@ -3,7 +3,6 @@
import asyncio import asyncio
import os import os
import random import random
import socketio
import logging import logging
import caribou import caribou
import sqlite3 import sqlite3
@ -11,7 +10,6 @@ import magic
from datetime import datetime from datetime import datetime
from aiohttp import web from aiohttp import web
from pathlib import Path from pathlib import Path
from library.Utils import load_file
from library.Emitter import Emitter from library.Emitter import Emitter
from library.config import Config from library.config import Config
from library.encoder import Encoder from library.encoder import Encoder
@ -21,17 +19,13 @@ from library.HttpSocket import HttpSocket
from library.HttpAPI import HttpAPI from library.HttpAPI import HttpAPI
from aiocron import crontab from aiocron import crontab
LOG = logging.getLogger('app') LOG = logging.getLogger("app")
MIME = magic.Magic(mime=True) MIME = magic.Magic(mime=True)
class main: class main:
config: Config = None config: Config = None
app: web.Application = None app: web.Application = None
sio: socketio.AsyncServer = None
routes: web.RouteTableDef = None
loop: asyncio.AbstractEventLoop = None
appLoader: str = None
http: HttpAPI = None http: HttpAPI = None
socket: HttpSocket = None socket: HttpSocket = None
@ -42,11 +36,11 @@ class main:
self.encoder = Encoder() self.encoder = Encoder()
self.checkDirectories() self.checkDirectories()
caribou.upgrade(self.config.db_file, os.path.join(self.rootPath, 'migrations')) caribou.upgrade(self.config.db_file, os.path.join(self.rootPath, "migrations"))
connection = sqlite3.connect(database=self.config.db_file, isolation_level=None) connection = sqlite3.connect(database=self.config.db_file, isolation_level=None)
connection.row_factory = sqlite3.Row connection.row_factory = sqlite3.Row
connection.execute('PRAGMA journal_mode=wal') connection.execute("PRAGMA journal_mode=wal")
emitter = Emitter() emitter = Emitter()
@ -56,96 +50,91 @@ class main:
self.http = HttpAPI(queue=queue, emitter=emitter, encoder=self.encoder) self.http = HttpAPI(queue=queue, emitter=emitter, encoder=self.encoder)
self.socket = HttpSocket(queue=queue, emitter=emitter, encoder=self.encoder) self.socket = HttpSocket(queue=queue, emitter=emitter, encoder=self.encoder)
WebhookFile = os.path.join(self.config.config_path, 'webhooks.json') WebhookFile = os.path.join(self.config.config_path, "webhooks.json")
if os.path.exists(WebhookFile): if os.path.exists(WebhookFile):
emitter.add_emitter(Webhooks(WebhookFile).emit) emitter.add_emitter(Webhooks(WebhookFile).emit)
def checkDirectories(self) -> None: def checkDirectories(self) -> None:
try: try:
LOG.debug(f'Checking download folder at [{self.config.download_path}]') LOG.debug(f"Checking download folder at '{self.config.download_path}'.")
if not os.path.exists(self.config.download_path): if not os.path.exists(self.config.download_path):
LOG.info(f'Creating download folder at [{self.config.download_path}]') LOG.info(f"Creating download folder at '{self.config.download_path}'.")
os.makedirs(self.config.download_path, exist_ok=True) os.makedirs(self.config.download_path, exist_ok=True)
except OSError as e: except OSError as e:
LOG.error(f'Could not create download folder at [{self.config.download_path}]') LOG.error(
f"Could not create download folder at '{self.config.download_path}'."
)
raise e raise e
try: try:
LOG.debug(f'Checking temp folder at [{self.config.temp_path}]') LOG.debug(f"Checking temp folder at '{self.config.temp_path}'.")
if not os.path.exists(self.config.temp_path): if not os.path.exists(self.config.temp_path):
LOG.info(f'Creating temp folder at [{self.config.temp_path}]') LOG.info(f"Creating temp folder at '{self.config.temp_path}'.")
os.makedirs(self.config.temp_path, exist_ok=True) os.makedirs(self.config.temp_path, exist_ok=True)
except OSError as e: except OSError as e:
LOG.error(f'Could not create temp folder at [{self.config.temp_path}]') LOG.error(f"Could not create temp folder at '{self.config.temp_path}'.")
raise e raise e
try: try:
LOG.debug(f'Checking config folder at [{self.config.config_path}]') LOG.debug(f"Checking config folder at '{self.config.config_path}'.")
if not os.path.exists(self.config.config_path): if not os.path.exists(self.config.config_path):
LOG.info(f'Creating config folder at [{self.config.config_path}]') LOG.info(f"Creating config folder at '{self.config.config_path}'.")
os.makedirs(self.config.config_path, exist_ok=True) os.makedirs(self.config.config_path, exist_ok=True)
except OSError as e: except OSError as e:
LOG.error(f'Could not create config folder at [{self.config.config_path}]') LOG.error(f"Could not create config folder at '{self.config.config_path}'.")
raise e raise e
try: try:
LOG.debug(f'Checking database file at [{self.config.db_file}]') LOG.debug(f"Checking database file at '{self.config.db_file}'.")
if not os.path.exists(self.config.db_file): if not os.path.exists(self.config.db_file):
LOG.info(f'Creating database file at [{self.config.db_file}]') LOG.info(f"Creating database file at '{self.config.db_file}'.")
with open(self.config.db_file, 'w') as _: with open(self.config.db_file, "w") as _:
pass pass
except OSError as e: except OSError as e:
LOG.error(f'Could not create database file at [{self.config.db_file}]') LOG.error(f"Could not create database file at '{self.config.db_file}'.")
raise e raise e
def load_tasks(self): async def cron_runner(self, task: dict):
tasks_file: str = os.path.join(self.config.config_path, 'tasks.json')
if not os.path.exists(tasks_file):
LOG.info(f'No tasks file found at {tasks_file}. Skipping Tasks.')
return
try: try:
(tasks, status, error) = load_file(tasks_file, list) taskName = task.get("name", task.get("url"))
if status is False: timeNow = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
raise Exception(error) LOG.info(f"Started 'Task: {taskName}' at '{timeNow}'.")
await self.socket.add(
url=task.get("url"),
quality=task.get("quality", "best"),
format=task.get("format", "any"),
folder=task.get("folder"),
ytdlp_cookies=task.get("ytdlp_cookies"),
ytdlp_config=task.get("ytdlp_config"),
output_template=task.get("output_template"),
)
timeNow = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
LOG.info(f"Completed 'Task: {taskName}' at '{timeNow}'.")
except Exception as e: except Exception as e:
LOG.error(f"Could not load tasks file '{tasks_file}'. Error message '{str(e)}'. Skipping Tasks.") timeNow = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
return LOG.error(
f"Failed 'Task: {taskName}' at '{timeNow}'. Error message '{str(e)}'."
)
for task in tasks: def load_tasks(self):
if not task.get('url'): for task in self.config.tasks:
LOG.warning(f'Invalid task {task}.') if not task.get("url"):
LOG.warning(f"Invalid task '{task}'. No URL found.")
continue continue
cron_timer: str = task.get('timer', f'{random.randint(1,59)} */1 * * *') cron_timer: str = task.get("timer", f"{random.randint(1,59)} */1 * * *")
async def cron_runner(task: dict):
taskName = task.get("name", task.get("url"))
timeNow = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
LOG.info(f'Started [Task: {taskName}] at [{timeNow}].')
await self.add(
url=task.get('url'),
quality=task.get('quality', 'best'),
format=task.get('format', 'any'),
folder=task.get('folder'),
ytdlp_cookies=task.get('ytdlp_cookies'),
ytdlp_config=task.get('ytdlp_config'),
output_template=task.get('output_template')
)
timeNow = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
LOG.info(f'Completed [Task: {taskName}] at [{timeNow}].')
crontab( crontab(
spec=cron_timer, spec=cron_timer,
func=cron_runner, func=self.cron_runner,
args=(task,), args=(task,),
start=True, start=True,
loop=self.loop loop=asyncio.get_event_loop(),
) )
LOG.info(f'Added [Task: {task.get("name", task.get("url"))}] executed every [{cron_timer}].') LOG.info(
f"Added 'Task: {task.get('name', task.get('url'))}' to be executed every '{cron_timer}'."
)
def start(self): def start(self):
self.socket.attach(self.app) self.socket.attach(self.app)
@ -153,7 +142,7 @@ class main:
self.load_tasks() self.load_tasks()
start: str = f'YTPTube v{self.config.version} - listening on http://{self.config.host}:{self.config.port}' start: str = f"YTPTube v{self.config.version} - started on http://{self.config.host}:{self.config.port}"
web.run_app( web.run_app(
self.app, self.app,
host=self.config.host, host=self.config.host,
@ -161,10 +150,10 @@ class main:
reuse_port=True, reuse_port=True,
loop=asyncio.get_event_loop(), loop=asyncio.get_event_loop(),
access_log=None, access_log=None,
print=lambda _: LOG.info(start) print=lambda _: LOG.info(start),
) )
if __name__ == '__main__': if __name__ == "__main__":
logging.basicConfig(level=logging.DEBUG) logging.basicConfig(level=logging.DEBUG)
main().start() main().start()