feat(events): expose typed payload snapshots

This commit is contained in:
jxxghp
2026-08-25 22:32:33 +08:00
parent c8146ee395
commit 3ba4027d84
13 changed files with 1023 additions and 54 deletions
+81 -20
View File
@@ -37,6 +37,21 @@ class EventErrorBehavior(StrEnum):
NOTIFY = "notify"
class EventPayloadMode(StrEnum):
"""描述事件 payload 是运行时对象还是稳定快照。"""
RUNTIME = "runtime"
SNAPSHOT = "snapshot"
EXTENSION = "extension"
class EventValidationMode(StrEnum):
"""描述契约校验失败后的兼容策略。"""
DIAGNOSTIC = "diagnostic"
STRICT = "strict"
@dataclass(frozen=True, slots=True)
class EventContract:
"""冻结一个事件的 payload 与运行语义。"""
@@ -49,6 +64,11 @@ class EventContract:
delivery: EventDelivery
error_behavior: EventErrorBehavior
ordering: str
input_model: type[BaseModel] | None = None
output_model: type[BaseModel] | None = None
schema_version: int = 1
payload_mode: EventPayloadMode = EventPayloadMode.RUNTIME
validation_mode: EventValidationMode = EventValidationMode.DIAGNOSTIC
sensitive_fields: tuple[str, ...] = ()
legacy_reason: str | None = None
@@ -66,13 +86,13 @@ _PAYLOAD_MODELS: dict[EventType | ChainEventType, type[BaseModel]] = {
EventType.SubscribeAdded: event_schemas.SubscribeAddedEventData,
EventType.SubscribeDeleted: event_schemas.SubscribeDeletedEventData,
EventType.SubscribeModified: event_schemas.SubscribeModifiedEventData,
EventType.DownloadAdded: event_schemas.DownloadAddedEventData,
EventType.TransferComplete: event_schemas.TransferResultEventData,
EventType.TransferFailed: event_schemas.TransferResultEventData,
EventType.SubtitleTransferComplete: event_schemas.TransferResultEventData,
EventType.SubtitleTransferFailed: event_schemas.TransferResultEventData,
EventType.AudioTransferComplete: event_schemas.TransferResultEventData,
EventType.AudioTransferFailed: event_schemas.TransferResultEventData,
EventType.DownloadAdded: event_schemas.DownloadAddedContractData,
EventType.TransferComplete: event_schemas.TransferResultContractData,
EventType.TransferFailed: event_schemas.TransferResultContractData,
EventType.SubtitleTransferComplete: event_schemas.TransferResultContractData,
EventType.SubtitleTransferFailed: event_schemas.TransferResultContractData,
EventType.AudioTransferComplete: event_schemas.TransferResultContractData,
EventType.AudioTransferFailed: event_schemas.TransferResultContractData,
EventType.HistoryDeleted: event_schemas.HistoryDeletedEventData,
EventType.DownloadFileDeleted: event_schemas.DownloadFileDeletedEventData,
EventType.DownloadDeleted: event_schemas.DownloadDeletedEventData,
@@ -81,7 +101,7 @@ _PAYLOAD_MODELS: dict[EventType | ChainEventType, type[BaseModel]] = {
EventType.NoticeMessage: event_schemas.NoticeMessageEventData,
EventType.SubscribeComplete: event_schemas.SubscribeCompleteEventData,
EventType.SystemError: event_schemas.SystemErrorEventData,
EventType.MetadataScrape: event_schemas.MetadataScrapeEventData,
EventType.MetadataScrape: event_schemas.MetadataScrapeContractData,
EventType.ModuleReload: event_schemas.EmptyEventData,
EventType.MessageAction: event_schemas.MessageActionEventData,
EventType.WorkflowExecute: event_schemas.WorkflowExecuteEventData,
@@ -93,15 +113,15 @@ _PAYLOAD_MODELS: dict[EventType | ChainEventType, type[BaseModel]] = {
ChainEventType.TransferRenameBuild: event_schemas.TransferRenameBuildEventData,
ChainEventType.TransferIntercept: event_schemas.TransferInterceptEventData,
ChainEventType.TransferOverwriteCheck: event_schemas.TransferOverwriteCheckEventData,
ChainEventType.ResourceSelection: event_schemas.ResourceSelectionEventData,
ChainEventType.ResourceDownload: event_schemas.ResourceDownloadEventData,
ChainEventType.ResourceSelection: event_schemas.ResourceSelectionContractData,
ChainEventType.ResourceDownload: event_schemas.ResourceDownloadContractData,
ChainEventType.DiscoverSource: event_schemas.DiscoverSourceEventData,
ChainEventType.MediaRecognizeConvert: event_schemas.MediaRecognizeConvertEventData,
ChainEventType.RecommendSource: event_schemas.RecommendSourceEventData,
ChainEventType.StorageOperSelection: event_schemas.StorageOperSelectionEventData,
ChainEventType.AgentLLMProvider: event_schemas.AgentLLMProviderEventData,
ChainEventType.SubscribeEpisodesRefresh: event_schemas.SubscribeEpisodesRefreshEventData,
ChainEventType.SubscribeCompletionCheck: event_schemas.SubscribeCompletionCheckEventData,
ChainEventType.SubscribeCompletionCheck: event_schemas.SubscribeCompletionCheckContractData,
ChainEventType.NameRecognize: event_schemas.NameRecognizeEventData,
ChainEventType.MusicNameRecognize: event_schemas.MusicNameRecognizeEventData,
ChainEventType.MediaRecognize: event_schemas.MediaRecognizeEventData,
@@ -109,6 +129,32 @@ _PAYLOAD_MODELS: dict[EventType | ChainEventType, type[BaseModel]] = {
ChainEventType.WorkflowExecution: ActionContext,
}
_SNAPSHOT_EVENTS = {
EventType.DownloadAdded,
EventType.TransferComplete,
EventType.TransferFailed,
EventType.SubtitleTransferComplete,
EventType.SubtitleTransferFailed,
EventType.AudioTransferComplete,
EventType.AudioTransferFailed,
EventType.MetadataScrape,
ChainEventType.ResourceSelection,
ChainEventType.ResourceDownload,
ChainEventType.SubscribeCompletionCheck,
}
_INPUT_MODELS = {
ChainEventType.ResourceSelection: event_schemas.ResourceSelectionInputContractData,
ChainEventType.ResourceDownload: event_schemas.ResourceDownloadInputContractData,
ChainEventType.SubscribeCompletionCheck: event_schemas.SubscribeCompletionCheckInputContractData,
}
_OUTPUT_MODELS = {
ChainEventType.ResourceSelection: event_schemas.ResourceSelectionOutputContractData,
ChainEventType.ResourceDownload: event_schemas.ResourceDownloadOutputContractData,
ChainEventType.SubscribeCompletionCheck: event_schemas.SubscribeCompletionCheckOutputContractData,
}
_DURABLE_REQUIRED = {
EventType.SubscribeAdded,
EventType.SubscribeModified,
@@ -150,6 +196,14 @@ def _build_contract(event_type: EventType | ChainEventType) -> EventContract:
EventErrorBehavior.STOP_CHAIN if is_chain else EventErrorBehavior.NOTIFY
),
ordering="priority_serial" if is_chain else "priority_queue",
input_model=_INPUT_MODELS.get(event_type, payload_model),
output_model=_OUTPUT_MODELS.get(event_type),
payload_mode=(
EventPayloadMode.SNAPSHOT
if event_type in _SNAPSHOT_EVENTS
else EventPayloadMode.RUNTIME
),
validation_mode=EventValidationMode.DIAGNOSTIC,
sensitive_fields=_SENSITIVE_FIELDS.get(event_type, ()),
legacy_reason=(
None
@@ -195,16 +249,23 @@ def validate_event_payload(
) -> tuple[str, ...]:
"""在发送边界诊断首批 typed payload,保持原对象和插件 dict 形状不变。"""
try:
model = get_event_contract(event_type).payload_model
contract = get_event_contract(event_type)
models = tuple(
model
for model in (contract.input_model, contract.output_model)
if model is not None
)
except KeyError:
# 动态插件在旧 ABI 下可能使用宿主枚举之外的字符串事件。
return ()
if model is None or payload is None:
if not models or payload is None:
return ()
if isinstance(payload, model):
return ()
try:
model.model_validate(payload)
except (ValidationError, TypeError, ValueError) as error:
return (str(error),)
return ()
problems: list[str] = []
for model in models:
if isinstance(payload, model):
continue
try:
model.model_validate(payload)
except (ValidationError, TypeError, ValueError) as error:
problems.append(f"{model.__name__}: {error}")
return tuple(problems)
+116
View File
@@ -0,0 +1,116 @@
"""事件 payload 的只读契约快照访问。"""
from __future__ import annotations
from dataclasses import dataclass
from typing import Any
from pydantic import BaseModel, ValidationError
from app.runtime.event.contracts import get_event_contract, normalize_event_type
from app.schemas.types import ChainEventType, EventType
@dataclass(frozen=True, slots=True)
class EventPayloadSnapshot:
"""保存原始 payload 及按登记契约解析出的独立快照。"""
event_type: EventType | ChainEventType | str
raw: Any
payload: BaseModel | None
input: BaseModel | None
output: BaseModel | None
known: bool
schema_version: int | None
payload_mode: str | None
validation_mode: str | None
errors: tuple[str, ...] = ()
@property
def valid(self) -> bool:
"""返回已知事件的全部登记模型是否都通过校验。"""
return self.known and not self.errors
def _model_snapshot(
model: type[BaseModel] | None,
payload: Any,
label: str,
cache: dict[type[BaseModel], BaseModel | None],
errors: list[str],
) -> BaseModel | None:
"""按单个模型构造快照,同一模型在一次访问中只解析一次。"""
if model is None:
return None
if model in cache:
return cache[model]
if isinstance(payload, model):
cache[model] = payload.model_copy(deep=True)
return cache[model]
try:
cache[model] = model.model_validate(payload)
except (ValidationError, TypeError, ValueError) as error:
cache[model] = None
errors.append(f"{label}:{model.__name__}: {error}")
return cache[model]
def snapshot_event_data(
event_type: EventType | ChainEventType | str,
event_data: Any,
) -> EventPayloadSnapshot:
"""按事件注册契约解析 payload,未知自定义事件保持原始数据并返回未登记状态。"""
normalized = normalize_event_type(event_type)
try:
contract = get_event_contract(normalized)
except KeyError:
return EventPayloadSnapshot(
event_type=normalized,
raw=event_data,
payload=None,
input=None,
output=None,
known=False,
schema_version=None,
payload_mode=None,
validation_mode=None,
)
errors: list[str] = []
cache: dict[type[BaseModel], BaseModel | None] = {}
payload_snapshot = _model_snapshot(
contract.payload_model,
event_data,
"payload",
cache,
errors,
)
input_snapshot = _model_snapshot(
contract.input_model,
event_data,
"input",
cache,
errors,
)
output_snapshot = _model_snapshot(
contract.output_model,
event_data,
"output",
cache,
errors,
)
return EventPayloadSnapshot(
event_type=normalized,
raw=event_data,
payload=payload_snapshot,
input=input_snapshot,
output=output_snapshot,
known=True,
schema_version=contract.schema_version,
payload_mode=contract.payload_mode.value,
validation_mode=contract.validation_mode.value,
errors=tuple(errors),
)
__all__ = ["EventPayloadSnapshot", "snapshot_event_data"]
+18 -9
View File
@@ -7,26 +7,29 @@ import uuid
from contextvars import ContextVar
from dataclasses import dataclass
from queue import Empty, PriorityQueue
from typing import Callable, Dict, List, Optional, Tuple, Union, Any, Type
from typing import TYPE_CHECKING, Any, Callable, Dict, List, Optional, Tuple, Type, Union
from app.runtime.config import global_vars
from app.runtime.thread import ThreadHelper
from app.runtime.log import logger
from app.schemas.event import ChainEventData
from app.schemas.types import ChainEventType, EventType
from app.runtime.rate import ExponentialBackoffRateLimiter
from app.foundation.singleton import Singleton
from app.runtime.config import global_vars
from app.runtime.correlation import get_correlation_id
from app.runtime.event.binding import (
EventBindingResolver,
EventHandlerBinding,
HandlerInstanceResolver,
)
from app.runtime.event.contracts import normalize_event_type, validate_event_payload
from app.runtime.event.dispatch import EventDispatcher
from app.runtime.event.errors import EventErrorNotifier, EventErrorPolicy
from app.runtime.event.registry import EventRegistry
from app.runtime.event.contracts import normalize_event_type, validate_event_payload
from app.runtime.correlation import get_correlation_id
from app.runtime.log import logger
from app.runtime.observability import record_metric
from app.runtime.rate import ExponentialBackoffRateLimiter
from app.runtime.thread import ThreadHelper
from app.schemas.event import ChainEventData
from app.schemas.types import ChainEventType, EventType
if TYPE_CHECKING:
from app.runtime.event.snapshot import EventPayloadSnapshot
DEFAULT_EVENT_PRIORITY = 10 # 事件的默认优先级
MIN_EVENT_CONSUMER_THREADS = 1 # 最小事件消费者线程数
@@ -84,6 +87,12 @@ class Event:
event_name = getattr(self.event_type, "value", self.event_type)
return f"<{event_kind}: {event_name}, ID: {self.event_id}, Priority: {self.priority}>"
def snapshot(self) -> "EventPayloadSnapshot":
"""返回插件可安全读取的类型化 payload 快照,原始 event_data 保持不变。"""
from app.runtime.event.snapshot import snapshot_event_data
return snapshot_event_data(self.event_type, self.event_data)
def __lt__(self, other):
"""
定义事件对象的比较规则,基于优先级比较