mirror of https://github.com/blackjack4494/yt-dlc
[websockets] Add `WebSocketFragmentFD` (#399)
Necessary for #392 Co-authored by: nao20010128nao, pukkandanpull/310/head
parent
ff0f78e1fe
commit
e36d50c5dd
@ -1,2 +1,3 @@
|
|||||||
mutagen
|
mutagen
|
||||||
pycryptodome
|
pycryptodome
|
||||||
|
websockets
|
||||||
|
@ -0,0 +1,59 @@
|
|||||||
|
import os
|
||||||
|
import signal
|
||||||
|
import asyncio
|
||||||
|
import threading
|
||||||
|
|
||||||
|
try:
|
||||||
|
import websockets
|
||||||
|
has_websockets = True
|
||||||
|
except ImportError:
|
||||||
|
has_websockets = False
|
||||||
|
|
||||||
|
from .common import FileDownloader
|
||||||
|
from .external import FFmpegFD
|
||||||
|
|
||||||
|
|
||||||
|
class FFmpegSinkFD(FileDownloader):
|
||||||
|
""" A sink to ffmpeg for downloading fragments in any form """
|
||||||
|
|
||||||
|
def real_download(self, filename, info_dict):
|
||||||
|
info_copy = info_dict.copy()
|
||||||
|
info_copy['url'] = '-'
|
||||||
|
|
||||||
|
async def call_conn(proc, stdin):
|
||||||
|
try:
|
||||||
|
await self.real_connection(stdin, info_dict)
|
||||||
|
except (BrokenPipeError, OSError):
|
||||||
|
pass
|
||||||
|
finally:
|
||||||
|
try:
|
||||||
|
stdin.flush()
|
||||||
|
stdin.close()
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
os.kill(os.getpid(), signal.SIGINT)
|
||||||
|
|
||||||
|
class FFmpegStdinFD(FFmpegFD):
|
||||||
|
@classmethod
|
||||||
|
def get_basename(cls):
|
||||||
|
return FFmpegFD.get_basename()
|
||||||
|
|
||||||
|
def on_process_started(self, proc, stdin):
|
||||||
|
thread = threading.Thread(target=asyncio.run, daemon=True, args=(call_conn(proc, stdin), ))
|
||||||
|
thread.start()
|
||||||
|
|
||||||
|
return FFmpegStdinFD(self.ydl, self.params or {}).download(filename, info_copy)
|
||||||
|
|
||||||
|
async def real_connection(self, sink, info_dict):
|
||||||
|
""" Override this in subclasses """
|
||||||
|
raise NotImplementedError('This method must be implemented by subclasses')
|
||||||
|
|
||||||
|
|
||||||
|
class WebSocketFragmentFD(FFmpegSinkFD):
|
||||||
|
async def real_connection(self, sink, info_dict):
|
||||||
|
async with websockets.connect(info_dict['url'], extra_headers=info_dict.get('http_headers', {})) as ws:
|
||||||
|
while True:
|
||||||
|
recv = await ws.recv()
|
||||||
|
if isinstance(recv, str):
|
||||||
|
recv = recv.encode('utf8')
|
||||||
|
sink.write(recv)
|
Loading…
Reference in New Issue