Compare commits

..

1 Commits

Author SHA1 Message Date
acgnhik
c9409a1570 release: 1.6.0 2022-04-09 23:01:32 +08:00
31 changed files with 184 additions and 538 deletions

View File

@@ -40,6 +40,11 @@ jobs:
- name: Unzip Python archive
run: Expand-Archive -LiteralPath "python.zip" -DestinationPath "build\python"
- name: Enter build directory
run: |
Set-Location -Path "build"
ls
- name: Rename ffmpeg directory
working-directory: build
run: Rename-Item -Path $($env:FFMPEG_ARCHIVE_NAME).Substring(0, $($env:FFMPEG_ARCHIVE_NAME).Length - 4) "ffmpeg"
@@ -78,9 +83,11 @@ jobs:
working-directory: build
run: Copy-Item "${{ github.workspace }}\run.bat" -Destination ".\run.bat"
- name: Copy run.ps1
- name: Exit build directory
working-directory: build
run: Copy-Item "${{ github.workspace }}\run.ps1" -Destination ".\run.ps1"
run: |
ls
Set-Location -Path ".."
- name: Zip files
run: |

View File

@@ -1,17 +1,5 @@
# 更新日志
## 1.6.2
- 忽略 Windows 注册表 JavaScript MIME 设置 (issue #12, 27)
- 修复 HLS 录制出错 (issue #39, 41)
- 修 bug (issue #47)
- Windows 绿色版默认主机绑定 0.0.0.0 并加上 api key
## 1.6.1
- 修复 bug (issue #37, 38, 40)
- 接收到错误的数据自动换线路 (issue #43)
## 1.6.0
- 更新 Pushplus 消息推送 url (issue #26)

View File

@@ -129,7 +129,7 @@ sudo docker run \
例如:`blrec --host 0.0.0.0 --port 8000`
### 网络安全
### 安全保障
指定 `SSL` 证书使用 **https** 协议并指定 `api key` 可防止被恶意访问和泄漏设置里的敏感信息
@@ -141,16 +141,6 @@ sudo docker run \
如果在不信任的环境下,请使用浏览器的隐式模式访问。
### 关于 api-key
api key 可以使用数字和字母,长度限制为最短 8 最长 80。
3 次尝试内 api key 正确客户端 ip 会自动加入白名单3 次错误后则 ip 会被加入黑名单,黑名单后请求会被拒绝 (403)。
黑名单和白名单数以及同时尝试连接的 ip 数量限制各为 100黑名单或白名单到达限制后不再接受除了白名单内的其它 ip 。
只有重启才会清空黑名单和白名单。
## 作为 ASGI 应用运行
uvicorn blrec.web:app

27
run.bat
View File

@@ -3,33 +3,20 @@ chcp 65001
set PATH=.\ffmpeg\bin;.\python;%PATH%
@REM 不使用代理
REM 不使用代理
set no_proxy=*
@REM 主机和端口绑定,可以按需修改。
set host=0.0.0.0
REM 默认本地主机和端口绑定
set host=localhost
set port=2233
@REM 关于 api key
@REM api key 可以使用数字和字母,长度限制为最短 8 最长 80
@REM 3 次尝试内 api key 正确客户端 ip 会自动加入白名单3 次错误后则 ip 会被加入黑名单,黑名单后请求会被拒绝 (403)。
@REM 黑名单和白名单数以及同时尝试连接的 ip 数量限制各为 100黑名单或白名单到达限制后不再接受除了白名单内的其它 ip 。
@REM 只有重启才会清空黑名单和白名单。
@REM 浏览器第一次访问会弹对话框要求输入 api key。
@REM 输入的 api key 会被保存在浏览器的 local storage下次使用同一浏览器不用再次输入。
@REM 请自行修改 api key不要使用默认的 api key。
set api_key=bili2233
REM 服务器主机和端口绑定,去掉注释并按照自己的情况修改。
REM set host=0.0.0.0
REM set port=80
set DEFAULT_LOG_DIR=日志文件
set DEFAULT_OUT_DIR=录播文件
python -m blrec -c settings.toml --open --host %host% --port %port% --api-key %api_key%
python -m blrec -c settings.toml --open --host %host% --port %port%
pause

27
run.ps1
View File

@@ -1,27 +0,0 @@
chcp 65001
$env:PATH = ".\ffmpeg\bin;.\python;" + $env:PATH
# 不使用代理
$env:no_proxy = "*"
# 主机和端口绑定,可以按需修改。
$env:host = "0.0.0.0"
$env:port = 2233
# 关于 api key
# api key 可以使用数字和字母,长度限制为最短 8 最长 80。
# 3 次尝试内 api key 正确客户端 ip 会自动加入白名单3 次错误后则 ip 会被加入黑名单,黑名单后请求会被拒绝 (403)。
# 黑名单和白名单数以及同时尝试连接的 ip 数量限制各为 100黑名单或白名单到达限制后不再接受除了白名单内的其它 ip 。
# 只有重启才会清空黑名单和白名单。
# 浏览器第一次访问会弹对话框要求输入 api key。
# 输入的 api key 会被保存在浏览器的 local storage下次使用同一浏览器不用再次输入。
# 请自行修改 api key不要使用默认的 api key。
$env:api_key = "bili2233"
$env:DEFAULT_LOG_DIR = "日志文件"
$env:DEFAULT_OUT_DIR = "录播文件"
python -m blrec -c settings.toml --open --host $env:host --port $env:port --api-key $env:api_key
pause

View File

@@ -1,4 +1,4 @@
__prog__ = 'blrec'
__version__ = '1.6.2'
__version__ = '1.6.0'
__github__ = 'https://github.com/acgnhiki/blrec'

View File

@@ -15,9 +15,7 @@ from tenacity import (
from .api import AppApi, WebApi
from .models import LiveStatus, RoomInfo, UserInfo
from .typing import (
ApiPlatform, StreamFormat, QualityNumber, StreamCodec, ResponseData
)
from .typing import StreamFormat, QualityNumber, StreamCodec, ResponseData
from .exceptions import (
LiveRoomHidden, LiveRoomLocked, LiveRoomEncrypted, NoStreamAvailable,
NoStreamFormatAvailable, NoStreamCodecAvailable, NoStreamQualityAvailable,
@@ -179,14 +177,12 @@ class Live:
async def get_live_stream_urls(
self,
qn: QualityNumber = 10000,
*,
api_platform: ApiPlatform = 'android',
stream_format: StreamFormat = 'flv',
stream_codec: StreamCodec = 'avc',
) -> List[str]:
if api_platform == 'android':
try:
info = await self._appapi.get_room_play_info(self._room_id, qn)
else:
except Exception:
info = await self._webapi.get_room_play_info(self._room_id, qn)
self._check_room_play_info(info)

View File

@@ -1,10 +1,7 @@
from typing import Any, Dict, Literal, Mapping
ApiPlatform = Literal[
'web',
'android',
]
Danmaku = Mapping[str, Any]
QualityNumber = Literal[
20000, # 4K
@@ -29,5 +26,3 @@ StreamCodec = Literal[
JsonResponse = Dict[str, Any]
ResponseData = Dict[str, Any]
Danmaku = Mapping[str, Any]

View File

@@ -34,11 +34,10 @@ from .stream_analyzer import StreamProfile
from .statistics import StatisticsCalculator
from ..event.event_emitter import EventListener, EventEmitter
from ..bili.live import Live
from ..bili.typing import ApiPlatform, StreamFormat, QualityNumber
from ..bili.typing import StreamFormat, QualityNumber
from ..bili.helpers import get_quality_name
from ..flv.data_analyser import MetaData
from ..flv.stream_processor import StreamProcessor, BaseOutputFileManager
from ..utils.io import wait_for
from ..utils.mixins import AsyncCooperationMixin, AsyncStoppableMixin
from ..path import escape_path
from ..logging.room_id import aio_task_with_room_id
@@ -103,8 +102,7 @@ class BaseStreamRecorder(
self._quality_number = quality_number
self._real_stream_format: Optional[StreamFormat] = None
self._real_quality_number: Optional[QualityNumber] = None
self._api_platform: ApiPlatform = 'android'
self._use_alternative_stream: bool = False
self._use_candidate_stream: bool = False
self.buffer_size = buffer_size or io.DEFAULT_BUFFER_SIZE # bytes
self.read_timeout = read_timeout or 3 # seconds
self.disconnection_timeout = disconnection_timeout or 600 # seconds
@@ -270,8 +268,7 @@ class BaseStreamRecorder(
self._stream_url = ''
self._stream_host = ''
self._stream_profile = {}
self._api_platform = 'android'
self._use_alternative_stream = False
self._use_candidate_stream = False
self._connection_recovered.clear()
self._thread = Thread(
target=self._run, name=f'StreamRecorder::{self._live.room_id}'
@@ -290,12 +287,6 @@ class BaseStreamRecorder(
def _run(self) -> None:
raise NotImplementedError()
def _rotate_api_platform(self) -> None:
if self._api_platform == 'android':
self._api_platform = 'web'
else:
self._api_platform = 'android'
@retry(
reraise=True,
retry=retry_if_exception_type((
@@ -309,16 +300,11 @@ class BaseStreamRecorder(
fmt = self._real_stream_format or self.stream_format
logger.info(
f'Getting the live stream url... qn: {qn}, format: {fmt}, '
f'api platform: {self._api_platform}, '
f'use alternative stream: {self._use_alternative_stream}'
f'use_candidate_stream: {self._use_candidate_stream}'
)
try:
urls = self._run_coroutine(
self._live.get_live_stream_urls(
qn,
api_platform=self._api_platform,
stream_format=fmt,
)
self._live.get_live_stream_urls(qn, fmt)
)
except NoStreamQualityAvailable:
logger.info(
@@ -346,10 +332,6 @@ class BaseStreamRecorder(
except NoStreamCodecAvailable as e:
logger.warning(repr(e))
raise TryAgain
except Exception as e:
logger.warning(f'Failed to get live stream urls: {repr(e)}')
self._rotate_api_platform()
raise TryAgain
else:
logger.info(
f'Adopted the stream format ({fmt}) and quality ({qn})'
@@ -357,19 +339,17 @@ class BaseStreamRecorder(
self._real_quality_number = qn
self._real_stream_format = fmt
if not self._use_alternative_stream:
if not self._use_candidate_stream:
url = urls[0]
else:
try:
url = urls[1]
except IndexError:
self._use_alternative_stream = False
self._rotate_api_platform()
logger.info(
'No alternative stream url available, will using the primary'
f' stream url from {self._api_platform} api instead.'
'No candidate stream url available, '
'will using the primary stream url instead.'
)
raise TryAgain
url = urls[0]
logger.info(f"Got live stream url: '{url}'")
return url
@@ -461,14 +441,8 @@ B站直播录像
class StreamProxy(io.RawIOBase):
def __init__(
self,
stream: io.BufferedIOBase,
*,
read_timeout: Optional[float] = None,
) -> None:
def __init__(self, stream: io.BufferedIOBase) -> None:
self._stream = stream
self._read_timmeout = read_timeout
self._offset = 0
self._size_updates = Subject()
@@ -491,14 +465,7 @@ class StreamProxy(io.RawIOBase):
return True
def read(self, size: int = -1) -> bytes:
if self._stream.closed:
raise EOFError
if self._read_timmeout:
data = wait_for(
self._stream.read, args=(size, ), timeout=self._read_timmeout
)
else:
data = self._stream.read(size)
data = self._stream.read(size)
self._offset += len(data)
self._size_updates.on_next(len(data))
return data
@@ -507,14 +474,7 @@ class StreamProxy(io.RawIOBase):
return self._offset
def readinto(self, b: Any) -> int:
if self._stream.closed:
raise EOFError
if self._read_timmeout:
n = wait_for(
self._stream.readinto, args=(b, ), timeout=self._read_timmeout
)
else:
n = self._stream.readinto(b)
n = self._stream.readinto(b)
self._offset += n
self._size_updates.on_next(n)
return n

View File

@@ -1,3 +0,0 @@
class FailedToFetchSegments(Exception):
pass

View File

@@ -84,7 +84,6 @@ class FLVStreamRecorder(
analyse_data=True,
dedup_join=True,
save_extra_metadata=True,
backup_timestamp=True,
)
def update_size(size: int) -> None:
@@ -106,7 +105,6 @@ class FLVStreamRecorder(
except Exception as e:
self._handle_exception(e)
finally:
self._stopped = True
if self._stream_processor is not None:
self._stream_processor.finalize()
self._stream_processor = None
@@ -152,6 +150,7 @@ class FLVStreamRecorder(
except Exception as e:
logger.exception(e)
self._handle_exception(e)
self._stopped = True
def _streaming_loop(self) -> None:
url = self._get_live_stream_url()
@@ -177,12 +176,12 @@ class FLVStreamRecorder(
except requests.exceptions.ConnectionError as e:
logger.warning(repr(e))
self._wait_for_connection_error()
except (FlvDataError, FlvStreamCorruptedError) as e:
except FlvDataError as e:
logger.warning(repr(e))
self._use_candidate_stream = not self._use_candidate_stream
url = self._get_live_stream_url()
except FlvStreamCorruptedError as e:
logger.warning(repr(e))
if not self._use_alternative_stream:
self._use_alternative_stream = True
else:
self._rotate_api_platform()
url = self._get_live_stream_url()
def _streaming(self, url: str) -> None:

View File

@@ -3,12 +3,12 @@ import time
import errno
import logging
from queue import Queue, Empty
from threading import Thread, Event, Lock, Condition
from threading import Thread, Event, Lock
from datetime import datetime
from contextlib import suppress
from urllib.parse import urlparse
from typing import List, Set, Optional
from typing import Set, Optional
import urllib3
import requests
@@ -30,12 +30,10 @@ from tenacity import (
from .stream_remuxer import StreamRemuxer
from .stream_analyzer import ffprobe, StreamProfile
from .base_stream_recorder import BaseStreamRecorder, StreamProxy
from .exceptions import FailedToFetchSegments
from .retry import wait_exponential_for_same_exceptions, before_sleep_log
from ..bili.live import Live
from ..bili.typing import StreamFormat, QualityNumber
from ..flv.stream_processor import StreamProcessor
from ..flv.exceptions import FlvDataError, FlvStreamCorruptedError
from ..utils.mixins import (
AsyncCooperationMixin, AsyncStoppableMixin, SupportDebugMixin
)
@@ -83,9 +81,6 @@ class HLSStreamRecorder(
duration_limit=duration_limit,
)
self._init_for_debug(self._live.room_id)
self._init_section_data: Optional[bytes] = None
self._ready_to_fetch_segments = Condition()
self._failed_to_fetch_segments = Event()
self._stream_analysed_lock = Lock()
self._last_segment_uris: Set[str] = set()
@@ -100,54 +95,52 @@ class HLSStreamRecorder(
)
self._playlist_debug_file = open(path, 'wt', encoding='utf-8')
self._session = requests.Session()
self._session.headers.update(self._live.headers)
with StreamRemuxer(self._live.room_id) as self._stream_remuxer:
with requests.Session() as self._session:
self._session.headers.update(self._live.headers)
self._stream_remuxer = StreamRemuxer(self._live.room_id)
self._segment_queue: Queue[Segment] = Queue(maxsize=1000)
self._segment_data_queue: Queue[bytes] = Queue(maxsize=100)
self._stream_host_available = Event()
self._segment_queue: Queue[Segment] = Queue(maxsize=1000)
self._segment_data_queue: Queue[bytes] = Queue(maxsize=100)
self._stream_host_available = Event()
self._segment_fetcher_thread = Thread(
target=self._run_segment_fetcher,
name=f'SegmentFetcher::{self._live.room_id}',
daemon=True,
)
self._segment_fetcher_thread.start()
self._segment_fetcher_thread = Thread(
target=self._run_segment_fetcher,
name=f'SegmentFetcher::{self._live.room_id}',
daemon=True,
)
self._segment_fetcher_thread.start()
self._segment_data_feeder_thread = Thread(
target=self._run_segment_data_feeder,
name=f'SegmentDataFeeder::{self._live.room_id}',
daemon=True,
)
self._segment_data_feeder_thread.start()
self._segment_data_feeder_thread = Thread(
target=self._run_segment_data_feeder,
name=f'SegmentDataFeeder::{self._live.room_id}',
daemon=True,
)
self._segment_data_feeder_thread.start()
self._stream_processor_thread = Thread(
target=self._run_stream_processor,
name=f'StreamProcessor::{self._live.room_id}',
daemon=True,
)
self._stream_processor_thread.start()
self._stream_processor_thread = Thread(
target=self._run_stream_processor,
name=f'StreamProcessor::{self._live.room_id}',
daemon=True,
)
self._stream_processor_thread.start()
try:
self._main_loop()
finally:
if self._stream_processor is not None:
self._stream_processor.cancel()
self._stream_processor_thread.join(timeout=10)
self._segment_fetcher_thread.join(timeout=10)
self._segment_data_feeder_thread.join(timeout=10)
self._stream_remuxer.stop()
self._stream_remuxer.raise_for_exception()
self._last_segment_uris.clear()
del self._segment_queue
del self._segment_data_queue
try:
self._main_loop()
finally:
if self._stream_processor is not None:
self._stream_processor.cancel()
self._segment_fetcher_thread.join(timeout=10)
self._segment_data_feeder_thread.join(timeout=10)
self._last_segment_uris.clear()
del self._segment_queue
del self._segment_data_queue
except TryAgain:
pass
except Exception as e:
self._handle_exception(e)
finally:
self._stopped = True
with suppress(Exception):
self._stream_processor_thread.join(timeout=10)
with suppress(Exception):
self._playlist_debug_file.close()
self._emit_event('stream_recording_stopped')
@@ -208,8 +201,6 @@ class HLSStreamRecorder(
except requests.exceptions.ConnectionError as e:
logger.warning(repr(e))
self._wait_for_connection_error()
except FailedToFetchSegments:
url = self._get_live_stream_url()
except RetryError as e:
logger.warning(repr(e))
@@ -221,14 +212,6 @@ class HLSStreamRecorder(
self._stream_analysed = False
while not self._stopped:
if self._failed_to_fetch_segments.is_set():
with self._segment_queue.mutex:
self._segment_queue.queue.clear()
with self._ready_to_fetch_segments:
self._ready_to_fetch_segments.notify_all()
self._failed_to_fetch_segments.clear()
raise FailedToFetchSegments()
content = self._fetch_playlist(url)
playlist = m3u8.loads(content, uri=url)
@@ -269,10 +252,8 @@ class HLSStreamRecorder(
if playlist.is_endlist:
logger.debug('playlist ended')
self._run_coroutine(self._live.update_room_info())
if not self._live.is_living():
self._stopped = True
break
self._stopped = True
break
time.sleep(1)
@@ -291,8 +272,6 @@ class HLSStreamRecorder(
assert self._stream_remuxer is not None
init_section = None
self._init_section_data = None
num_of_continuously_failed = 0
self._failed_to_fetch_segments.clear()
while not self._stopped:
try:
@@ -328,13 +307,6 @@ class HLSStreamRecorder(
except requests.exceptions.HTTPError as e:
logger.warning(f'Failed to fetch segment: {repr(e)}')
if e.response.status_code in (403, 404, 599):
num_of_continuously_failed += 1
if num_of_continuously_failed >= 3:
self._failed_to_fetch_segments.set()
with self._ready_to_fetch_segments:
self._ready_to_fetch_segments.wait()
num_of_continuously_failed = 0
self._failed_to_fetch_segments.clear()
break
except requests.exceptions.ConnectionError as e:
logger.warning(repr(e))
@@ -343,7 +315,6 @@ class HLSStreamRecorder(
logger.warning(repr(e))
break
else:
num_of_continuously_failed = 0
break
def _run_segment_data_feeder(self) -> None:
@@ -358,8 +329,6 @@ class HLSStreamRecorder(
def _segment_data_feeder(self) -> None:
assert self._stream_remuxer is not None
MAX_SEGMENT_DATA_CACHE = 3
segment_data_cache: List[bytes] = []
bytes_io = io.BytesIO()
segment_count = 0
@@ -390,41 +359,10 @@ class HLSStreamRecorder(
bytes_io = io.BytesIO()
segment_count = 0
self._stream_analysed = True
try:
if self._stream_remuxer.stopped:
self._stream_remuxer.start()
while True:
ready = self._stream_remuxer.wait(timeout=1)
if self._stopped:
return
if ready:
break
if segment_data_cache:
if self._init_section_data:
self._stream_remuxer.input.write(
self._init_section_data
)
for cached_data in segment_data_cache:
if cached_data == self._init_section_data:
continue
self._stream_remuxer.input.write(cached_data)
self._stream_remuxer.input.write(data)
except BrokenPipeError as e:
if not self._stopped:
logger.warning(repr(e))
else:
logger.debug(repr(e))
except ValueError as e:
if not self._stopped:
logger.warning(repr(e))
else:
logger.debug(repr(e))
segment_data_cache.append(data)
if len(segment_data_cache) > MAX_SEGMENT_DATA_CACHE:
segment_data_cache.pop(0)
except BrokenPipeError:
return
def _run_stream_processor(self) -> None:
logger.debug('Stream processor thread started')
@@ -449,49 +387,19 @@ class HLSStreamRecorder(
analyse_data=True,
dedup_join=True,
save_extra_metadata=True,
backup_timestamp=True,
)
self._stream_processor.size_updates.subscribe(update_size)
try:
while not self._stopped:
while True:
ready = self._stream_remuxer.wait(timeout=1)
if self._stopped:
return
if ready:
break
self._stream_host_available.wait()
self._stream_processor.set_metadata(self._make_metadata())
try:
self._stream_processor.process_stream(
StreamProxy(
self._stream_remuxer.output,
read_timeout=10,
) # type: ignore
)
except BrokenPipeError as e:
logger.debug(repr(e))
except TimeoutError as e:
logger.debug(repr(e))
self._stream_remuxer.stop()
except FlvDataError as e:
logger.warning(repr(e))
self._stream_remuxer.stop()
except FlvStreamCorruptedError as e:
logger.warning(repr(e))
self._stream_remuxer.stop()
except ValueError as e:
logger.warning(repr(e))
self._stream_remuxer.stop()
self._stream_host_available.wait()
self._stream_processor.set_metadata(self._make_metadata())
self._stream_processor.process_stream(
StreamProxy(self._stream_remuxer.output), # type: ignore
)
except Exception as e:
if not self._stopped:
logger.exception(e)
self._handle_exception(e)
else:
logger.debug(repr(e))
finally:
self._stream_processor.finalize()
self._progress_bar = None

View File

@@ -1,16 +1,14 @@
import re
import os
import io
import errno
import shlex
import logging
from threading import Thread, Condition
from threading import Thread, Event
from subprocess import Popen, PIPE, CalledProcessError
from typing import Optional, cast
from typing import List, Optional, cast
from ..utils.mixins import StoppableMixin, SupportDebugMixin
from ..utils.io import wait_for
logger = logging.getLogger(__name__)
@@ -19,21 +17,15 @@ logger = logging.getLogger(__name__)
__all__ = 'StreamRemuxer',
class FFmpegError(Exception):
pass
class StreamRemuxer(StoppableMixin, SupportDebugMixin):
_ERROR_PATTERN = re.compile(
r'\b(error|failed|missing|invalid|corrupt)\b', re.IGNORECASE
)
def __init__(self, room_id: int, bufsize: int = 1024 * 1024) -> None:
super().__init__()
self._room_id = room_id
self._bufsize = bufsize
self._exception: Optional[Exception] = None
self._ready = Condition()
self._subprocess_setup = Event()
self._MAX_ERROR_MESSAGES = 10
self._error_messages: List[str] = []
self._env = None
self._init_for_debug(room_id)
@@ -58,22 +50,15 @@ class StreamRemuxer(StoppableMixin, SupportDebugMixin):
def __enter__(self): # type: ignore
self.start()
self.wait()
self.wait_for_subprocess()
return self
def __exit__(self, exc_type, value, traceback): # type: ignore
self.stop()
self.raise_for_exception()
def wait(self, timeout: Optional[float] = None) -> bool:
with self._ready:
return self._ready.wait(timeout=timeout)
def restart(self) -> None:
logger.debug('Restarting stream remuxer...')
self.stop()
self.start()
logger.debug('Restarted stream remuxer')
def wait_for_subprocess(self) -> None:
self._subprocess_setup.wait()
def raise_for_exception(self) -> None:
if not self.exception:
@@ -100,17 +85,12 @@ class StreamRemuxer(StoppableMixin, SupportDebugMixin):
def _run(self) -> None:
logger.debug('Started stream remuxer')
self._exception = None
self._error_messages.clear()
self._subprocess_setup.clear()
try:
self._run_subprocess()
except BrokenPipeError as e:
logger.debug(repr(e))
except FFmpegError as e:
if not self._stopped:
logger.warning(repr(e))
else:
logger.debug(repr(e))
except TimeoutError as e:
logger.debug(repr(e))
except BrokenPipeError:
pass
except Exception as e:
# OSError: [Errno 22] Invalid argument
# https://stackoverflow.com/questions/23688492/oserror-errno-22-invalid-argument-in-subprocess
@@ -120,43 +100,43 @@ class StreamRemuxer(StoppableMixin, SupportDebugMixin):
self._exception = e
logger.exception(e)
finally:
self._stopped = True
logger.debug('Stopped stream remuxer')
def _run_subprocess(self) -> None:
cmd = 'ffmpeg -xerror -i pipe:0 -c copy -copyts -f flv pipe:1'
cmd = 'ffmpeg -i pipe:0 -c copy -f flv pipe:1'
args = shlex.split(cmd)
with Popen(
args, stdin=PIPE, stdout=PIPE, stderr=PIPE,
bufsize=self._bufsize, env=self._env,
) as self._subprocess:
with self._ready:
self._ready.notify_all()
self._subprocess_setup.set()
assert self._subprocess.stderr is not None
with io.TextIOWrapper(
self._subprocess.stderr,
encoding='utf-8',
errors='backslashreplace'
) as stderr:
while not self._stopped:
line = wait_for(stderr.readline, timeout=10)
if not line:
if self._subprocess.poll() is not None:
break
else:
continue
if self._debug:
logger.debug('ffmpeg: %s', line)
self._check_error(line)
while not self._stopped:
data = self._subprocess.stderr.readline()
if not data:
if self._subprocess.poll() is not None:
break
else:
continue
line = data.decode('utf-8', errors='backslashreplace')
if self._debug:
logger.debug('ffmpeg: %s', line)
self._check_error(line)
if not self._stopped and self._subprocess.returncode not in (0, 255):
# 255: Exiting normally, received signal 2.
raise CalledProcessError(self._subprocess.returncode, cmd=cmd)
raise CalledProcessError(
self._subprocess.returncode,
cmd=cmd,
output='\n'.join(self._error_messages),
)
def _check_error(self, line: str) -> None:
match = self._ERROR_PATTERN.search(line)
if not match:
if 'error' not in line.lower() and 'failed' not in line.lower():
return
raise FFmpegError(line)
logger.warning(f'ffmpeg error: {line}')
self._error_messages.append(line)
if len(self._error_messages) > self._MAX_ERROR_MESSAGES:
self._error_messages.remove(self._error_messages[0])

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

View File

@@ -10,6 +10,6 @@
<body>
<app-root></app-root>
<noscript>Please enable JavaScript to continue using this application.</noscript>
<script src="runtime.459a2fcc68ef4fa7.js" type="module"></script><script src="polyfills.4b08448aee19bb22.js" type="module"></script><script src="main.8a8c73fae6ff9291.js" type="module"></script>
<script src="runtime.23c91f03d62c595a.js" type="module"></script><script src="polyfills.4b08448aee19bb22.js" type="module"></script><script src="main.8a8c73fae6ff9291.js" type="module"></script>
</body></html>

View File

@@ -1,6 +1,6 @@
{
"configVersion": 1,
"timestamp": 1650082907617,
"timestamp": 1649386979751,
"index": "/index.html",
"assetGroups": [
{
@@ -14,15 +14,15 @@
"/103.5b5d2a6e5a8a7479.js",
"/146.92e3b29c4c754544.js",
"/45.c90c3cea2bf1a66e.js",
"/66.8d16a032cbce41ed.js",
"/694.dbcbd5e1953ea55d.js",
"/869.ac675e78fa0ea7cf.js",
"/66.d8b06f1fef317761.js",
"/694.92a3e0c2fc842a42.js",
"/869.95d68b28a4188d76.js",
"/common.858f777e9296e6f2.js",
"/index.html",
"/main.8a8c73fae6ff9291.js",
"/manifest.webmanifest",
"/polyfills.4b08448aee19bb22.js",
"/runtime.459a2fcc68ef4fa7.js",
"/runtime.23c91f03d62c595a.js",
"/styles.1f581691b230dc4d.css"
],
"patterns": []
@@ -1637,9 +1637,9 @@
"/103.5b5d2a6e5a8a7479.js": "cc0240f217015b6d4ddcc14f31fcc42e1c1c282a",
"/146.92e3b29c4c754544.js": "3824de681dd1f982ea69a065cdf54d7a1e781f4d",
"/45.c90c3cea2bf1a66e.js": "e5bfb8cf3803593e6b8ea14c90b3d3cb6a066764",
"/66.8d16a032cbce41ed.js": "a473089370c2fe27f96a778acf1e709dc5770b31",
"/694.dbcbd5e1953ea55d.js": "41973de76799b085188903f46cc6974f20395741",
"/869.ac675e78fa0ea7cf.js": "f45052016cb5201d5784b3f261e719d96bd1b153",
"/66.d8b06f1fef317761.js": "43676d9dc886b5624dadecc50f17d4972b183d2d",
"/694.92a3e0c2fc842a42.js": "f8f093029b9996b3db0c4e738bf9f8573fba8392",
"/869.95d68b28a4188d76.js": "cd1add38c89b1df3c0783b74c931b51839f1c530",
"/assets/animal/panda.js": "fec2868bb3053dd2da45f96bbcb86d5116ed72b1",
"/assets/animal/panda.svg": "bebd302cdc601e0ead3a6d2710acf8753f3d83b1",
"/assets/fill/.gitkeep": "da39a3ee5e6b4b0d3255bfef95601890afd80709",
@@ -3234,11 +3234,11 @@
"/assets/twotone/warning.js": "fb2d7ea232f3a99bf8f080dbc94c65699232ac01",
"/assets/twotone/warning.svg": "8c7a2d3e765a2e7dd58ac674870c6655cecb0068",
"/common.858f777e9296e6f2.js": "b68ca68e1e214a2537d96935c23410126cc564dd",
"/index.html": "be9d178866ccd58dfdc2a4c1b375ab030302163b",
"/index.html": "114f00ffcd1f7fa5aaaa7f2fcf3109f26c77c715",
"/main.8a8c73fae6ff9291.js": "41a5a5a8fb5cda4cfa0e28532812594816257122",
"/manifest.webmanifest": "62c1cb8c5ad2af551a956b97013ab55ce77dd586",
"/polyfills.4b08448aee19bb22.js": "8e73f2d42cc13ca353cea5c886d930bd6da08d0d",
"/runtime.459a2fcc68ef4fa7.js": "a68b948e588e75f8cb2fb5315ac41623e7c6ed1e",
"/runtime.23c91f03d62c595a.js": "0819f1120ed1e37c2ad069ef949147450c951069",
"/styles.1f581691b230dc4d.css": "6f5befbbad57c2b2e80aae855139744b8010d150"
},
"navigationUrls": [

View File

@@ -1 +1 @@
(()=>{"use strict";var e,v={},m={};function r(e){var i=m[e];if(void 0!==i)return i.exports;var t=m[e]={exports:{}};return v[e].call(t.exports,t,t.exports,r),t.exports}r.m=v,e=[],r.O=(i,t,o,f)=>{if(!t){var a=1/0;for(n=0;n<e.length;n++){for(var[t,o,f]=e[n],c=!0,l=0;l<t.length;l++)(!1&f||a>=f)&&Object.keys(r.O).every(b=>r.O[b](t[l]))?t.splice(l--,1):(c=!1,f<a&&(a=f));if(c){e.splice(n--,1);var d=o();void 0!==d&&(i=d)}}return i}f=f||0;for(var n=e.length;n>0&&e[n-1][2]>f;n--)e[n]=e[n-1];e[n]=[t,o,f]},r.n=e=>{var i=e&&e.__esModule?()=>e.default:()=>e;return r.d(i,{a:i}),i},r.d=(e,i)=>{for(var t in i)r.o(i,t)&&!r.o(e,t)&&Object.defineProperty(e,t,{enumerable:!0,get:i[t]})},r.f={},r.e=e=>Promise.all(Object.keys(r.f).reduce((i,t)=>(r.f[t](e,i),i),[])),r.u=e=>(592===e?"common":e)+"."+{45:"c90c3cea2bf1a66e",66:"8d16a032cbce41ed",103:"5b5d2a6e5a8a7479",146:"92e3b29c4c754544",592:"858f777e9296e6f2",694:"dbcbd5e1953ea55d",869:"ac675e78fa0ea7cf"}[e]+".js",r.miniCssF=e=>{},r.o=(e,i)=>Object.prototype.hasOwnProperty.call(e,i),(()=>{var e={},i="blrec:";r.l=(t,o,f,n)=>{if(e[t])e[t].push(o);else{var a,c;if(void 0!==f)for(var l=document.getElementsByTagName("script"),d=0;d<l.length;d++){var u=l[d];if(u.getAttribute("src")==t||u.getAttribute("data-webpack")==i+f){a=u;break}}a||(c=!0,(a=document.createElement("script")).type="module",a.charset="utf-8",a.timeout=120,r.nc&&a.setAttribute("nonce",r.nc),a.setAttribute("data-webpack",i+f),a.src=r.tu(t)),e[t]=[o];var s=(g,b)=>{a.onerror=a.onload=null,clearTimeout(p);var _=e[t];if(delete e[t],a.parentNode&&a.parentNode.removeChild(a),_&&_.forEach(h=>h(b)),g)return g(b)},p=setTimeout(s.bind(null,void 0,{type:"timeout",target:a}),12e4);a.onerror=s.bind(null,a.onerror),a.onload=s.bind(null,a.onload),c&&document.head.appendChild(a)}}})(),r.r=e=>{"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},(()=>{var e;r.tu=i=>(void 0===e&&(e={createScriptURL:t=>t},"undefined"!=typeof trustedTypes&&trustedTypes.createPolicy&&(e=trustedTypes.createPolicy("angular#bundler",e))),e.createScriptURL(i))})(),r.p="",(()=>{var e={666:0};r.f.j=(o,f)=>{var n=r.o(e,o)?e[o]:void 0;if(0!==n)if(n)f.push(n[2]);else if(666!=o){var a=new Promise((u,s)=>n=e[o]=[u,s]);f.push(n[2]=a);var c=r.p+r.u(o),l=new Error;r.l(c,u=>{if(r.o(e,o)&&(0!==(n=e[o])&&(e[o]=void 0),n)){var s=u&&("load"===u.type?"missing":u.type),p=u&&u.target&&u.target.src;l.message="Loading chunk "+o+" failed.\n("+s+": "+p+")",l.name="ChunkLoadError",l.type=s,l.request=p,n[1](l)}},"chunk-"+o,o)}else e[o]=0},r.O.j=o=>0===e[o];var i=(o,f)=>{var l,d,[n,a,c]=f,u=0;if(n.some(p=>0!==e[p])){for(l in a)r.o(a,l)&&(r.m[l]=a[l]);if(c)var s=c(r)}for(o&&o(f);u<n.length;u++)r.o(e,d=n[u])&&e[d]&&e[d][0](),e[n[u]]=0;return r.O(s)},t=self.webpackChunkblrec=self.webpackChunkblrec||[];t.forEach(i.bind(null,0)),t.push=i.bind(null,t.push.bind(t))})()})();
(()=>{"use strict";var e,v={},m={};function r(e){var i=m[e];if(void 0!==i)return i.exports;var t=m[e]={exports:{}};return v[e].call(t.exports,t,t.exports,r),t.exports}r.m=v,e=[],r.O=(i,t,f,o)=>{if(!t){var a=1/0;for(n=0;n<e.length;n++){for(var[t,f,o]=e[n],c=!0,l=0;l<t.length;l++)(!1&o||a>=o)&&Object.keys(r.O).every(b=>r.O[b](t[l]))?t.splice(l--,1):(c=!1,o<a&&(a=o));if(c){e.splice(n--,1);var d=f();void 0!==d&&(i=d)}}return i}o=o||0;for(var n=e.length;n>0&&e[n-1][2]>o;n--)e[n]=e[n-1];e[n]=[t,f,o]},r.n=e=>{var i=e&&e.__esModule?()=>e.default:()=>e;return r.d(i,{a:i}),i},r.d=(e,i)=>{for(var t in i)r.o(i,t)&&!r.o(e,t)&&Object.defineProperty(e,t,{enumerable:!0,get:i[t]})},r.f={},r.e=e=>Promise.all(Object.keys(r.f).reduce((i,t)=>(r.f[t](e,i),i),[])),r.u=e=>(592===e?"common":e)+"."+{45:"c90c3cea2bf1a66e",66:"d8b06f1fef317761",103:"5b5d2a6e5a8a7479",146:"92e3b29c4c754544",592:"858f777e9296e6f2",694:"92a3e0c2fc842a42",869:"95d68b28a4188d76"}[e]+".js",r.miniCssF=e=>{},r.o=(e,i)=>Object.prototype.hasOwnProperty.call(e,i),(()=>{var e={},i="blrec:";r.l=(t,f,o,n)=>{if(e[t])e[t].push(f);else{var a,c;if(void 0!==o)for(var l=document.getElementsByTagName("script"),d=0;d<l.length;d++){var u=l[d];if(u.getAttribute("src")==t||u.getAttribute("data-webpack")==i+o){a=u;break}}a||(c=!0,(a=document.createElement("script")).type="module",a.charset="utf-8",a.timeout=120,r.nc&&a.setAttribute("nonce",r.nc),a.setAttribute("data-webpack",i+o),a.src=r.tu(t)),e[t]=[f];var s=(g,b)=>{a.onerror=a.onload=null,clearTimeout(p);var _=e[t];if(delete e[t],a.parentNode&&a.parentNode.removeChild(a),_&&_.forEach(h=>h(b)),g)return g(b)},p=setTimeout(s.bind(null,void 0,{type:"timeout",target:a}),12e4);a.onerror=s.bind(null,a.onerror),a.onload=s.bind(null,a.onload),c&&document.head.appendChild(a)}}})(),r.r=e=>{"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},(()=>{var e;r.tu=i=>(void 0===e&&(e={createScriptURL:t=>t},"undefined"!=typeof trustedTypes&&trustedTypes.createPolicy&&(e=trustedTypes.createPolicy("angular#bundler",e))),e.createScriptURL(i))})(),r.p="",(()=>{var e={666:0};r.f.j=(f,o)=>{var n=r.o(e,f)?e[f]:void 0;if(0!==n)if(n)o.push(n[2]);else if(666!=f){var a=new Promise((u,s)=>n=e[f]=[u,s]);o.push(n[2]=a);var c=r.p+r.u(f),l=new Error;r.l(c,u=>{if(r.o(e,f)&&(0!==(n=e[f])&&(e[f]=void 0),n)){var s=u&&("load"===u.type?"missing":u.type),p=u&&u.target&&u.target.src;l.message="Loading chunk "+f+" failed.\n("+s+": "+p+")",l.name="ChunkLoadError",l.type=s,l.request=p,n[1](l)}},"chunk-"+f,f)}else e[f]=0},r.O.j=f=>0===e[f];var i=(f,o)=>{var l,d,[n,a,c]=o,u=0;if(n.some(p=>0!==e[p])){for(l in a)r.o(a,l)&&(r.m[l]=a[l]);if(c)var s=c(r)}for(f&&f(o);u<n.length;u++)r.o(e,d=n[u])&&e[d]&&e[d][0](),e[n[u]]=0;return r.O(s)},t=self.webpackChunkblrec=self.webpackChunkblrec||[];t.forEach(i.bind(null,0)),t.push=i.bind(null,t.push.bind(t))})()})();

View File

@@ -37,15 +37,8 @@ __all__ = 'FlvParser', 'FlvDumper'
class FlvParser:
def __init__(
self,
stream: RandomIO,
backup_timestamp: bool = False,
restore_timestamp: bool = False,
) -> None:
def __init__(self, stream: RandomIO) -> None:
self._stream = stream
self._backup_timestamp = backup_timestamp
self._restore_timestamp = restore_timestamp
self._reader = StructReader(stream)
def parse_header(self) -> FlvHeader:
@@ -115,42 +108,19 @@ class FlvParser:
def parse_flv_tag_header(self, data: bytes) -> FlvTagHeader:
reader = StructReader(BytesIO(data))
flag = reader.read_ui8()
filtered = bool(flag & 0b0010_0000)
if filtered:
raise FlvDataError('Unsupported Filtered FLV Tag', data)
tag_type = TagType(flag & 0b0001_1111)
data_size = reader.read_ui24()
timestamp = reader.read_ui24()
timestamp_extended = reader.read_ui8()
timestamp = timestamp_extended << 24 | timestamp
stream_id = reader.read_ui24()
if self._backup_timestamp:
return FlvTagHeader(
filtered=filtered,
tag_type=tag_type,
data_size=data_size,
timestamp=timestamp_extended << 24 | timestamp,
stream_id=timestamp,
)
elif self._restore_timestamp:
return FlvTagHeader(
filtered=filtered,
tag_type=tag_type,
data_size=data_size,
timestamp=stream_id,
stream_id=stream_id,
)
else:
return FlvTagHeader(
filtered=filtered,
tag_type=tag_type,
data_size=data_size,
timestamp=timestamp_extended << 24 | timestamp,
stream_id=stream_id,
)
return FlvTagHeader(
filtered, tag_type, data_size, timestamp, stream_id
)
def parse_audio_tag_header(self, data: bytes) -> AudioTagHeader:
reader = StructReader(BytesIO(data))

View File

@@ -16,21 +16,9 @@ logger = logging.getLogger(__name__)
class FlvReader:
def __init__(
self,
stream: RandomIO,
*,
backup_timestamp: bool = False,
restore_timestamp: bool = False,
) -> None:
def __init__(self, stream: RandomIO) -> None:
self._stream = stream
self._backup_timestamp = backup_timestamp
self._restore_timestamp = restore_timestamp
self._parser = FlvParser(
stream,
backup_timestamp=backup_timestamp,
restore_timestamp=restore_timestamp,
)
self._parser = FlvParser(stream)
def read_header(self) -> FlvHeader:
header = self._parser.parse_header()

View File

@@ -116,8 +116,8 @@ class FlvTagHeader:
filtered: bool
tag_type: TagType
data_size: int = attr.ib(validator=[non_negative_integer_validator])
timestamp: int
stream_id: int
timestamp: int = attr.ib(validator=[non_negative_integer_validator])
stream_id: int = attr.ib(validator=[non_negative_integer_validator])
@attr.s(auto_attribs=True, slots=True, frozen=True)

View File

@@ -31,9 +31,9 @@ from .exceptions import (
CutStream,
)
from .common import (
is_metadata_tag, parse_metadata, is_audio_tag, is_video_tag,
is_video_sequence_header, is_audio_sequence_header,
enrich_metadata, update_metadata, is_data_tag, read_tags_in_duration,
is_audio_tag, is_video_tag, is_metadata_tag, parse_metadata,
is_audio_data_tag, is_video_data_tag, enrich_metadata, update_metadata,
is_data_tag, read_tags_in_duration,
)
from ..path import extra_metadata_path
if TYPE_CHECKING:
@@ -63,8 +63,6 @@ class StreamProcessor:
analyse_data: bool = False,
dedup_join: bool = False,
save_extra_metadata: bool = False,
backup_timestamp: bool = False,
restore_timestamp: bool = False,
) -> None:
self._file_manager = file_manager
self._parameters_checker = ParametersChecker()
@@ -81,8 +79,6 @@ class StreamProcessor:
self._analyse_data = analyse_data
self._dedup_join = dedup_join
self._save_x_metadata = save_extra_metadata
self._backup_timestamp = backup_timestamp
self._restore_timestamp = restore_timestamp
self._cancelled: bool = False
self._finalized: bool = False
@@ -227,12 +223,7 @@ class StreamProcessor:
def _process_stream(self, stream: RandomIO) -> None:
logger.debug(f'Processing the {self._stream_count}th stream...')
self._in_reader = FlvReaderWithTimestampFix(
stream,
backup_timestamp=self._backup_timestamp,
restore_timestamp=self._restore_timestamp,
)
self._in_reader = FlvReaderWithTimestampFix(stream)
flv_header = self._read_header()
self._has_audio = flv_header.has_audio()
@@ -561,31 +552,18 @@ class StreamProcessor:
return header
def _ensure_ts_correct(self, tag: FlvTag) -> None:
if not tag.timestamp + self._delta < 0:
if not is_audio_data_tag(tag) or not is_video_data_tag(tag):
return
logger.warning(
f'Incorrect timestamp: {tag.timestamp + self._delta}\n'
f'last output tag: {self._last_tags[0]}\n'
f'current tag: {tag}'
)
if tag.is_audio_tag() or tag.is_video_tag():
self._delta = (
self._last_tags[0].timestamp +
self._in_reader.calc_interval(tag) - tag.timestamp
)
logger.debug(f'Updated delta: {self._delta}')
elif tag.is_script_tag():
self._delta = (
self._last_tags[0].timestamp - tag.timestamp
)
logger.debug(f'Updated delta: {self._delta}')
else:
pass
if tag.timestamp + self._delta < 0:
self._delta = -tag.timestamp
logger.warning('Incorrect timestamp: {}, new delta: {}'.format(
tag, self._delta
))
def _correct_ts(self, tag: FlvTag, delta: int) -> FlvTag:
if delta == 0 and tag.timestamp >= 0:
return tag
return tag.evolve(timestamp=tag.timestamp + delta)
return tag.evolve(timestamp=max(0, tag.timestamp + delta))
def _calc_delta_duplicated(self, last_duplicated_tag: FlvTag) -> int:
return self._last_tags[0].timestamp - last_duplicated_tag.timestamp
@@ -786,17 +764,8 @@ class RobustFlvReader(FlvReader):
class FlvReaderWithTimestampFix(RobustFlvReader):
def __init__(
self,
stream: RandomIO,
backup_timestamp: bool = False,
restore_timestamp: bool = False,
) -> None:
super().__init__(
stream,
backup_timestamp=backup_timestamp,
restore_timestamp=restore_timestamp,
)
def __init__(self, stream: RandomIO) -> None:
super().__init__(stream)
self._last_tag: Optional[FlvTag] = None
self._last_video_tag: Optional[VideoTag] = None
self._last_audio_tag: Optional[AudioTag] = None
@@ -873,17 +842,11 @@ class FlvReaderWithTimestampFix(RobustFlvReader):
if is_video_tag(tag):
if self._last_video_tag is None:
return False
if is_video_sequence_header(self._last_video_tag):
return tag.timestamp < self._last_video_tag.timestamp
else:
return tag.timestamp <= self._last_video_tag.timestamp
return tag.timestamp <= self._last_video_tag.timestamp
elif is_audio_tag(tag):
if self._last_audio_tag is None:
return False
if is_audio_sequence_header(self._last_audio_tag):
return tag.timestamp < self._last_audio_tag.timestamp
else:
return tag.timestamp <= self._last_audio_tag.timestamp
return tag.timestamp <= self._last_audio_tag.timestamp
else:
return False

View File

@@ -1,29 +0,0 @@
import atexit
from concurrent.futures import ThreadPoolExecutor
from concurrent.futures import TimeoutError as _TimeoutError
from typing import Callable, Any, Iterable, Mapping, TypeVar
_T = TypeVar('_T')
_executor = None
def wait_for(
func: Callable[..., _T],
*,
args: Iterable[Any] = [],
kwargs: Mapping[str, Any] = {},
timeout: float
) -> _T:
global _executor
if _executor is None:
_executor = ThreadPoolExecutor(thread_name_prefix='wait_for')
atexit.register(_executor.shutdown)
future = _executor.submit(func, *args, **kwargs)
try:
return future.result(timeout=timeout)
except _TimeoutError:
raise TimeoutError(timeout, func, args, kwargs) from None

View File

@@ -1,7 +1,6 @@
import os
from abc import ABC, abstractmethod
import asyncio
import threading
from typing import Awaitable, TypeVar, final
@@ -9,28 +8,24 @@ class SwitchableMixin(ABC):
def __init__(self) -> None:
super().__init__()
self._enabled = False
self._enabled_lock = threading.Lock()
@property
def enabled(self) -> bool:
with self._enabled_lock:
return self._enabled
return self._enabled
@final
def enable(self) -> None:
with self._enabled_lock:
if self._enabled:
return
self._enabled = True
self._do_enable()
if self._enabled:
return
self._enabled = True
self._do_enable()
@final
def disable(self) -> None:
with self._enabled_lock:
if not self._enabled:
return
self._enabled = False
self._do_disable()
if not self._enabled:
return
self._enabled = False
self._do_disable()
@abstractmethod
def _do_enable(self) -> None:
@@ -45,28 +40,24 @@ class StoppableMixin(ABC):
def __init__(self) -> None:
super().__init__()
self._stopped = True
self._stopped_lock = threading.Lock()
@property
def stopped(self) -> bool:
with self._stopped_lock:
return self._stopped
return self._stopped
@final
def start(self) -> None:
with self._stopped_lock:
if not self._stopped:
return
self._stopped = False
self._do_start()
if not self._stopped:
return
self._stopped = False
self._do_start()
@final
def stop(self) -> None:
with self._stopped_lock:
if self._stopped:
return
self._stopped = True
self._do_stop()
if self._stopped:
return
self._stopped = True
self._do_stop()
@abstractmethod
def _do_start(self) -> None:

View File

@@ -6,7 +6,6 @@ from fastapi import FastAPI, status, Request, Depends
from fastapi.responses import JSONResponse
from fastapi.middleware.cors import CORSMiddleware
from fastapi.staticfiles import StaticFiles
from starlette.responses import Response
from pydantic import ValidationError
from pkg_resources import resource_filename
@@ -144,20 +143,6 @@ class WebAppFiles(StaticFiles):
path = 'index.html'
return await super().lookup_path(path)
def file_response(self, full_path: str, *args, **kwargs) -> Response: # type: ignore # noqa
# ignore MIME types from Windows registry
# workaround for https://github.com/acgnhiki/blrec/issues/12
response = super().file_response(full_path, *args, **kwargs)
if full_path.endswith('.js'):
js_media_type = 'application/javascript'
if response.media_type != js_media_type:
response.media_type = js_media_type
headers = response.headers
headers['content-type'] = js_media_type
response.raw_headers = headers.raw
del response._headers
return response
directory = resource_filename(__name__, '../data/webapp')
api.mount('/', WebAppFiles(directory=directory, html=True), name='webapp')

View File

@@ -1,10 +1,8 @@
<nz-modal
nzTitle="修改文件路径模板"
nzCentered
[nzFooter]="modalFooter"
[(nzVisible)]="visible"
[nzOkDisabled]="control.invalid || control.value.trim() === value"
(nzOnOk)="handleConfirm()"
(nzOnCancel)="handleCancel()"
>
<ng-container *nzModalContent>
<form nz-form [formGroup]="settingsForm">

View File

@@ -28,7 +28,7 @@ export function toBitRateString(
let unit: string;
if (bitrate <= 0) {
return '0' + spacer + 'kbps';
return '0 kbps/s';
}
if (bitrate < 1e6) {
@@ -60,7 +60,7 @@ export function toByteRateString(
let unit: string;
if (rate <= 0) {
return '0' + spacer + 'B/s';
return '0B/s';
}
if (rate < 1e3) {

View File

@@ -86,7 +86,7 @@
<span class="label">下载速度</span>
<app-wave-graph [value]="data.task_status.dl_rate"></app-wave-graph>
<span class="value">
{{ data.task_status.dl_rate * 8 | datarate: { bitrate: true } }}
{{ data.task_status.dl_rate | datarate: { bitrate: true } }}
</span>
</li>
<li class="info-item">