Merge pull request #83 from arabcoders/dev

upgrade webhooks to support both form and json request.
This commit is contained in:
Abdulmohsen 2024-03-29 00:18:05 +03:00 committed by GitHub
commit 33931a8660
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 113 additions and 54 deletions

View file

@ -17,6 +17,7 @@
"dotenv", "dotenv",
"finaldir", "finaldir",
"getpid", "getpid",
"httpx",
"libcurl", "libcurl",
"libx", "libx",
"mpegts", "mpegts",

View file

@ -13,6 +13,7 @@ aiocron = ">=1.8"
python-dotenv = ">=1.0.1" python-dotenv = ">=1.0.1"
python-magic = ">=0.4.27" python-magic = ">=0.4.27"
debugpy = ">=1.8.1" debugpy = ">=1.8.1"
httpx = "*"
[dev-packages] [dev-packages]

47
Pipfile.lock generated
View file

@ -1,7 +1,7 @@
{ {
"_meta": { "_meta": {
"hash": { "hash": {
"sha256": "6f7d266189a02316b0a7f6d1958351963de5dfa35b6e3860c1da5da875629843" "sha256": "b8c580658ddc19c696fab69a5317e66edf477f5560957ab549357b25b1b70212"
}, },
"pipfile-spec": 6, "pipfile-spec": 6,
"requires": { "requires": {
@ -115,6 +115,14 @@
"markers": "python_version >= '3.7'", "markers": "python_version >= '3.7'",
"version": "==1.3.1" "version": "==1.3.1"
}, },
"anyio": {
"hashes": [
"sha256:048e05d0f6caeed70d731f3db756d35dcc1f35747c8c403364a8332c630441b8",
"sha256:f75253795a87df48568485fd18cdd2a3fa5c4f7c5be8e5e36637733fce06fed6"
],
"markers": "python_version >= '3.8'",
"version": "==4.3.0"
},
"argparse": { "argparse": {
"hashes": [ "hashes": [
"sha256:62b089a55be1d8949cd2bc7e0df0bddb9e028faefc8c32038cc84862aefdd6e4", "sha256:62b089a55be1d8949cd2bc7e0df0bddb9e028faefc8c32038cc84862aefdd6e4",
@ -349,11 +357,11 @@
}, },
"croniter": { "croniter": {
"hashes": [ "hashes": [
"sha256:78bf110a2c7dbbfdd98b926318ae6c64a731a4c637c7befe3685755110834746", "sha256:28763ad39c404e159140874f08010cfd8a18f4c2a7cea1ce73e9506a4380cfc1",
"sha256:8bff16c9af4ef1fb6f05416973b8f7cb54997c02f2f8365251f9bf1dded91866" "sha256:84dc95b2eb6760144cc01eca65a6b9cc1619c93b2dc37d8a27f4319b3eb740de"
], ],
"markers": "python_version >= '2.6' and python_version not in '3.0, 3.1, 3.2, 3.3'", "markers": "python_version >= '2.6' and python_version not in '3.0, 3.1, 3.2, 3.3'",
"version": "==2.0.2" "version": "==2.0.3"
}, },
"debugpy": { "debugpy": {
"hashes": [ "hashes": [
@ -475,6 +483,23 @@
"markers": "python_version >= '3.7'", "markers": "python_version >= '3.7'",
"version": "==0.14.0" "version": "==0.14.0"
}, },
"httpcore": {
"hashes": [
"sha256:34a38e2f9291467ee3b44e89dd52615370e152954ba21721378a87b2960f7a61",
"sha256:421f18bac248b25d310f3cacd198d55b8e6125c107797b609ff9b7a6ba7991b5"
],
"markers": "python_version >= '3.8'",
"version": "==1.0.5"
},
"httpx": {
"hashes": [
"sha256:71d5465162c13681bff01ad59b2cc68dd838ea1f10e51574bac27103f00c91a5",
"sha256:a0cb88a46f32dc874e04ee956e4c2764aba2aa228f650b06788ba6bda2962ab5"
],
"index": "pypi",
"markers": "python_version >= '3.8'",
"version": "==0.27.0"
},
"humanfriendly": { "humanfriendly": {
"hashes": [ "hashes": [
"sha256:1697e1a8a8f550fd43c2865cd84542fc175a61dcb779b6fee18cf6b6ccba1477", "sha256:1697e1a8a8f550fd43c2865cd84542fc175a61dcb779b6fee18cf6b6ccba1477",
@ -669,12 +694,12 @@
}, },
"python-socketio": { "python-socketio": {
"hashes": [ "hashes": [
"sha256:bbcbd758ed8c183775cb2853ba001361e2fa018babf5cbe11a5b77e91c2ec2a2", "sha256:ae6a1de5c5209ca859dc574dccc8931c4be17ee003e74ce3b8d1306162bb4a37",
"sha256:f1a0228b8b1fbdbd93fbbedd821ebce0ef54b2b5bf6e98fcf710deaa7c574259" "sha256:b9f22a8ff762d7a6e123d16a43ddb1a27d50f07c3c88ea999334f2f89b0ad52b"
], ],
"index": "pypi", "index": "pypi",
"markers": "python_version >= '3.8'", "markers": "python_version >= '3.8'",
"version": "==5.11.1" "version": "==5.11.2"
}, },
"pytz": { "pytz": {
"hashes": [ "hashes": [
@ -707,6 +732,14 @@
"markers": "python_version >= '2.7' and python_version not in '3.0, 3.1, 3.2, 3.3'", "markers": "python_version >= '2.7' and python_version not in '3.0, 3.1, 3.2, 3.3'",
"version": "==1.16.0" "version": "==1.16.0"
}, },
"sniffio": {
"hashes": [
"sha256:2f6da418d1f1e0fddd844478f41680e794e6051915791a034ff65e5f100525a2",
"sha256:f4324edc670a0f49750a81b895f35c3adb843cca46f0530f79fc1babb23789dc"
],
"markers": "python_version >= '3.7'",
"version": "==1.3.1"
},
"tzlocal": { "tzlocal": {
"hashes": [ "hashes": [
"sha256:49816ef2fe65ea8ac19d19aa7a1ae0551c834303d5014c6d5a62e4cbda8047b8", "sha256:49816ef2fe65ea8ac19d19aa7a1ae0551c834303d5014c6d5a62e4cbda8047b8",

View file

@ -259,6 +259,8 @@ The `config/webhooks.json`, is a json file, which can be used to add webhook end
"request":{ "request":{
// (url: string) - REQUIRED- The webhook url // (url: string) - REQUIRED- The webhook url
"url": "https://mysecert.webhook.com/endpoint", "url": "https://mysecert.webhook.com/endpoint",
// (type: string) - OPTIONAL - The request type, it can be json or form.
"type": "json",
// (method: string) - OPTIONAL - The request method, it can be POST or PUT // (method: string) - OPTIONAL - The request method, it can be POST or PUT
"method": "POST", "method": "POST",
// (headers: dictionary) - OPTIONAL - Extra headers to include. // (headers: dictionary) - OPTIONAL - Extra headers to include.

View file

@ -4,7 +4,6 @@ from datetime import datetime, timezone
from email.utils import formatdate from email.utils import formatdate
import json import json
from sqlite3 import Connection from sqlite3 import Connection
from Utils import calcDownloadPath
from Config import Config from Config import Config
from Download import Download from Download import Download
from ItemDTO import ItemDTO from ItemDTO import ItemDTO

View file

@ -3,8 +3,8 @@ import json
import logging import logging
import os import os
from ItemDTO import ItemDTO from ItemDTO import ItemDTO
from aiohttp import client import httpx
from version import APP_VERSION
LOG = logging.getLogger('Webhooks') LOG = logging.getLogger('Webhooks')
@ -49,20 +49,43 @@ class Webhooks:
req: dict = target.get('request') req: dict = target.get('request')
try: try:
LOG.info(f"Sending {event=} {item.id=} to [{target.get('name')}]") LOG.info(f"Sending {event=} {item.id=} to [{target.get('name')}]")
async with client.ClientSession() as session: async with httpx.AsyncClient() as client:
headers = req.get('headers', {}) if 'headers' in req else {} request_type = req.get('type', 'json')
async with session.request(method=req.get('method', 'POST'), url=req.get('url'), json=item.__dict__, headers=headers) as response:
respData = {
'url': req.get('url'),
'status': response.status,
'text': await response.text()
}
msg = f"[{target.get('name')}] Response to [{event=} {item.id=}] [status: {response.status}]."
if respData.get('text'):
msg += f" [Body: {respData.get('text','??')}]"
LOG.info(msg)
return respData reqBody = {
'method': req.get('method', 'POST'),
'url': req.get('url'),
'headers': {
'User-Agent': f"YTPTube/{APP_VERSION}"
},
}
if req.get('headers', None):
reqBody['headers'].update(req.get('headers'))
match(request_type):
case 'json':
reqBody['json'] = item.__dict__
reqBody['headers']['Content-Type'] = 'application/json'
case _:
reqBody['data'] = item.__dict__
reqBody['headers']['Content-Type'] = 'application/x-www-form-urlencoded'
response = await client.request(**reqBody)
respData = {
'url': req.get('url'),
'status': response.status_code,
'text': response.text
}
msg = f"[{target.get('name')}] Response to [{event=} {item.id=}] [status: {response.status_code}]."
if respData.get('text'):
msg += f" [Body: {respData.get('text','??')}]"
LOG.info(msg)
return respData
except Exception as e: except Exception as e:
return { return {
'url': req.get('url'), 'url': req.get('url'),

View file

@ -11,7 +11,6 @@ from Utils import ObjectSerializer, Notifier
from aiohttp import web, client from aiohttp import web, client
from aiohttp.web import Request, Response from aiohttp.web import Request, Response
from Webhooks import Webhooks from Webhooks import Webhooks
from Download import Download
from player.M3u8 import M3u8 from player.M3u8 import M3u8
from player.Segments import Segments from player.Segments import Segments
import socketio import socketio
@ -416,11 +415,16 @@ class Main:
if not file: if not file:
raise web.HTTPBadRequest(reason='file is required.') raise web.HTTPBadRequest(reason='file is required.')
return web.Response( try:
text=await M3u8(url=f"{self.config.url_host}{self.config.url_prefix}").make_stream( text = await M3u8(url=f"{self.config.url_host}{self.config.url_prefix}").make_stream(
download_path=self.config.download_path, download_path=self.config.download_path,
file=file file=file
), )
except Exception as e:
return web.HTTPNotFound(reason=str(e))
return web.Response(
text=text,
headers={ headers={
'Content-Type': 'application/x-mpegURL', 'Content-Type': 'application/x-mpegURL',
'Cache-Control': 'no-cache', 'Cache-Control': 'no-cache',
@ -443,14 +447,14 @@ class Main:
raise web.HTTPBadRequest(reason='segment is required') raise web.HTTPBadRequest(reason='segment is required')
segmenter = Segments( segmenter = Segments(
segment_index=int(segment), index=int(segment),
segment_duration=float('{:.6f}'.format(float(sd if sd else M3u8.segment_duration))), duration=float('{:.6f}'.format(float(sd if sd else M3u8.duration))),
vconvert=True if vc == 1 else False, vconvert=True if vc == 1 else False,
aconvert=True if ac == 1 else False aconvert=True if ac == 1 else False
) )
return web.Response( return web.Response(
body=await segmenter.stream(download_path=self.config.download_path, file=file), body=await segmenter.stream(path=self.config.download_path, file=file),
headers={ headers={
'Content-Type': 'video/mpegts', 'Content-Type': 'video/mpegts',
'Cache-Control': 'no-cache', 'Cache-Control': 'no-cache',

View file

@ -9,12 +9,12 @@ class M3u8:
ok_vcodecs: tuple = ('h264', 'x264', 'avc',) ok_vcodecs: tuple = ('h264', 'x264', 'avc',)
ok_acodecs: tuple = ('aac', 'mp3',) ok_acodecs: tuple = ('aac', 'mp3',)
segment_duration: float = 10.000000
url: str = None url: str = None
duration: float = 6.000000
def __init__(self, url: str, segment_duration: float = 6.000000): def __init__(self, url: str, segment_duration: float = None):
self.url = url self.url = url
self.segment_duration = float(segment_duration) self.duration = float(segment_duration) if segment_duration is not None else self.duration
async def make_stream(self, download_path: str, file: str): async def make_stream(self, download_path: str, file: str):
realFile: str = calcDownloadPath(basePath=download_path, folder=file, createPath=False) realFile: str = calcDownloadPath(basePath=download_path, folder=file, createPath=False)
@ -35,12 +35,12 @@ class M3u8:
m3u8 = "#EXTM3U\n" m3u8 = "#EXTM3U\n"
m3u8 += "#EXT-X-VERSION:3\n" m3u8 += "#EXT-X-VERSION:3\n"
m3u8 += f"#EXT-X-TARGETDURATION:{int(self.segment_duration)}\n" m3u8 += f"#EXT-X-TARGETDURATION:{int(self.duration)}\n"
m3u8 += "#EXT-X-MEDIA-SEQUENCE:0\n" m3u8 += "#EXT-X-MEDIA-SEQUENCE:0\n"
m3u8 += "#EXT-X-PLAYLIST-TYPE:VOD\n" m3u8 += "#EXT-X-PLAYLIST-TYPE:VOD\n"
segmentSize: float = '{:.6f}'.format(self.segment_duration) segmentSize: float = '{:.6f}'.format(self.duration)
splits: int = math.ceil(duration / self.segment_duration) splits: int = math.ceil(duration / self.duration)
segmentParams: dict = {} segmentParams: dict = {}
@ -54,7 +54,7 @@ class M3u8:
for i in range(splits): for i in range(splits):
if (i + 1) == splits: if (i + 1) == splits:
segmentParams.update({'sd': '{:.6f}'.format(duration - (i * self.segment_duration))}) segmentParams.update({'sd': '{:.6f}'.format(duration - (i * self.duration))})
m3u8 += f"#EXTINF:{segmentParams['sd']}, nodesc\n" m3u8 += f"#EXTINF:{segmentParams['sd']}, nodesc\n"
else: else:
m3u8 += f"#EXTINF:{segmentSize}, nodesc\n" m3u8 += f"#EXTINF:{segmentSize}, nodesc\n"

View file

@ -2,7 +2,6 @@ import asyncio
import hashlib import hashlib
import logging import logging
import os import os
import subprocess
import tempfile import tempfile
from Utils import calcDownloadPath from Utils import calcDownloadPath
@ -10,19 +9,19 @@ LOG = logging.getLogger('segments')
class Segments: class Segments:
segment_duration: int duration: int
segment_index: int index: int
vconvert: bool vconvert: bool
aconvert: bool aconvert: bool
def __init__(self, segment_index: int, segment_duration: float, vconvert: bool, aconvert: bool): def __init__(self, index: int, duration: float, vconvert: bool, aconvert: bool):
self.segment_duration = float(segment_duration) self.index = int(index)
self.segment_index = int(segment_index) self.duration = float(duration)
self.vconvert = bool(vconvert) self.vconvert = bool(vconvert)
self.aconvert = bool(aconvert) self.aconvert = bool(aconvert)
async def stream(self, download_path: str, file: str) -> bytes: async def stream(self, path: str, file: str) -> bytes:
realFile: str = calcDownloadPath(basePath=download_path, folder=file, createPath=False) realFile: str = calcDownloadPath(basePath=path, folder=file, createPath=False)
if not os.path.exists(realFile): if not os.path.exists(realFile):
raise Exception(f"File {realFile} does not exist.") raise Exception(f"File {realFile} does not exist.")
@ -33,10 +32,10 @@ class Segments:
if not os.path.exists(tmpFile): if not os.path.exists(tmpFile):
os.symlink(realFile, tmpFile) os.symlink(realFile, tmpFile)
if self.segment_index == 0: if self.index == 0:
startTime: float = '{:.6f}'.format(0) startTime: float = '{:.6f}'.format(0)
else: else:
startTime: float = '{:.6f}'.format((self.segment_duration * self.segment_index)) startTime: float = '{:.6f}'.format((self.duration * self.index))
fargs = [] fargs = []
fargs.append('-xerror') fargs.append('-xerror')
@ -47,7 +46,7 @@ class Segments:
fargs.append('-ss') fargs.append('-ss')
fargs.append(str(startTime if startTime else '0.00000')) fargs.append(str(startTime if startTime else '0.00000'))
fargs.append('-t') fargs.append('-t')
fargs.append(str('{:.6f}'.format(self.segment_duration))) fargs.append(str('{:.6f}'.format(self.duration)))
fargs.append('-copyts') fargs.append('-copyts')
@ -99,7 +98,7 @@ class Segments:
fargs.append('mpegts') fargs.append('mpegts')
fargs.append('pipe:1') fargs.append('pipe:1')
LOG.debug(f"Streaming '{realFile}' segment '{self.segment_index}'. " + " ".join(fargs)) LOG.debug(f"Streaming '{realFile}' segment '{self.index}'. " + " ".join(fargs))
proc = await asyncio.subprocess.create_subprocess_exec( proc = await asyncio.subprocess.create_subprocess_exec(
'ffmpeg', *fargs, 'ffmpeg', *fargs,
@ -111,7 +110,7 @@ class Segments:
data, err = await proc.communicate() data, err = await proc.communicate()
if 0 != proc.returncode: if 0 != proc.returncode:
LOG.error(f'Failed to stream {realFile} segment {self.segment_index}. {err.decode("utf-8")}') LOG.error(f'Failed to stream {realFile} segment {self.index}. {err.decode("utf-8")}')
raise Exception(f'Failed to stream {realFile} segment {self.segment_index}.') raise Exception(f'Failed to stream {realFile} segment {self.index}.')
return data return data

View file

@ -4,11 +4,8 @@ Python wrapper for ffprobe command line tool. ffprobe must exist in the path.
import asyncio import asyncio
import functools import functools
import json import json
import logging
import operator import operator
import os import os
import pipes
import subprocess
class FFProbeError(Exception): class FFProbeError(Exception):