Compare commits

...

22 Commits

Author SHA1 Message Date
acgnhik
2b0537086e release: 1.8.1 2022-07-20 18:52:13 +08:00
acgnhik
a4d61c1f29 fix: fix bug 2022-07-20 15:41:35 +08:00
acgnhik
17ea3ffc04 release: 1.8.0
close #55
close #57
close #60
solve #64
solve #67
solve #73
solve #84
fix #92
2022-07-17 16:41:32 +08:00
acgnhik
7f366463d4 feat: sort url info by host 2022-07-17 11:31:38 +08:00
acgnhik
fdd93e70a5 refactor: refactor operators 2022-07-17 11:17:28 +08:00
acgnhik
1370cecea8 fix: fix KeyError: 'Title' 2022-07-14 11:03:23 +08:00
acgnhik
18cd928369 release: 1.8.0-alpha.5 2022-07-11 13:56:20 +08:00
acgnhik
7eea5c80b0 chore: update readme 2022-07-11 13:48:40 +08:00
acgnhik
3b523ca11a refactor: add sort operator 2022-07-11 13:05:26 +08:00
acgnhiki
133409a81a Merge pull request #87 from happytommyl/master
Change the working directory
2022-07-05 13:58:34 +08:00
acgnhik
613718e56b fix: fix AssertionError: self._last_metadata is not None 2022-07-05 13:54:20 +08:00
happytommyl
955a531655 Change the working directory
Change the working directory to that of the batch file so that it can be run globaly from the terminal
2022-07-05 02:26:12 +08:00
happytommyl
1ea9d10fce Change the working directory
Change the working directory to that of the batch file so that it can be run globaly from the terminal
2022-07-05 02:17:56 +08:00
acgnhik
66999ef587 fix: fix ValueError
`ValueError: All strings must be XML compatible: Unicode or ASCII, no NULL bytes or control characters`
2022-07-03 10:56:53 +08:00
acgnhik
e5277510c2 refactor: workaround for negative timestamp 2022-07-02 23:03:08 +08:00
acgnhik
e546b47e29 feat: adjust stime of Danmaku when stream recording interrupted 2022-07-02 16:46:44 +08:00
acgnhik
261a2993be refactor: refactor RequestExceptionHandler 2022-06-26 11:00:15 +08:00
acgnhik
44b031ef70 fix: fix HLS recording stopped unexpectedly 2022-06-26 10:55:46 +08:00
acgnhik
8d720cecc5 refactor: refactor operators 2022-06-25 14:27:34 +08:00
acgnhik
ae6058146d release: 1.8.0-alpha.4 2022-06-19 13:39:51 +08:00
acgnhik
765df3ead9 perf: postprocessing one video only at the same time 2022-06-19 13:32:56 +08:00
acgnhik
b74c25ebb6 perf: adjust max_workers 2022-06-19 13:28:45 +08:00
44 changed files with 872 additions and 140 deletions

View File

@@ -1,5 +1,30 @@
# 更新日志
## 1.8.1
修 bug
## 1.8.0
- 重构直播流录制
- 重构弹幕客户端
- 修复了一些 bug
- 优先使用 web api
- 添加直播流时间相关元数据
- 支持 Liquid 模板自定义通知消息
- 同一时间只处理一个录播文件
- 流录制中断重新调整弹幕时间
- 对流主机进行排序
## 1.8.0-alpha.5
- 流录制中断重新调整弹幕时间
- 重构并修正了一些问题
## 1.8.0-alpha.4
- 改善录制多个直播间出现卡顿
## 1.8.0-alpha.3
- 重构并修正了一些问题

View File

@@ -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
[![JetBrains Logo (Main) logo](https://resources.jetbrains.com/storage/products/company/brand/logos/jb_beam.svg)](https://jb.gg/OpenSource)
## 其它相关工具或项目
| 名称 | 链接 | 简介 |
| --- | --- | --- |
| 录播姬 | [官网](https://rec.danmuji.org/) | 简单易用成熟稳定的 B 站直播录制工具 |
| rclone | [官网](https://rclone.org/) | 可以挂载网盘用于存放录播文件 |
| alist | [官网](https://alist-doc.nn.ci/) | 网盘文件浏览、播放 |
| filebrowser | [官网](https://filebrowser.org/) | 服务器文件管理 |

View File

@@ -1,6 +1,6 @@
@echo off
chcp 65001
cd %~dp0
set PATH=.\ffmpeg\bin;.\python;%PATH%
@REM 不使用代理

View File

@@ -1,5 +1,5 @@
chcp 65001
Set-Location $PSScriptRoot
$env:PATH = ".\ffmpeg\bin;.\python;" + $env:PATH
# 不使用代理

View File

@@ -1,3 +1,3 @@
__prog__ = 'blrec'
__version__ = '1.8.0-alpha.3'
__version__ = '1.8.1'
__github__ = 'https://github.com/acgnhiki/blrec'

View File

@@ -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]

View File

@@ -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

View File

@@ -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,

View File

@@ -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'))

View File

@@ -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',

View File

@@ -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

View File

@@ -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:

View File

@@ -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:

View 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)

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -1,7 +1,6 @@
from __future__ import annotations
import asyncio
import time
import logging
from datetime import datetime
from typing import Iterator, Optional

View File

@@ -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')

View File

@@ -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:

View File

@@ -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:

View File

@@ -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)

View File

@@ -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)

View File

@@ -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

View 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),
)

View File

@@ -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:

View 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

View File

@@ -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',
)

View File

@@ -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

View File

@@ -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

View File

@@ -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

View File

@@ -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:

View File

@@ -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

View File

@@ -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)

View File

@@ -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

View File

@@ -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(

View File

@@ -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:

View File

@@ -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:

View File

@@ -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

View 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

View File

@@ -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

View File

@@ -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)):

View File

@@ -4,7 +4,7 @@ import asyncio
import logging
from contextlib import suppress
from pathlib import PurePath
from typing import Any, Awaitable, Dict, Iterator, List, Optional, Union
from typing import Any, Awaitable, Dict, Final, Iterator, List, Optional, Union
from reactivex.scheduler import ThreadPoolScheduler
@@ -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
@@ -48,6 +49,8 @@ class Postprocessor(
AsyncCooperationMixin,
SupportDebugMixin,
):
_worker_semaphore: Final = asyncio.Semaphore(value=1)
def __init__(
self,
live: Live,
@@ -127,43 +130,59 @@ class Postprocessor(
@aio_task_with_room_id
async def _worker(self) -> None:
while True:
self._status = PostprocessorStatus.WAITING
self._postprocessing_path = None
self._postprocessing_progress = None
video_path = await self._queue.get()
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
async with self._worker_semaphore:
logger.debug(f'Postprocessing... {video_path}')
try:
if self.remux_to_mp4:
self._status = PostprocessorStatus.REMUXING
result_path = await self._remux_flv_to_mp4(video_path)
elif self.inject_extra_metadata:
self._status = PostprocessorStatus.INJECTING
result_path = await self._inject_extra_metadata(video_path)
else:
result_path = video_path
if not await self._is_vaild_flv_file(video_path):
logger.warning(f'The flv file may be invalid: {video_path}')
if not self._debug:
await discard_file(extra_metadata_path(video_path), 'DEBUG')
try:
if self.remux_to_mp4:
self._status = PostprocessorStatus.REMUXING
result_path = await self._remux_flv_to_mp4(video_path)
elif self.inject_extra_metadata:
self._status = PostprocessorStatus.INJECTING
result_path = await self._inject_extra_metadata(video_path)
else:
result_path = video_path
self._completed_files.append(result_path)
await self._emit('video_postprocessing_completed', self, result_path)
except Exception as exc:
submit_exception(exc)
finally:
self._queue.task_done()
if not self._debug:
await discard_file(extra_metadata_path(video_path), 'DEBUG')
self._completed_files.append(result_path)
await self._emit(
'video_postprocessing_completed', self, result_path
)
except Exception as exc:
submit_exception(exc)
finally:
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)}")
@@ -200,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

View File

@@ -19,7 +19,7 @@ def wait_for(
) -> _T:
global _executor
if _executor is None:
_executor = ThreadPoolExecutor(thread_name_prefix='wait_for')
_executor = ThreadPoolExecutor(max_workers=200, thread_name_prefix='wait_for')
atexit.register(_executor.shutdown)
future = _executor.submit(func, *args, **kwargs)