69 lines
2.2 KiB
Python
69 lines
2.2 KiB
Python
"""
|
||
Skill 聊天总线 - 各 Skill 一起执行时的共享聊天与交互
|
||
"""
|
||
|
||
import time
|
||
import logging
|
||
from typing import Dict, Any, List, Optional
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
class SkillChatBus:
|
||
"""Skill 间共享消息总线,同一轮执行内可读写"""
|
||
|
||
def __init__(self):
|
||
self._messages: List[Dict[str, Any]] = []
|
||
self._session_id: Optional[str] = None
|
||
|
||
def clear(self, session_id: str = None):
|
||
"""清空当前会话,开始新一轮(如新一条自然语言命令)"""
|
||
self._messages.clear()
|
||
self._session_id = session_id or str(int(time.time() * 1000))
|
||
|
||
def append(
|
||
self,
|
||
from_skill: str,
|
||
message: str,
|
||
to_skill: Optional[str] = None,
|
||
data: Optional[Dict[str, Any]] = None,
|
||
):
|
||
"""Skill 或执行器发一条消息,其它 Skill 可读"""
|
||
self._messages.append({
|
||
"from_skill": from_skill,
|
||
"to_skill": to_skill,
|
||
"message": message,
|
||
"data": data or {},
|
||
"ts": time.time(),
|
||
"index": len(self._messages),
|
||
})
|
||
logger.debug(f"[SkillBus] {from_skill} -> {to_skill or 'all'}: {message[:50]}")
|
||
|
||
def get_messages(
|
||
self,
|
||
since_index: int = 0,
|
||
from_skill: Optional[str] = None,
|
||
to_skill: Optional[str] = None,
|
||
) -> List[Dict[str, Any]]:
|
||
"""取消息列表,可按发送方/接收方过滤"""
|
||
out = []
|
||
for m in self._messages:
|
||
if m["index"] < since_index:
|
||
continue
|
||
if from_skill and m.get("from_skill") != from_skill:
|
||
continue
|
||
if to_skill and m.get("to_skill") and m.get("to_skill") != to_skill:
|
||
continue
|
||
out.append(m)
|
||
return out
|
||
|
||
def get_last(self, from_skill: Optional[str] = None) -> Optional[Dict[str, Any]]:
|
||
"""取最后一条(可选指定来自哪个 Skill)"""
|
||
for m in reversed(self._messages):
|
||
if from_skill is None or m.get("from_skill") == from_skill:
|
||
return m
|
||
return None
|
||
|
||
def session_id(self) -> Optional[str]:
|
||
return self._session_id
|