Compare commits
19 Commits
v1.8.0-alp
...
v1.8.1
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2b0537086e | ||
|
|
a4d61c1f29 | ||
|
|
17ea3ffc04 | ||
|
|
7f366463d4 | ||
|
|
fdd93e70a5 | ||
|
|
1370cecea8 | ||
|
|
18cd928369 | ||
|
|
7eea5c80b0 | ||
|
|
3b523ca11a | ||
|
|
133409a81a | ||
|
|
613718e56b | ||
|
|
955a531655 | ||
|
|
1ea9d10fce | ||
|
|
66999ef587 | ||
|
|
e5277510c2 | ||
|
|
e546b47e29 | ||
|
|
261a2993be | ||
|
|
44b031ef70 | ||
|
|
8d720cecc5 |
21
CHANGELOG.md
21
CHANGELOG.md
@@ -1,5 +1,26 @@
|
||||
# 更新日志
|
||||
|
||||
## 1.8.1
|
||||
|
||||
修 bug
|
||||
|
||||
## 1.8.0
|
||||
|
||||
- 重构直播流录制
|
||||
- 重构弹幕客户端
|
||||
- 修复了一些 bug
|
||||
- 优先使用 web api
|
||||
- 添加直播流时间相关元数据
|
||||
- 支持 Liquid 模板自定义通知消息
|
||||
- 同一时间只处理一个录播文件
|
||||
- 流录制中断重新调整弹幕时间
|
||||
- 对流主机进行排序
|
||||
|
||||
## 1.8.0-alpha.5
|
||||
|
||||
- 流录制中断重新调整弹幕时间
|
||||
- 重构并修正了一些问题
|
||||
|
||||
## 1.8.0-alpha.4
|
||||
|
||||
- 改善录制多个直播间出现卡顿
|
||||
|
||||
52
README.md
52
README.md
@@ -26,7 +26,7 @@
|
||||
- 事件通知(支持邮箱、`ServerChan`、`pushplus`)
|
||||
- `Webhook`(可配合 `REST API` 实现录制控制,录制完成后压制、上传等自定义需求)
|
||||
|
||||
## 先决条件
|
||||
## 前提条件
|
||||
|
||||
Python 3.8+
|
||||
ffmpeg、 ffprobe
|
||||
@@ -42,7 +42,9 @@
|
||||
|
||||
- 免安装绿色版
|
||||
|
||||
Windows 64 位系统用户也可以用打包好的免安装绿色版,下载后解压运行 `run.bat` 即可。
|
||||
支持 Windows 10+ 或 Windows Server 2016+,下载后解压运行 `run.bat` 或 `run.ps1` 。
|
||||
|
||||
不是官方或最新系统可能需要安装系统更新或缺少的 `C` 或 `C++` 运行时库
|
||||
|
||||
下载
|
||||
|
||||
@@ -192,6 +194,42 @@ api key 可以使用数字和字母,长度限制为最短 8 最长 80。
|
||||
|
||||
---
|
||||
|
||||
## 开发
|
||||
|
||||
1. 克隆代码
|
||||
|
||||
`git clone https://github.com/acgnhiki/blrec.git`
|
||||
|
||||
2. 进入项目目录
|
||||
|
||||
`cd blrec`
|
||||
|
||||
3. 创建虚拟环境
|
||||
|
||||
`python3 -m venv .venv`
|
||||
|
||||
4. 激活虚拟环境
|
||||
|
||||
`source .venv/bin/activate`
|
||||
|
||||
5. 以可编辑方式安装
|
||||
|
||||
`pip install -e .[dev]`
|
||||
|
||||
6. 修改代码
|
||||
|
||||
……
|
||||
|
||||
7. 运行 blrec
|
||||
|
||||
`blrec`
|
||||
|
||||
8. 退出虚拟环境
|
||||
|
||||
`deactivate`
|
||||
|
||||
---
|
||||
|
||||
## 常见问题
|
||||
|
||||
[FAQ](FAQ.md)
|
||||
@@ -205,3 +243,13 @@ api key 可以使用数字和字母,长度限制为最短 8 最长 80。
|
||||
## Thanks
|
||||
|
||||
[](https://jb.gg/OpenSource)
|
||||
|
||||
|
||||
## 其它相关工具或项目
|
||||
|
||||
| 名称 | 链接 | 简介 |
|
||||
| --- | --- | --- |
|
||||
| 录播姬 | [官网](https://rec.danmuji.org/) | 简单易用成熟稳定的 B 站直播录制工具 |
|
||||
| rclone | [官网](https://rclone.org/) | 可以挂载网盘用于存放录播文件 |
|
||||
| alist | [官网](https://alist-doc.nn.ci/) | 网盘文件浏览、播放 |
|
||||
| filebrowser | [官网](https://filebrowser.org/) | 服务器文件管理 |
|
||||
|
||||
2
run.bat
2
run.bat
@@ -1,6 +1,6 @@
|
||||
@echo off
|
||||
chcp 65001
|
||||
|
||||
cd %~dp0
|
||||
set PATH=.\ffmpeg\bin;.\python;%PATH%
|
||||
|
||||
@REM 不使用代理
|
||||
|
||||
2
run.ps1
2
run.ps1
@@ -1,5 +1,5 @@
|
||||
chcp 65001
|
||||
|
||||
Set-Location $PSScriptRoot
|
||||
$env:PATH = ".\ffmpeg\bin;.\python;" + $env:PATH
|
||||
|
||||
# 不使用代理
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
__prog__ = 'blrec'
|
||||
__version__ = '1.8.0-alpha.4'
|
||||
__version__ = '1.8.1'
|
||||
__github__ = 'https://github.com/acgnhiki/blrec'
|
||||
|
||||
@@ -3,7 +3,7 @@ import json
|
||||
import logging
|
||||
import re
|
||||
import time
|
||||
from typing import Dict, List, cast
|
||||
from typing import Any, Dict, List, cast
|
||||
|
||||
import aiohttp
|
||||
from jsonpath import jsonpath
|
||||
@@ -234,7 +234,26 @@ class Live:
|
||||
if qn not in accept_qn or codec['current_qn'] != qn:
|
||||
raise NoStreamQualityAvailable(stream_format, stream_codec, qn)
|
||||
|
||||
urls = [i['host'] + codec['base_url'] + i['extra'] for i in codec['url_info']]
|
||||
def sort_by_host(info: Any) -> int:
|
||||
host = info['host']
|
||||
if match := re.search(r'gotcha(\d+)', host):
|
||||
num = match.group(1)
|
||||
if num == '04':
|
||||
return 0
|
||||
if num == '09':
|
||||
return 1
|
||||
if num == '08':
|
||||
return 2
|
||||
return int(num)
|
||||
elif re.search(r'cn-[a-z]+-[a-z]+', host):
|
||||
return 1000
|
||||
elif 'mcdn' in host:
|
||||
return 2000
|
||||
else:
|
||||
return 10000
|
||||
|
||||
url_info = sorted(codec['url_info'], key=sort_by_host)
|
||||
urls = [i['host'] + codec['base_url'] + i['extra'] for i in url_info]
|
||||
|
||||
if not select_alternative:
|
||||
return urls[0]
|
||||
|
||||
@@ -116,8 +116,11 @@ class DanmakuDumper(
|
||||
self, video_path: str, record_start_time: int
|
||||
) -> None:
|
||||
with self._lock:
|
||||
self._delta: float = 0
|
||||
self._record_start_time: int = record_start_time
|
||||
self._timebase: int = self._record_start_time * 1000
|
||||
self._stream_recording_interrupted: bool = False
|
||||
self._path = danmaku_path(video_path)
|
||||
self._record_start_time = record_start_time
|
||||
self._files.append(self._path)
|
||||
self._start_dumping()
|
||||
|
||||
@@ -126,6 +129,20 @@ class DanmakuDumper(
|
||||
await self._stop_dumping()
|
||||
self._path = None
|
||||
|
||||
async def on_stream_recording_interrupted(self, duration: float) -> None:
|
||||
logger.debug(f'Stream recording interrupted, {duration}')
|
||||
self._duration = duration
|
||||
self._stream_recording_recovered = asyncio.Condition()
|
||||
self._stream_recording_interrupted = True
|
||||
|
||||
async def on_stream_recording_recovered(self, timestamp: int) -> None:
|
||||
logger.debug(f'Stream recording recovered, {timestamp}')
|
||||
self._timebase = timestamp * 1000
|
||||
self._delta = self._duration * 1000
|
||||
self._stream_recording_interrupted = False
|
||||
async with self._stream_recording_recovered:
|
||||
self._stream_recording_recovered.notify_all()
|
||||
|
||||
def _start_dumping(self) -> None:
|
||||
self._create_dump_task()
|
||||
|
||||
@@ -172,6 +189,7 @@ class DanmakuDumper(
|
||||
async def _dumping_loop(self, writer: DanmakuWriter) -> None:
|
||||
while True:
|
||||
msg = await self._receiver.get_message()
|
||||
|
||||
if isinstance(msg, DanmuMsg):
|
||||
await writer.write_danmu(self._make_danmu(msg))
|
||||
self._statistics.submit(1)
|
||||
@@ -193,6 +211,13 @@ class DanmakuDumper(
|
||||
else:
|
||||
logger.warning('Unsupported message type:', repr(msg))
|
||||
|
||||
if self._stream_recording_interrupted:
|
||||
logger.debug(
|
||||
f'Last message before stream recording interrupted: {repr(msg)}'
|
||||
)
|
||||
async with self._stream_recording_recovered:
|
||||
await self._stream_recording_recovered.wait()
|
||||
|
||||
def _make_metadata(self) -> Metadata:
|
||||
return Metadata(
|
||||
user_name=self._live.user_info.name,
|
||||
@@ -259,4 +284,4 @@ class DanmakuDumper(
|
||||
)
|
||||
|
||||
def _calc_stime(self, timestamp: int) -> float:
|
||||
return max((timestamp - self._record_start_time * 1000), 0) / 1000
|
||||
return (max(timestamp - self._timebase, 0) + self._delta) / 1000
|
||||
|
||||
@@ -6,6 +6,7 @@ from reactivex.scheduler import NewThreadScheduler
|
||||
from ..bili.live import Live
|
||||
from ..bili.typing import QualityNumber
|
||||
from ..flv import operators as flv_ops
|
||||
from ..utils.mixins import SupportDebugMixin
|
||||
from .stream_recorder_impl import StreamRecorderImpl
|
||||
|
||||
__all__ = ('FLVStreamRecorderImpl',)
|
||||
@@ -14,7 +15,7 @@ __all__ = ('FLVStreamRecorderImpl',)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class FLVStreamRecorderImpl(StreamRecorderImpl):
|
||||
class FLVStreamRecorderImpl(StreamRecorderImpl, SupportDebugMixin):
|
||||
def __init__(
|
||||
self,
|
||||
live: Live,
|
||||
@@ -40,6 +41,7 @@ class FLVStreamRecorderImpl(StreamRecorderImpl):
|
||||
filesize_limit=filesize_limit,
|
||||
duration_limit=duration_limit,
|
||||
)
|
||||
self._init_for_debug(live.room_id)
|
||||
|
||||
def _run(self) -> None:
|
||||
self._subscription = (
|
||||
@@ -47,11 +49,12 @@ class FLVStreamRecorderImpl(StreamRecorderImpl):
|
||||
.pipe(
|
||||
self._stream_url_resolver,
|
||||
self._stream_fetcher,
|
||||
self._recording_monitor,
|
||||
self._dl_statistics,
|
||||
self._stream_parser,
|
||||
self._request_exception_handler,
|
||||
self._connection_error_handler,
|
||||
flv_ops.process(),
|
||||
self._request_exception_handler,
|
||||
flv_ops.process(sort_tags=True, trace=self._debug),
|
||||
self._cutter,
|
||||
self._limiter,
|
||||
self._join_point_extractor,
|
||||
|
||||
@@ -57,8 +57,9 @@ class HLSStreamRecorderImpl(StreamRecorderImpl):
|
||||
NewThreadScheduler(self._thread_factory('PlaylistFetcher'))
|
||||
),
|
||||
self._playlist_fetcher,
|
||||
self._request_exception_handler,
|
||||
self._recording_monitor,
|
||||
self._connection_error_handler,
|
||||
self._request_exception_handler,
|
||||
self._playlist_resolver,
|
||||
ops.observe_on(
|
||||
NewThreadScheduler(self._thread_factory('SegmentFetcher'))
|
||||
|
||||
@@ -4,6 +4,7 @@ from .hls_prober import HLSProber, StreamProfile
|
||||
from .playlist_fetcher import PlaylistFetcher
|
||||
from .playlist_resolver import PlaylistResolver
|
||||
from .progress_bar import ProgressBar
|
||||
from .recording_monitor import RecordingMonitor
|
||||
from .request_exception_handler import RequestExceptionHandler
|
||||
from .segment_fetcher import InitSectionData, SegmentData, SegmentFetcher
|
||||
from .segment_remuxer import SegmentRemuxer
|
||||
@@ -21,6 +22,7 @@ __all__ = (
|
||||
'PlaylistFetcher',
|
||||
'PlaylistResolver',
|
||||
'ProgressBar',
|
||||
'RecordingMonitor',
|
||||
'RequestExceptionHandler',
|
||||
'SegmentData',
|
||||
'SegmentFetcher',
|
||||
|
||||
@@ -28,6 +28,7 @@ class ExceptionHandler(AsyncCooperationMixin):
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
def on_error(exc: Exception) -> None:
|
||||
logger.exception(repr(exc))
|
||||
self._submit_exception(exc)
|
||||
try:
|
||||
raise exc
|
||||
|
||||
@@ -5,6 +5,7 @@ import logging
|
||||
from typing import List, Optional, Union
|
||||
|
||||
from reactivex import Observable, Subject, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ...utils.ffprobe import StreamProfile, ffprobe
|
||||
from .segment_fetcher import InitSectionData, SegmentData
|
||||
@@ -39,6 +40,9 @@ class HLSProber:
|
||||
observer: abc.ObserverBase[Union[InitSectionData, SegmentData]],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
self._reset()
|
||||
|
||||
def on_next(item: Union[InitSectionData, SegmentData]) -> None:
|
||||
@@ -59,10 +63,17 @@ class HLSProber:
|
||||
|
||||
observer.on_next(item)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
self._reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
def _do_probe(self) -> None:
|
||||
|
||||
@@ -29,6 +29,7 @@ class PlaylistResolver:
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
last_seg_uris: OrderedSet[str] = OrderedSet()
|
||||
|
||||
def on_next(playlist: m3u8.M3U8) -> None:
|
||||
|
||||
69
src/blrec/core/operators/recording_monitor.py
Normal file
69
src/blrec/core/operators/recording_monitor.py
Normal file
@@ -0,0 +1,69 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Final, Optional, TypeVar
|
||||
|
||||
from reactivex import Observable, Subject, abc
|
||||
|
||||
from ...bili.live import Live
|
||||
from ...flv import operators as flv_ops
|
||||
from ...utils.mixins import AsyncCooperationMixin
|
||||
|
||||
__all__ = ('RecordingMonitor',)
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_T = TypeVar('_T')
|
||||
|
||||
|
||||
class RecordingMonitor(AsyncCooperationMixin):
|
||||
def __init__(self, live: Live, analyser: flv_ops.Analyser) -> None:
|
||||
super().__init__()
|
||||
self._live = live
|
||||
self._analyser = analyser
|
||||
self._interrupted: Subject[float] = Subject()
|
||||
self._recovered: Subject[int] = Subject()
|
||||
|
||||
@property
|
||||
def interrupted(self) -> Observable[float]:
|
||||
return self._interrupted
|
||||
|
||||
@property
|
||||
def recovered(self) -> Observable[int]:
|
||||
return self._recovered
|
||||
|
||||
def __call__(self, source: Observable[_T]) -> Observable[_T]:
|
||||
return self._monitor(source)
|
||||
|
||||
def _monitor(self, source: Observable[_T]) -> Observable[_T]:
|
||||
CRITERIA: Final[int] = 1
|
||||
recording: bool = False
|
||||
failed_count: int = 0
|
||||
|
||||
def subscribe(
|
||||
observer: abc.ObserverBase[_T],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
def on_next(item: _T) -> None:
|
||||
nonlocal recording, failed_count
|
||||
recording = True
|
||||
if failed_count >= CRITERIA:
|
||||
ts = self._run_coroutine(self._live.get_timestamp())
|
||||
self._recovered.on_next(ts)
|
||||
failed_count = 0
|
||||
observer.on_next(item)
|
||||
|
||||
def on_error(exc: Exception) -> None:
|
||||
nonlocal failed_count
|
||||
if recording:
|
||||
failed_count += 1
|
||||
if failed_count == CRITERIA:
|
||||
self._interrupted.on_next(self._analyser.duration)
|
||||
observer.on_error(exc)
|
||||
|
||||
return source.subscribe(
|
||||
on_next, on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return Observable(subscribe)
|
||||
@@ -2,6 +2,7 @@ from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
import time
|
||||
from typing import Optional, TypeVar
|
||||
|
||||
import requests
|
||||
@@ -19,9 +20,12 @@ _T = TypeVar('_T')
|
||||
|
||||
|
||||
class RequestExceptionHandler:
|
||||
def __init__(self) -> None:
|
||||
self._last_retry_time = time.monotonic()
|
||||
|
||||
def __call__(self, source: Observable[_T]) -> Observable[_T]:
|
||||
return self._handle(source).pipe(
|
||||
utils_ops.retry(delay=1, should_retry=self._should_retry)
|
||||
utils_ops.retry(should_retry=self._should_retry)
|
||||
)
|
||||
|
||||
def _handle(self, source: Observable[_T]) -> Observable[_T]:
|
||||
@@ -32,19 +36,20 @@ class RequestExceptionHandler:
|
||||
def on_error(exc: Exception) -> None:
|
||||
try:
|
||||
raise exc
|
||||
except requests.exceptions.RequestException: # XXX: ConnectionError
|
||||
logger.warning(repr(exc))
|
||||
except urllib3.exceptions.HTTPError:
|
||||
logger.warning(repr(exc))
|
||||
except asyncio.exceptions.TimeoutError:
|
||||
logger.warning(repr(exc))
|
||||
except requests.exceptions.Timeout:
|
||||
logger.warning(repr(exc))
|
||||
except requests.exceptions.HTTPError:
|
||||
logger.warning(repr(exc))
|
||||
except urllib3.exceptions.TimeoutError:
|
||||
logger.warning(repr(exc))
|
||||
except urllib3.exceptions.ProtocolError:
|
||||
# ProtocolError('Connection broken: IncompleteRead(
|
||||
logger.warning(repr(exc))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if self._should_retry(exc):
|
||||
if time.monotonic() - self._last_retry_time < 1:
|
||||
time.sleep(1)
|
||||
self._last_retry_time = time.monotonic()
|
||||
|
||||
observer.on_error(exc)
|
||||
|
||||
return source.subscribe(
|
||||
@@ -57,11 +62,9 @@ class RequestExceptionHandler:
|
||||
if isinstance(
|
||||
exc,
|
||||
(
|
||||
requests.exceptions.RequestException, # XXX: ConnectionError
|
||||
urllib3.exceptions.HTTPError,
|
||||
asyncio.exceptions.TimeoutError,
|
||||
requests.exceptions.Timeout,
|
||||
requests.exceptions.HTTPError,
|
||||
urllib3.exceptions.TimeoutError,
|
||||
urllib3.exceptions.ProtocolError,
|
||||
),
|
||||
):
|
||||
return True
|
||||
|
||||
@@ -55,6 +55,7 @@ class SegmentFetcher:
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
init_section: Optional[InitializationSection] = None
|
||||
|
||||
def on_next(seg: m3u8.Segment) -> None:
|
||||
@@ -76,7 +77,9 @@ class SegmentFetcher:
|
||||
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
nonlocal init_section
|
||||
disposed = True
|
||||
init_section = None
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
|
||||
@@ -21,7 +21,8 @@ logging.getLogger(urllib3.__name__).setLevel(logging.WARNING)
|
||||
|
||||
|
||||
class SegmentRemuxer:
|
||||
_MAX_SEGMENT_DATA_CACHE: Final = 3
|
||||
_SEGMENT_DATA_CACHE: Final = 10
|
||||
_MAX_SEGMENT_DATA_CACHE: Final = 15
|
||||
|
||||
def __init__(self, live: Live) -> None:
|
||||
self._live = live
|
||||
@@ -47,6 +48,12 @@ class SegmentRemuxer:
|
||||
segment_data_cache: List[bytes] = []
|
||||
self._stream_remuxer.stop()
|
||||
|
||||
def reset() -> None:
|
||||
nonlocal init_section_data, segment_data_cache
|
||||
init_section_data = None
|
||||
segment_data_cache = []
|
||||
self._stream_remuxer.stop()
|
||||
|
||||
def write(data: bytes) -> int:
|
||||
return wait_for(
|
||||
self._stream_remuxer.input.write,
|
||||
@@ -56,6 +63,7 @@ class SegmentRemuxer:
|
||||
|
||||
def on_next(data: Union[InitSectionData, SegmentData]) -> None:
|
||||
nonlocal init_section_data
|
||||
nonlocal segment_data_cache
|
||||
|
||||
if isinstance(data, InitSectionData):
|
||||
init_section_data = data.payload
|
||||
@@ -64,8 +72,8 @@ class SegmentRemuxer:
|
||||
|
||||
try:
|
||||
if self._stream_remuxer.stopped:
|
||||
self._stream_remuxer.start()
|
||||
while True:
|
||||
self._stream_remuxer.start()
|
||||
ready = self._stream_remuxer.wait(timeout=1)
|
||||
if disposed:
|
||||
return
|
||||
@@ -86,14 +94,22 @@ class SegmentRemuxer:
|
||||
except Exception as e:
|
||||
logger.warning(f'Failed to write data to stream remuxer: {repr(e)}')
|
||||
self._stream_remuxer.stop()
|
||||
if len(segment_data_cache) >= self._MAX_SEGMENT_DATA_CACHE:
|
||||
segment_data_cache = segment_data_cache[
|
||||
-self._MAX_SEGMENT_DATA_CACHE + 1 :
|
||||
]
|
||||
else:
|
||||
if len(segment_data_cache) >= self._SEGMENT_DATA_CACHE:
|
||||
segment_data_cache = segment_data_cache[
|
||||
-self._SEGMENT_DATA_CACHE + 1 :
|
||||
]
|
||||
|
||||
segment_data_cache.append(data.payload)
|
||||
if len(segment_data_cache) > self._MAX_SEGMENT_DATA_CACHE:
|
||||
segment_data_cache.pop(0)
|
||||
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import time
|
||||
import logging
|
||||
from datetime import datetime
|
||||
from typing import Iterator, Optional
|
||||
|
||||
@@ -254,6 +254,12 @@ class StreamRecorder(
|
||||
async def on_video_file_completed(self, path: str) -> None:
|
||||
await self._emit('video_file_completed', path)
|
||||
|
||||
async def on_stream_recording_interrupted(self, timestamp: int) -> None:
|
||||
await self._emit('stream_recording_interrupted', timestamp)
|
||||
|
||||
async def on_stream_recording_recovered(self, timestamp: int) -> None:
|
||||
await self._emit('stream_recording_recovered', timestamp)
|
||||
|
||||
async def on_stream_recording_completed(self) -> None:
|
||||
await self._emit('stream_recording_completed')
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import logging
|
||||
from abc import ABC, abstractmethod
|
||||
from datetime import datetime
|
||||
from threading import Thread
|
||||
from typing import Any, Iterator, List, Optional, Tuple, Union
|
||||
|
||||
@@ -14,6 +15,7 @@ from ..event.event_emitter import EventEmitter, EventListener
|
||||
from ..flv import operators as flv_ops
|
||||
from ..flv.metadata_dumper import MetadataDumper
|
||||
from ..flv.operators import StreamProfile
|
||||
from ..flv.utils import format_timestamp
|
||||
from ..logging.room_id import aio_task_with_room_id
|
||||
from ..utils.mixins import AsyncCooperationMixin, AsyncStoppableMixin
|
||||
from . import operators as core_ops
|
||||
@@ -35,6 +37,12 @@ class StreamRecorderEventListener(EventListener):
|
||||
async def on_video_file_completed(self, path: str) -> None:
|
||||
...
|
||||
|
||||
async def on_stream_recording_interrupted(self, duratin: float) -> None:
|
||||
...
|
||||
|
||||
async def on_stream_recording_recovered(self, timestamp: int) -> None:
|
||||
...
|
||||
|
||||
async def on_stream_recording_completed(self) -> None:
|
||||
...
|
||||
|
||||
@@ -87,6 +95,7 @@ class StreamRecorderImpl(
|
||||
self._path_provider = PathProvider(live, out_dir, path_template)
|
||||
self._dumper = flv_ops.Dumper(self._path_provider, buffer_size)
|
||||
self._rec_statistics = core_ops.SizedStatistics()
|
||||
self._recording_monitor = core_ops.RecordingMonitor(live, self._analyser)
|
||||
|
||||
self._prober: Union[flv_ops.Prober, core_ops.HLSProber]
|
||||
self._dl_statistics: Union[core_ops.StreamStatistics, core_ops.SizedStatistics]
|
||||
@@ -135,6 +144,19 @@ class StreamRecorderImpl(
|
||||
self._dumper.file_opened.subscribe(on_file_opened)
|
||||
self._dumper.file_closed.subscribe(on_file_closed)
|
||||
|
||||
def on_recording_interrupted(duration: float) -> None:
|
||||
duration_string = format_timestamp(int(duration * 1000))
|
||||
logger.info(f'Recording interrupted, current duration: {duration_string}')
|
||||
self._emit_event('stream_recording_interrupted', duration)
|
||||
|
||||
def on_recording_recovered(timestamp: int) -> None:
|
||||
datetime_string = datetime.fromtimestamp(timestamp).isoformat()
|
||||
logger.info(f'Recording recovered, current date time {(datetime_string)}')
|
||||
self._emit_event('stream_recording_recovered', timestamp)
|
||||
|
||||
self._recording_monitor.interrupted.subscribe(on_recording_interrupted)
|
||||
self._recording_monitor.recovered.subscribe(on_recording_recovered)
|
||||
|
||||
@property
|
||||
def stream_url(self) -> str:
|
||||
return self._stream_url_resolver.stream_url
|
||||
@@ -328,6 +350,7 @@ class StreamRecorderImpl(
|
||||
|
||||
def _dispose(self) -> None:
|
||||
self._subscription.dispose()
|
||||
del self._subscription
|
||||
self._on_completed()
|
||||
|
||||
def _on_completed(self) -> None:
|
||||
|
||||
@@ -1,20 +1,18 @@
|
||||
from __future__ import annotations
|
||||
import html
|
||||
|
||||
import asyncio
|
||||
import html
|
||||
import logging
|
||||
import unicodedata
|
||||
from datetime import datetime, timezone, timedelta
|
||||
from typing import AsyncIterator, Final, List, Any
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any, AsyncIterator, Final, List
|
||||
|
||||
from lxml import etree
|
||||
import aiofiles
|
||||
import attr
|
||||
from lxml import etree
|
||||
|
||||
from .models import Danmu, GiftSendRecord, GuardBuyRecord, Metadata, SuperChatRecord
|
||||
from .typing import Element
|
||||
from .models import (
|
||||
Metadata, Danmu, GiftSendRecord, GuardBuyRecord, SuperChatRecord
|
||||
)
|
||||
|
||||
|
||||
__all__ = 'DanmakuReader', 'DanmakuWriter'
|
||||
|
||||
@@ -52,12 +50,16 @@ class DanmakuReader: # TODO rewrite
|
||||
room_title=self._tree.xpath('/i/metadata/room_title')[0].text,
|
||||
area=self._tree.xpath('/i/metadata/area')[0].text,
|
||||
parent_area=self._tree.xpath('/i/metadata/parent_area')[0].text,
|
||||
live_start_time=int(datetime.fromisoformat(
|
||||
self._tree.xpath('/i/metadata/live_start_time')[0].text
|
||||
).timestamp()),
|
||||
record_start_time=int(datetime.fromisoformat(
|
||||
self._tree.xpath('/i/metadata/record_start_time')[0].text
|
||||
).timestamp()),
|
||||
live_start_time=int(
|
||||
datetime.fromisoformat(
|
||||
self._tree.xpath('/i/metadata/live_start_time')[0].text
|
||||
).timestamp()
|
||||
),
|
||||
record_start_time=int(
|
||||
datetime.fromisoformat(
|
||||
self._tree.xpath('/i/metadata/record_start_time')[0].text
|
||||
).timestamp()
|
||||
),
|
||||
recorder=self._tree.xpath('/i/metadata/recorder')[0].text,
|
||||
)
|
||||
|
||||
@@ -83,7 +85,9 @@ class DanmakuReader: # TODO rewrite
|
||||
|
||||
|
||||
class DanmakuWriter:
|
||||
_XML_HEAD: Final[str] = """\
|
||||
_XML_HEAD: Final[
|
||||
str
|
||||
] = """\
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<i>
|
||||
<chatserver>chat.bilibili.com</chatserver>
|
||||
@@ -191,10 +195,7 @@ class DanmakuWriter:
|
||||
value_serializer=record_value_serializer,
|
||||
)
|
||||
elem = etree.Element('sc', attrib=attrib)
|
||||
try:
|
||||
elem.text = record.message
|
||||
except ValueError:
|
||||
elem.text = remove_control_characters(record.message)
|
||||
elem.text = record.message
|
||||
return ' ' + etree.tostring(elem, encoding='utf8').decode() + '\n'
|
||||
|
||||
|
||||
@@ -206,8 +207,8 @@ def record_value_serializer(
|
||||
if attribute.name == 'cointype':
|
||||
return '金瓜子' if value == 'gold' else '银瓜子'
|
||||
if not isinstance(value, str):
|
||||
return str(value)
|
||||
return value
|
||||
value = str(value)
|
||||
return remove_control_characters(value)
|
||||
|
||||
|
||||
def remove_control_characters(s: str) -> str:
|
||||
|
||||
@@ -9,7 +9,16 @@ from . import scriptdata
|
||||
from .avc import extract_resolution
|
||||
from .io import FlvReader
|
||||
from .io_protocols import RandomIO
|
||||
from .models import AudioTag, AVCPacketType, FlvTag, ScriptTag, TagType, VideoTag
|
||||
from .models import (
|
||||
AudioTag,
|
||||
AVCPacketType,
|
||||
CodecID,
|
||||
FlvTag,
|
||||
FrameType,
|
||||
ScriptTag,
|
||||
TagType,
|
||||
VideoTag,
|
||||
)
|
||||
from .utils import OffsetRepositor
|
||||
|
||||
|
||||
@@ -155,8 +164,31 @@ def is_video_nalu_keyframe(tag: FlvTag) -> TypeGuard[VideoTag]:
|
||||
return is_video_tag(tag) and tag.is_keyframe() and tag.is_avc_nalu()
|
||||
|
||||
|
||||
def is_avc_end_sequence(tag: FlvTag) -> TypeGuard[VideoTag]:
|
||||
return is_video_tag(tag) and tag.is_avc_end()
|
||||
|
||||
|
||||
def is_avc_end_sequence_tag(value: Any) -> TypeGuard[VideoTag]:
|
||||
return isinstance(value, FlvTag) and is_avc_end_sequence(value)
|
||||
|
||||
|
||||
def create_avc_end_sequence_tag(offset: int = 0, timestamp: int = 0) -> VideoTag:
|
||||
return VideoTag(
|
||||
offset=offset,
|
||||
filtered=False,
|
||||
tag_type=TagType.VIDEO,
|
||||
data_size=5,
|
||||
timestamp=timestamp,
|
||||
stream_id=timestamp,
|
||||
frame_type=FrameType.KEY_FRAME,
|
||||
codec_id=CodecID.AVC,
|
||||
avc_packet_type=AVCPacketType.AVC_END_OF_SEQENCE,
|
||||
composition_time=0,
|
||||
)
|
||||
|
||||
|
||||
def parse_scriptdata(script_tag: ScriptTag) -> scriptdata.ScriptData:
|
||||
assert script_tag.body is not None
|
||||
assert script_tag.body
|
||||
return scriptdata.load(script_tag.body)
|
||||
|
||||
|
||||
@@ -253,8 +285,8 @@ class Resolution:
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def from_aac_sequence_header(cls, tag: VideoTag) -> Resolution:
|
||||
def from_avc_sequence_header(cls, tag: VideoTag) -> Resolution:
|
||||
assert tag.avc_packet_type == AVCPacketType.AVC_SEQUENCE_HEADER
|
||||
assert tag.body is not None
|
||||
assert tag.body
|
||||
width, height = extract_resolution(tag.body)
|
||||
return cls(width, height)
|
||||
|
||||
@@ -197,6 +197,9 @@ class FlvDumper:
|
||||
self._writer.write_ui32(size)
|
||||
|
||||
def dump_tag(self, tag: FlvTag) -> None:
|
||||
if tag.timestamp < 0:
|
||||
raise FlvDataError(f'Incorrect timestamp: {tag.timestamp}', tag)
|
||||
|
||||
self.dump_flv_tag_header(tag)
|
||||
|
||||
if tag.is_audio_tag():
|
||||
@@ -210,10 +213,10 @@ class FlvDumper:
|
||||
else:
|
||||
raise FlvDataError(f'Unsupported tag type: {tag.tag_type}')
|
||||
|
||||
if tag.body is None:
|
||||
self._stream.seek(tag.tag_end_offset)
|
||||
else:
|
||||
if tag.body:
|
||||
self._writer.write(tag.body)
|
||||
else:
|
||||
self._stream.seek(tag.tag_end_offset)
|
||||
|
||||
def dump_flv_tag_header(self, tag: FlvTag) -> None:
|
||||
self._writer.write_ui8((int(tag.filtered) << 5) | tag.tag_type.value)
|
||||
|
||||
@@ -16,7 +16,7 @@ def get_metadata(path: str) -> Dict[str, Any]:
|
||||
with open(path, mode='rb') as file:
|
||||
reader = FlvReader(file)
|
||||
reader.read_header()
|
||||
if (tag := find_metadata_tag(reversed(list(read_tags(reader, 5))))):
|
||||
if (tag := find_metadata_tag(list(read_tags(reader, 5)))):
|
||||
return parse_metadata(tag)
|
||||
raise EOFError
|
||||
|
||||
|
||||
56
src/blrec/flv/metadata_analysis.py
Normal file
56
src/blrec/flv/metadata_analysis.py
Normal file
@@ -0,0 +1,56 @@
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
|
||||
import attr
|
||||
from reactivex import Observable
|
||||
from reactivex import operators as ops
|
||||
|
||||
from ..path import extra_metadata_path
|
||||
from . import operators as flv_ops
|
||||
from .operators.helpers import from_file
|
||||
|
||||
__all__ = 'AnalysingProgress', 'analyse_metadata'
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@attr.s(auto_attribs=True, slots=True, frozen=True)
|
||||
class AnalysingProgress:
|
||||
count: int
|
||||
total: int
|
||||
|
||||
|
||||
def analyse_metadata(
|
||||
path: str, *, show_progress: bool = False
|
||||
) -> Observable[AnalysingProgress]:
|
||||
filesize = os.path.getsize(path)
|
||||
filename = os.path.basename(path)
|
||||
|
||||
def dump_metadata() -> None:
|
||||
try:
|
||||
metadata = analyser.make_metadata()
|
||||
data = attr.asdict(metadata, filter=lambda a, v: v is not None)
|
||||
file_path = extra_metadata_path(path)
|
||||
with open(file_path, 'wt', encoding='utf8') as file:
|
||||
json.dump(data, file)
|
||||
except Exception as e:
|
||||
logger.error(f'Failed to dump metadata: {e}')
|
||||
else:
|
||||
logger.debug(f"Successfully dumped metadata to file: '{path}'")
|
||||
|
||||
analyser = flv_ops.Analyser()
|
||||
return from_file(path).pipe(
|
||||
analyser,
|
||||
flv_ops.ProgressBar(
|
||||
desc='Analysing',
|
||||
postfix=filename,
|
||||
total=filesize,
|
||||
disable=not show_progress,
|
||||
),
|
||||
ops.map(lambda i: len(i)), # type: ignore
|
||||
ops.scan(lambda acc, x: acc + x, 0), # type: ignore
|
||||
ops.map(lambda s: AnalysingProgress(s, filesize)), # type: ignore
|
||||
ops.do_action(on_completed=dump_metadata),
|
||||
)
|
||||
@@ -22,11 +22,12 @@ class MetadataDumper(SwitchableMixin):
|
||||
joinpoint_extractor: flv_ops.JoinPointExtractor,
|
||||
) -> None:
|
||||
super().__init__()
|
||||
|
||||
self._dumper = dumper
|
||||
self._analyser = analyser
|
||||
self._joinpoint_extractor = joinpoint_extractor
|
||||
self._reset()
|
||||
|
||||
def _reset(self) -> None:
|
||||
self._last_metadata: Optional[flv_ops.MetaData] = None
|
||||
self._last_join_points: Optional[List[flv_ops.JoinPoint]] = None
|
||||
|
||||
@@ -40,6 +41,7 @@ class MetadataDumper(SwitchableMixin):
|
||||
self._file_closed_subscription = self._dumper.file_closed.subscribe(
|
||||
self._dump_metadata
|
||||
)
|
||||
self._reset()
|
||||
logger.debug('Enabled metadata dumper')
|
||||
|
||||
def _do_disable(self) -> None:
|
||||
@@ -49,6 +51,7 @@ class MetadataDumper(SwitchableMixin):
|
||||
self._join_points_subscription.dispose()
|
||||
with suppress(Exception):
|
||||
self._file_closed_subscription.dispose()
|
||||
self._reset()
|
||||
logger.debug('Disabled metadata dumper')
|
||||
|
||||
def _update_metadata(self, metadata: Optional[flv_ops.MetaData]) -> None:
|
||||
@@ -61,13 +64,19 @@ class MetadataDumper(SwitchableMixin):
|
||||
path = extra_metadata_path(video_path)
|
||||
logger.debug(f"Dumping metadata to file: '{path}'")
|
||||
|
||||
assert self._last_metadata is not None
|
||||
assert self._last_join_points is not None
|
||||
if self._last_metadata is not None:
|
||||
data = attr.asdict(self._last_metadata, filter=lambda a, v: v is not None)
|
||||
else:
|
||||
data = {}
|
||||
logger.warning('The metadata may be lost duo to something went wrong')
|
||||
|
||||
data = attr.asdict(self._last_metadata, filter=lambda a, v: v is not None)
|
||||
data['joinpoints'] = list(
|
||||
map(lambda p: p.to_metadata_value(), self._last_join_points)
|
||||
)
|
||||
if self._last_join_points is not None:
|
||||
data['joinpoints'] = list(
|
||||
map(lambda p: p.to_metadata_value(), self._last_join_points)
|
||||
)
|
||||
else:
|
||||
data['joinpoints'] = []
|
||||
logger.warning('The joinpoints may be lost duo to something went wrong')
|
||||
|
||||
try:
|
||||
with open(path, 'wt', encoding='utf8') as file:
|
||||
|
||||
@@ -152,7 +152,7 @@ _T = TypeVar('_T', bound='FlvTag')
|
||||
@attr.s(auto_attribs=True, slots=True, frozen=True, kw_only=True)
|
||||
class FlvTag(ABC, FlvTagHeader):
|
||||
offset: int = attr.ib(validator=[non_negative_integer_validator])
|
||||
body: Optional[bytes] = attr.ib(default=None, repr=cksum)
|
||||
body: bytes = attr.ib(default=b'', repr=cksum)
|
||||
|
||||
def __len__(self) -> int:
|
||||
return self.tag_size
|
||||
|
||||
@@ -11,6 +11,7 @@ from .parse import parse
|
||||
from .probe import Prober, StreamProfile
|
||||
from .process import process
|
||||
from .progress import ProgressBar
|
||||
from .sort import sort
|
||||
from .split import split
|
||||
|
||||
__all__ = (
|
||||
@@ -33,6 +34,7 @@ __all__ = (
|
||||
'Prober',
|
||||
'process',
|
||||
'ProgressBar',
|
||||
'sort',
|
||||
'split',
|
||||
'StreamProfile',
|
||||
)
|
||||
|
||||
@@ -128,7 +128,11 @@ class Analyser:
|
||||
self._video_analysed = False
|
||||
|
||||
@property
|
||||
def metadatas(self) -> Observable[MetaData]:
|
||||
def duration(self) -> float:
|
||||
return self._last_timestamp / 1000
|
||||
|
||||
@property
|
||||
def metadatas(self) -> Observable[Optional[MetaData]]:
|
||||
return self._metadatas
|
||||
|
||||
def __call__(self, source: FLVStream) -> FLVStream:
|
||||
@@ -237,9 +241,12 @@ class Analyser:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
stream_index: int = -1
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
self._reset()
|
||||
stream_index: int = -1
|
||||
|
||||
def push_metadata() -> None:
|
||||
try:
|
||||
metadata = self.make_metadata()
|
||||
@@ -270,7 +277,10 @@ class Analyser:
|
||||
observer.on_error(e)
|
||||
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
push_metadata()
|
||||
self._reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, on_error, on_completed, scheduler=scheduler
|
||||
@@ -318,7 +328,7 @@ class Analyser:
|
||||
self._keyframe_timestamps.append(tag.timestamp)
|
||||
self._keyframe_filepositions.append(self.calc_file_size())
|
||||
if tag.is_avc_header():
|
||||
self._resolution = Resolution.from_aac_sequence_header(tag)
|
||||
self._resolution = Resolution.from_avc_sequence_header(tag)
|
||||
logger.debug(f'Resolution: {self._resolution}')
|
||||
else:
|
||||
pass
|
||||
|
||||
@@ -85,6 +85,9 @@ def concat(
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
delta: int = 0
|
||||
action: ACTION = ACTION.NOOP
|
||||
last_tags: List[FlvTag] = []
|
||||
@@ -93,6 +96,22 @@ def concat(
|
||||
last_audio_sequence_header: Optional[AudioTag] = None
|
||||
last_video_sequence_header: Optional[VideoTag] = None
|
||||
|
||||
def reset() -> None:
|
||||
nonlocal delta
|
||||
nonlocal action
|
||||
nonlocal last_tags
|
||||
nonlocal gathered_tags
|
||||
nonlocal last_flv_header
|
||||
nonlocal last_audio_sequence_header
|
||||
nonlocal last_video_sequence_header
|
||||
delta = 0
|
||||
action = ACTION.NOOP
|
||||
last_tags = []
|
||||
gathered_tags = []
|
||||
last_flv_header = None
|
||||
last_audio_sequence_header = None
|
||||
last_video_sequence_header = None
|
||||
|
||||
def update_last_tags(tag: FlvTag) -> None:
|
||||
nonlocal last_audio_sequence_header, last_video_sequence_header
|
||||
last_tags.append(tag)
|
||||
@@ -182,7 +201,7 @@ def concat(
|
||||
return tag.evolve(timestamp=tag.timestamp + delta)
|
||||
|
||||
def make_join_point_tag(next_tag: FlvTag, seamless: bool) -> ScriptTag:
|
||||
assert next_tag.body is not None
|
||||
assert next_tag.body
|
||||
join_point = JoinPoint(
|
||||
seamless=seamless,
|
||||
timestamp=float(next_tag.timestamp),
|
||||
@@ -315,10 +334,17 @@ def concat(
|
||||
do_concat()
|
||||
observer.on_error(e)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, on_error, on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return _concat
|
||||
@@ -340,11 +366,19 @@ class JoinPointExtractor:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
stream_index: int = -1
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
stream_index: int = -1
|
||||
join_points: List[JoinPoint] = []
|
||||
join_point_tag: Optional[ScriptTag] = None
|
||||
|
||||
def reset() -> None:
|
||||
nonlocal stream_index, join_points, join_point_tag
|
||||
stream_index = -1
|
||||
join_points = []
|
||||
join_point_tag = None
|
||||
|
||||
def push_join_points() -> None:
|
||||
self._join_points.on_next(join_points.copy())
|
||||
|
||||
@@ -381,7 +415,10 @@ class JoinPointExtractor:
|
||||
observer.on_error(e)
|
||||
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
push_join_points()
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, on_error, on_completed, scheduler=scheduler
|
||||
@@ -402,11 +439,18 @@ class JoinPointExtractor:
|
||||
) -> JoinPoint:
|
||||
script_data = parse_scriptdata(join_point_tag)
|
||||
join_point_data = cast(JoinPointData, script_data['value'])
|
||||
assert next_tag.body is not None, next_tag
|
||||
assert next_tag.body, next_tag
|
||||
join_point = JoinPoint(
|
||||
seamless=join_point_data['seamless'],
|
||||
timestamp=next_tag.timestamp,
|
||||
crc32=join_point_data['crc32'],
|
||||
)
|
||||
logger.debug(f'Extracted join point: {join_point}; next tag: {next_tag}')
|
||||
if cksum(next_tag.body) != join_point_data['crc32']:
|
||||
logger.warning(
|
||||
f'Timestamp of extracted join point may be incorrect\n'
|
||||
f'join point data: {join_point_data}\n'
|
||||
f'join point tag: {join_point_tag}\n'
|
||||
f'next tag: {next_tag}\n'
|
||||
)
|
||||
return join_point
|
||||
|
||||
@@ -2,6 +2,7 @@ import logging
|
||||
from typing import Callable, Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import is_script_tag, is_sequence_header
|
||||
from ..models import FlvHeader, FlvTag
|
||||
@@ -20,9 +21,17 @@ def correct() -> Callable[[FLVStream], FLVStream]:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
delta: Optional[int] = None
|
||||
first_data_tag: Optional[FlvTag] = None
|
||||
|
||||
def reset() -> None:
|
||||
nonlocal delta, first_data_tag
|
||||
delta = None
|
||||
first_data_tag = None
|
||||
|
||||
def correct_ts(tag: FlvTag, delta: int) -> FlvTag:
|
||||
if delta == 0:
|
||||
return tag
|
||||
@@ -72,10 +81,17 @@ def correct() -> Callable[[FLVStream], FLVStream]:
|
||||
tag = correct_ts(tag, delta)
|
||||
observer.on_next(tag)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return _correct
|
||||
|
||||
@@ -4,6 +4,7 @@ import logging
|
||||
from typing import Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import (
|
||||
is_audio_sequence_header,
|
||||
@@ -58,6 +59,9 @@ class Cutter:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
def on_next(item: FLVStreamItem) -> None:
|
||||
if isinstance(item, FlvHeader):
|
||||
self._reset()
|
||||
@@ -71,10 +75,17 @@ class Cutter:
|
||||
self._triggered = False
|
||||
observer.on_next(item)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
self._reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
def _update_flv_header(self, header: FlvHeader) -> None:
|
||||
|
||||
@@ -2,6 +2,7 @@ import logging
|
||||
from typing import Callable, List, Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..models import FlvHeader
|
||||
from .typing import FLVStream, FLVStreamItem
|
||||
@@ -19,9 +20,17 @@ def defragment(min_tags: int = 10) -> Callable[[FLVStream], FLVStream]:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
gathering: bool = False
|
||||
gathered_items: List[FLVStreamItem] = []
|
||||
|
||||
def reset() -> None:
|
||||
nonlocal gathering, gathered_items
|
||||
gathering = False
|
||||
gathered_items = []
|
||||
|
||||
def on_next(item: FLVStreamItem) -> None:
|
||||
nonlocal gathering
|
||||
|
||||
@@ -47,10 +56,17 @@ def defragment(min_tags: int = 10) -> Callable[[FLVStream], FLVStream]:
|
||||
else:
|
||||
observer.on_next(item)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return _defragment
|
||||
|
||||
@@ -90,9 +90,6 @@ class Dumper:
|
||||
self._timestamp_updates.on_next(0)
|
||||
else:
|
||||
if self._flv_writer is not None:
|
||||
# XXX: negative timestamp will cause
|
||||
# `struct.error: ubyte format requires 0 <= number <= 255`
|
||||
assert item.timestamp >= 0, item
|
||||
size = self._flv_writer.write_tag(item)
|
||||
self._size_updates.on_next(size)
|
||||
self._timestamp_updates.on_next(item.timestamp)
|
||||
|
||||
@@ -3,6 +3,7 @@ import math
|
||||
from typing import Callable, Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import (
|
||||
is_audio_tag,
|
||||
@@ -27,6 +28,9 @@ def fix() -> Callable[[FLVStream], FLVStream]:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
delta: int = 0
|
||||
last_tag: Optional[FlvTag] = None
|
||||
last_audio_tag: Optional[AudioTag] = None
|
||||
@@ -73,40 +77,59 @@ def fix() -> Callable[[FLVStream], FLVStream]:
|
||||
|
||||
def update_delta(tag: FlvTag) -> None:
|
||||
nonlocal delta
|
||||
|
||||
if is_video_tag(tag):
|
||||
assert last_video_tag is not None
|
||||
delta = (
|
||||
last_video_tag.timestamp - tag.timestamp + video_frame_interval
|
||||
)
|
||||
elif is_audio_tag(tag):
|
||||
assert last_audio_tag is not None
|
||||
delta = (
|
||||
last_audio_tag.timestamp - tag.timestamp + sound_sample_interval
|
||||
)
|
||||
|
||||
assert last_tag is not None
|
||||
delta = last_tag.timestamp + delta - tag.timestamp + calc_interval(tag)
|
||||
if tag.timestamp + delta <= last_tag.timestamp:
|
||||
if is_video_tag(tag):
|
||||
delta = (
|
||||
last_tag.timestamp - tag.timestamp + video_frame_interval
|
||||
)
|
||||
elif is_audio_tag(tag):
|
||||
delta = (
|
||||
last_tag.timestamp - tag.timestamp + sound_sample_interval
|
||||
)
|
||||
|
||||
def correct_ts(tag: FlvTag) -> FlvTag:
|
||||
if delta == 0:
|
||||
return tag
|
||||
return tag.evolve(timestamp=tag.timestamp + delta)
|
||||
|
||||
def calc_interval(tag: FlvTag) -> int:
|
||||
if is_audio_tag(tag):
|
||||
return sound_sample_interval
|
||||
elif is_video_tag(tag):
|
||||
return video_frame_interval
|
||||
else:
|
||||
logger.warning(f'Unexpected tag type: {tag}')
|
||||
return min(sound_sample_interval, video_frame_interval)
|
||||
|
||||
def is_ts_rebounded(tag: FlvTag) -> bool:
|
||||
if is_audio_tag(tag):
|
||||
if last_audio_tag is None:
|
||||
return False
|
||||
return tag.timestamp < last_audio_tag.timestamp
|
||||
if last_audio_tag.is_aac_header():
|
||||
return tag.timestamp + delta < last_audio_tag.timestamp
|
||||
else:
|
||||
return tag.timestamp + delta <= last_audio_tag.timestamp
|
||||
elif is_video_tag(tag):
|
||||
if last_video_tag is None:
|
||||
return False
|
||||
return tag.timestamp < last_video_tag.timestamp
|
||||
if last_video_tag.is_avc_header():
|
||||
return tag.timestamp + delta < last_video_tag.timestamp
|
||||
else:
|
||||
return tag.timestamp + delta <= last_video_tag.timestamp
|
||||
else:
|
||||
return False
|
||||
|
||||
def is_ts_incontinuous(tag: FlvTag) -> bool:
|
||||
tolerance = 1
|
||||
if last_tag is None:
|
||||
return False
|
||||
return tag.timestamp - last_tag.timestamp > max(
|
||||
sound_sample_interval, video_frame_interval
|
||||
return (
|
||||
tag.timestamp + delta - last_tag.timestamp
|
||||
> max(sound_sample_interval, video_frame_interval) + tolerance
|
||||
)
|
||||
|
||||
def on_next(item: FLVStreamItem) -> None:
|
||||
@@ -127,8 +150,9 @@ def fix() -> Callable[[FLVStream], FLVStream]:
|
||||
update_delta(tag)
|
||||
logger.warning(
|
||||
f'Timestamp rebounded, updated delta: {delta}\n'
|
||||
f'last audio tag: {last_audio_tag}\n'
|
||||
f'last tag: {last_tag}\n'
|
||||
f'last video tag: {last_video_tag}\n'
|
||||
f'last audio tag: {last_audio_tag}\n'
|
||||
f'current tag: {tag}'
|
||||
)
|
||||
elif is_ts_incontinuous(tag):
|
||||
@@ -136,17 +160,26 @@ def fix() -> Callable[[FLVStream], FLVStream]:
|
||||
logger.warning(
|
||||
f'Timestamp incontinuous, updated delta: {delta}\n'
|
||||
f'last tag: {last_tag}\n'
|
||||
f'last video tag: {last_video_tag}\n'
|
||||
f'last audio tag: {last_audio_tag}\n'
|
||||
f'current tag: {tag}'
|
||||
)
|
||||
|
||||
update_last_tags(tag)
|
||||
tag = correct_ts(tag)
|
||||
update_last_tags(tag)
|
||||
observer.on_next(tag)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return _fix
|
||||
|
||||
@@ -4,6 +4,7 @@ import logging
|
||||
from typing import Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import (
|
||||
is_audio_sequence_header,
|
||||
@@ -51,6 +52,9 @@ class Limiter:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
def on_next(item: FLVStreamItem) -> None:
|
||||
if isinstance(item, FlvHeader):
|
||||
self._reset()
|
||||
@@ -62,10 +66,17 @@ class Limiter:
|
||||
self._insert_header_and_tags(observer)
|
||||
observer.on_next(item)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
self._reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
def _insert_header_and_tags(
|
||||
|
||||
@@ -5,7 +5,9 @@ from typing import Callable, Optional
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import create_avc_end_sequence_tag, is_avc_end_sequence
|
||||
from ..io import FlvReader
|
||||
from ..models import FlvTag
|
||||
from .typing import FLVStream, FLVStreamItem
|
||||
|
||||
__all__ = ('parse',)
|
||||
@@ -30,6 +32,7 @@ def parse(
|
||||
subscription = SerialDisposable()
|
||||
|
||||
def on_next(stream: io.RawIOBase) -> None:
|
||||
tag: Optional[FlvTag] = None
|
||||
try:
|
||||
try:
|
||||
reader = FlvReader(
|
||||
@@ -42,6 +45,11 @@ def parse(
|
||||
tag = reader.read_tag()
|
||||
observer.on_next(tag)
|
||||
finally:
|
||||
if tag is not None and not is_avc_end_sequence(tag):
|
||||
tag = create_avc_end_sequence_tag(
|
||||
offset=tag.next_tag_offset, timestamp=tag.timestamp
|
||||
)
|
||||
observer.on_next(tag)
|
||||
stream.close()
|
||||
except EOFError as e:
|
||||
if complete_on_eof:
|
||||
|
||||
@@ -5,6 +5,7 @@ import logging
|
||||
from typing import List, Optional, cast
|
||||
|
||||
from reactivex import Observable, Subject, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ...utils.ffprobe import StreamProfile, ffprobe
|
||||
from ..common import find_aac_header_tag, find_avc_header_tag
|
||||
@@ -38,6 +39,9 @@ class Prober:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
self._reset()
|
||||
|
||||
def on_next(item: FLVStreamItem) -> None:
|
||||
@@ -58,10 +62,17 @@ class Prober:
|
||||
|
||||
observer.on_next(item)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
self._reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
def _do_probe(self) -> None:
|
||||
|
||||
@@ -1,9 +1,14 @@
|
||||
import logging
|
||||
from typing import Callable
|
||||
|
||||
from reactivex import operators as ops
|
||||
|
||||
from ..common import is_avc_end_sequence_tag
|
||||
from .concat import concat
|
||||
from .correct import correct
|
||||
from .defragment import defragment
|
||||
from .fix import fix
|
||||
from .sort import sort
|
||||
from .split import split
|
||||
from .typing import FLVStream
|
||||
|
||||
@@ -12,8 +17,28 @@ __all__ = ('process',)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def process() -> Callable[[FLVStream], FLVStream]:
|
||||
def process(
|
||||
sort_tags: bool = False, trace: bool = False
|
||||
) -> Callable[[FLVStream], FLVStream]:
|
||||
def _process(source: FLVStream) -> FLVStream:
|
||||
return source.pipe(defragment(), split(), fix(), concat())
|
||||
if sort_tags:
|
||||
return source.pipe(
|
||||
defragment(),
|
||||
split(),
|
||||
sort(trace=trace),
|
||||
ops.filter(lambda v: not is_avc_end_sequence_tag(v)), # type: ignore
|
||||
correct(),
|
||||
fix(),
|
||||
concat(),
|
||||
)
|
||||
else:
|
||||
return source.pipe(
|
||||
defragment(),
|
||||
split(),
|
||||
ops.filter(lambda v: not is_avc_end_sequence_tag(v)), # type: ignore
|
||||
correct(),
|
||||
fix(),
|
||||
concat(),
|
||||
)
|
||||
|
||||
return _process
|
||||
|
||||
134
src/blrec/flv/operators/sort.py
Normal file
134
src/blrec/flv/operators/sort.py
Normal file
@@ -0,0 +1,134 @@
|
||||
import logging
|
||||
from typing import Callable, List, Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import (
|
||||
find_aac_header_tag,
|
||||
find_avc_header_tag,
|
||||
find_metadata_tag,
|
||||
is_audio_tag,
|
||||
is_avc_end_sequence,
|
||||
is_script_tag,
|
||||
is_video_nalu_keyframe,
|
||||
is_video_tag,
|
||||
)
|
||||
from ..models import AudioTag, FlvHeader, FlvTag, ScriptTag, VideoTag
|
||||
from .typing import FLVStream, FLVStreamItem
|
||||
|
||||
__all__ = ('sort',)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def sort(trace: bool = False) -> Callable[[FLVStream], FLVStream]:
|
||||
"Sort tags in GOP by timestamp to ensure subsequent operators work as expected."
|
||||
|
||||
def _sort(source: FLVStream) -> FLVStream:
|
||||
def subscribe(
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
gop_tags: List[FlvTag] = []
|
||||
|
||||
def reset() -> None:
|
||||
nonlocal gop_tags
|
||||
gop_tags = []
|
||||
|
||||
def push_gop_tags() -> None:
|
||||
if not gop_tags:
|
||||
return
|
||||
|
||||
if trace:
|
||||
logger.debug(
|
||||
'Tags in GOP:\n'
|
||||
f'Number of tags: {len(gop_tags)}\n'
|
||||
f'Total size of tags: {sum(map(len, gop_tags))}\n'
|
||||
f'The first tag is {gop_tags[0]}\n'
|
||||
f'The last tag is {gop_tags[-1]}'
|
||||
)
|
||||
|
||||
if len(gop_tags) < 10:
|
||||
avc_header_tag = find_avc_header_tag(gop_tags)
|
||||
aac_header_tag = find_aac_header_tag(gop_tags)
|
||||
if avc_header_tag is not None and aac_header_tag is not None:
|
||||
if (metadata_tag := find_metadata_tag(gop_tags)) is not None:
|
||||
observer.on_next(metadata_tag)
|
||||
observer.on_next(avc_header_tag)
|
||||
observer.on_next(aac_header_tag)
|
||||
gop_tags.clear()
|
||||
return
|
||||
|
||||
script_tags: List[ScriptTag] = []
|
||||
video_tags: List[VideoTag] = []
|
||||
audio_tags: List[AudioTag] = []
|
||||
for tag in gop_tags:
|
||||
if is_video_tag(tag):
|
||||
video_tags.append(tag)
|
||||
elif is_audio_tag(tag):
|
||||
audio_tags.append(tag)
|
||||
elif is_script_tag(tag):
|
||||
script_tags.append(tag)
|
||||
|
||||
sorted_tags: List[FlvTag] = []
|
||||
i = len(audio_tags) - 1
|
||||
for video_tag in reversed(video_tags):
|
||||
sorted_tags.insert(0, video_tag)
|
||||
while i >= 0 and audio_tags[i].timestamp >= video_tag.timestamp:
|
||||
sorted_tags.insert(1, audio_tags[i])
|
||||
i -= 1
|
||||
|
||||
for tag in script_tags:
|
||||
observer.on_next(tag)
|
||||
|
||||
for tag in sorted_tags:
|
||||
observer.on_next(tag)
|
||||
|
||||
gop_tags.clear()
|
||||
|
||||
def on_next(item: FLVStreamItem) -> None:
|
||||
if isinstance(item, FlvHeader) or is_avc_end_sequence(item):
|
||||
push_gop_tags()
|
||||
observer.on_next(item)
|
||||
return
|
||||
|
||||
if is_video_nalu_keyframe(item):
|
||||
push_gop_tags()
|
||||
gop_tags.append(item)
|
||||
else:
|
||||
gop_tags.append(item)
|
||||
|
||||
def on_completed() -> None:
|
||||
push_gop_tags()
|
||||
observer.on_completed()
|
||||
|
||||
def on_error(exc: Exception) -> None:
|
||||
push_gop_tags()
|
||||
observer.on_error(exc)
|
||||
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
if gop_tags:
|
||||
logger.debug(
|
||||
'Remaining tags:\n'
|
||||
f'Number of tags: {len(gop_tags)}\n'
|
||||
f'Total size of tags: {sum(map(len, gop_tags))}\n'
|
||||
f'The first tag is {gop_tags[0]}\n'
|
||||
f'The last tag is {gop_tags[-1]}'
|
||||
)
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, on_error, on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return _sort
|
||||
@@ -2,10 +2,15 @@ import logging
|
||||
from typing import Callable, Optional
|
||||
|
||||
from reactivex import Observable, abc
|
||||
from reactivex.disposable import CompositeDisposable, Disposable, SerialDisposable
|
||||
|
||||
from ..common import is_audio_sequence_header, is_metadata_tag, is_video_sequence_header
|
||||
from ..common import (
|
||||
is_audio_sequence_header,
|
||||
is_metadata_tag,
|
||||
is_video_sequence_header,
|
||||
parse_metadata,
|
||||
)
|
||||
from ..models import AudioTag, FlvHeader, ScriptTag, VideoTag
|
||||
from .correct import correct
|
||||
from .typing import FLVStream, FLVStreamItem
|
||||
|
||||
__all__ = ('split',)
|
||||
@@ -21,6 +26,9 @@ def split() -> Callable[[FLVStream], FLVStream]:
|
||||
observer: abc.ObserverBase[FLVStreamItem],
|
||||
scheduler: Optional[abc.SchedulerBase] = None,
|
||||
) -> abc.DisposableBase:
|
||||
disposed = False
|
||||
subscription = SerialDisposable()
|
||||
|
||||
changed: bool = False
|
||||
last_flv_header: Optional[FlvHeader] = None
|
||||
last_metadata_tag: Optional[ScriptTag] = None
|
||||
@@ -62,7 +70,11 @@ def split() -> Callable[[FLVStream], FLVStream]:
|
||||
tag = item
|
||||
|
||||
if is_metadata_tag(tag):
|
||||
logger.debug(f'Metadata tag: {tag}')
|
||||
metadata = parse_metadata(tag)
|
||||
logger.debug(f'Metadata tag: {tag}, metadata: {metadata}')
|
||||
if last_metadata_tag is not None:
|
||||
last_metadata_tag = tag
|
||||
return
|
||||
last_metadata_tag = tag
|
||||
elif is_audio_sequence_header(tag):
|
||||
logger.debug(f'Audio sequence header: {tag}')
|
||||
@@ -91,10 +103,17 @@ def split() -> Callable[[FLVStream], FLVStream]:
|
||||
|
||||
observer.on_next(tag)
|
||||
|
||||
return source.subscribe(
|
||||
def dispose() -> None:
|
||||
nonlocal disposed
|
||||
disposed = True
|
||||
reset()
|
||||
|
||||
subscription.disposable = source.subscribe(
|
||||
on_next, observer.on_error, observer.on_completed, scheduler=scheduler
|
||||
)
|
||||
|
||||
return Observable(subscribe).pipe(correct())
|
||||
return CompositeDisposable(subscription, Disposable(dispose))
|
||||
|
||||
return Observable(subscribe)
|
||||
|
||||
return _split
|
||||
|
||||
@@ -21,7 +21,11 @@ async def make_metadata_file(flv_path: str) -> str:
|
||||
|
||||
async def _make_metadata_content(flv_path: str) -> str:
|
||||
metadata = await get_metadata(flv_path)
|
||||
extra_metadata = await get_extra_metadata(flv_path)
|
||||
try:
|
||||
extra_metadata = await get_extra_metadata(flv_path)
|
||||
except Exception as e:
|
||||
logger.warning(f'Failed to get extra metadata: {repr(e)}')
|
||||
extra_metadata = {}
|
||||
|
||||
comment = cast(str, metadata.get('Comment', ''))
|
||||
chapters = ''
|
||||
@@ -29,8 +33,11 @@ async def _make_metadata_content(flv_path: str) -> str:
|
||||
if join_points := extra_metadata.get('joinpoints'):
|
||||
join_points = list(map(JoinPoint.from_metadata_value, join_points))
|
||||
comment += '\n\n' + make_comment_for_joinpoints(join_points)
|
||||
duration = int(cast(float, metadata['duration']) * 1000)
|
||||
chapters = _make_chapters(join_points, duration)
|
||||
last_timestamp = int(
|
||||
cast(float, extra_metadata.get('duration') or metadata.get('duration'))
|
||||
* 1000
|
||||
)
|
||||
chapters = _make_chapters(join_points, last_timestamp)
|
||||
|
||||
comment = '\\\n'.join(comment.splitlines())
|
||||
|
||||
@@ -48,13 +55,13 @@ Comment={comment}
|
||||
"""
|
||||
|
||||
|
||||
def _make_chapters(join_points: Iterable[JoinPoint], duration: int) -> str:
|
||||
def _make_chapters(join_points: Iterable[JoinPoint], last_timestamp: int) -> str:
|
||||
join_points = filter(lambda p: not p.seamless, join_points)
|
||||
timestamps = list(map(lambda p: p.timestamp, join_points))
|
||||
if not timestamps:
|
||||
return ''
|
||||
timestamps.insert(0, 0)
|
||||
timestamps.append(duration)
|
||||
timestamps.append(last_timestamp)
|
||||
|
||||
result = ''
|
||||
for i in range(1, len(timestamps)):
|
||||
|
||||
@@ -13,6 +13,7 @@ from ..core import Recorder, RecorderEventListener
|
||||
from ..event.event_emitter import EventEmitter, EventListener
|
||||
from ..exception import exception_callback, submit_exception
|
||||
from ..flv.helpers import is_valid_flv_file
|
||||
from ..flv.metadata_analysis import analyse_metadata
|
||||
from ..flv.metadata_injection import InjectingProgress, inject_metadata
|
||||
from ..logging.room_id import aio_task_with_room_id
|
||||
from ..path import extra_metadata_path
|
||||
@@ -141,9 +142,7 @@ class Postprocessor(
|
||||
logger.debug(f'Postprocessing... {video_path}')
|
||||
|
||||
if not await self._is_vaild_flv_file(video_path):
|
||||
logger.warning(f'Invalid flv file: {video_path}')
|
||||
self._queue.task_done()
|
||||
continue
|
||||
logger.warning(f'The flv file may be invalid: {video_path}')
|
||||
|
||||
try:
|
||||
if self.remux_to_mp4:
|
||||
@@ -168,9 +167,22 @@ class Postprocessor(
|
||||
self._queue.task_done()
|
||||
|
||||
async def _inject_extra_metadata(self, path: str) -> str:
|
||||
logger.info(f"Injecting metadata for '{path}' ...")
|
||||
try:
|
||||
metadata = await get_extra_metadata(path)
|
||||
logger.info(f"Injecting metadata for '{path}' ...")
|
||||
try:
|
||||
metadata = await get_extra_metadata(path)
|
||||
except Exception as e:
|
||||
logger.warning(f'Failed to get extra metadata: {repr(e)}')
|
||||
logger.info(f"Analysing metadata for '{path}' ...")
|
||||
await self._analyse_metadata(path)
|
||||
metadata = await get_extra_metadata(path)
|
||||
else:
|
||||
if 'keyframes' not in metadata:
|
||||
logger.warning('The keyframes metadata lost')
|
||||
logger.info(f"Analysing metadata for '{path}' ...")
|
||||
await self._analyse_metadata(path)
|
||||
new_metadata = await get_extra_metadata(path)
|
||||
metadata.update(new_metadata)
|
||||
await self._inject_metadata(path, metadata)
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to inject metadata for '{path}': {repr(e)}")
|
||||
@@ -207,6 +219,19 @@ class Postprocessor(
|
||||
|
||||
return result_path
|
||||
|
||||
def _analyse_metadata(self, path: str) -> Awaitable[None]:
|
||||
future: asyncio.Future[None] = asyncio.Future()
|
||||
self._postprocessing_path = path
|
||||
|
||||
subscription = analyse_metadata(path, show_progress=True).subscribe(
|
||||
on_error=lambda e: future.set_exception(e),
|
||||
on_completed=lambda: future.set_result(None),
|
||||
scheduler=self._scheduler,
|
||||
)
|
||||
future.add_done_callback(lambda f: subscription.dispose())
|
||||
|
||||
return future
|
||||
|
||||
def _inject_metadata(self, path: str, metadata: Dict[str, Any]) -> Awaitable[None]:
|
||||
future: asyncio.Future[None] = asyncio.Future()
|
||||
self._postprocessing_path = path
|
||||
|
||||
Reference in New Issue
Block a user