diff --git a/pyproject.toml b/pyproject.toml index 473913b..ed6b00d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -16,7 +16,7 @@ dependencies = [ ] [project.optional-dependencies] -dev = ["pytest>=8.0", "pytest-asyncio>=0.23"] +dev = ["pytest>=8.0", "pytest-asyncio>=0.23", "fakeredis>=2.0"] [tool.setuptools.packages.find] where = ["src"] diff --git a/src/nmfs_agents/tools/event_bus.py b/src/nmfs_agents/tools/event_bus.py index 37f3864..4507f27 100644 --- a/src/nmfs_agents/tools/event_bus.py +++ b/src/nmfs_agents/tools/event_bus.py @@ -46,8 +46,10 @@ class EventBus: try: self._r.xgroup_create(self._stream, group, id="0", mkstream=True) except Exception as e: - if "BUSYGROUP" not in str(e): - logger.warning("xgroup_create %s: %s", group, e) + if "BUSYGROUP" in str(e): + logger.debug("consumer group %s 已存在", group) + else: + raise def publish(self, event: AgentEvent) -> str: data = { @@ -69,16 +71,12 @@ class EventBus: count: int = 10, block_ms: int = 0, ) -> list[AgentEvent]: - try: - results = self._r.xreadgroup( - consumer_group, consumer_name, - {self._stream: ">"}, - count=count, - block=block_ms or None, - ) - except Exception as e: - logger.warning("EventBus read_new 失败: %s", e) - return [] + results = self._r.xreadgroup( + consumer_group, consumer_name, + {self._stream: ">"}, + count=count, + block=block_ms or None, + ) events: list[AgentEvent] = [] if results: for _stream, messages in results: