Коммуникация между модулями
Note
В этой главе требуется ErisPulse 2.8.0+.
В ErisPulse между модулями реализована трехуровневая модель коммуникации, расположенная по принципу "точка-точка → направленный → широковещательный":
| Уровень | API | Семантика | Типичные сценарии |
|---|---|---|---|
| RPC | await sdk.module.call("Chat", "get_history", ...) |
Точечный запрос-ответ с контрактом / аудитом / таймаутом | Вызов возможностей другого модуля (получение истории, перевод, возврат средств) |
| Направленный событие | await lifecycle.emit("message_received", {...}, to="Chat") |
Доставка только модулям, зарегистрировавшим хуки жизненного цикла | Уведомление о состоянии сверху вниз (например, "получено новое сообщение") |
| Рассылка | await lifecycle.emit("config.updated", {...}) |
События жизненного цикла, видимые по всей системе | Горячая перенастройка конфигурации, изменение состояния модуля |
{!--< tips >!--}
Правило выбора: для получения значения используйте call, для уведомления конкретного модуля — emit(..., to=...), для уведомления всех — emit(...).
{!--< /tips >!--}
RPC: module.call
result = await sdk.module.call("Chat", "get_history", session_id, n=20)
Отличие от прямого доступа к атрибуту sdk.module.Chat.get_history(...) (без изменений):
module.call() |
Прямой доступ к атрибуту | |
|---|---|---|
| Целевой модуль не зарегистрирован / не активен | Бросает ModuleNotAvailableError |
Бросает AttributeError |
| Модуль с ленивой загрузкой | Автоматически активируется (модуль на основе событий использует блокировку активации) | Асинхронная инициализация модуля бросает RuntimeError |
current_owner |
Относится к целевому модулю (его внутренние wait_reply / отправка / логи корректно привязаны) | Сохраняется вызывающей стороной |
| Таймаут | По умолчанию 30 секунд (timeout= переопределяет, None — без ограничения) |
Нет |
| Аудит в рамках scope | Вызов проходит через выходной барьер actions.<вызывающий>.call |
Нет |
| Проверка контракта | Белый список meta.services |
Нет |
Иерархия исключений
ModuleError # Базовый класс исключений модульной системы
└── ModuleCallError # Базовый класс для вызова между модулями (содержит атрибуты module / method)
├── ModuleNotAvailableError # Целевой модуль не зарегистрирован / не активен / активация не удалась
├── ServiceNotProvidedError # Метод не входит в белый список services / приватный метод / не существует
└── ModuleCallTimeoutError # Таймаут асинхронного метода
Все исключения входят в иерархию ErisPulseError, их можно перехватить с помощью from ErisPulse.Core import ModuleCallError.
Контракт сервиса: meta.services
Провайдер сервиса объявляет белый список предоставляемых методов в get_meta() (симметрично commands):
from ErisPulse.Core.Bases import BaseModule, ModuleMeta
class ChatModule(BaseModule):
@staticmethod
def get_meta() -> ModuleMeta:
return ModuleMeta(
name="聊天",
services=[
"get_history", # Простая форма
{"name": "translate", "description": "把文本翻译成指定语言"}, # С описанием
],
)
async def get_history(self, session_id, n=20): ...
async def translate(self, text, target_lang): ...
def _internal_helper(self): ... # Методы с подчеркиванием всегда запрещены для внешнего вызова
Разработчикам не нужно ничего делать по умолчанию:
- Если
servicesне объявлен → все публичные методы по умолчанию доступны дляmodule.call()(как и прямой доступ к атрибуту), без каких-либо объявления - После объявления → доступ сужается до белого списка, вызов за пределами списка бросает
ServiceNotProvidedError— используется для маркировки "это те методы, которые публично доступны" - Основной контроль за ограничениями находится на стороне пользователя:
scope.actionsопределяет, "кто может вызывать кого" (см. аудит ниже),servicesмодуля автора — это только объявление интерфейса, две системы не заменяют друг друга
Описание сервиса: добавьте человекопонятное описание для каждого сервиса — если не нужно, ничего не пишите, описание автоматически берется из первой строки docstring (рамка уже требует стиля docstring):
async def translate(self, text, target_lang):
"""把文本翻译成指定语言""" # ← Эта строка автоматически становится описанием сервиса
...
Для тонкой настройки (переопределение docstring / многоязычность) используйте словарную форму объявления description (поддержка i18n-словарей):
services=[
{"name": "translate", "description": "把文本翻译成指定语言"},
{"name": "summarize", "description": {"i18n": "Chat.meta.svc.summarize", "default": "摘要对话"}},
]
Служебный каталог: services()
sdk.module.services()
# {'Chat': [{'name': 'get_history', 'signature': '(session_id, n=20)',
# 'description': '把文本翻译成指定语言'}]}
sdk.module.services("Chat") # Только для запроса указанного модуля
- В каталоге перечислены только модули, явно объявившие
meta.services(модули без объявления не попадают в каталог) - Каждый сервис снабжен строкой сигнатуры метода (извлекается с помощью
inspect.signature) и описанием - Также включены в топологию:
sdk.module.get_topology()содержит для каждого модуля полеservices
{!--< tips >!--}
Маршрут к MCP: каталог сервисов (имя + сигнатура + описание) соответствует форме MCP tool —
каждый сервис естественным образом имеет вид {"name", "description", "parameters"}.
В будущем фреймворк сможет напрямую выставить services() как конечную точку MCP server, позволяя AI обнаруживать и вызывать возможности модуля;
scope.actions.call аудит естественным образом становится безопасным барьером для вызова AI.
{!--< /tips >!--}
Аудит исходящих вызовов: кто может вызывать кого
Каждый вызов module.call() проходит через барьер исходящих вызовов в контексте вызывающего модуля:
[ErisPulse.scope.actions.CallerModule.call]
deny = ["Chat.get_history"] # Запретить CallerModule вызывать Chat.get_history
# allow = ["Chat.get_*"] # Или белый список: разрешить только вызовы Chat.get_* сервисов
- Формат
name—<目标模块>.<方法名>, поддерживается точное совпадение, шаблоны иre:регулярные выражения - Вызовы из фреймворка (без контекста owner, например, запуск скриптов) не подпадают под аудит
- Запрещённый вызов бросает
ModuleCallError(в логах TRACEcore.module.call_denied)
Способы конфигурирования см. в разделе Scope по исходящему измерению.
Направленные события: параметр to в lifecycle.emit
Жизненные события поддерживают направленную доставку: при указании параметра to с именем целевого владельца (owner) событие доставляется только хукам, зарегистрированным с этим
владельцем (хуки, зарегистрированные в on_load модуля, автоматически привязываются к нему),
другие модули и обработчики с * не получают уведомления.
from ErisPulse.Core.lifecycle import lifecycle
# Отправитель: событие доставляется только хукам, зарегистрированным в модуле Chat
await lifecycle.emit("message_received", {"text": "hi", "from": "u1"}, to="Chat")
# Подписчик (внутри модуля Chat): регистрирует хук с тем же именем, owner автоматически записывается при регистрации
@lifecycle.on("message_received")
async def on_message_received(data): ...
@lifecycle.on("message") # Префиксные хуки также работают (фильтруются по owner)
async def on_any(data): ...
Семантические детали:
- Если у целевого владельца нет зарегистрированных хуков → событие тихо отбрасывается (не отправляется туда, где его нет),
можно заранее проверить с помощью
lifecycle.has_handlers("message_received") - При передаче
dataв виде dict автоматически добавляется_trace_id(не перезаписывает уже существующее значение), интегрируется с системой трассировки - Рассылка и направленные события используют одну и ту же систему регистрации хуков:
emit(...)безto— рассылка по всей системе, сto— событие видно только целевому модулю emit_sync/submit_event(совместимые API) также поддерживают параметрto=
Note
Направленные события — это лёгкие уведомления, не выполняют проверку цели и ленивую активацию; если нужна проверка существования цели,
аудит контракта или возврат значения, используйте RPC: module.call.
Ленивая загрузка и вызов
module.call() прозрачно активирует модули с ленивой загрузкой:
- Модули на основе событий (
activate_onобъявлен) → проходят через блокировку активации_activate(), после активации stub-триггер автоматически удаляется - Обычные модули с ленивой загрузкой → синхронная инициализация или обычный путь загрузки (идемпотентно)
- Ошибка активации →
ModuleNotAvailableError
То есть: вызывающая сторона не должна заботиться о том, загружен ли целевой модуль, и не нужно ждать какого-либо события для его активации.
Направленные события (lifecycle.emit(..., to=...)) не выполняют ленивую активацию — если целевой модуль не загружен, у него нет хуков,
событие тихо отбрасывается; если нужно гарантировать доставку, используйте module.call().
Воспроизведение при холодном запуске
Новый или перезапущенный модуль пропустил часть диалога — get_load_strategy(replay=...) позволяет фреймворку воспроизвести
самому модулю последние сообщения из почтового ящика сессии после готовности:
from ErisPulse.loaders import ModuleLoadStrategy
class MyAIModule(BaseModule):
@staticmethod
def get_load_strategy():
return ModuleLoadStrategy(
lazy_load=False,
priority=100,
replay="5m", # Воспроизвести последние 5 минут (можно также "1h" / "300" секунд)
)
async def on_load(self, event):
@message.on_message()
async def handle(e):
if e.get("replayed"):
# Синтезированное событие: только восстановление контекста, без побочных эффектов, таких как отправка
...
Семантические детали:
- Источник данных — почтовый ящик сессии (
sdk.transcript.recent()), воспроизведение выполняется в фоне после загрузки модуля, не блокируя запуск - Синтезированное событие помечается флагом
replayed: True, содержит полныеplatform / detail_type / user_id / alt_message, доставляется только обработчикам этого модуля — другие модули не затрагиваются воспроизведением - При отсутствии почтового ящика, отсутствии записей или недопустимом значении
replay(предупреждениеreplay_invalid) воспроизведение пропускается
Идемпотентная дедупликация событий
После переподключения websocket-платформы часто повторно отправляются одни и те же события (с одинаковым event["id"]) — входная точка доставки по id выполняет
дедупликацию с помощью LRU (емкость 4096), одинаковые события доставляются только один раз.
[ErisPulse.framework]
event_dedupe = true # По умолчанию включено; в тестовой среде можно отключить, если синтезировать события с фиксированным id
При регистрации адаптера (начало жизненного цикла нового подключения) автоматически сбрасывается кэш дедупликации.
Связанные документы
- Система взаимодействия сессий - wait_reply / таймеры / многопоточные ожидания / взаимное исключение сессий
- Scope - полная конфигурация аудита по исходящему измерению
- Система принадлежности (owner) - как контекст owner проходит через межмодульные вызовы
- Система ленивой загрузки - ленивая загрузка и активация на основе событий (activate_on)
- Управление жизненным циклом - механизмы шин событий на уровне рассылки