Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c9409a1570 |
11
.github/workflows/portable.yml
vendored
11
.github/workflows/portable.yml
vendored
@@ -40,6 +40,11 @@ jobs:
|
||||
- name: Unzip Python archive
|
||||
run: Expand-Archive -LiteralPath "python.zip" -DestinationPath "build\python"
|
||||
|
||||
- name: Enter build directory
|
||||
run: |
|
||||
Set-Location -Path "build"
|
||||
ls
|
||||
|
||||
- name: Rename ffmpeg directory
|
||||
working-directory: build
|
||||
run: Rename-Item -Path $($env:FFMPEG_ARCHIVE_NAME).Substring(0, $($env:FFMPEG_ARCHIVE_NAME).Length - 4) "ffmpeg"
|
||||
@@ -78,9 +83,11 @@ jobs:
|
||||
working-directory: build
|
||||
run: Copy-Item "${{ github.workspace }}\run.bat" -Destination ".\run.bat"
|
||||
|
||||
- name: Copy run.ps1
|
||||
- name: Exit build directory
|
||||
working-directory: build
|
||||
run: Copy-Item "${{ github.workspace }}\run.ps1" -Destination ".\run.ps1"
|
||||
run: |
|
||||
ls
|
||||
Set-Location -Path ".."
|
||||
|
||||
- name: Zip files
|
||||
run: |
|
||||
|
||||
12
CHANGELOG.md
12
CHANGELOG.md
@@ -1,17 +1,5 @@
|
||||
# 更新日志
|
||||
|
||||
## 1.6.2
|
||||
|
||||
- 忽略 Windows 注册表 JavaScript MIME 设置 (issue #12, 27)
|
||||
- 修复 HLS 录制出错 (issue #39, 41)
|
||||
- 修 bug (issue #47)
|
||||
- Windows 绿色版默认主机绑定 0.0.0.0 并加上 api key
|
||||
|
||||
## 1.6.1
|
||||
|
||||
- 修复 bug (issue #37, 38, 40)
|
||||
- 接收到错误的数据自动换线路 (issue #43)
|
||||
|
||||
## 1.6.0
|
||||
|
||||
- 更新 Pushplus 消息推送 url (issue #26)
|
||||
|
||||
12
README.md
12
README.md
@@ -129,7 +129,7 @@ sudo docker run \
|
||||
|
||||
例如:`blrec --host 0.0.0.0 --port 8000`
|
||||
|
||||
### 网络安全
|
||||
### 安全保障
|
||||
|
||||
指定 `SSL` 证书使用 **https** 协议并指定 `api key` 可防止被恶意访问和泄漏设置里的敏感信息
|
||||
|
||||
@@ -141,16 +141,6 @@ sudo docker run \
|
||||
|
||||
如果在不信任的环境下,请使用浏览器的隐式模式访问。
|
||||
|
||||
### 关于 api-key
|
||||
|
||||
api key 可以使用数字和字母,长度限制为最短 8 最长 80。
|
||||
|
||||
3 次尝试内 api key 正确客户端 ip 会自动加入白名单,3 次错误后则 ip 会被加入黑名单,黑名单后请求会被拒绝 (403)。
|
||||
|
||||
黑名单和白名单数以及同时尝试连接的 ip 数量限制各为 100,黑名单或白名单到达限制后不再接受除了白名单内的其它 ip 。
|
||||
|
||||
只有重启才会清空黑名单和白名单。
|
||||
|
||||
## 作为 ASGI 应用运行
|
||||
|
||||
uvicorn blrec.web:app
|
||||
|
||||
27
run.bat
27
run.bat
@@ -3,33 +3,20 @@ chcp 65001
|
||||
|
||||
set PATH=.\ffmpeg\bin;.\python;%PATH%
|
||||
|
||||
@REM 不使用代理
|
||||
REM 不使用代理
|
||||
set no_proxy=*
|
||||
|
||||
@REM 主机和端口绑定,可以按需修改。
|
||||
set host=0.0.0.0
|
||||
REM 默认本地主机和端口绑定
|
||||
set host=localhost
|
||||
set port=2233
|
||||
|
||||
@REM 关于 api key
|
||||
|
||||
@REM api key 可以使用数字和字母,长度限制为最短 8 最长 80。
|
||||
|
||||
@REM 3 次尝试内 api key 正确客户端 ip 会自动加入白名单,3 次错误后则 ip 会被加入黑名单,黑名单后请求会被拒绝 (403)。
|
||||
|
||||
@REM 黑名单和白名单数以及同时尝试连接的 ip 数量限制各为 100,黑名单或白名单到达限制后不再接受除了白名单内的其它 ip 。
|
||||
|
||||
@REM 只有重启才会清空黑名单和白名单。
|
||||
|
||||
@REM 浏览器第一次访问会弹对话框要求输入 api key。
|
||||
|
||||
@REM 输入的 api key 会被保存在浏览器的 local storage,下次使用同一浏览器不用再次输入。
|
||||
|
||||
@REM 请自行修改 api key,不要使用默认的 api key。
|
||||
set api_key=bili2233
|
||||
REM 服务器主机和端口绑定,去掉注释并按照自己的情况修改。
|
||||
REM set host=0.0.0.0
|
||||
REM set port=80
|
||||
|
||||
set DEFAULT_LOG_DIR=日志文件
|
||||
set DEFAULT_OUT_DIR=录播文件
|
||||
|
||||
python -m blrec -c settings.toml --open --host %host% --port %port% --api-key %api_key%
|
||||
python -m blrec -c settings.toml --open --host %host% --port %port%
|
||||
|
||||
pause
|
||||
|
||||
27
run.ps1
27
run.ps1
@@ -1,27 +0,0 @@
|
||||
chcp 65001
|
||||
|
||||
$env:PATH = ".\ffmpeg\bin;.\python;" + $env:PATH
|
||||
|
||||
# 不使用代理
|
||||
$env:no_proxy = "*"
|
||||
|
||||
# 主机和端口绑定,可以按需修改。
|
||||
$env:host = "0.0.0.0"
|
||||
$env:port = 2233
|
||||
|
||||
# 关于 api key
|
||||
# api key 可以使用数字和字母,长度限制为最短 8 最长 80。
|
||||
# 3 次尝试内 api key 正确客户端 ip 会自动加入白名单,3 次错误后则 ip 会被加入黑名单,黑名单后请求会被拒绝 (403)。
|
||||
# 黑名单和白名单数以及同时尝试连接的 ip 数量限制各为 100,黑名单或白名单到达限制后不再接受除了白名单内的其它 ip 。
|
||||
# 只有重启才会清空黑名单和白名单。
|
||||
# 浏览器第一次访问会弹对话框要求输入 api key。
|
||||
# 输入的 api key 会被保存在浏览器的 local storage,下次使用同一浏览器不用再次输入。
|
||||
# 请自行修改 api key,不要使用默认的 api key。
|
||||
$env:api_key = "bili2233"
|
||||
|
||||
$env:DEFAULT_LOG_DIR = "日志文件"
|
||||
$env:DEFAULT_OUT_DIR = "录播文件"
|
||||
|
||||
python -m blrec -c settings.toml --open --host $env:host --port $env:port --api-key $env:api_key
|
||||
|
||||
pause
|
||||
@@ -1,4 +1,4 @@
|
||||
|
||||
__prog__ = 'blrec'
|
||||
__version__ = '1.6.2'
|
||||
__version__ = '1.6.0'
|
||||
__github__ = 'https://github.com/acgnhiki/blrec'
|
||||
|
||||
@@ -15,9 +15,7 @@ from tenacity import (
|
||||
|
||||
from .api import AppApi, WebApi
|
||||
from .models import LiveStatus, RoomInfo, UserInfo
|
||||
from .typing import (
|
||||
ApiPlatform, StreamFormat, QualityNumber, StreamCodec, ResponseData
|
||||
)
|
||||
from .typing import StreamFormat, QualityNumber, StreamCodec, ResponseData
|
||||
from .exceptions import (
|
||||
LiveRoomHidden, LiveRoomLocked, LiveRoomEncrypted, NoStreamAvailable,
|
||||
NoStreamFormatAvailable, NoStreamCodecAvailable, NoStreamQualityAvailable,
|
||||
@@ -179,14 +177,12 @@ class Live:
|
||||
async def get_live_stream_urls(
|
||||
self,
|
||||
qn: QualityNumber = 10000,
|
||||
*,
|
||||
api_platform: ApiPlatform = 'android',
|
||||
stream_format: StreamFormat = 'flv',
|
||||
stream_codec: StreamCodec = 'avc',
|
||||
) -> List[str]:
|
||||
if api_platform == 'android':
|
||||
try:
|
||||
info = await self._appapi.get_room_play_info(self._room_id, qn)
|
||||
else:
|
||||
except Exception:
|
||||
info = await self._webapi.get_room_play_info(self._room_id, qn)
|
||||
|
||||
self._check_room_play_info(info)
|
||||
|
||||
@@ -1,10 +1,7 @@
|
||||
from typing import Any, Dict, Literal, Mapping
|
||||
|
||||
|
||||
ApiPlatform = Literal[
|
||||
'web',
|
||||
'android',
|
||||
]
|
||||
Danmaku = Mapping[str, Any]
|
||||
|
||||
QualityNumber = Literal[
|
||||
20000, # 4K
|
||||
@@ -29,5 +26,3 @@ StreamCodec = Literal[
|
||||
|
||||
JsonResponse = Dict[str, Any]
|
||||
ResponseData = Dict[str, Any]
|
||||
|
||||
Danmaku = Mapping[str, Any]
|
||||
|
||||
@@ -34,11 +34,10 @@ from .stream_analyzer import StreamProfile
|
||||
from .statistics import StatisticsCalculator
|
||||
from ..event.event_emitter import EventListener, EventEmitter
|
||||
from ..bili.live import Live
|
||||
from ..bili.typing import ApiPlatform, StreamFormat, QualityNumber
|
||||
from ..bili.typing import StreamFormat, QualityNumber
|
||||
from ..bili.helpers import get_quality_name
|
||||
from ..flv.data_analyser import MetaData
|
||||
from ..flv.stream_processor import StreamProcessor, BaseOutputFileManager
|
||||
from ..utils.io import wait_for
|
||||
from ..utils.mixins import AsyncCooperationMixin, AsyncStoppableMixin
|
||||
from ..path import escape_path
|
||||
from ..logging.room_id import aio_task_with_room_id
|
||||
@@ -103,8 +102,7 @@ class BaseStreamRecorder(
|
||||
self._quality_number = quality_number
|
||||
self._real_stream_format: Optional[StreamFormat] = None
|
||||
self._real_quality_number: Optional[QualityNumber] = None
|
||||
self._api_platform: ApiPlatform = 'android'
|
||||
self._use_alternative_stream: bool = False
|
||||
self._use_candidate_stream: bool = False
|
||||
self.buffer_size = buffer_size or io.DEFAULT_BUFFER_SIZE # bytes
|
||||
self.read_timeout = read_timeout or 3 # seconds
|
||||
self.disconnection_timeout = disconnection_timeout or 600 # seconds
|
||||
@@ -270,8 +268,7 @@ class BaseStreamRecorder(
|
||||
self._stream_url = ''
|
||||
self._stream_host = ''
|
||||
self._stream_profile = {}
|
||||
self._api_platform = 'android'
|
||||
self._use_alternative_stream = False
|
||||
self._use_candidate_stream = False
|
||||
self._connection_recovered.clear()
|
||||
self._thread = Thread(
|
||||
target=self._run, name=f'StreamRecorder::{self._live.room_id}'
|
||||
@@ -290,12 +287,6 @@ class BaseStreamRecorder(
|
||||
def _run(self) -> None:
|
||||
raise NotImplementedError()
|
||||
|
||||
def _rotate_api_platform(self) -> None:
|
||||
if self._api_platform == 'android':
|
||||
self._api_platform = 'web'
|
||||
else:
|
||||
self._api_platform = 'android'
|
||||
|
||||
@retry(
|
||||
reraise=True,
|
||||
retry=retry_if_exception_type((
|
||||
@@ -309,16 +300,11 @@ class BaseStreamRecorder(
|
||||
fmt = self._real_stream_format or self.stream_format
|
||||
logger.info(
|
||||
f'Getting the live stream url... qn: {qn}, format: {fmt}, '
|
||||
f'api platform: {self._api_platform}, '
|
||||
f'use alternative stream: {self._use_alternative_stream}'
|
||||
f'use_candidate_stream: {self._use_candidate_stream}'
|
||||
)
|
||||
try:
|
||||
urls = self._run_coroutine(
|
||||
self._live.get_live_stream_urls(
|
||||
qn,
|
||||
api_platform=self._api_platform,
|
||||
stream_format=fmt,
|
||||
)
|
||||
self._live.get_live_stream_urls(qn, fmt)
|
||||
)
|
||||
except NoStreamQualityAvailable:
|
||||
logger.info(
|
||||
@@ -346,10 +332,6 @@ class BaseStreamRecorder(
|
||||
except NoStreamCodecAvailable as e:
|
||||
logger.warning(repr(e))
|
||||
raise TryAgain
|
||||
except Exception as e:
|
||||
logger.warning(f'Failed to get live stream urls: {repr(e)}')
|
||||
self._rotate_api_platform()
|
||||
raise TryAgain
|
||||
else:
|
||||
logger.info(
|
||||
f'Adopted the stream format ({fmt}) and quality ({qn})'
|
||||
@@ -357,19 +339,17 @@ class BaseStreamRecorder(
|
||||
self._real_quality_number = qn
|
||||
self._real_stream_format = fmt
|
||||
|
||||
if not self._use_alternative_stream:
|
||||
if not self._use_candidate_stream:
|
||||
url = urls[0]
|
||||
else:
|
||||
try:
|
||||
url = urls[1]
|
||||
except IndexError:
|
||||
self._use_alternative_stream = False
|
||||
self._rotate_api_platform()
|
||||
logger.info(
|
||||
'No alternative stream url available, will using the primary'
|
||||
f' stream url from {self._api_platform} api instead.'
|
||||
'No candidate stream url available, '
|
||||
'will using the primary stream url instead.'
|
||||
)
|
||||
raise TryAgain
|
||||
url = urls[0]
|
||||
logger.info(f"Got live stream url: '{url}'")
|
||||
|
||||
return url
|
||||
@@ -461,14 +441,8 @@ B站直播录像
|
||||
|
||||
|
||||
class StreamProxy(io.RawIOBase):
|
||||
def __init__(
|
||||
self,
|
||||
stream: io.BufferedIOBase,
|
||||
*,
|
||||
read_timeout: Optional[float] = None,
|
||||
) -> None:
|
||||
def __init__(self, stream: io.BufferedIOBase) -> None:
|
||||
self._stream = stream
|
||||
self._read_timmeout = read_timeout
|
||||
self._offset = 0
|
||||
self._size_updates = Subject()
|
||||
|
||||
@@ -491,14 +465,7 @@ class StreamProxy(io.RawIOBase):
|
||||
return True
|
||||
|
||||
def read(self, size: int = -1) -> bytes:
|
||||
if self._stream.closed:
|
||||
raise EOFError
|
||||
if self._read_timmeout:
|
||||
data = wait_for(
|
||||
self._stream.read, args=(size, ), timeout=self._read_timmeout
|
||||
)
|
||||
else:
|
||||
data = self._stream.read(size)
|
||||
data = self._stream.read(size)
|
||||
self._offset += len(data)
|
||||
self._size_updates.on_next(len(data))
|
||||
return data
|
||||
@@ -507,14 +474,7 @@ class StreamProxy(io.RawIOBase):
|
||||
return self._offset
|
||||
|
||||
def readinto(self, b: Any) -> int:
|
||||
if self._stream.closed:
|
||||
raise EOFError
|
||||
if self._read_timmeout:
|
||||
n = wait_for(
|
||||
self._stream.readinto, args=(b, ), timeout=self._read_timmeout
|
||||
)
|
||||
else:
|
||||
n = self._stream.readinto(b)
|
||||
n = self._stream.readinto(b)
|
||||
self._offset += n
|
||||
self._size_updates.on_next(n)
|
||||
return n
|
||||
|
||||
@@ -1,3 +0,0 @@
|
||||
|
||||
class FailedToFetchSegments(Exception):
|
||||
pass
|
||||
@@ -84,7 +84,6 @@ class FLVStreamRecorder(
|
||||
analyse_data=True,
|
||||
dedup_join=True,
|
||||
save_extra_metadata=True,
|
||||
backup_timestamp=True,
|
||||
)
|
||||
|
||||
def update_size(size: int) -> None:
|
||||
@@ -106,7 +105,6 @@ class FLVStreamRecorder(
|
||||
except Exception as e:
|
||||
self._handle_exception(e)
|
||||
finally:
|
||||
self._stopped = True
|
||||
if self._stream_processor is not None:
|
||||
self._stream_processor.finalize()
|
||||
self._stream_processor = None
|
||||
@@ -152,6 +150,7 @@ class FLVStreamRecorder(
|
||||
except Exception as e:
|
||||
logger.exception(e)
|
||||
self._handle_exception(e)
|
||||
self._stopped = True
|
||||
|
||||
def _streaming_loop(self) -> None:
|
||||
url = self._get_live_stream_url()
|
||||
@@ -177,12 +176,12 @@ class FLVStreamRecorder(
|
||||
except requests.exceptions.ConnectionError as e:
|
||||
logger.warning(repr(e))
|
||||
self._wait_for_connection_error()
|
||||
except (FlvDataError, FlvStreamCorruptedError) as e:
|
||||
except FlvDataError as e:
|
||||
logger.warning(repr(e))
|
||||
self._use_candidate_stream = not self._use_candidate_stream
|
||||
url = self._get_live_stream_url()
|
||||
except FlvStreamCorruptedError as e:
|
||||
logger.warning(repr(e))
|
||||
if not self._use_alternative_stream:
|
||||
self._use_alternative_stream = True
|
||||
else:
|
||||
self._rotate_api_platform()
|
||||
url = self._get_live_stream_url()
|
||||
|
||||
def _streaming(self, url: str) -> None:
|
||||
|
||||
@@ -3,12 +3,12 @@ import time
|
||||
import errno
|
||||
import logging
|
||||
from queue import Queue, Empty
|
||||
from threading import Thread, Event, Lock, Condition
|
||||
from threading import Thread, Event, Lock
|
||||
from datetime import datetime
|
||||
from contextlib import suppress
|
||||
from urllib.parse import urlparse
|
||||
|
||||
from typing import List, Set, Optional
|
||||
from typing import Set, Optional
|
||||
|
||||
import urllib3
|
||||
import requests
|
||||
@@ -30,12 +30,10 @@ from tenacity import (
|
||||
from .stream_remuxer import StreamRemuxer
|
||||
from .stream_analyzer import ffprobe, StreamProfile
|
||||
from .base_stream_recorder import BaseStreamRecorder, StreamProxy
|
||||
from .exceptions import FailedToFetchSegments
|
||||
from .retry import wait_exponential_for_same_exceptions, before_sleep_log
|
||||
from ..bili.live import Live
|
||||
from ..bili.typing import StreamFormat, QualityNumber
|
||||
from ..flv.stream_processor import StreamProcessor
|
||||
from ..flv.exceptions import FlvDataError, FlvStreamCorruptedError
|
||||
from ..utils.mixins import (
|
||||
AsyncCooperationMixin, AsyncStoppableMixin, SupportDebugMixin
|
||||
)
|
||||
@@ -83,9 +81,6 @@ class HLSStreamRecorder(
|
||||
duration_limit=duration_limit,
|
||||
)
|
||||
self._init_for_debug(self._live.room_id)
|
||||
self._init_section_data: Optional[bytes] = None
|
||||
self._ready_to_fetch_segments = Condition()
|
||||
self._failed_to_fetch_segments = Event()
|
||||
self._stream_analysed_lock = Lock()
|
||||
self._last_segment_uris: Set[str] = set()
|
||||
|
||||
@@ -100,54 +95,52 @@ class HLSStreamRecorder(
|
||||
)
|
||||
self._playlist_debug_file = open(path, 'wt', encoding='utf-8')
|
||||
|
||||
self._session = requests.Session()
|
||||
self._session.headers.update(self._live.headers)
|
||||
with StreamRemuxer(self._live.room_id) as self._stream_remuxer:
|
||||
with requests.Session() as self._session:
|
||||
self._session.headers.update(self._live.headers)
|
||||
|
||||
self._stream_remuxer = StreamRemuxer(self._live.room_id)
|
||||
self._segment_queue: Queue[Segment] = Queue(maxsize=1000)
|
||||
self._segment_data_queue: Queue[bytes] = Queue(maxsize=100)
|
||||
self._stream_host_available = Event()
|
||||
self._segment_queue: Queue[Segment] = Queue(maxsize=1000)
|
||||
self._segment_data_queue: Queue[bytes] = Queue(maxsize=100)
|
||||
self._stream_host_available = Event()
|
||||
|
||||
self._segment_fetcher_thread = Thread(
|
||||
target=self._run_segment_fetcher,
|
||||
name=f'SegmentFetcher::{self._live.room_id}',
|
||||
daemon=True,
|
||||
)
|
||||
self._segment_fetcher_thread.start()
|
||||
self._segment_fetcher_thread = Thread(
|
||||
target=self._run_segment_fetcher,
|
||||
name=f'SegmentFetcher::{self._live.room_id}',
|
||||
daemon=True,
|
||||
)
|
||||
self._segment_fetcher_thread.start()
|
||||
|
||||
self._segment_data_feeder_thread = Thread(
|
||||
target=self._run_segment_data_feeder,
|
||||
name=f'SegmentDataFeeder::{self._live.room_id}',
|
||||
daemon=True,
|
||||
)
|
||||
self._segment_data_feeder_thread.start()
|
||||
self._segment_data_feeder_thread = Thread(
|
||||
target=self._run_segment_data_feeder,
|
||||
name=f'SegmentDataFeeder::{self._live.room_id}',
|
||||
daemon=True,
|
||||
)
|
||||
self._segment_data_feeder_thread.start()
|
||||
|
||||
self._stream_processor_thread = Thread(
|
||||
target=self._run_stream_processor,
|
||||
name=f'StreamProcessor::{self._live.room_id}',
|
||||
daemon=True,
|
||||
)
|
||||
self._stream_processor_thread.start()
|
||||
self._stream_processor_thread = Thread(
|
||||
target=self._run_stream_processor,
|
||||
name=f'StreamProcessor::{self._live.room_id}',
|
||||
daemon=True,
|
||||
)
|
||||
self._stream_processor_thread.start()
|
||||
|
||||
try:
|
||||
self._main_loop()
|
||||
finally:
|
||||
if self._stream_processor is not None:
|
||||
self._stream_processor.cancel()
|
||||
self._stream_processor_thread.join(timeout=10)
|
||||
self._segment_fetcher_thread.join(timeout=10)
|
||||
self._segment_data_feeder_thread.join(timeout=10)
|
||||
self._stream_remuxer.stop()
|
||||
self._stream_remuxer.raise_for_exception()
|
||||
self._last_segment_uris.clear()
|
||||
del self._segment_queue
|
||||
del self._segment_data_queue
|
||||
try:
|
||||
self._main_loop()
|
||||
finally:
|
||||
if self._stream_processor is not None:
|
||||
self._stream_processor.cancel()
|
||||
self._segment_fetcher_thread.join(timeout=10)
|
||||
self._segment_data_feeder_thread.join(timeout=10)
|
||||
self._last_segment_uris.clear()
|
||||
del self._segment_queue
|
||||
del self._segment_data_queue
|
||||
except TryAgain:
|
||||
pass
|
||||
except Exception as e:
|
||||
self._handle_exception(e)
|
||||
finally:
|
||||
self._stopped = True
|
||||
with suppress(Exception):
|
||||
self._stream_processor_thread.join(timeout=10)
|
||||
with suppress(Exception):
|
||||
self._playlist_debug_file.close()
|
||||
self._emit_event('stream_recording_stopped')
|
||||
@@ -208,8 +201,6 @@ class HLSStreamRecorder(
|
||||
except requests.exceptions.ConnectionError as e:
|
||||
logger.warning(repr(e))
|
||||
self._wait_for_connection_error()
|
||||
except FailedToFetchSegments:
|
||||
url = self._get_live_stream_url()
|
||||
except RetryError as e:
|
||||
logger.warning(repr(e))
|
||||
|
||||
@@ -221,14 +212,6 @@ class HLSStreamRecorder(
|
||||
self._stream_analysed = False
|
||||
|
||||
while not self._stopped:
|
||||
if self._failed_to_fetch_segments.is_set():
|
||||
with self._segment_queue.mutex:
|
||||
self._segment_queue.queue.clear()
|
||||
with self._ready_to_fetch_segments:
|
||||
self._ready_to_fetch_segments.notify_all()
|
||||
self._failed_to_fetch_segments.clear()
|
||||
raise FailedToFetchSegments()
|
||||
|
||||
content = self._fetch_playlist(url)
|
||||
playlist = m3u8.loads(content, uri=url)
|
||||
|
||||
@@ -269,10 +252,8 @@ class HLSStreamRecorder(
|
||||
|
||||
if playlist.is_endlist:
|
||||
logger.debug('playlist ended')
|
||||
self._run_coroutine(self._live.update_room_info())
|
||||
if not self._live.is_living():
|
||||
self._stopped = True
|
||||
break
|
||||
self._stopped = True
|
||||
break
|
||||
|
||||
time.sleep(1)
|
||||
|
||||
@@ -291,8 +272,6 @@ class HLSStreamRecorder(
|
||||
assert self._stream_remuxer is not None
|
||||
init_section = None
|
||||
self._init_section_data = None
|
||||
num_of_continuously_failed = 0
|
||||
self._failed_to_fetch_segments.clear()
|
||||
|
||||
while not self._stopped:
|
||||
try:
|
||||
@@ -328,13 +307,6 @@ class HLSStreamRecorder(
|
||||
except requests.exceptions.HTTPError as e:
|
||||
logger.warning(f'Failed to fetch segment: {repr(e)}')
|
||||
if e.response.status_code in (403, 404, 599):
|
||||
num_of_continuously_failed += 1
|
||||
if num_of_continuously_failed >= 3:
|
||||
self._failed_to_fetch_segments.set()
|
||||
with self._ready_to_fetch_segments:
|
||||
self._ready_to_fetch_segments.wait()
|
||||
num_of_continuously_failed = 0
|
||||
self._failed_to_fetch_segments.clear()
|
||||
break
|
||||
except requests.exceptions.ConnectionError as e:
|
||||
logger.warning(repr(e))
|
||||
@@ -343,7 +315,6 @@ class HLSStreamRecorder(
|
||||
logger.warning(repr(e))
|
||||
break
|
||||
else:
|
||||
num_of_continuously_failed = 0
|
||||
break
|
||||
|
||||
def _run_segment_data_feeder(self) -> None:
|
||||
@@ -358,8 +329,6 @@ class HLSStreamRecorder(
|
||||
|
||||
def _segment_data_feeder(self) -> None:
|
||||
assert self._stream_remuxer is not None
|
||||
MAX_SEGMENT_DATA_CACHE = 3
|
||||
segment_data_cache: List[bytes] = []
|
||||
bytes_io = io.BytesIO()
|
||||
segment_count = 0
|
||||
|
||||
@@ -390,41 +359,10 @@ class HLSStreamRecorder(
|
||||
bytes_io = io.BytesIO()
|
||||
segment_count = 0
|
||||
self._stream_analysed = True
|
||||
|
||||
try:
|
||||
if self._stream_remuxer.stopped:
|
||||
self._stream_remuxer.start()
|
||||
while True:
|
||||
ready = self._stream_remuxer.wait(timeout=1)
|
||||
if self._stopped:
|
||||
return
|
||||
if ready:
|
||||
break
|
||||
if segment_data_cache:
|
||||
if self._init_section_data:
|
||||
self._stream_remuxer.input.write(
|
||||
self._init_section_data
|
||||
)
|
||||
for cached_data in segment_data_cache:
|
||||
if cached_data == self._init_section_data:
|
||||
continue
|
||||
self._stream_remuxer.input.write(cached_data)
|
||||
|
||||
self._stream_remuxer.input.write(data)
|
||||
except BrokenPipeError as e:
|
||||
if not self._stopped:
|
||||
logger.warning(repr(e))
|
||||
else:
|
||||
logger.debug(repr(e))
|
||||
except ValueError as e:
|
||||
if not self._stopped:
|
||||
logger.warning(repr(e))
|
||||
else:
|
||||
logger.debug(repr(e))
|
||||
|
||||
segment_data_cache.append(data)
|
||||
if len(segment_data_cache) > MAX_SEGMENT_DATA_CACHE:
|
||||
segment_data_cache.pop(0)
|
||||
except BrokenPipeError:
|
||||
return
|
||||
|
||||
def _run_stream_processor(self) -> None:
|
||||
logger.debug('Stream processor thread started')
|
||||
@@ -449,49 +387,19 @@ class HLSStreamRecorder(
|
||||
analyse_data=True,
|
||||
dedup_join=True,
|
||||
save_extra_metadata=True,
|
||||
backup_timestamp=True,
|
||||
)
|
||||
self._stream_processor.size_updates.subscribe(update_size)
|
||||
|
||||
try:
|
||||
while not self._stopped:
|
||||
while True:
|
||||
ready = self._stream_remuxer.wait(timeout=1)
|
||||
if self._stopped:
|
||||
return
|
||||
if ready:
|
||||
break
|
||||
|
||||
self._stream_host_available.wait()
|
||||
self._stream_processor.set_metadata(self._make_metadata())
|
||||
|
||||
try:
|
||||
self._stream_processor.process_stream(
|
||||
StreamProxy(
|
||||
self._stream_remuxer.output,
|
||||
read_timeout=10,
|
||||
) # type: ignore
|
||||
)
|
||||
except BrokenPipeError as e:
|
||||
logger.debug(repr(e))
|
||||
except TimeoutError as e:
|
||||
logger.debug(repr(e))
|
||||
self._stream_remuxer.stop()
|
||||
except FlvDataError as e:
|
||||
logger.warning(repr(e))
|
||||
self._stream_remuxer.stop()
|
||||
except FlvStreamCorruptedError as e:
|
||||
logger.warning(repr(e))
|
||||
self._stream_remuxer.stop()
|
||||
except ValueError as e:
|
||||
logger.warning(repr(e))
|
||||
self._stream_remuxer.stop()
|
||||
self._stream_host_available.wait()
|
||||
self._stream_processor.set_metadata(self._make_metadata())
|
||||
self._stream_processor.process_stream(
|
||||
StreamProxy(self._stream_remuxer.output), # type: ignore
|
||||
)
|
||||
except Exception as e:
|
||||
if not self._stopped:
|
||||
logger.exception(e)
|
||||
self._handle_exception(e)
|
||||
else:
|
||||
logger.debug(repr(e))
|
||||
finally:
|
||||
self._stream_processor.finalize()
|
||||
self._progress_bar = None
|
||||
|
||||
@@ -1,16 +1,14 @@
|
||||
import re
|
||||
import os
|
||||
import io
|
||||
import errno
|
||||
import shlex
|
||||
import logging
|
||||
from threading import Thread, Condition
|
||||
from threading import Thread, Event
|
||||
from subprocess import Popen, PIPE, CalledProcessError
|
||||
from typing import Optional, cast
|
||||
from typing import List, Optional, cast
|
||||
|
||||
|
||||
from ..utils.mixins import StoppableMixin, SupportDebugMixin
|
||||
from ..utils.io import wait_for
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -19,21 +17,15 @@ logger = logging.getLogger(__name__)
|
||||
__all__ = 'StreamRemuxer',
|
||||
|
||||
|
||||
class FFmpegError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
class StreamRemuxer(StoppableMixin, SupportDebugMixin):
|
||||
_ERROR_PATTERN = re.compile(
|
||||
r'\b(error|failed|missing|invalid|corrupt)\b', re.IGNORECASE
|
||||
)
|
||||
|
||||
def __init__(self, room_id: int, bufsize: int = 1024 * 1024) -> None:
|
||||
super().__init__()
|
||||
self._room_id = room_id
|
||||
self._bufsize = bufsize
|
||||
self._exception: Optional[Exception] = None
|
||||
self._ready = Condition()
|
||||
self._subprocess_setup = Event()
|
||||
self._MAX_ERROR_MESSAGES = 10
|
||||
self._error_messages: List[str] = []
|
||||
self._env = None
|
||||
|
||||
self._init_for_debug(room_id)
|
||||
@@ -58,22 +50,15 @@ class StreamRemuxer(StoppableMixin, SupportDebugMixin):
|
||||
|
||||
def __enter__(self): # type: ignore
|
||||
self.start()
|
||||
self.wait()
|
||||
self.wait_for_subprocess()
|
||||
return self
|
||||
|
||||
def __exit__(self, exc_type, value, traceback): # type: ignore
|
||||
self.stop()
|
||||
self.raise_for_exception()
|
||||
|
||||
def wait(self, timeout: Optional[float] = None) -> bool:
|
||||
with self._ready:
|
||||
return self._ready.wait(timeout=timeout)
|
||||
|
||||
def restart(self) -> None:
|
||||
logger.debug('Restarting stream remuxer...')
|
||||
self.stop()
|
||||
self.start()
|
||||
logger.debug('Restarted stream remuxer')
|
||||
def wait_for_subprocess(self) -> None:
|
||||
self._subprocess_setup.wait()
|
||||
|
||||
def raise_for_exception(self) -> None:
|
||||
if not self.exception:
|
||||
@@ -100,17 +85,12 @@ class StreamRemuxer(StoppableMixin, SupportDebugMixin):
|
||||
def _run(self) -> None:
|
||||
logger.debug('Started stream remuxer')
|
||||
self._exception = None
|
||||
self._error_messages.clear()
|
||||
self._subprocess_setup.clear()
|
||||
try:
|
||||
self._run_subprocess()
|
||||
except BrokenPipeError as e:
|
||||
logger.debug(repr(e))
|
||||
except FFmpegError as e:
|
||||
if not self._stopped:
|
||||
logger.warning(repr(e))
|
||||
else:
|
||||
logger.debug(repr(e))
|
||||
except TimeoutError as e:
|
||||
logger.debug(repr(e))
|
||||
except BrokenPipeError:
|
||||
pass
|
||||
except Exception as e:
|
||||
# OSError: [Errno 22] Invalid argument
|
||||
# https://stackoverflow.com/questions/23688492/oserror-errno-22-invalid-argument-in-subprocess
|
||||
@@ -120,43 +100,43 @@ class StreamRemuxer(StoppableMixin, SupportDebugMixin):
|
||||
self._exception = e
|
||||
logger.exception(e)
|
||||
finally:
|
||||
self._stopped = True
|
||||
logger.debug('Stopped stream remuxer')
|
||||
|
||||
def _run_subprocess(self) -> None:
|
||||
cmd = 'ffmpeg -xerror -i pipe:0 -c copy -copyts -f flv pipe:1'
|
||||
cmd = 'ffmpeg -i pipe:0 -c copy -f flv pipe:1'
|
||||
args = shlex.split(cmd)
|
||||
|
||||
with Popen(
|
||||
args, stdin=PIPE, stdout=PIPE, stderr=PIPE,
|
||||
bufsize=self._bufsize, env=self._env,
|
||||
) as self._subprocess:
|
||||
with self._ready:
|
||||
self._ready.notify_all()
|
||||
|
||||
self._subprocess_setup.set()
|
||||
assert self._subprocess.stderr is not None
|
||||
with io.TextIOWrapper(
|
||||
self._subprocess.stderr,
|
||||
encoding='utf-8',
|
||||
errors='backslashreplace'
|
||||
) as stderr:
|
||||
while not self._stopped:
|
||||
line = wait_for(stderr.readline, timeout=10)
|
||||
if not line:
|
||||
if self._subprocess.poll() is not None:
|
||||
break
|
||||
else:
|
||||
continue
|
||||
if self._debug:
|
||||
logger.debug('ffmpeg: %s', line)
|
||||
self._check_error(line)
|
||||
|
||||
while not self._stopped:
|
||||
data = self._subprocess.stderr.readline()
|
||||
if not data:
|
||||
if self._subprocess.poll() is not None:
|
||||
break
|
||||
else:
|
||||
continue
|
||||
line = data.decode('utf-8', errors='backslashreplace')
|
||||
if self._debug:
|
||||
logger.debug('ffmpeg: %s', line)
|
||||
self._check_error(line)
|
||||
|
||||
if not self._stopped and self._subprocess.returncode not in (0, 255):
|
||||
# 255: Exiting normally, received signal 2.
|
||||
raise CalledProcessError(self._subprocess.returncode, cmd=cmd)
|
||||
raise CalledProcessError(
|
||||
self._subprocess.returncode,
|
||||
cmd=cmd,
|
||||
output='\n'.join(self._error_messages),
|
||||
)
|
||||
|
||||
def _check_error(self, line: str) -> None:
|
||||
match = self._ERROR_PATTERN.search(line)
|
||||
if not match:
|
||||
if 'error' not in line.lower() and 'failed' not in line.lower():
|
||||
return
|
||||
raise FFmpegError(line)
|
||||
logger.warning(f'ffmpeg error: {line}')
|
||||
self._error_messages.append(line)
|
||||
if len(self._error_messages) > self._MAX_ERROR_MESSAGES:
|
||||
self._error_messages.remove(self._error_messages[0])
|
||||
|
||||
File diff suppressed because one or more lines are too long
1
src/blrec/data/webapp/66.d8b06f1fef317761.js
Normal file
1
src/blrec/data/webapp/66.d8b06f1fef317761.js
Normal file
File diff suppressed because one or more lines are too long
1
src/blrec/data/webapp/694.92a3e0c2fc842a42.js
Normal file
1
src/blrec/data/webapp/694.92a3e0c2fc842a42.js
Normal file
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@@ -10,6 +10,6 @@
|
||||
<body>
|
||||
<app-root></app-root>
|
||||
<noscript>Please enable JavaScript to continue using this application.</noscript>
|
||||
<script src="runtime.459a2fcc68ef4fa7.js" type="module"></script><script src="polyfills.4b08448aee19bb22.js" type="module"></script><script src="main.8a8c73fae6ff9291.js" type="module"></script>
|
||||
<script src="runtime.23c91f03d62c595a.js" type="module"></script><script src="polyfills.4b08448aee19bb22.js" type="module"></script><script src="main.8a8c73fae6ff9291.js" type="module"></script>
|
||||
|
||||
</body></html>
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"configVersion": 1,
|
||||
"timestamp": 1650082907617,
|
||||
"timestamp": 1649386979751,
|
||||
"index": "/index.html",
|
||||
"assetGroups": [
|
||||
{
|
||||
@@ -14,15 +14,15 @@
|
||||
"/103.5b5d2a6e5a8a7479.js",
|
||||
"/146.92e3b29c4c754544.js",
|
||||
"/45.c90c3cea2bf1a66e.js",
|
||||
"/66.8d16a032cbce41ed.js",
|
||||
"/694.dbcbd5e1953ea55d.js",
|
||||
"/869.ac675e78fa0ea7cf.js",
|
||||
"/66.d8b06f1fef317761.js",
|
||||
"/694.92a3e0c2fc842a42.js",
|
||||
"/869.95d68b28a4188d76.js",
|
||||
"/common.858f777e9296e6f2.js",
|
||||
"/index.html",
|
||||
"/main.8a8c73fae6ff9291.js",
|
||||
"/manifest.webmanifest",
|
||||
"/polyfills.4b08448aee19bb22.js",
|
||||
"/runtime.459a2fcc68ef4fa7.js",
|
||||
"/runtime.23c91f03d62c595a.js",
|
||||
"/styles.1f581691b230dc4d.css"
|
||||
],
|
||||
"patterns": []
|
||||
@@ -1637,9 +1637,9 @@
|
||||
"/103.5b5d2a6e5a8a7479.js": "cc0240f217015b6d4ddcc14f31fcc42e1c1c282a",
|
||||
"/146.92e3b29c4c754544.js": "3824de681dd1f982ea69a065cdf54d7a1e781f4d",
|
||||
"/45.c90c3cea2bf1a66e.js": "e5bfb8cf3803593e6b8ea14c90b3d3cb6a066764",
|
||||
"/66.8d16a032cbce41ed.js": "a473089370c2fe27f96a778acf1e709dc5770b31",
|
||||
"/694.dbcbd5e1953ea55d.js": "41973de76799b085188903f46cc6974f20395741",
|
||||
"/869.ac675e78fa0ea7cf.js": "f45052016cb5201d5784b3f261e719d96bd1b153",
|
||||
"/66.d8b06f1fef317761.js": "43676d9dc886b5624dadecc50f17d4972b183d2d",
|
||||
"/694.92a3e0c2fc842a42.js": "f8f093029b9996b3db0c4e738bf9f8573fba8392",
|
||||
"/869.95d68b28a4188d76.js": "cd1add38c89b1df3c0783b74c931b51839f1c530",
|
||||
"/assets/animal/panda.js": "fec2868bb3053dd2da45f96bbcb86d5116ed72b1",
|
||||
"/assets/animal/panda.svg": "bebd302cdc601e0ead3a6d2710acf8753f3d83b1",
|
||||
"/assets/fill/.gitkeep": "da39a3ee5e6b4b0d3255bfef95601890afd80709",
|
||||
@@ -3234,11 +3234,11 @@
|
||||
"/assets/twotone/warning.js": "fb2d7ea232f3a99bf8f080dbc94c65699232ac01",
|
||||
"/assets/twotone/warning.svg": "8c7a2d3e765a2e7dd58ac674870c6655cecb0068",
|
||||
"/common.858f777e9296e6f2.js": "b68ca68e1e214a2537d96935c23410126cc564dd",
|
||||
"/index.html": "be9d178866ccd58dfdc2a4c1b375ab030302163b",
|
||||
"/index.html": "114f00ffcd1f7fa5aaaa7f2fcf3109f26c77c715",
|
||||
"/main.8a8c73fae6ff9291.js": "41a5a5a8fb5cda4cfa0e28532812594816257122",
|
||||
"/manifest.webmanifest": "62c1cb8c5ad2af551a956b97013ab55ce77dd586",
|
||||
"/polyfills.4b08448aee19bb22.js": "8e73f2d42cc13ca353cea5c886d930bd6da08d0d",
|
||||
"/runtime.459a2fcc68ef4fa7.js": "a68b948e588e75f8cb2fb5315ac41623e7c6ed1e",
|
||||
"/runtime.23c91f03d62c595a.js": "0819f1120ed1e37c2ad069ef949147450c951069",
|
||||
"/styles.1f581691b230dc4d.css": "6f5befbbad57c2b2e80aae855139744b8010d150"
|
||||
},
|
||||
"navigationUrls": [
|
||||
|
||||
@@ -1 +1 @@
|
||||
(()=>{"use strict";var e,v={},m={};function r(e){var i=m[e];if(void 0!==i)return i.exports;var t=m[e]={exports:{}};return v[e].call(t.exports,t,t.exports,r),t.exports}r.m=v,e=[],r.O=(i,t,o,f)=>{if(!t){var a=1/0;for(n=0;n<e.length;n++){for(var[t,o,f]=e[n],c=!0,l=0;l<t.length;l++)(!1&f||a>=f)&&Object.keys(r.O).every(b=>r.O[b](t[l]))?t.splice(l--,1):(c=!1,f<a&&(a=f));if(c){e.splice(n--,1);var d=o();void 0!==d&&(i=d)}}return i}f=f||0;for(var n=e.length;n>0&&e[n-1][2]>f;n--)e[n]=e[n-1];e[n]=[t,o,f]},r.n=e=>{var i=e&&e.__esModule?()=>e.default:()=>e;return r.d(i,{a:i}),i},r.d=(e,i)=>{for(var t in i)r.o(i,t)&&!r.o(e,t)&&Object.defineProperty(e,t,{enumerable:!0,get:i[t]})},r.f={},r.e=e=>Promise.all(Object.keys(r.f).reduce((i,t)=>(r.f[t](e,i),i),[])),r.u=e=>(592===e?"common":e)+"."+{45:"c90c3cea2bf1a66e",66:"8d16a032cbce41ed",103:"5b5d2a6e5a8a7479",146:"92e3b29c4c754544",592:"858f777e9296e6f2",694:"dbcbd5e1953ea55d",869:"ac675e78fa0ea7cf"}[e]+".js",r.miniCssF=e=>{},r.o=(e,i)=>Object.prototype.hasOwnProperty.call(e,i),(()=>{var e={},i="blrec:";r.l=(t,o,f,n)=>{if(e[t])e[t].push(o);else{var a,c;if(void 0!==f)for(var l=document.getElementsByTagName("script"),d=0;d<l.length;d++){var u=l[d];if(u.getAttribute("src")==t||u.getAttribute("data-webpack")==i+f){a=u;break}}a||(c=!0,(a=document.createElement("script")).type="module",a.charset="utf-8",a.timeout=120,r.nc&&a.setAttribute("nonce",r.nc),a.setAttribute("data-webpack",i+f),a.src=r.tu(t)),e[t]=[o];var s=(g,b)=>{a.onerror=a.onload=null,clearTimeout(p);var _=e[t];if(delete e[t],a.parentNode&&a.parentNode.removeChild(a),_&&_.forEach(h=>h(b)),g)return g(b)},p=setTimeout(s.bind(null,void 0,{type:"timeout",target:a}),12e4);a.onerror=s.bind(null,a.onerror),a.onload=s.bind(null,a.onload),c&&document.head.appendChild(a)}}})(),r.r=e=>{"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},(()=>{var e;r.tu=i=>(void 0===e&&(e={createScriptURL:t=>t},"undefined"!=typeof trustedTypes&&trustedTypes.createPolicy&&(e=trustedTypes.createPolicy("angular#bundler",e))),e.createScriptURL(i))})(),r.p="",(()=>{var e={666:0};r.f.j=(o,f)=>{var n=r.o(e,o)?e[o]:void 0;if(0!==n)if(n)f.push(n[2]);else if(666!=o){var a=new Promise((u,s)=>n=e[o]=[u,s]);f.push(n[2]=a);var c=r.p+r.u(o),l=new Error;r.l(c,u=>{if(r.o(e,o)&&(0!==(n=e[o])&&(e[o]=void 0),n)){var s=u&&("load"===u.type?"missing":u.type),p=u&&u.target&&u.target.src;l.message="Loading chunk "+o+" failed.\n("+s+": "+p+")",l.name="ChunkLoadError",l.type=s,l.request=p,n[1](l)}},"chunk-"+o,o)}else e[o]=0},r.O.j=o=>0===e[o];var i=(o,f)=>{var l,d,[n,a,c]=f,u=0;if(n.some(p=>0!==e[p])){for(l in a)r.o(a,l)&&(r.m[l]=a[l]);if(c)var s=c(r)}for(o&&o(f);u<n.length;u++)r.o(e,d=n[u])&&e[d]&&e[d][0](),e[n[u]]=0;return r.O(s)},t=self.webpackChunkblrec=self.webpackChunkblrec||[];t.forEach(i.bind(null,0)),t.push=i.bind(null,t.push.bind(t))})()})();
|
||||
(()=>{"use strict";var e,v={},m={};function r(e){var i=m[e];if(void 0!==i)return i.exports;var t=m[e]={exports:{}};return v[e].call(t.exports,t,t.exports,r),t.exports}r.m=v,e=[],r.O=(i,t,f,o)=>{if(!t){var a=1/0;for(n=0;n<e.length;n++){for(var[t,f,o]=e[n],c=!0,l=0;l<t.length;l++)(!1&o||a>=o)&&Object.keys(r.O).every(b=>r.O[b](t[l]))?t.splice(l--,1):(c=!1,o<a&&(a=o));if(c){e.splice(n--,1);var d=f();void 0!==d&&(i=d)}}return i}o=o||0;for(var n=e.length;n>0&&e[n-1][2]>o;n--)e[n]=e[n-1];e[n]=[t,f,o]},r.n=e=>{var i=e&&e.__esModule?()=>e.default:()=>e;return r.d(i,{a:i}),i},r.d=(e,i)=>{for(var t in i)r.o(i,t)&&!r.o(e,t)&&Object.defineProperty(e,t,{enumerable:!0,get:i[t]})},r.f={},r.e=e=>Promise.all(Object.keys(r.f).reduce((i,t)=>(r.f[t](e,i),i),[])),r.u=e=>(592===e?"common":e)+"."+{45:"c90c3cea2bf1a66e",66:"d8b06f1fef317761",103:"5b5d2a6e5a8a7479",146:"92e3b29c4c754544",592:"858f777e9296e6f2",694:"92a3e0c2fc842a42",869:"95d68b28a4188d76"}[e]+".js",r.miniCssF=e=>{},r.o=(e,i)=>Object.prototype.hasOwnProperty.call(e,i),(()=>{var e={},i="blrec:";r.l=(t,f,o,n)=>{if(e[t])e[t].push(f);else{var a,c;if(void 0!==o)for(var l=document.getElementsByTagName("script"),d=0;d<l.length;d++){var u=l[d];if(u.getAttribute("src")==t||u.getAttribute("data-webpack")==i+o){a=u;break}}a||(c=!0,(a=document.createElement("script")).type="module",a.charset="utf-8",a.timeout=120,r.nc&&a.setAttribute("nonce",r.nc),a.setAttribute("data-webpack",i+o),a.src=r.tu(t)),e[t]=[f];var s=(g,b)=>{a.onerror=a.onload=null,clearTimeout(p);var _=e[t];if(delete e[t],a.parentNode&&a.parentNode.removeChild(a),_&&_.forEach(h=>h(b)),g)return g(b)},p=setTimeout(s.bind(null,void 0,{type:"timeout",target:a}),12e4);a.onerror=s.bind(null,a.onerror),a.onload=s.bind(null,a.onload),c&&document.head.appendChild(a)}}})(),r.r=e=>{"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},(()=>{var e;r.tu=i=>(void 0===e&&(e={createScriptURL:t=>t},"undefined"!=typeof trustedTypes&&trustedTypes.createPolicy&&(e=trustedTypes.createPolicy("angular#bundler",e))),e.createScriptURL(i))})(),r.p="",(()=>{var e={666:0};r.f.j=(f,o)=>{var n=r.o(e,f)?e[f]:void 0;if(0!==n)if(n)o.push(n[2]);else if(666!=f){var a=new Promise((u,s)=>n=e[f]=[u,s]);o.push(n[2]=a);var c=r.p+r.u(f),l=new Error;r.l(c,u=>{if(r.o(e,f)&&(0!==(n=e[f])&&(e[f]=void 0),n)){var s=u&&("load"===u.type?"missing":u.type),p=u&&u.target&&u.target.src;l.message="Loading chunk "+f+" failed.\n("+s+": "+p+")",l.name="ChunkLoadError",l.type=s,l.request=p,n[1](l)}},"chunk-"+f,f)}else e[f]=0},r.O.j=f=>0===e[f];var i=(f,o)=>{var l,d,[n,a,c]=o,u=0;if(n.some(p=>0!==e[p])){for(l in a)r.o(a,l)&&(r.m[l]=a[l]);if(c)var s=c(r)}for(f&&f(o);u<n.length;u++)r.o(e,d=n[u])&&e[d]&&e[d][0](),e[n[u]]=0;return r.O(s)},t=self.webpackChunkblrec=self.webpackChunkblrec||[];t.forEach(i.bind(null,0)),t.push=i.bind(null,t.push.bind(t))})()})();
|
||||
@@ -37,15 +37,8 @@ __all__ = 'FlvParser', 'FlvDumper'
|
||||
|
||||
|
||||
class FlvParser:
|
||||
def __init__(
|
||||
self,
|
||||
stream: RandomIO,
|
||||
backup_timestamp: bool = False,
|
||||
restore_timestamp: bool = False,
|
||||
) -> None:
|
||||
def __init__(self, stream: RandomIO) -> None:
|
||||
self._stream = stream
|
||||
self._backup_timestamp = backup_timestamp
|
||||
self._restore_timestamp = restore_timestamp
|
||||
self._reader = StructReader(stream)
|
||||
|
||||
def parse_header(self) -> FlvHeader:
|
||||
@@ -115,42 +108,19 @@ class FlvParser:
|
||||
|
||||
def parse_flv_tag_header(self, data: bytes) -> FlvTagHeader:
|
||||
reader = StructReader(BytesIO(data))
|
||||
|
||||
flag = reader.read_ui8()
|
||||
filtered = bool(flag & 0b0010_0000)
|
||||
if filtered:
|
||||
raise FlvDataError('Unsupported Filtered FLV Tag', data)
|
||||
|
||||
tag_type = TagType(flag & 0b0001_1111)
|
||||
data_size = reader.read_ui24()
|
||||
timestamp = reader.read_ui24()
|
||||
timestamp_extended = reader.read_ui8()
|
||||
timestamp = timestamp_extended << 24 | timestamp
|
||||
stream_id = reader.read_ui24()
|
||||
|
||||
if self._backup_timestamp:
|
||||
return FlvTagHeader(
|
||||
filtered=filtered,
|
||||
tag_type=tag_type,
|
||||
data_size=data_size,
|
||||
timestamp=timestamp_extended << 24 | timestamp,
|
||||
stream_id=timestamp,
|
||||
)
|
||||
elif self._restore_timestamp:
|
||||
return FlvTagHeader(
|
||||
filtered=filtered,
|
||||
tag_type=tag_type,
|
||||
data_size=data_size,
|
||||
timestamp=stream_id,
|
||||
stream_id=stream_id,
|
||||
)
|
||||
else:
|
||||
return FlvTagHeader(
|
||||
filtered=filtered,
|
||||
tag_type=tag_type,
|
||||
data_size=data_size,
|
||||
timestamp=timestamp_extended << 24 | timestamp,
|
||||
stream_id=stream_id,
|
||||
)
|
||||
return FlvTagHeader(
|
||||
filtered, tag_type, data_size, timestamp, stream_id
|
||||
)
|
||||
|
||||
def parse_audio_tag_header(self, data: bytes) -> AudioTagHeader:
|
||||
reader = StructReader(BytesIO(data))
|
||||
|
||||
@@ -16,21 +16,9 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class FlvReader:
|
||||
def __init__(
|
||||
self,
|
||||
stream: RandomIO,
|
||||
*,
|
||||
backup_timestamp: bool = False,
|
||||
restore_timestamp: bool = False,
|
||||
) -> None:
|
||||
def __init__(self, stream: RandomIO) -> None:
|
||||
self._stream = stream
|
||||
self._backup_timestamp = backup_timestamp
|
||||
self._restore_timestamp = restore_timestamp
|
||||
self._parser = FlvParser(
|
||||
stream,
|
||||
backup_timestamp=backup_timestamp,
|
||||
restore_timestamp=restore_timestamp,
|
||||
)
|
||||
self._parser = FlvParser(stream)
|
||||
|
||||
def read_header(self) -> FlvHeader:
|
||||
header = self._parser.parse_header()
|
||||
|
||||
@@ -116,8 +116,8 @@ class FlvTagHeader:
|
||||
filtered: bool
|
||||
tag_type: TagType
|
||||
data_size: int = attr.ib(validator=[non_negative_integer_validator])
|
||||
timestamp: int
|
||||
stream_id: int
|
||||
timestamp: int = attr.ib(validator=[non_negative_integer_validator])
|
||||
stream_id: int = attr.ib(validator=[non_negative_integer_validator])
|
||||
|
||||
|
||||
@attr.s(auto_attribs=True, slots=True, frozen=True)
|
||||
|
||||
@@ -31,9 +31,9 @@ from .exceptions import (
|
||||
CutStream,
|
||||
)
|
||||
from .common import (
|
||||
is_metadata_tag, parse_metadata, is_audio_tag, is_video_tag,
|
||||
is_video_sequence_header, is_audio_sequence_header,
|
||||
enrich_metadata, update_metadata, is_data_tag, read_tags_in_duration,
|
||||
is_audio_tag, is_video_tag, is_metadata_tag, parse_metadata,
|
||||
is_audio_data_tag, is_video_data_tag, enrich_metadata, update_metadata,
|
||||
is_data_tag, read_tags_in_duration,
|
||||
)
|
||||
from ..path import extra_metadata_path
|
||||
if TYPE_CHECKING:
|
||||
@@ -63,8 +63,6 @@ class StreamProcessor:
|
||||
analyse_data: bool = False,
|
||||
dedup_join: bool = False,
|
||||
save_extra_metadata: bool = False,
|
||||
backup_timestamp: bool = False,
|
||||
restore_timestamp: bool = False,
|
||||
) -> None:
|
||||
self._file_manager = file_manager
|
||||
self._parameters_checker = ParametersChecker()
|
||||
@@ -81,8 +79,6 @@ class StreamProcessor:
|
||||
self._analyse_data = analyse_data
|
||||
self._dedup_join = dedup_join
|
||||
self._save_x_metadata = save_extra_metadata
|
||||
self._backup_timestamp = backup_timestamp
|
||||
self._restore_timestamp = restore_timestamp
|
||||
|
||||
self._cancelled: bool = False
|
||||
self._finalized: bool = False
|
||||
@@ -227,12 +223,7 @@ class StreamProcessor:
|
||||
def _process_stream(self, stream: RandomIO) -> None:
|
||||
logger.debug(f'Processing the {self._stream_count}th stream...')
|
||||
|
||||
self._in_reader = FlvReaderWithTimestampFix(
|
||||
stream,
|
||||
backup_timestamp=self._backup_timestamp,
|
||||
restore_timestamp=self._restore_timestamp,
|
||||
)
|
||||
|
||||
self._in_reader = FlvReaderWithTimestampFix(stream)
|
||||
flv_header = self._read_header()
|
||||
self._has_audio = flv_header.has_audio()
|
||||
|
||||
@@ -561,31 +552,18 @@ class StreamProcessor:
|
||||
return header
|
||||
|
||||
def _ensure_ts_correct(self, tag: FlvTag) -> None:
|
||||
if not tag.timestamp + self._delta < 0:
|
||||
if not is_audio_data_tag(tag) or not is_video_data_tag(tag):
|
||||
return
|
||||
logger.warning(
|
||||
f'Incorrect timestamp: {tag.timestamp + self._delta}\n'
|
||||
f'last output tag: {self._last_tags[0]}\n'
|
||||
f'current tag: {tag}'
|
||||
)
|
||||
if tag.is_audio_tag() or tag.is_video_tag():
|
||||
self._delta = (
|
||||
self._last_tags[0].timestamp +
|
||||
self._in_reader.calc_interval(tag) - tag.timestamp
|
||||
)
|
||||
logger.debug(f'Updated delta: {self._delta}')
|
||||
elif tag.is_script_tag():
|
||||
self._delta = (
|
||||
self._last_tags[0].timestamp - tag.timestamp
|
||||
)
|
||||
logger.debug(f'Updated delta: {self._delta}')
|
||||
else:
|
||||
pass
|
||||
if tag.timestamp + self._delta < 0:
|
||||
self._delta = -tag.timestamp
|
||||
logger.warning('Incorrect timestamp: {}, new delta: {}'.format(
|
||||
tag, self._delta
|
||||
))
|
||||
|
||||
def _correct_ts(self, tag: FlvTag, delta: int) -> FlvTag:
|
||||
if delta == 0 and tag.timestamp >= 0:
|
||||
return tag
|
||||
return tag.evolve(timestamp=tag.timestamp + delta)
|
||||
return tag.evolve(timestamp=max(0, tag.timestamp + delta))
|
||||
|
||||
def _calc_delta_duplicated(self, last_duplicated_tag: FlvTag) -> int:
|
||||
return self._last_tags[0].timestamp - last_duplicated_tag.timestamp
|
||||
@@ -786,17 +764,8 @@ class RobustFlvReader(FlvReader):
|
||||
|
||||
|
||||
class FlvReaderWithTimestampFix(RobustFlvReader):
|
||||
def __init__(
|
||||
self,
|
||||
stream: RandomIO,
|
||||
backup_timestamp: bool = False,
|
||||
restore_timestamp: bool = False,
|
||||
) -> None:
|
||||
super().__init__(
|
||||
stream,
|
||||
backup_timestamp=backup_timestamp,
|
||||
restore_timestamp=restore_timestamp,
|
||||
)
|
||||
def __init__(self, stream: RandomIO) -> None:
|
||||
super().__init__(stream)
|
||||
self._last_tag: Optional[FlvTag] = None
|
||||
self._last_video_tag: Optional[VideoTag] = None
|
||||
self._last_audio_tag: Optional[AudioTag] = None
|
||||
@@ -873,17 +842,11 @@ class FlvReaderWithTimestampFix(RobustFlvReader):
|
||||
if is_video_tag(tag):
|
||||
if self._last_video_tag is None:
|
||||
return False
|
||||
if is_video_sequence_header(self._last_video_tag):
|
||||
return tag.timestamp < self._last_video_tag.timestamp
|
||||
else:
|
||||
return tag.timestamp <= self._last_video_tag.timestamp
|
||||
return tag.timestamp <= self._last_video_tag.timestamp
|
||||
elif is_audio_tag(tag):
|
||||
if self._last_audio_tag is None:
|
||||
return False
|
||||
if is_audio_sequence_header(self._last_audio_tag):
|
||||
return tag.timestamp < self._last_audio_tag.timestamp
|
||||
else:
|
||||
return tag.timestamp <= self._last_audio_tag.timestamp
|
||||
return tag.timestamp <= self._last_audio_tag.timestamp
|
||||
else:
|
||||
return False
|
||||
|
||||
|
||||
@@ -1,29 +0,0 @@
|
||||
import atexit
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from concurrent.futures import TimeoutError as _TimeoutError
|
||||
from typing import Callable, Any, Iterable, Mapping, TypeVar
|
||||
|
||||
|
||||
_T = TypeVar('_T')
|
||||
|
||||
|
||||
_executor = None
|
||||
|
||||
|
||||
def wait_for(
|
||||
func: Callable[..., _T],
|
||||
*,
|
||||
args: Iterable[Any] = [],
|
||||
kwargs: Mapping[str, Any] = {},
|
||||
timeout: float
|
||||
) -> _T:
|
||||
global _executor
|
||||
if _executor is None:
|
||||
_executor = ThreadPoolExecutor(thread_name_prefix='wait_for')
|
||||
atexit.register(_executor.shutdown)
|
||||
|
||||
future = _executor.submit(func, *args, **kwargs)
|
||||
try:
|
||||
return future.result(timeout=timeout)
|
||||
except _TimeoutError:
|
||||
raise TimeoutError(timeout, func, args, kwargs) from None
|
||||
@@ -1,7 +1,6 @@
|
||||
import os
|
||||
from abc import ABC, abstractmethod
|
||||
import asyncio
|
||||
import threading
|
||||
from typing import Awaitable, TypeVar, final
|
||||
|
||||
|
||||
@@ -9,28 +8,24 @@ class SwitchableMixin(ABC):
|
||||
def __init__(self) -> None:
|
||||
super().__init__()
|
||||
self._enabled = False
|
||||
self._enabled_lock = threading.Lock()
|
||||
|
||||
@property
|
||||
def enabled(self) -> bool:
|
||||
with self._enabled_lock:
|
||||
return self._enabled
|
||||
return self._enabled
|
||||
|
||||
@final
|
||||
def enable(self) -> None:
|
||||
with self._enabled_lock:
|
||||
if self._enabled:
|
||||
return
|
||||
self._enabled = True
|
||||
self._do_enable()
|
||||
if self._enabled:
|
||||
return
|
||||
self._enabled = True
|
||||
self._do_enable()
|
||||
|
||||
@final
|
||||
def disable(self) -> None:
|
||||
with self._enabled_lock:
|
||||
if not self._enabled:
|
||||
return
|
||||
self._enabled = False
|
||||
self._do_disable()
|
||||
if not self._enabled:
|
||||
return
|
||||
self._enabled = False
|
||||
self._do_disable()
|
||||
|
||||
@abstractmethod
|
||||
def _do_enable(self) -> None:
|
||||
@@ -45,28 +40,24 @@ class StoppableMixin(ABC):
|
||||
def __init__(self) -> None:
|
||||
super().__init__()
|
||||
self._stopped = True
|
||||
self._stopped_lock = threading.Lock()
|
||||
|
||||
@property
|
||||
def stopped(self) -> bool:
|
||||
with self._stopped_lock:
|
||||
return self._stopped
|
||||
return self._stopped
|
||||
|
||||
@final
|
||||
def start(self) -> None:
|
||||
with self._stopped_lock:
|
||||
if not self._stopped:
|
||||
return
|
||||
self._stopped = False
|
||||
self._do_start()
|
||||
if not self._stopped:
|
||||
return
|
||||
self._stopped = False
|
||||
self._do_start()
|
||||
|
||||
@final
|
||||
def stop(self) -> None:
|
||||
with self._stopped_lock:
|
||||
if self._stopped:
|
||||
return
|
||||
self._stopped = True
|
||||
self._do_stop()
|
||||
if self._stopped:
|
||||
return
|
||||
self._stopped = True
|
||||
self._do_stop()
|
||||
|
||||
@abstractmethod
|
||||
def _do_start(self) -> None:
|
||||
|
||||
@@ -6,7 +6,6 @@ from fastapi import FastAPI, status, Request, Depends
|
||||
from fastapi.responses import JSONResponse
|
||||
from fastapi.middleware.cors import CORSMiddleware
|
||||
from fastapi.staticfiles import StaticFiles
|
||||
from starlette.responses import Response
|
||||
from pydantic import ValidationError
|
||||
from pkg_resources import resource_filename
|
||||
|
||||
@@ -144,20 +143,6 @@ class WebAppFiles(StaticFiles):
|
||||
path = 'index.html'
|
||||
return await super().lookup_path(path)
|
||||
|
||||
def file_response(self, full_path: str, *args, **kwargs) -> Response: # type: ignore # noqa
|
||||
# ignore MIME types from Windows registry
|
||||
# workaround for https://github.com/acgnhiki/blrec/issues/12
|
||||
response = super().file_response(full_path, *args, **kwargs)
|
||||
if full_path.endswith('.js'):
|
||||
js_media_type = 'application/javascript'
|
||||
if response.media_type != js_media_type:
|
||||
response.media_type = js_media_type
|
||||
headers = response.headers
|
||||
headers['content-type'] = js_media_type
|
||||
response.raw_headers = headers.raw
|
||||
del response._headers
|
||||
return response
|
||||
|
||||
|
||||
directory = resource_filename(__name__, '../data/webapp')
|
||||
api.mount('/', WebAppFiles(directory=directory, html=True), name='webapp')
|
||||
|
||||
@@ -1,10 +1,8 @@
|
||||
<nz-modal
|
||||
nzTitle="修改文件路径模板"
|
||||
nzCentered
|
||||
[nzFooter]="modalFooter"
|
||||
[(nzVisible)]="visible"
|
||||
[nzOkDisabled]="control.invalid || control.value.trim() === value"
|
||||
(nzOnOk)="handleConfirm()"
|
||||
(nzOnCancel)="handleCancel()"
|
||||
>
|
||||
<ng-container *nzModalContent>
|
||||
<form nz-form [formGroup]="settingsForm">
|
||||
|
||||
@@ -28,7 +28,7 @@ export function toBitRateString(
|
||||
let unit: string;
|
||||
|
||||
if (bitrate <= 0) {
|
||||
return '0' + spacer + 'kbps';
|
||||
return '0 kbps/s';
|
||||
}
|
||||
|
||||
if (bitrate < 1e6) {
|
||||
@@ -60,7 +60,7 @@ export function toByteRateString(
|
||||
let unit: string;
|
||||
|
||||
if (rate <= 0) {
|
||||
return '0' + spacer + 'B/s';
|
||||
return '0B/s';
|
||||
}
|
||||
|
||||
if (rate < 1e3) {
|
||||
|
||||
@@ -86,7 +86,7 @@
|
||||
<span class="label">下载速度</span>
|
||||
<app-wave-graph [value]="data.task_status.dl_rate"></app-wave-graph>
|
||||
<span class="value">
|
||||
{{ data.task_status.dl_rate * 8 | datarate: { bitrate: true } }}
|
||||
{{ data.task_status.dl_rate | datarate: { bitrate: true } }}
|
||||
</span>
|
||||
</li>
|
||||
<li class="info-item">
|
||||
|
||||
Reference in New Issue
Block a user