Compare commits
11 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6c1ed1bbb7 | ||
|
|
996dca3e94 | ||
|
|
d54721d553 | ||
|
|
97b1b3cd02 | ||
|
|
b80019a258 | ||
|
|
74dd739ec7 | ||
|
|
486e2ba552 | ||
|
|
60a8f23a14 | ||
|
|
34d9aa63ef | ||
|
|
6469881220 | ||
|
|
50971eeb0e |
4
.github/workflows/portable.yml
vendored
4
.github/workflows/portable.yml
vendored
@@ -8,7 +8,7 @@ on:
|
||||
env:
|
||||
FFMPEG_ARCHIVE_URL: https://github.com/BtbN/FFmpeg-Builds/releases/download/latest/ffmpeg-master-latest-win64-lgpl-shared.zip
|
||||
FFMPEG_ARCHIVE_NAME: ffmpeg-master-latest-win64-lgpl-shared.zip
|
||||
PYTHON_ARCHIVE_URL: https://www.python.org/ftp/python/3.10.4/python-3.10.4-embed-amd64.zip
|
||||
PYTHON_ARCHIVE_URL: https://www.python.org/ftp/python/3.11.1/python-3.11.1-embed-amd64.zip
|
||||
|
||||
jobs:
|
||||
|
||||
@@ -23,7 +23,7 @@ jobs:
|
||||
- name: Setup Python
|
||||
uses: actions/setup-python@v3
|
||||
with:
|
||||
python-version: "3.10"
|
||||
python-version: "3.11.1"
|
||||
|
||||
- name: Download ffmpeg archive
|
||||
run: Invoke-WebRequest -Uri $($env:FFMPEG_ARCHIVE_URL) -OutFile ffmpeg.zip
|
||||
|
||||
@@ -1,5 +1,12 @@
|
||||
# 更新日志
|
||||
|
||||
## 1.13.0
|
||||
|
||||
- 支持 Python 3.11
|
||||
- 改进直播监控
|
||||
- 优化在 Linux 下的内存占用
|
||||
- docker 时区设置为默认 `Asia/Shanghai`
|
||||
|
||||
## 1.12.0
|
||||
|
||||
- 支持自定义 Telegram bot api 地址
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
# syntax=docker/dockerfile:1
|
||||
|
||||
FROM python:3.10-slim-buster
|
||||
FROM python:3.11-slim-buster
|
||||
|
||||
WORKDIR /app
|
||||
VOLUME ["/cfg", "/log", "/rec"]
|
||||
@@ -18,6 +18,7 @@ RUN apt-get update && \
|
||||
ENV DEFAULT_SETTINGS_FILE=/cfg/settings.toml
|
||||
ENV DEFAULT_LOG_DIR=/log
|
||||
ENV DEFAULT_OUT_DIR=/rec
|
||||
ENV TZ="Asia/Shanghai"
|
||||
|
||||
EXPOSE 2233
|
||||
ENTRYPOINT ["blrec", "--host", "0.0.0.0"]
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
# syntax=docker/dockerfile:1
|
||||
|
||||
FROM python:3.10-slim-buster
|
||||
FROM python:3.11-slim-buster
|
||||
|
||||
WORKDIR /app
|
||||
VOLUME ["/cfg", "/log", "/rec"]
|
||||
@@ -20,6 +20,7 @@ RUN sed -i "s/deb.debian.org/mirrors.tuna.tsinghua.edu.cn/g" /etc/apt/sources.li
|
||||
ENV DEFAULT_SETTINGS_FILE=/cfg/settings.toml
|
||||
ENV DEFAULT_LOG_DIR=/log
|
||||
ENV DEFAULT_OUT_DIR=/rec
|
||||
ENV TZ="Asia/Shanghai"
|
||||
|
||||
EXPOSE 2233
|
||||
ENTRYPOINT ["blrec", "--host", "0.0.0.0"]
|
||||
|
||||
@@ -23,7 +23,7 @@
|
||||
- 支持按文件大小或时长分割文件
|
||||
- 支持转换 `flv` 为 `mp4` 格式(需要安装 `ffmpeg`)
|
||||
- 硬盘空间检测并支持空间不足自动删除旧录播文件。
|
||||
- 事件通知(支持邮箱、`ServerChan`、`pushplus`)
|
||||
- 事件通知(支持邮箱、`ServerChan`、`PushDeer`、`pushplus`、`Telegram`、`Bark` )
|
||||
- `Webhook`(可配合 `REST API` 实现录制控制,录制完成后压制、上传等自定义需求)
|
||||
|
||||
## 前提条件
|
||||
@@ -82,6 +82,7 @@
|
||||
- 默认设置文件位置: `ENV DEFAULT_SETTINGS_FILE=/cfg/settings.toml`
|
||||
- 默认日志存放目录: `ENV DEFAULT_LOG_DIR=/log`
|
||||
- 默认录播存放目录: `ENV DEFAULT_OUT_DIR=/rec`
|
||||
- 默认时区: `ENV TZ="Asia/Shanghai"`
|
||||
|
||||
### 默认参数运行
|
||||
|
||||
@@ -240,11 +241,6 @@ api key 可以使用数字和字母,长度限制为最短 8 最长 80。
|
||||
|
||||
---
|
||||
|
||||
## Thanks
|
||||
|
||||
[](https://jb.gg/OpenSource)
|
||||
|
||||
|
||||
## 其它相关工具或项目
|
||||
|
||||
| 名称 | 链接 | 简介 |
|
||||
|
||||
10
setup.cfg
10
setup.cfg
@@ -38,13 +38,13 @@ install_requires =
|
||||
python-liquid >= 1.2.1, < 2.0.0
|
||||
typing-extensions >= 3.10.0.0
|
||||
ordered-set >= 4.1.0, < 5.0.0
|
||||
fastapi >= 0.70.0, < 0.71.0
|
||||
fastapi >= 0.88.0, < 0.89.0
|
||||
email_validator >= 1.1.3, < 2.0.0
|
||||
click < 8.1.0
|
||||
typer >= 0.4.1, < 0.5.0
|
||||
typer >= 0.7.0, < 0.8.0
|
||||
aiohttp >= 3.8.1, < 4.0.0
|
||||
requests >= 2.24.0, < 3.0.0
|
||||
aiofiles >= 0.8.0, < 0.9.0
|
||||
aiofiles >= 22.1.0, < 23.0.0
|
||||
tenacity >= 8.0.1, < 9.0.0
|
||||
colorama >= 0.4.4, < 0.5.0
|
||||
humanize >= 3.13.1, < 4.0.0
|
||||
@@ -52,14 +52,14 @@ install_requires =
|
||||
attrs >= 21.2.0, < 22.0.0
|
||||
lxml >= 4.6.4, < 5.0.0
|
||||
toml >= 0.10.2, < 0.11.0
|
||||
m3u8 >= 1.0.0, < 2.0.0
|
||||
m3u8 >= 3.3.0, < 4.0.0
|
||||
av >= 10.0.0, < 11.0.0
|
||||
jsonpath == 0.82
|
||||
psutil >= 5.8.0, < 6.0.0
|
||||
reactivex >= 4.0.0, < 5.0.0
|
||||
bitarray >= 2.2.5, < 3.0.0
|
||||
brotli >= 1.0.9, < 2.0.0
|
||||
uvicorn[standard] >= 0.15.0, < 0.16.0
|
||||
uvicorn[standard] >= 0.20.0, < 0.21.0
|
||||
|
||||
[options.extras_require]
|
||||
dev =
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
__prog__ = 'blrec'
|
||||
__version__ = '1.12.0'
|
||||
__version__ = '1.13.0'
|
||||
__github__ = 'https://github.com/acgnhiki/blrec'
|
||||
|
||||
@@ -1,43 +1,36 @@
|
||||
import os
|
||||
import logging
|
||||
import asyncio
|
||||
import logging
|
||||
import os
|
||||
from typing import Iterator, List, Optional
|
||||
|
||||
import attr
|
||||
import psutil
|
||||
|
||||
from . import __prog__, __version__
|
||||
from .flv.operators import MetaData, StreamProfile
|
||||
from .disk_space import SpaceMonitor, SpaceReclaimer
|
||||
from .bili.helpers import ensure_room_id
|
||||
from .disk_space import SpaceMonitor, SpaceReclaimer
|
||||
from .event.event_submitters import SpaceEventSubmitter
|
||||
from .exception import ExceptionHandler, ExistsError, exception_callback
|
||||
from .flv.operators import MetaData, StreamProfile
|
||||
from .notification import (
|
||||
BarkNotifier,
|
||||
EmailNotifier,
|
||||
PushdeerNotifier,
|
||||
PushplusNotifier,
|
||||
ServerchanNotifier,
|
||||
TelegramNotifier,
|
||||
)
|
||||
from .setting import Settings, SettingsIn, SettingsManager, SettingsOut, TaskOptions
|
||||
from .setting.typing import KeySetOfSettings
|
||||
from .task import (
|
||||
DanmakuFileDetail,
|
||||
RecordTaskManager,
|
||||
TaskData,
|
||||
TaskParam,
|
||||
VideoFileDetail,
|
||||
DanmakuFileDetail,
|
||||
)
|
||||
from .exception import ExistsError, ExceptionHandler, exception_callback
|
||||
from .event.event_submitters import SpaceEventSubmitter
|
||||
from .setting import (
|
||||
SettingsManager,
|
||||
Settings,
|
||||
SettingsIn,
|
||||
SettingsOut,
|
||||
TaskOptions,
|
||||
)
|
||||
from .setting.typing import KeySetOfSettings
|
||||
from .notification import (
|
||||
EmailNotifier,
|
||||
ServerchanNotifier,
|
||||
PushdeerNotifier,
|
||||
PushplusNotifier,
|
||||
TelegramNotifier,
|
||||
BarkNotifier,
|
||||
)
|
||||
from .webhook import WebHookEmitter
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -105,6 +98,7 @@ class Application:
|
||||
await self.exit()
|
||||
|
||||
async def launch(self) -> None:
|
||||
logger.info('Launching Application...')
|
||||
self._setup()
|
||||
logger.debug(f'Default umask {os.umask(000)}')
|
||||
logger.info(f'Launched Application v{__version__}')
|
||||
@@ -112,10 +106,12 @@ class Application:
|
||||
task.add_done_callback(exception_callback)
|
||||
|
||||
async def exit(self) -> None:
|
||||
logger.info('Exiting Application...')
|
||||
await self._exit()
|
||||
logger.info('Exited Application')
|
||||
|
||||
async def abort(self) -> None:
|
||||
logger.info('Aborting Application...')
|
||||
await self._exit(force=True)
|
||||
logger.info('Aborted Application')
|
||||
|
||||
@@ -128,6 +124,7 @@ class Application:
|
||||
logger.info('Restarting Application...')
|
||||
await self.exit()
|
||||
await self.launch()
|
||||
logger.info('Restarted Application')
|
||||
|
||||
def has_task(self, room_id: int) -> bool:
|
||||
return self._task_manager.has_task(room_id)
|
||||
@@ -136,9 +133,7 @@ class Application:
|
||||
room_id = await ensure_room_id(room_id)
|
||||
|
||||
if self._task_manager.has_task(room_id):
|
||||
raise ExistsError(
|
||||
f'a task for the room {room_id} is already existed'
|
||||
)
|
||||
raise ExistsError(f'a task for the room {room_id} is already existed')
|
||||
|
||||
settings = self._settings_manager.find_task_settings(room_id)
|
||||
if not settings:
|
||||
@@ -214,9 +209,7 @@ class Application:
|
||||
await self._settings_manager.mark_task_recorder_enabled(room_id)
|
||||
logger.info(f'Successfully enabled recorder for task {room_id}')
|
||||
|
||||
async def disable_task_recorder(
|
||||
self, room_id: int, force: bool = False
|
||||
) -> None:
|
||||
async def disable_task_recorder(self, room_id: int, force: bool = False) -> None:
|
||||
logger.info(f'Disabling recorder for task {room_id}...')
|
||||
await self._task_manager.disable_task_recorder(room_id, force)
|
||||
await self._settings_manager.mark_task_recorder_disabled(room_id)
|
||||
@@ -249,9 +242,7 @@ class Application:
|
||||
def get_task_stream_profile(self, room_id: int) -> StreamProfile:
|
||||
return self._task_manager.get_task_stream_profile(room_id)
|
||||
|
||||
def get_task_video_file_details(
|
||||
self, room_id: int
|
||||
) -> Iterator[VideoFileDetail]:
|
||||
def get_task_video_file_details(self, room_id: int) -> Iterator[VideoFileDetail]:
|
||||
yield from self._task_manager.get_task_video_file_details(room_id)
|
||||
|
||||
def get_task_danmaku_file_details(
|
||||
@@ -291,9 +282,7 @@ class Application:
|
||||
async def change_task_options(
|
||||
self, room_id: int, options: TaskOptions
|
||||
) -> TaskOptions:
|
||||
return await self._settings_manager.change_task_options(
|
||||
room_id, options
|
||||
)
|
||||
return await self._settings_manager.change_task_options(room_id, options)
|
||||
|
||||
def _setup(self) -> None:
|
||||
self._setup_logger()
|
||||
@@ -320,9 +309,7 @@ class Application:
|
||||
self._space_event_submitter = SpaceEventSubmitter(self._space_monitor)
|
||||
|
||||
def _setup_space_reclaimer(self) -> None:
|
||||
self._space_reclaimer = SpaceReclaimer(
|
||||
self._space_monitor, self._out_dir,
|
||||
)
|
||||
self._space_reclaimer = SpaceReclaimer(self._space_monitor, self._out_dir)
|
||||
self._settings_manager.apply_space_reclaimer_settings()
|
||||
self._space_reclaimer.enable()
|
||||
|
||||
|
||||
@@ -1,13 +1,17 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import random
|
||||
from contextlib import suppress
|
||||
|
||||
from blrec.exception import exception_callback
|
||||
from blrec.logging.room_id import aio_task_with_room_id
|
||||
|
||||
from .danmaku_client import DanmakuClient, DanmakuListener, DanmakuCommand
|
||||
from .live import Live
|
||||
from .typing import Danmaku
|
||||
from .models import LiveStatus, RoomInfo
|
||||
from ..event.event_emitter import EventListener, EventEmitter
|
||||
from ..event.event_emitter import EventEmitter, EventListener
|
||||
from ..utils.mixins import SwitchableMixin
|
||||
|
||||
from .danmaku_client import DanmakuClient, DanmakuCommand, DanmakuListener
|
||||
from .live import Live
|
||||
from .models import LiveStatus, RoomInfo
|
||||
from .typing import Danmaku
|
||||
|
||||
__all__ = 'LiveMonitor', 'LiveEventListener'
|
||||
|
||||
@@ -37,9 +41,7 @@ class LiveEventListener(EventListener):
|
||||
...
|
||||
|
||||
|
||||
class LiveMonitor(
|
||||
EventEmitter[LiveEventListener], DanmakuListener, SwitchableMixin
|
||||
):
|
||||
class LiveMonitor(EventEmitter[LiveEventListener], DanmakuListener, SwitchableMixin):
|
||||
def __init__(self, danmaku_client: DanmakuClient, live: Live) -> None:
|
||||
super().__init__()
|
||||
self._danmaku_client = danmaku_client
|
||||
@@ -49,18 +51,49 @@ class LiveMonitor(
|
||||
self._previous_status = self._live.room_info.live_status
|
||||
if self._live.is_living():
|
||||
self._status_count = 2
|
||||
self._stream_available = True
|
||||
else:
|
||||
self._status_count = 0
|
||||
self._stream_available = False
|
||||
|
||||
def _do_enable(self) -> None:
|
||||
self._init_status()
|
||||
self._danmaku_client.add_listener(self)
|
||||
self._start_polling()
|
||||
logger.debug('Enabled live monitor')
|
||||
|
||||
def _do_disable(self) -> None:
|
||||
self._danmaku_client.remove_listener(self)
|
||||
asyncio.create_task(self._stop_polling())
|
||||
asyncio.create_task(self._stop_checking())
|
||||
logger.debug('Disabled live monitor')
|
||||
|
||||
def _start_polling(self) -> None:
|
||||
self._polling_task = asyncio.create_task(self._poll_live_status())
|
||||
self._polling_task.add_done_callback(exception_callback)
|
||||
logger.debug('Started polling live status')
|
||||
|
||||
async def _stop_polling(self) -> None:
|
||||
self._polling_task.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await self._polling_task
|
||||
del self._polling_task
|
||||
logger.debug('Stopped polling live status')
|
||||
|
||||
def _start_checking(self) -> None:
|
||||
self._checking_task = asyncio.create_task(self._check_if_stream_available())
|
||||
self._checking_task.add_done_callback(exception_callback)
|
||||
logger.debug('Started checking if stream available')
|
||||
|
||||
async def _stop_checking(self) -> None:
|
||||
if not hasattr(self, '_checking_task'):
|
||||
return
|
||||
self._checking_task.cancel()
|
||||
with suppress(asyncio.CancelledError):
|
||||
await self._checking_task
|
||||
del self._checking_task
|
||||
logger.debug('Stopped checking if stream available')
|
||||
|
||||
async def on_client_reconnected(self) -> None:
|
||||
# check the live status after the client reconnected and simulate
|
||||
# events if necessary.
|
||||
@@ -89,8 +122,10 @@ class LiveMonitor(
|
||||
danmu_cmd = danmu['cmd']
|
||||
|
||||
if danmu_cmd == DanmakuCommand.LIVE.value:
|
||||
await self._live.update_room_info()
|
||||
await self._handle_status_change(LiveStatus.LIVE)
|
||||
elif danmu_cmd == DanmakuCommand.PREPARING.value:
|
||||
await self._live.update_room_info()
|
||||
if danmu.get('round', None) == 1:
|
||||
await self._handle_status_change(LiveStatus.ROUND)
|
||||
else:
|
||||
@@ -100,41 +135,57 @@ class LiveMonitor(
|
||||
await self._emit('room_changed', self._live.room_info)
|
||||
|
||||
async def _handle_status_change(self, current_status: LiveStatus) -> None:
|
||||
logger.debug('Live status changed from {} to {}'.format(
|
||||
self._previous_status.name, current_status.name
|
||||
))
|
||||
|
||||
await self._live.update_room_info()
|
||||
if (s := self._live.room_info.live_status) != current_status:
|
||||
logger.warning(
|
||||
'Updated live status {} is inconsistent with '
|
||||
'current live status {}'.format(s.name, current_status.name)
|
||||
logger.debug(
|
||||
'Live status changed from {} to {}'.format(
|
||||
self._previous_status.name, current_status.name
|
||||
)
|
||||
|
||||
await self._emit(
|
||||
'live_status_changed', current_status, self._previous_status
|
||||
)
|
||||
|
||||
await self._emit('live_status_changed', current_status, self._previous_status)
|
||||
|
||||
if current_status != LiveStatus.LIVE:
|
||||
self._status_count = 0
|
||||
self._stream_available = False
|
||||
await self._emit('live_ended', self._live)
|
||||
else:
|
||||
self._status_count += 1
|
||||
|
||||
if self._status_count == 1:
|
||||
assert self._previous_status != LiveStatus.LIVE
|
||||
self._start_checking()
|
||||
await self._emit('live_began', self._live)
|
||||
elif self._status_count == 2:
|
||||
assert self._previous_status == LiveStatus.LIVE
|
||||
await self._emit('live_stream_available', self._live)
|
||||
if not self._stream_available:
|
||||
self._stream_available = True
|
||||
await self._stop_checking()
|
||||
await self._emit('live_stream_available', self._live)
|
||||
elif self._status_count > 2:
|
||||
assert self._previous_status == LiveStatus.LIVE
|
||||
await self._emit('live_stream_reset', self._live)
|
||||
else:
|
||||
pass
|
||||
|
||||
logger.debug('Number of sequential LIVE status: {}'.format(
|
||||
self._status_count
|
||||
))
|
||||
logger.debug('Number of sequential LIVE status: {}'.format(self._status_count))
|
||||
|
||||
self._previous_status = current_status
|
||||
|
||||
@aio_task_with_room_id
|
||||
async def _poll_live_status(self) -> None:
|
||||
while True:
|
||||
await asyncio.sleep(600 + random.randrange(-60, 60))
|
||||
await self._live.update_room_info()
|
||||
current_status = self._live.room_info.live_status
|
||||
if current_status != self._previous_status:
|
||||
await self._handle_status_change(current_status)
|
||||
|
||||
@aio_task_with_room_id
|
||||
async def _check_if_stream_available(self) -> None:
|
||||
while not self._stream_available:
|
||||
try:
|
||||
await self._live.get_live_stream_url()
|
||||
except Exception:
|
||||
await asyncio.sleep(1)
|
||||
else:
|
||||
self._stream_available = True
|
||||
await self._emit('live_stream_available', self._live)
|
||||
|
||||
@@ -1,13 +1,13 @@
|
||||
import logging
|
||||
from typing import Optional
|
||||
|
||||
from reactivex import operators as ops
|
||||
from reactivex.scheduler import NewThreadScheduler
|
||||
|
||||
from blrec.bili.live import Live
|
||||
from blrec.bili.typing import QualityNumber
|
||||
from blrec.hls import operators as hls_ops
|
||||
from blrec.hls.metadata_dumper import MetadataDumper
|
||||
from blrec.utils import operators as utils_ops
|
||||
|
||||
from . import operators as core_ops
|
||||
from .stream_recorder_impl import StreamRecorderImpl
|
||||
@@ -84,16 +84,14 @@ class HLSRawStreamRecorderImpl(StreamRecorderImpl):
|
||||
self._stream_param_holder.get_stream_params() # type: ignore
|
||||
.pipe(
|
||||
self._stream_url_resolver,
|
||||
ops.subscribe_on(
|
||||
NewThreadScheduler(self._thread_factory('PlaylistDownloader'))
|
||||
),
|
||||
self._playlist_fetcher,
|
||||
self._recording_monitor,
|
||||
self._connection_error_handler,
|
||||
self._request_exception_handler,
|
||||
self._playlist_dumper,
|
||||
ops.observe_on(
|
||||
NewThreadScheduler(self._thread_factory('SegmentDownloader'))
|
||||
utils_ops.observe_on_new_thread(
|
||||
queue_size=60,
|
||||
thread_name=f'SegmentDownloader::{self._live.room_id}',
|
||||
),
|
||||
self._segment_fetcher,
|
||||
self._dl_statistics,
|
||||
@@ -103,5 +101,10 @@ class HLSRawStreamRecorderImpl(StreamRecorderImpl):
|
||||
self._progress_bar,
|
||||
self._exception_handler,
|
||||
)
|
||||
.subscribe(on_completed=self._on_completed)
|
||||
.subscribe(
|
||||
on_completed=self._on_completed,
|
||||
scheduler=NewThreadScheduler(
|
||||
self._thread_factory('HLSRawStreamRecorder')
|
||||
),
|
||||
)
|
||||
)
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
import logging
|
||||
from typing import Optional
|
||||
|
||||
from reactivex import operators as ops
|
||||
from reactivex.scheduler import NewThreadScheduler
|
||||
|
||||
from blrec.bili.live import Live
|
||||
@@ -9,6 +8,7 @@ from blrec.bili.typing import QualityNumber
|
||||
from blrec.flv import operators as flv_ops
|
||||
from blrec.flv.metadata_dumper import MetadataDumper
|
||||
from blrec.hls import operators as hls_ops
|
||||
from blrec.utils import operators as utils_ops
|
||||
|
||||
from . import operators as core_ops
|
||||
from .stream_recorder_impl import StreamRecorderImpl
|
||||
@@ -130,22 +130,19 @@ class HLSStreamRecorderImpl(StreamRecorderImpl):
|
||||
self._stream_param_holder.get_stream_params() # type: ignore
|
||||
.pipe(
|
||||
self._stream_url_resolver,
|
||||
ops.subscribe_on(
|
||||
NewThreadScheduler(self._thread_factory('PlaylistFetcher'))
|
||||
),
|
||||
self._playlist_fetcher,
|
||||
self._recording_monitor,
|
||||
self._connection_error_handler,
|
||||
self._request_exception_handler,
|
||||
self._playlist_resolver,
|
||||
ops.observe_on(
|
||||
NewThreadScheduler(self._thread_factory('SegmentFetcher'))
|
||||
utils_ops.observe_on_new_thread(
|
||||
queue_size=60, thread_name=f'SegmentFetcher::{self._live.room_id}'
|
||||
),
|
||||
self._segment_fetcher,
|
||||
self._dl_statistics,
|
||||
self._prober,
|
||||
ops.observe_on(
|
||||
NewThreadScheduler(self._thread_factory('StreamRecorder'))
|
||||
utils_ops.observe_on_new_thread(
|
||||
queue_size=10, thread_name=f'StreamRecorder::{self._live.room_id}'
|
||||
),
|
||||
self._segment_remuxer,
|
||||
self._segment_parser,
|
||||
@@ -160,5 +157,8 @@ class HLSStreamRecorderImpl(StreamRecorderImpl):
|
||||
self._progress_bar,
|
||||
self._exception_handler,
|
||||
)
|
||||
.subscribe(on_completed=self._on_completed)
|
||||
.subscribe(
|
||||
on_completed=self._on_completed,
|
||||
scheduler=NewThreadScheduler(self._thread_factory('HLSStreamRecorder')),
|
||||
)
|
||||
)
|
||||
|
||||
@@ -14,7 +14,6 @@ from blrec.bili.models import RoomInfo
|
||||
from blrec.bili.typing import QualityNumber, StreamFormat
|
||||
from blrec.event.event_emitter import EventEmitter, EventListener
|
||||
from blrec.flv.operators import MetaData, StreamProfile
|
||||
from blrec.logging.room_id import aio_task_with_room_id
|
||||
from blrec.setting.typing import RecordingMode
|
||||
from blrec.utils.mixins import AsyncStoppableMixin
|
||||
|
||||
@@ -457,8 +456,6 @@ class Recorder(
|
||||
await self._prepare()
|
||||
if self._stream_available:
|
||||
await self._stream_recorder.start()
|
||||
else:
|
||||
asyncio.create_task(self._guard())
|
||||
|
||||
logger.info('Started recording')
|
||||
await self._emit('recording_started', self)
|
||||
@@ -489,29 +486,6 @@ class Recorder(
|
||||
self._danmaku_dumper.clear_files()
|
||||
self._stream_recorder.clear_files()
|
||||
|
||||
@aio_task_with_room_id
|
||||
async def _guard(self, timeout: float = 60) -> None:
|
||||
await asyncio.sleep(timeout)
|
||||
if not self._recording:
|
||||
return
|
||||
if self._stream_available:
|
||||
return
|
||||
logger.debug(
|
||||
f'Stream not available in {timeout} seconds, the event maybe lost.'
|
||||
)
|
||||
|
||||
await self._live.update_info()
|
||||
if self._live.is_living():
|
||||
logger.debug('The live is living now')
|
||||
self._stream_available = True
|
||||
if self._stream_recorder.stopped:
|
||||
await self._stream_recorder.start()
|
||||
else:
|
||||
logger.debug('The live has ended before streaming')
|
||||
self._stream_available = False
|
||||
if not self._stream_recorder.stopped:
|
||||
await self.stop()
|
||||
|
||||
def _print_waiting_message(self) -> None:
|
||||
logger.info('Waiting... until the live starts')
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ from blrec.bili.typing import QualityNumber, StreamFormat
|
||||
from blrec.event.event_emitter import EventEmitter
|
||||
from blrec.flv.operators import MetaData, StreamProfile
|
||||
from blrec.setting.typing import RecordingMode
|
||||
from blrec.utils.libc import malloc_trim
|
||||
from blrec.utils.mixins import AsyncStoppableMixin
|
||||
|
||||
from .flv_stream_recorder_impl import FLVStreamRecorderImpl
|
||||
@@ -255,6 +256,7 @@ class StreamRecorder(
|
||||
|
||||
async def _do_stop(self) -> None:
|
||||
await self._impl.stop()
|
||||
malloc_trim(0)
|
||||
|
||||
async def on_video_file_created(self, path: str, record_start_time: int) -> None:
|
||||
await self._emit('video_file_created', path, record_start_time)
|
||||
|
||||
@@ -148,6 +148,8 @@ class Postprocessor(
|
||||
if video_path.endswith('.flv'):
|
||||
if not await self._is_vaild_flv_file(video_path):
|
||||
logger.warning(f'The flv file may be invalid: {video_path}')
|
||||
if os.path.getsize(video_path) < 1024**2:
|
||||
continue
|
||||
if self.remux_to_mp4:
|
||||
self._status = PostprocessorStatus.REMUXING
|
||||
(
|
||||
|
||||
@@ -7,6 +7,8 @@ from typing import TYPE_CHECKING, Dict, Iterator, Optional
|
||||
import aiohttp
|
||||
from tenacity import retry, retry_if_exception_type, stop_after_delay, wait_exponential
|
||||
|
||||
from blrec.utils.libc import malloc_trim
|
||||
|
||||
from ..bili.exceptions import ApiRequestError
|
||||
from ..exception import NotFoundError, submit_exception
|
||||
from ..flv.operators import MetaData, StreamProfile
|
||||
@@ -17,8 +19,8 @@ if TYPE_CHECKING:
|
||||
from ..setting import SettingsManager
|
||||
|
||||
from ..setting import (
|
||||
DanmakuSettings,
|
||||
BiliApiSettings,
|
||||
DanmakuSettings,
|
||||
HeaderSettings,
|
||||
OutputSettings,
|
||||
PostprocessingSettings,
|
||||
@@ -52,12 +54,14 @@ class RecordTaskManager:
|
||||
logger.info('Load all tasks complete')
|
||||
|
||||
async def destroy_all_tasks(self) -> None:
|
||||
logger.info('Destroying all tasks...')
|
||||
if not self._tasks:
|
||||
return
|
||||
await asyncio.wait([t.destroy() for t in self._tasks.values() if t.ready])
|
||||
logger.debug('Destroying all tasks...')
|
||||
for task in self._tasks.values():
|
||||
if not task.ready:
|
||||
continue
|
||||
await task.destroy()
|
||||
self._tasks.clear()
|
||||
logger.info('Successfully destroyed all task')
|
||||
malloc_trim(0)
|
||||
logger.debug('Successfully destroyed all task')
|
||||
|
||||
def has_task(self, room_id: int) -> bool:
|
||||
return room_id in self._tasks
|
||||
@@ -110,72 +114,110 @@ class RecordTaskManager:
|
||||
logger.info(f'Successfully added task {settings.room_id}')
|
||||
|
||||
async def remove_task(self, room_id: int) -> None:
|
||||
logger.debug(f'Removing task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.disable_recorder(force=True)
|
||||
await task.disable_monitor()
|
||||
await task.destroy()
|
||||
del self._tasks[room_id]
|
||||
malloc_trim(0)
|
||||
logger.debug(f'Removed task {room_id}')
|
||||
|
||||
async def remove_all_tasks(self) -> None:
|
||||
coros = [self.remove_task(i) for i, t in self._tasks.items() if t.ready]
|
||||
if coros:
|
||||
await asyncio.wait(coros)
|
||||
logger.debug('Removing all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.remove_task(room_id)
|
||||
malloc_trim(0)
|
||||
logger.debug('Removed all tasks')
|
||||
|
||||
async def start_task(self, room_id: int) -> None:
|
||||
logger.debug(f'Starting task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.update_info()
|
||||
await task.enable_monitor()
|
||||
await task.enable_recorder()
|
||||
logger.debug(f'Started task {room_id}')
|
||||
|
||||
async def stop_task(self, room_id: int, force: bool = False) -> None:
|
||||
logger.debug(f'Stopping task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.disable_recorder(force)
|
||||
await task.disable_monitor()
|
||||
logger.debug(f'Stopped task {room_id}')
|
||||
|
||||
async def start_all_tasks(self) -> None:
|
||||
await self.update_all_task_infos()
|
||||
await self.enable_all_task_monitors()
|
||||
await self.enable_all_task_recorders()
|
||||
logger.debug('Starting all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.start_task(room_id)
|
||||
logger.debug('Started all tasks')
|
||||
|
||||
async def stop_all_tasks(self, force: bool = False) -> None:
|
||||
await self.disable_all_task_recorders(force)
|
||||
await self.disable_all_task_monitors()
|
||||
logger.debug('Stopping all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.stop_task(room_id, force=force)
|
||||
logger.debug('Stopped all tasks')
|
||||
|
||||
async def enable_task_monitor(self, room_id: int) -> None:
|
||||
logger.debug(f'Enabling live monitor for task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.enable_monitor()
|
||||
logger.debug(f'Enabled live monitor for task {room_id}')
|
||||
|
||||
async def disable_task_monitor(self, room_id: int) -> None:
|
||||
logger.debug(f'Disabling live monitor for task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.disable_monitor()
|
||||
logger.debug(f'Disabled live monitor for task {room_id}')
|
||||
|
||||
async def enable_all_task_monitors(self) -> None:
|
||||
coros = [t.enable_monitor() for t in self._tasks.values() if t.ready]
|
||||
if coros:
|
||||
await asyncio.wait(coros)
|
||||
logger.debug('Enabling live monitor for all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.enable_task_monitor(room_id)
|
||||
logger.debug('Enabled live monitor for all tasks')
|
||||
|
||||
async def disable_all_task_monitors(self) -> None:
|
||||
coros = [t.disable_monitor() for t in self._tasks.values() if t.ready]
|
||||
if coros:
|
||||
await asyncio.wait(coros)
|
||||
logger.debug('Disabling live monitor for all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.disable_task_monitor(room_id)
|
||||
logger.debug('Disabled live monitor for all tasks')
|
||||
|
||||
async def enable_task_recorder(self, room_id: int) -> None:
|
||||
logger.debug(f'Enabling recorder for task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.enable_recorder()
|
||||
logger.debug(f'Enabled recorder for task {room_id}')
|
||||
|
||||
async def disable_task_recorder(self, room_id: int, force: bool = False) -> None:
|
||||
logger.debug(f'Disabling recorder for task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.disable_recorder(force)
|
||||
logger.debug(f'Disabled recorder for task {room_id}')
|
||||
|
||||
async def enable_all_task_recorders(self) -> None:
|
||||
coros = [t.enable_recorder() for t in self._tasks.values() if t.ready]
|
||||
if coros:
|
||||
await asyncio.wait(coros)
|
||||
logger.debug('Enabling recorder for all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.enable_task_recorder(room_id)
|
||||
logger.debug('Enabled recorder for all tasks')
|
||||
|
||||
async def disable_all_task_recorders(self, force: bool = False) -> None:
|
||||
coros = [t.disable_recorder(force) for t in self._tasks.values() if t.ready]
|
||||
if coros:
|
||||
await asyncio.wait(coros)
|
||||
logger.debug('Disabling recorder for all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.disable_task_recorder(room_id, force=force)
|
||||
logger.debug('Disabled recorder for all tasks')
|
||||
|
||||
def get_task_data(self, room_id: int) -> TaskData:
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
@@ -216,15 +258,18 @@ class RecordTaskManager:
|
||||
return task.cut_stream()
|
||||
|
||||
async def update_task_info(self, room_id: int) -> None:
|
||||
logger.debug(f'Updating info for task {room_id}...')
|
||||
task = self._get_task(room_id, check_ready=True)
|
||||
await task.update_info(raise_exception=True)
|
||||
logger.debug(f'Updated info for task {room_id}')
|
||||
|
||||
async def update_all_task_infos(self) -> None:
|
||||
coros = [
|
||||
t.update_info(raise_exception=True) for t in self._tasks.values() if t.ready
|
||||
]
|
||||
if coros:
|
||||
await asyncio.wait(coros)
|
||||
logger.debug('Updating info for all tasks...')
|
||||
for room_id, task in self._tasks.items():
|
||||
if not task.ready:
|
||||
continue
|
||||
await self.update_task_info(room_id)
|
||||
logger.debug('Updated info for all tasks')
|
||||
|
||||
def apply_task_bili_api_settings(
|
||||
self, room_id: int, settings: BiliApiSettings
|
||||
|
||||
16
src/blrec/utils/libc.py
Normal file
16
src/blrec/utils/libc.py
Normal file
@@ -0,0 +1,16 @@
|
||||
from ctypes import cdll
|
||||
from ctypes.util import find_library
|
||||
|
||||
lib_name = find_library('c')
|
||||
if not lib_name:
|
||||
libc = None
|
||||
else:
|
||||
libc = cdll.LoadLibrary(lib_name)
|
||||
|
||||
|
||||
def malloc_trim(pad: int) -> bool:
|
||||
"""Release free memory from the heap"""
|
||||
assert pad >= 0, 'pad must be >= 0'
|
||||
if libc is None:
|
||||
return False
|
||||
return libc.malloc_trim(pad) == 1
|
||||
@@ -1,4 +1,5 @@
|
||||
from .replace import replace
|
||||
from .retry import retry
|
||||
from .observe_on import observe_on_new_thread
|
||||
|
||||
__all__ = ('replace', 'retry')
|
||||
__all__ = ('replace', 'retry', 'observe_on_new_thread')
|
||||
|
||||
54
src/blrec/utils/operators/observe_on.py
Normal file
54
src/blrec/utils/operators/observe_on.py
Normal file
@@ -0,0 +1,54 @@
|
||||
from queue import Queue
|
||||
from threading import Thread
|
||||
from typing import Any, Callable, Optional, TypeVar
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
|
||||
_T = TypeVar('_T')
|
||||
|
||||
|
||||
def observe_on_new_thread(
|
||||
queue_size: Optional[int] = None, thread_name: Optional[str] = None
|
||||
) -> Callable[[Observable[_T]], Observable[_T]]:
|
||||
def observe_on(source: Observable[_T]) -> Observable[_T]:
|
||||
def subscribe(
|
||||
observer: abc.ObserverBase[_T],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
queue: Queue[Callable[..., Any]] = Queue(maxsize=queue_size or 0)
|
||||
|
||||
def run() -> None:
|
||||
while not disposed:
|
||||
queue.get()()
|
||||
|
||||
thread = Thread(target=run, name=thread_name, daemon=True)
|
||||
thread.start()
|
||||
|
||||
def on_next(value: _T) -> None:
|
||||
queue.put(lambda: observer.on_next(value))
|
||||
|
||||
def on_error(exc: Exception) -> None:
|
||||
queue.put(lambda: observer.on_error(exc))
|
||||
|
||||
def on_completed() -> None:
|
||||
queue.put(lambda: observer.on_completed)
|
||||
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
queue.put(lambda: None)
|
||||
thread.join()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, on_error, on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return observe_on
|
||||
@@ -137,12 +137,12 @@ api.include_router(update.router)
|
||||
|
||||
|
||||
class WebAppFiles(StaticFiles):
|
||||
async def lookup_path(
|
||||
def lookup_path(
|
||||
self, path: str
|
||||
) -> Tuple[str, Optional[os.stat_result]]:
|
||||
if path == '404.html':
|
||||
path = 'index.html'
|
||||
return await super().lookup_path(path)
|
||||
return super().lookup_path(path)
|
||||
|
||||
def file_response(self, full_path: str, *args, **kwargs) -> Response: # type: ignore # noqa
|
||||
# ignore MIME types from Windows registry
|
||||
|
||||
@@ -1,11 +1,10 @@
|
||||
import logging
|
||||
import secrets
|
||||
from typing import Optional, Set, Dict
|
||||
from typing import Dict, Optional, Set
|
||||
|
||||
from fastapi import status, Request, Header
|
||||
from fastapi import Header, Request, status
|
||||
from fastapi.exceptions import HTTPException
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@@ -21,34 +20,26 @@ attempting_clients: Dict[str, int] = {}
|
||||
|
||||
|
||||
async def authenticate(
|
||||
request: Request,
|
||||
x_api_key: Optional[str] = Header(None),
|
||||
request: Request, x_api_key: Optional[str] = Header(None)
|
||||
) -> None:
|
||||
assert api_key, 'api_key is required'
|
||||
|
||||
if not x_api_key:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_401_UNAUTHORIZED,
|
||||
detail='No api key',
|
||||
status_code=status.HTTP_401_UNAUTHORIZED, detail='No api key'
|
||||
)
|
||||
|
||||
assert request.client is not None, 'client should not be None'
|
||||
client_ip = request.client.host
|
||||
assert client_ip, 'client_ip is required'
|
||||
|
||||
if client_ip in blacklist:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_403_FORBIDDEN,
|
||||
detail='Blacklisted',
|
||||
)
|
||||
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail='Blacklisted')
|
||||
if client_ip not in whitelist:
|
||||
if (
|
||||
len(whitelist) >= MAX_WHITELIST or
|
||||
len(blacklist) >= MAX_BLACKLIST
|
||||
):
|
||||
if len(whitelist) >= MAX_WHITELIST or len(blacklist) >= MAX_BLACKLIST:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_403_FORBIDDEN,
|
||||
detail='Max clients allowed in whitelist or blacklist '
|
||||
'will exceeded',
|
||||
detail='Max clients allowed in whitelist or blacklist ' 'will exceeded',
|
||||
)
|
||||
if len(attempting_clients) >= MAX_ATTEMPTING_CLIENTS:
|
||||
raise HTTPException(
|
||||
@@ -71,8 +62,7 @@ async def authenticate(
|
||||
if client_ip in whitelist:
|
||||
whitelist.remove(client_ip)
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_401_UNAUTHORIZED,
|
||||
detail='API key is invalid',
|
||||
status_code=status.HTTP_401_UNAUTHORIZED, detail='API key is invalid'
|
||||
)
|
||||
|
||||
if client_ip in attempting_clients:
|
||||
|
||||
Reference in New Issue
Block a user