fix(persistence): separate retryable rejection errors (#6423)

This commit is contained in:
InfinityPacer
2026-08-23 20:10:59 +08:00
committed by GitHub
parent 46d5c8538e
commit c6ffc309d4
8 changed files with 79 additions and 48 deletions
+2 -2
View File
@@ -57,7 +57,7 @@ from app.api.dependencies.plugin import (
from app.adapters.external.server import MoviePilotServerHelper from app.adapters.external.server import MoviePilotServerHelper
from app.adapters.external.market import PluginHelper from app.adapters.external.market import PluginHelper
from app.adapters.system.plugin.package import PluginPackageManager from app.adapters.system.plugin.package import PluginPackageManager
from app.schemas.exception import DatabaseWorkerOverloadedError from app.schemas.exception import PersistenceUnavailableError
from app.runtime.log import logger from app.runtime.log import logger
from app.schemas.types import SystemConfigKey from app.schemas.types import SystemConfigKey
from app.api.context import get_background_task_registry, resolve_background_task_registry from app.api.context import get_background_task_registry, resolve_background_task_registry
@@ -893,7 +893,7 @@ async def save_plugin_folders(
folders, folders,
) )
return _SchemaResponse(success=True) return _SchemaResponse(success=True)
except DatabaseWorkerOverloadedError: except PersistenceUnavailableError:
raise raise
except Exception as e: except Exception as e:
logger.error(f"[文件夹API] 保存文件夹配置失败: {str(e)}") logger.error(f"[文件夹API] 保存文件夹配置失败: {str(e)}")
+5 -6
View File
@@ -12,10 +12,7 @@ from app.application.database import (
AsyncDatabaseExecutor, AsyncDatabaseExecutor,
) )
from app.schemas.agent import AgentChatSessionDetail, AgentChatSessionSummary from app.schemas.agent import AgentChatSessionDetail, AgentChatSessionSummary
from app.schemas.exception import ( from app.schemas.exception import AgentChatPersistenceUnavailableError
DatabaseWorkerClosedError,
DatabaseWorkerOverloadedError,
)
from app.runtime.observability import record_metric from app.runtime.observability import record_metric
@@ -364,14 +361,16 @@ class AgentChatPersistenceService:
"""在线程 worker 内完成同步写入并丢弃仓储对象返回值。""" """在线程 worker 内完成同步写入并丢弃仓储对象返回值。"""
# 同时限制全局和单会话等待量,避免一个热点会话占满总 admission 后饿死其他会话。 # 同时限制全局和单会话等待量,避免一个热点会话占满总 admission 后饿死其他会话。
if self._closing: if self._closing:
raise DatabaseWorkerClosedError("AgentChat 持久化服务当前不可接收任务") raise AgentChatPersistenceUnavailableError(
"AgentChat 持久化服务当前不可接收任务"
)
session_pending = self._pending_by_session.get(session_id, 0) session_pending = self._pending_by_session.get(session_id, 0)
if ( if (
self._pending_writes >= self._capacity self._pending_writes >= self._capacity
or session_pending >= self._session_capacity or session_pending >= self._session_capacity
): ):
record_metric("agent.chat.persistence.rejected") record_metric("agent.chat.persistence.rejected")
raise DatabaseWorkerOverloadedError( raise AgentChatPersistenceUnavailableError(
f"AgentChat 写入容量已用尽(全局上限 {self._capacity}" f"AgentChat 写入容量已用尽(全局上限 {self._capacity}"
f"单会话上限 {self._session_capacity}" f"单会话上限 {self._session_capacity}"
) )
+6 -6
View File
@@ -7,7 +7,7 @@ from collections.abc import Awaitable, Callable
from dataclasses import dataclass, field from dataclasses import dataclass, field
from typing import Any, Optional from typing import Any, Optional
from app.schemas.exception import DatabaseWorkerOverloadedError from app.schemas.exception import PersistenceUnavailableError
from app.application.plugin.lifecycle import plugin_lifecycle from app.application.plugin.lifecycle import plugin_lifecycle
from app.runtime.log import logger from app.runtime.log import logger
@@ -196,7 +196,7 @@ class PluginInstallCommand:
message=str(err), message=str(err),
package_installed=False, package_installed=False,
) )
if isinstance(err, DatabaseWorkerOverloadedError): if isinstance(err, PersistenceUnavailableError):
raise raise
return result return result
if not package_installed: if not package_installed:
@@ -228,7 +228,7 @@ class PluginInstallCommand:
package_installed=True, package_installed=True,
installed_list_persisted=state.installed_list_touched, installed_list_persisted=state.installed_list_touched,
) )
if isinstance(err, DatabaseWorkerOverloadedError): if isinstance(err, PersistenceUnavailableError):
raise raise
return result return result
@@ -247,7 +247,7 @@ class PluginInstallCommand:
installed_list_persisted=installed_list_persisted, installed_list_persisted=installed_list_persisted,
runtime_touched=True, runtime_touched=True,
) )
if isinstance(err, DatabaseWorkerOverloadedError): if isinstance(err, PersistenceUnavailableError):
raise raise
return result return result
@@ -267,7 +267,7 @@ class PluginInstallCommand:
runtime_touched=True, runtime_touched=True,
registrations_touched=True, registrations_touched=True,
) )
if isinstance(err, DatabaseWorkerOverloadedError): if isinstance(err, PersistenceUnavailableError):
raise raise
return result return result
@@ -412,7 +412,7 @@ class PluginInstallCommand:
failure_stage=failure_stage, failure_stage=failure_stage,
rollback=rollback, rollback=rollback,
) )
if isinstance(err, DatabaseWorkerOverloadedError): if isinstance(err, PersistenceUnavailableError):
raise raise
return result return result
+6 -11
View File
@@ -15,8 +15,7 @@ from app.adapters.web.plugin.routes import FastAPIDynamicRouteRegistry
from app.adapters.web.health import install_health_routes from app.adapters.web.health import install_health_routes
from app.application.plugin.routes import configure_plugin_routes from app.application.plugin.routes import configure_plugin_routes
from app.schemas.exception import ( from app.schemas.exception import (
DatabaseWorkerClosedError, PersistenceUnavailableError,
DatabaseWorkerOverloadedError,
) )
from app.adapters.web.security.access import ( from app.adapters.web.security.access import (
configure_token_codec, configure_token_codec,
@@ -235,11 +234,11 @@ async def localized_http_exception_handler(
) )
async def database_worker_overloaded_handler( async def persistence_unavailable_handler(
request: Request, request: Request,
_exc: DatabaseWorkerClosedError | DatabaseWorkerOverloadedError, _exc: PersistenceUnavailableError,
) -> JSONResponse: ) -> JSONResponse:
"""数据库 worker 暂不可用映射为可重试的 503 响应。""" """持久化能力暂不可用映射为可重试的 503 响应。"""
return await localized_http_exception_handler( return await localized_http_exception_handler(
request, request,
HTTPException( HTTPException(
@@ -335,12 +334,8 @@ def create_app() -> FastAPI:
_app.add_exception_handler(HTTPException, localized_http_exception_handler) _app.add_exception_handler(HTTPException, localized_http_exception_handler)
_app.add_exception_handler( _app.add_exception_handler(
DatabaseWorkerOverloadedError, PersistenceUnavailableError,
database_worker_overloaded_handler, persistence_unavailable_handler,
)
_app.add_exception_handler(
DatabaseWorkerClosedError,
database_worker_overloaded_handler,
) )
_app.add_exception_handler( _app.add_exception_handler(
RequestValidationError, RequestValidationError,
+10 -2
View File
@@ -58,14 +58,22 @@ class StorageQueryError(Exception):
pass pass
class DatabaseWorkerClosedError(RuntimeError): class PersistenceUnavailableError(RuntimeError):
"""持久化执行能力暂时拒绝新操作,调用方可在稍后重试。"""
class DatabaseWorkerClosedError(PersistenceUnavailableError):
"""数据库执行器尚未启动或已经停止。""" """数据库执行器尚未启动或已经停止。"""
class DatabaseWorkerOverloadedError(RuntimeError): class DatabaseWorkerOverloadedError(PersistenceUnavailableError):
"""数据库执行器的运行与排队容量已经用尽。""" """数据库执行器的运行与排队容量已经用尽。"""
class AgentChatPersistenceUnavailableError(PersistenceUnavailableError):
"""AgentChat 持久化服务因关闭或自身容量限制拒绝新写入。"""
class TMDbException(Exception): class TMDbException(Exception):
""" """
用于表示TheMovieDB数据源请求失败的跨层异常契约。 用于表示TheMovieDB数据源请求失败的跨层异常契约。
+4 -4
View File
@@ -12,7 +12,7 @@ import pytest
from sqlalchemy import delete, select from sqlalchemy import delete, select
from app.schemas.exception import ( from app.schemas.exception import (
DatabaseWorkerClosedError, AgentChatPersistenceUnavailableError,
DatabaseWorkerOverloadedError, DatabaseWorkerOverloadedError,
) )
from app.application.messaging.chat import AgentChatPersistenceService, AgentChatService from app.application.messaging.chat import AgentChatPersistenceService, AgentChatService
@@ -302,7 +302,7 @@ async def test_agent_chat_persistence_bounds_session_waiters_and_releases_cancel
messages=[], messages=[],
) )
) )
with pytest.raises(DatabaseWorkerOverloadedError): with pytest.raises(AgentChatPersistenceUnavailableError):
await third await third
second.cancel() second.cancel()
with pytest.raises(asyncio.CancelledError): with pytest.raises(asyncio.CancelledError):
@@ -347,7 +347,7 @@ async def test_agent_chat_persistence_session_admission_is_fair() -> None:
) )
) )
await asyncio.sleep(0) await asyncio.sleep(0)
with pytest.raises(DatabaseWorkerOverloadedError): with pytest.raises(AgentChatPersistenceUnavailableError):
await service.async_save_agent_messages( await service.async_save_agent_messages(
session_id="hot-session", user_id="1", messages=[] session_id="hot-session", user_id="1", messages=[]
) )
@@ -393,7 +393,7 @@ async def test_agent_chat_persistence_shutdown_drains_active_writes() -> None:
shutdown = asyncio.create_task(service.shutdown()) shutdown = asyncio.create_task(service.shutdown())
await asyncio.sleep(0) await asyncio.sleep(0)
assert not shutdown.done() assert not shutdown.done()
with pytest.raises(DatabaseWorkerClosedError): with pytest.raises(AgentChatPersistenceUnavailableError):
await service.async_save_agent_messages( await service.async_save_agent_messages(
session_id="new-session", user_id="1", messages=[] session_id="new-session", user_id="1", messages=[]
) )
+31 -12
View File
@@ -18,14 +18,16 @@ from app.api.response import (
ResponseAPIRouter, ResponseAPIRouter,
) )
from app.factory import ( from app.factory import (
database_worker_overloaded_handler,
localized_http_exception_handler, localized_http_exception_handler,
localized_unhandled_exception_handler, localized_unhandled_exception_handler,
localized_validation_exception_handler, localized_validation_exception_handler,
persistence_unavailable_handler,
) )
from app.schemas.exception import ( from app.schemas.exception import (
AgentChatPersistenceUnavailableError,
DatabaseWorkerClosedError, DatabaseWorkerClosedError,
DatabaseWorkerOverloadedError, DatabaseWorkerOverloadedError,
PersistenceUnavailableError,
) )
from app.runtime.localization import LocaleHelper from app.runtime.localization import LocaleHelper
from app.runtime.config import settings from app.runtime.config import settings
@@ -67,12 +69,8 @@ def api_app() -> FastAPI:
app.router.route_class = ResponseAPIRoute app.router.route_class = ResponseAPIRoute
app.add_exception_handler(HTTPException, localized_http_exception_handler) app.add_exception_handler(HTTPException, localized_http_exception_handler)
app.add_exception_handler( app.add_exception_handler(
DatabaseWorkerOverloadedError, PersistenceUnavailableError,
database_worker_overloaded_handler, persistence_unavailable_handler,
)
app.add_exception_handler(
DatabaseWorkerClosedError,
database_worker_overloaded_handler,
) )
from fastapi.exceptions import RequestValidationError from fastapi.exceptions import RequestValidationError
@@ -137,6 +135,11 @@ def api_app() -> FastAPI:
"""模拟数据库 worker 在关闭态拒绝新任务。""" """模拟数据库 worker 在关闭态拒绝新任务。"""
raise DatabaseWorkerClosedError("worker closed") raise DatabaseWorkerClosedError("worker closed")
@app.get("/agent-chat-persistence-unavailable")
async def get_agent_chat_persistence_unavailable() -> None:
"""模拟 AgentChat 自身 admission 拒绝新写入。"""
raise AgentChatPersistenceUnavailableError("agent persistence full")
@app.get("/native", response_model=None) @app.get("/native", response_model=None)
async def get_native_response() -> dict[str, bool]: async def get_native_response() -> dict[str, bool]:
"""返回显式旁路的原生 JSON 协议。""" """返回显式旁路的原生 JSON 协议。"""
@@ -247,14 +250,30 @@ async def test_database_worker_closed_is_retryable_service_unavailable(
} }
def test_create_app_registers_closed_database_worker_handler() -> None: async def test_agent_chat_persistence_rejection_is_retryable_service_unavailable(
"""生产组合根必须为 worker 关闭态登记 503 处理器。""" api_app: FastAPI,
) -> None:
"""AgentChat 自身 admission 拒绝也应返回可重试的 503。"""
async with make_client(api_app) as client:
response = await client.get("/agent-chat-persistence-unavailable")
assert response.status_code == 503
assert response.headers["retry-after"] == "1"
assert response.json() == {
"success": False,
"message": "服务当前繁忙,请稍后重试",
"data": None,
}
def test_create_app_registers_persistence_unavailable_handler() -> None:
"""生产组合根必须为持久化暂不可用登记统一 503 处理器。"""
from app.factory import create_app from app.factory import create_app
app = create_app() app = create_app()
assert app.exception_handlers[DatabaseWorkerClosedError] is ( assert app.exception_handlers[PersistenceUnavailableError] is (
database_worker_overloaded_handler persistence_unavailable_handler
) )
@@ -283,7 +302,7 @@ async def test_database_worker_overload_preserves_retry_after_for_native_protoco
"client": ("testclient", 123), "client": ("testclient", 123),
"root_path": "", "root_path": "",
} }
response = await database_worker_overloaded_handler( response = await persistence_unavailable_handler(
Request(scope), Request(scope),
DatabaseWorkerOverloadedError("worker full"), DatabaseWorkerOverloadedError("worker full"),
) )
+15 -5
View File
@@ -3,7 +3,11 @@ from unittest.mock import AsyncMock, Mock, patch
import pytest import pytest
from app.schemas.exception import DatabaseWorkerOverloadedError from app.schemas.exception import (
DatabaseWorkerClosedError,
DatabaseWorkerOverloadedError,
PersistenceUnavailableError,
)
from app.application.plugin.install import PluginInstallCommand from app.application.plugin.install import PluginInstallCommand
@@ -249,17 +253,23 @@ async def test_persistence_exception_after_write_restores_installed_list():
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_database_worker_overload_rolls_back_and_reaches_api_boundary(): @pytest.mark.parametrize(
"""配置 worker 背压完成补偿后继续抛出,交由 API 映射为 503。""" "error_type",
[DatabaseWorkerClosedError, DatabaseWorkerOverloadedError],
)
async def test_persistence_unavailable_rolls_back_and_reaches_api_boundary(
error_type: type[PersistenceUnavailableError],
) -> None:
"""持久化能力暂不可用时完成补偿并交由 API 映射为 503。"""
checkpoint = object() checkpoint = object()
rollback = AsyncMock() rollback = AsyncMock()
command = _command( command = _command(
checkpointer=AsyncMock(return_value=checkpoint), checkpointer=AsyncMock(return_value=checkpoint),
writer=AsyncMock(side_effect=DatabaseWorkerOverloadedError("worker full")), writer=AsyncMock(side_effect=error_type("persistence unavailable")),
rollback=rollback, rollback=rollback,
) )
with pytest.raises(DatabaseWorkerOverloadedError): with pytest.raises(error_type):
await command.execute( await command.execute(
plugin_id="DemoPlugin", plugin_id="DemoPlugin",
repo_url="https://github.com/demo/plugins", repo_url="https://github.com/demo/plugins",