mirror of
https://hubproxy.babadafafafafa.cn/https://github.com/jxxghp/MoviePilot.git
synced 2026-09-20 08:03:34 +08:00
fix(plugin): 加载入口补生命周期锁,手工重载带上分身 (#6679)
This commit is contained in:
@@ -44,11 +44,17 @@ def refresh_plugin_registrations(plugin_id: str) -> None:
|
||||
|
||||
|
||||
def reload_plugin_runtime(plugin_id: str) -> PluginRuntimeStatus:
|
||||
"""重载插件实例并重新注册其命令、定时任务和 API。"""
|
||||
"""重载插件实例树并重新注册其命令、定时任务和 API。
|
||||
|
||||
分身与本体共用同一份源码,只重载本体会让分身继续持有旧模块与旧类对象,改完源码
|
||||
点重载后新旧代码在同一进程内并存;分身的 API、调度与命令注册也不会刷新。文件监听
|
||||
那条重载路径本来就是按实例树走的,手工重载必须与它一致。
|
||||
"""
|
||||
plugin_manager = get_plugin_manager()
|
||||
with plugin_manager.mutation(f"重载插件 {plugin_id}"):
|
||||
runtime_status = plugin_manager.reload_plugin(plugin_id)
|
||||
refresh_plugin_registrations(plugin_id)
|
||||
runtime_status = plugin_manager.reload_plugin_tree(plugin_id)
|
||||
for reload_target in plugin_manager.get_plugin_reload_targets(plugin_id):
|
||||
refresh_plugin_registrations(reload_target)
|
||||
return runtime_status
|
||||
|
||||
|
||||
|
||||
@@ -106,75 +106,82 @@ class PluginLifecycle:
|
||||
self,
|
||||
plugin_id: Optional[str] = None,
|
||||
) -> dict[str, PluginRuntimeStatus]:
|
||||
"""加载并初始化插件,返回每个目标的明确运行结果。"""
|
||||
loadable_plugins = self._loadable_plugins()
|
||||
results: dict[str, PluginRuntimeStatus] = {}
|
||||
if plugin_id:
|
||||
self._runtime_status_writer(plugin_id, PluginRuntimeStatus.READY)
|
||||
"""加载并初始化插件,返回每个目标的明确运行结果。
|
||||
|
||||
def check_module(module: Any) -> bool:
|
||||
"""判断模块是否具备宿主插件最小生命周期钩子。"""
|
||||
return hasattr(module, "init_plugin") and hasattr(module, "plugin_name")
|
||||
与 stop/quiesce 共用同一把可重入锁。加载要构造实例、建库并注册事件订阅,
|
||||
停止取的却是进入时的运行表快照:两者交叠时,本次加载的实例落在快照之外而
|
||||
停止仍会把它从运行表里抹掉,它注册的定时任务、线程和事件订阅却留在原处,
|
||||
此后没有任何句柄能再停掉它。
|
||||
"""
|
||||
with self._lifecycle_lock:
|
||||
loadable_plugins = self._loadable_plugins()
|
||||
results: dict[str, PluginRuntimeStatus] = {}
|
||||
if plugin_id:
|
||||
self._runtime_status_writer(plugin_id, PluginRuntimeStatus.READY)
|
||||
|
||||
plugins = self._load_plugins(plugin_id, loadable_plugins, check_module)
|
||||
plugins.sort(key=lambda item: getattr(item, "plugin_order", 0))
|
||||
for plugin in plugins:
|
||||
current_id = plugin.__name__
|
||||
if plugin_id and current_id.casefold() != plugin_id.casefold():
|
||||
continue
|
||||
try:
|
||||
if not self._auth_checker(plugin):
|
||||
def check_module(module: Any) -> bool:
|
||||
"""判断模块是否具备宿主插件最小生命周期钩子。"""
|
||||
return hasattr(module, "init_plugin") and hasattr(module, "plugin_name")
|
||||
|
||||
plugins = self._load_plugins(plugin_id, loadable_plugins, check_module)
|
||||
plugins.sort(key=lambda item: getattr(item, "plugin_order", 0))
|
||||
for plugin in plugins:
|
||||
current_id = plugin.__name__
|
||||
if plugin_id and current_id.casefold() != plugin_id.casefold():
|
||||
continue
|
||||
try:
|
||||
if not self._auth_checker(plugin):
|
||||
self._remove_classification(current_id)
|
||||
if current_id in self._classes:
|
||||
self._classes[current_id] = plugin
|
||||
status = PluginRuntimeStatus.BLOCKED_BY_POLICY
|
||||
self._runtime_status_writer(plugin_id or current_id, status)
|
||||
results[plugin_id or current_id] = status
|
||||
continue
|
||||
self._remove_classification(current_id)
|
||||
if current_id in self._classes:
|
||||
self._classes[current_id] = plugin
|
||||
status = PluginRuntimeStatus.BLOCKED_BY_POLICY
|
||||
self._classes[current_id] = plugin
|
||||
with bind_plugin_instance(current_id):
|
||||
instance = plugin()
|
||||
instance.init_plugin(self._plugin_config(current_id))
|
||||
self._ensure_database(current_id, instance)
|
||||
enabled = bool(instance.get_state())
|
||||
if enabled:
|
||||
self._refresh_classification_safely(current_id, instance)
|
||||
else:
|
||||
self._remove_classification(current_id)
|
||||
self._quiesced_hooks.pop(current_id, None)
|
||||
self._running[current_id] = instance
|
||||
self._logger.info(
|
||||
f"加载插件:{current_id} 版本:{instance.plugin_version}"
|
||||
)
|
||||
if enabled:
|
||||
self._enable_events(plugin)
|
||||
else:
|
||||
self._disable_events(plugin)
|
||||
status = PluginRuntimeStatus.ACTIVE
|
||||
self._runtime_status_writer(plugin_id or current_id, status)
|
||||
results[plugin_id or current_id] = status
|
||||
continue
|
||||
self._remove_classification(current_id)
|
||||
self._classes[current_id] = plugin
|
||||
with bind_plugin_instance(current_id):
|
||||
instance = plugin()
|
||||
instance.init_plugin(self._plugin_config(current_id))
|
||||
self._ensure_database(current_id, instance)
|
||||
enabled = bool(instance.get_state())
|
||||
if enabled:
|
||||
self._refresh_classification_safely(current_id, instance)
|
||||
else:
|
||||
except Exception as error: # noqa: BLE001
|
||||
self._remove_classification(current_id)
|
||||
self._quiesced_hooks.pop(current_id, None)
|
||||
self._running[current_id] = instance
|
||||
self._logger.info(
|
||||
f"加载插件:{current_id} 版本:{instance.plugin_version}"
|
||||
)
|
||||
if enabled:
|
||||
self._enable_events(plugin)
|
||||
else:
|
||||
self._disable_events(plugin)
|
||||
status = PluginRuntimeStatus.ACTIVE
|
||||
self._runtime_status_writer(plugin_id or current_id, status)
|
||||
results[plugin_id or current_id] = status
|
||||
except Exception as error: # noqa: BLE001
|
||||
self._remove_classification(current_id)
|
||||
status = PluginRuntimeStatus.LOAD_FAILED
|
||||
self._runtime_status_writer(plugin_id or current_id, status)
|
||||
results[plugin_id or current_id] = status
|
||||
# 建库发生在进入运行态之前:失败的插件不会出现在 _running 里,卸载路径
|
||||
# 因此够不到它,句柄只能在这里释放
|
||||
self._release_databases((current_id,))
|
||||
self._logger.error(
|
||||
f"加载插件 {current_id} 出错:{error} - {traceback.format_exc()}"
|
||||
)
|
||||
if plugin_id and not any(
|
||||
result_id.casefold() == plugin_id.casefold()
|
||||
for result_id in results
|
||||
):
|
||||
self._remove_classification(plugin_id)
|
||||
status = PluginRuntimeStatus.LOAD_FAILED
|
||||
self._runtime_status_writer(plugin_id or current_id, status)
|
||||
results[plugin_id or current_id] = status
|
||||
# 建库发生在进入运行态之前:失败的插件不会出现在 _running 里,卸载路径
|
||||
# 因此够不到它,句柄只能在这里释放
|
||||
self._release_databases((current_id,))
|
||||
self._logger.error(
|
||||
f"加载插件 {current_id} 出错:{error} - {traceback.format_exc()}"
|
||||
)
|
||||
if plugin_id and not any(
|
||||
result_id.casefold() == plugin_id.casefold()
|
||||
for result_id in results
|
||||
):
|
||||
self._remove_classification(plugin_id)
|
||||
status = PluginRuntimeStatus.LOAD_FAILED
|
||||
self._runtime_status_writer(plugin_id, status)
|
||||
results[plugin_id] = status
|
||||
self._clear_tools()
|
||||
return results
|
||||
self._runtime_status_writer(plugin_id, status)
|
||||
results[plugin_id] = status
|
||||
self._clear_tools()
|
||||
return results
|
||||
|
||||
@staticmethod
|
||||
def _declaration(instance: Any, hook_name: str) -> Any:
|
||||
|
||||
119
tests/test_plugin_lifecycle_concurrency.py
Normal file
119
tests/test_plugin_lifecycle_concurrency.py
Normal file
@@ -0,0 +1,119 @@
|
||||
"""插件生命周期入口并发交错时的运行实例归属测试。"""
|
||||
|
||||
import threading
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
from app.runtime.extensions.plugin.lifecycle import PluginLifecycle
|
||||
|
||||
# 停止一侧在停机 hook 里等待加载一侧的预算。start 不持锁时加载会立刻完成并置位,
|
||||
# 等待即刻返回;start 持锁时加载被挡在锁外,这里必然等满。两种实现下交错顺序都是
|
||||
# 确定的,用例不依赖线程被调度的先后。
|
||||
_INTERLEAVE_BUDGET = 0.3
|
||||
# 线程握手与回收的宽松预算,仅用于避免实现回归时把整个测试进程挂死
|
||||
_HANDSHAKE_BUDGET = 5.0
|
||||
|
||||
|
||||
class _SeedInstance:
|
||||
"""代表上一轮加载留下的运行实例,并在停机 hook 中把执行权让给加载线程。"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
stopped: list,
|
||||
quiesce_entered: threading.Event,
|
||||
start_finished: threading.Event,
|
||||
) -> None:
|
||||
"""记录停机事实,并保存与加载线程握手用的两个事件。"""
|
||||
self._stopped = stopped
|
||||
self._quiesce_entered = quiesce_entered
|
||||
self._start_finished = start_finished
|
||||
|
||||
def stop_service(self) -> None:
|
||||
"""登记自身已被停止,并在停止事务中途给加载线程留出执行窗口。"""
|
||||
self._stopped.append(self)
|
||||
self._quiesce_entered.set()
|
||||
self._start_finished.wait(_INTERLEAVE_BUDGET)
|
||||
|
||||
|
||||
def _plugin_class(constructed: list) -> type:
|
||||
"""构造满足最小生命周期合同、并登记自身实例化事实的插件类。"""
|
||||
|
||||
class DemoPlugin:
|
||||
"""提供并发交错用例需要的最小插件行为。"""
|
||||
|
||||
plugin_name = "演示插件"
|
||||
plugin_version = "1.0.0"
|
||||
|
||||
def init_plugin(self, _config) -> None:
|
||||
"""登记本实例已完成初始化,等价于其后台服务已经拉起。"""
|
||||
constructed.append(self)
|
||||
|
||||
@staticmethod
|
||||
def get_state() -> bool:
|
||||
"""并发用例只关心启用态插件,恒定返回启用。"""
|
||||
return True
|
||||
|
||||
return DemoPlugin
|
||||
|
||||
|
||||
def _lifecycle(*, plugin_class: type, running: dict) -> PluginLifecycle:
|
||||
"""构造隔离事件、模块清理和数据库的生命周期实例。"""
|
||||
return PluginLifecycle(
|
||||
classes={"DemoPlugin": plugin_class},
|
||||
running=running,
|
||||
load_plugins=lambda _plugin_id, _loadable, _check: [plugin_class],
|
||||
loadable_plugins=lambda: ["DemoPlugin"],
|
||||
plugin_config=lambda _plugin_id: {},
|
||||
auth_checker=lambda _plugin: True,
|
||||
clear_modules=MagicMock(),
|
||||
clear_tools=MagicMock(),
|
||||
enable_events=MagicMock(),
|
||||
disable_events=MagicMock(),
|
||||
runtime_status_writer=lambda _plugin_id, _status: None,
|
||||
database=lambda: MagicMock(),
|
||||
log=MagicMock(),
|
||||
event_sender=MagicMock(),
|
||||
)
|
||||
|
||||
|
||||
def test_concurrent_start_does_not_orphan_the_instance_it_loads() -> None:
|
||||
"""加载与停止交错时,被加载出来的实例不能既不在运行表也没被停止。
|
||||
|
||||
停止取的是进入时的运行表快照,卸载却按插件 ID 清运行表。start 不与 stop 互斥
|
||||
时,交错窗口里加载出来的实例会被这次卸载一并抹掉,而它的 stop_service 从未被
|
||||
调用——它注册的定时任务、线程和事件订阅继续在跑,宿主却再也拿不到句柄停它。
|
||||
"""
|
||||
constructed: list = []
|
||||
stopped: list = []
|
||||
quiesce_entered = threading.Event()
|
||||
start_finished = threading.Event()
|
||||
seed = _SeedInstance(
|
||||
stopped=stopped,
|
||||
quiesce_entered=quiesce_entered,
|
||||
start_finished=start_finished,
|
||||
)
|
||||
running: dict = {"DemoPlugin": seed}
|
||||
lifecycle = _lifecycle(plugin_class=_plugin_class(constructed), running=running)
|
||||
|
||||
def load() -> None:
|
||||
"""等停止进入停机 hook 之后再发起加载,制造确定的交错顺序。"""
|
||||
assert quiesce_entered.wait(_HANDSHAKE_BUDGET)
|
||||
lifecycle.start("DemoPlugin")
|
||||
start_finished.set()
|
||||
|
||||
loader = threading.Thread(target=load, name="plugin-start", daemon=True)
|
||||
loader.start()
|
||||
try:
|
||||
lifecycle.stop("DemoPlugin")
|
||||
finally:
|
||||
loader.join(_HANDSHAKE_BUDGET)
|
||||
|
||||
assert not loader.is_alive()
|
||||
assert stopped == [seed]
|
||||
assert len(constructed) == 1
|
||||
orphaned = [
|
||||
instance
|
||||
for instance in constructed
|
||||
if instance not in running.values() and instance not in stopped
|
||||
]
|
||||
assert orphaned == []
|
||||
@@ -3,14 +3,18 @@
|
||||
import importlib.util
|
||||
import sys
|
||||
import threading
|
||||
from contextlib import nullcontext
|
||||
from pathlib import Path
|
||||
from types import ModuleType, SimpleNamespace
|
||||
from unittest.mock import MagicMock
|
||||
|
||||
import pytest
|
||||
|
||||
import app.application.plugin.management as plugin_management
|
||||
from app.runtime.event.registry import EventRegistry
|
||||
from app.runtime.extensions.plugin.clone import PluginCloneService
|
||||
from app.runtime.extensions.plugin.loader import PluginLoader
|
||||
from app.runtime.extensions.plugin.manager import PluginManager
|
||||
from app.runtime.extensions.plugin.storage import (
|
||||
PluginInstanceDirectory,
|
||||
PluginInstanceStore,
|
||||
@@ -735,3 +739,80 @@ def test_clone_service_rolls_back_descriptor_and_config_after_load_failure():
|
||||
assert instances == {}
|
||||
assert configs == {"DemoPlugin": {"enabled": True}}
|
||||
assert removed == ["DemoPluginbroken"]
|
||||
|
||||
|
||||
def _tree_manager(reloaded: list[str]) -> MagicMock:
|
||||
"""构造带一个分身的插件管理器替身,重载树语义取自真实实现。
|
||||
|
||||
只替换重载动作本身,实例树的解析仍走 PluginManager 的真实方法,避免用例把
|
||||
「分身也被重载」断言在一个自己编造的树上。
|
||||
"""
|
||||
directory, records = _make_directory()
|
||||
records["DemoPluginwork"] = PluginInstance(
|
||||
instance_id="DemoPluginwork",
|
||||
source_plugin_id="DemoPlugin",
|
||||
plugin_name="工作实例",
|
||||
)
|
||||
storage, _written = _make_storage({})
|
||||
manager = MagicMock()
|
||||
manager._plugin_instance_store = PluginInstanceStore(
|
||||
storage=lambda: storage,
|
||||
directory=lambda: directory,
|
||||
)
|
||||
manager._plugin_quiesce_lock = threading.RLock()
|
||||
manager.mutation.side_effect = lambda _operation: nullcontext()
|
||||
manager.reload_plugin.side_effect = lambda plugin_id: (
|
||||
reloaded.append(plugin_id) or PluginRuntimeStatus.ACTIVE
|
||||
)
|
||||
manager.get_plugin_source_id.side_effect = lambda plugin_id: (
|
||||
PluginManager.get_plugin_source_id(manager, plugin_id)
|
||||
)
|
||||
manager.reload_plugin_tree.side_effect = lambda plugin_id: (
|
||||
PluginManager.reload_plugin_tree(manager, plugin_id)
|
||||
)
|
||||
manager.get_plugin_reload_targets.side_effect = lambda plugin_id: (
|
||||
PluginManager.get_plugin_reload_targets(manager, plugin_id)
|
||||
)
|
||||
return manager
|
||||
|
||||
|
||||
def test_manual_reload_covers_source_plugin_and_its_clones(monkeypatch):
|
||||
"""手工重载本体要连分身一起换掉旧类对象,并逐个刷新它们的注册。
|
||||
|
||||
只重载本体时,分身继续持有旧模块与旧类对象,改完源码点重载后新旧代码会在同一
|
||||
进程内并存,分身的 API、调度与命令注册也不会刷新。
|
||||
"""
|
||||
reloaded: list[str] = []
|
||||
refreshed: list[str] = []
|
||||
manager = _tree_manager(reloaded)
|
||||
monkeypatch.setattr(plugin_management, "get_plugin_manager", lambda: manager)
|
||||
monkeypatch.setattr(
|
||||
plugin_management,
|
||||
"refresh_plugin_registrations",
|
||||
refreshed.append,
|
||||
)
|
||||
|
||||
status = plugin_management.reload_plugin_runtime("DemoPlugin")
|
||||
|
||||
assert status is PluginRuntimeStatus.ACTIVE
|
||||
assert reloaded == ["DemoPlugin", "DemoPluginwork"]
|
||||
assert refreshed == ["DemoPlugin", "DemoPluginwork"]
|
||||
|
||||
|
||||
def test_manual_reload_of_a_clone_rebuilds_the_whole_instance_tree(monkeypatch):
|
||||
"""从分身发起的重载要回到源码本体,再带上同源的全部分身。"""
|
||||
reloaded: list[str] = []
|
||||
refreshed: list[str] = []
|
||||
manager = _tree_manager(reloaded)
|
||||
monkeypatch.setattr(plugin_management, "get_plugin_manager", lambda: manager)
|
||||
monkeypatch.setattr(
|
||||
plugin_management,
|
||||
"refresh_plugin_registrations",
|
||||
refreshed.append,
|
||||
)
|
||||
|
||||
status = plugin_management.reload_plugin_runtime("DemoPluginwork")
|
||||
|
||||
assert status is PluginRuntimeStatus.ACTIVE
|
||||
assert reloaded == ["DemoPlugin", "DemoPluginwork"]
|
||||
assert refreshed == ["DemoPlugin", "DemoPluginwork"]
|
||||
|
||||
Reference in New Issue
Block a user