add chat state

This commit is contained in:
Worm 2026-02-18 13:32:28 +08:00
parent 789f1d2d67
commit ead9c0d3e5
4 changed files with 123 additions and 5 deletions

View file

@ -500,6 +500,13 @@ PicoClaw 将数据存储在您配置的工作区中(默认:`~/.picoclaw/work
- 对安全敏感的 XMPP 渠道设置较短 TTL`1h`),降低明文历史泄露风险 - 对安全敏感的 XMPP 渠道设置较短 TTL`1h`),降低明文历史泄露风险
- 对相对不敏感或需要长期上下文的渠道保留 `"false"`,完全不自动清理 - 对相对不敏感或需要长期上下文的渠道保留 `"false"`,完全不自动清理
对于 XMPP 渠道Agent 会默认:
- 在处理用户消息时发送「正在输入 / active」状态XEP0085方便前端展示输入指示
- 在收到带 `<request xmlns="urn:xmpp:receipts"/>` 的消息时自动发送 `<received/>` 回执XEP0184
这两个行为都是**默认开启且向后兼容**的:旧版 `config.json` 不需要增加任何字段即可使用,新版客户端若不支持相应 XEP 也会直接忽略这些附加元素。
### 心跳 / 周期性任务 (Heartbeat) ### 心跳 / 周期性任务 (Heartbeat)
PicoClaw 可以自动执行周期性任务。在工作区创建 `HEARTBEAT.md` 文件: PicoClaw 可以自动执行周期性任务。在工作区创建 `HEARTBEAT.md` 文件:

View file

@ -275,17 +275,19 @@ func (al *AgentLoop) processMessage(ctx context.Context, msg bus.InboundMessage)
"session_key": msg.SessionKey, "session_key": msg.SessionKey,
}) })
// Route system messages to processSystemMessage
if msg.Channel == "system" { if msg.Channel == "system" {
return al.processSystemMessage(ctx, msg) return al.processSystemMessage(ctx, msg)
} }
// Check for commands
if response, handled := al.handleCommand(ctx, msg); handled { if response, handled := al.handleCommand(ctx, msg); handled {
return response, nil return response, nil
} }
// Process as user message if al.channelManager != nil && msg.Channel != "" && msg.ChatID != "" && !constants.IsInternalChannel(msg.Channel) {
al.channelManager.NotifyTyping(ctx, msg.Channel, msg.ChatID, true)
defer al.channelManager.NotifyTyping(ctx, msg.Channel, msg.ChatID, false)
}
return al.runAgentLoop(ctx, processOptions{ return al.runAgentLoop(ctx, processOptions{
SessionKey: msg.SessionKey, SessionKey: msg.SessionKey,
Channel: msg.Channel, Channel: msg.Channel,

View file

@ -17,6 +17,11 @@ import (
"github.com/sipeed/picoclaw/pkg/logger" "github.com/sipeed/picoclaw/pkg/logger"
) )
type TypingChannel interface {
Channel
SendTyping(ctx context.Context, chatID string, composing bool) error
}
type Manager struct { type Manager struct {
channels map[string]Channel channels map[string]Channel
bus *bus.MessageBus bus *bus.MessageBus
@ -339,6 +344,20 @@ func (m *Manager) UnregisterChannel(name string) {
delete(m.channels, name) delete(m.channels, name)
} }
func (m *Manager) NotifyTyping(ctx context.Context, channelName, chatID string, composing bool) {
m.mu.RLock()
channel, exists := m.channels[channelName]
m.mu.RUnlock()
if !exists {
return
}
typingChannel, ok := channel.(TypingChannel)
if !ok {
return
}
_ = typingChannel.SendTyping(ctx, chatID, composing)
}
func (m *Manager) SendToChannel(ctx context.Context, channelName, chatID, content string) error { func (m *Manager) SendToChannel(ctx context.Context, channelName, chatID, content string) error {
m.mu.RLock() m.mu.RLock()
channel, exists := m.channels[channelName] channel, exists := m.channels[channelName]

View file

@ -47,6 +47,9 @@ type XMPPChannel struct {
lastTime time.Time lastTime time.Time
} }
const chatStatesNS = "http://jabber.org/protocol/chatstates"
const receiptsNS = "urn:xmpp:receipts"
func NewXMPPChannel(cfg config.XMPPConfig, messageBus *bus.MessageBus) (*XMPPChannel, error) { func NewXMPPChannel(cfg config.XMPPConfig, messageBus *bus.MessageBus) (*XMPPChannel, error) {
if cfg.JID == "" { if cfg.JID == "" {
return nil, fmt.Errorf("xmpp jid is required") return nil, fmt.Errorf("xmpp jid is required")
@ -269,12 +272,24 @@ func (c *XMPPChannel) sendChatMessage(to jid.JID, body string, desc string, medi
payload = xmlstream.MultiReader(children...) payload = xmlstream.MultiReader(children...)
} }
stateElem := xmlstream.Wrap(
nil,
xml.StartElement{
Name: xml.Name{Local: "active"},
Attr: []xml.Attr{
{Name: xml.Name{Local: "xmlns"}, Value: chatStatesNS},
},
},
)
fullPayload := xmlstream.MultiReader(payload, stateElem)
st := stanza.Message{ st := stanza.Message{
Type: stanza.ChatMessage, Type: stanza.ChatMessage,
To: to, To: to,
} }
return c.sendStanza(st.Wrap(payload)) return c.sendStanza(st.Wrap(fullPayload))
} }
func (c *XMPPChannel) sendStanza(r xml.TokenReader) error { func (c *XMPPChannel) sendStanza(r xml.TokenReader) error {
@ -302,7 +317,10 @@ func (c *XMPPChannel) sendStanza(r xml.TokenReader) error {
func (c *XMPPChannel) handleIncomingMessage(msg stanza.Message, t xmlstream.TokenReadEncoder) error { func (c *XMPPChannel) handleIncomingMessage(msg stanza.Message, t xmlstream.TokenReadEncoder) error {
var payload struct { var payload struct {
Body string `xml:"body"` Body string `xml:"body"`
Request *struct {
XMLName xml.Name `xml:"urn:xmpp:receipts request"`
} `xml:"urn:xmpp:receipts request"`
} }
d := xml.NewTokenDecoder(t) d := xml.NewTokenDecoder(t)
@ -312,9 +330,16 @@ func (c *XMPPChannel) handleIncomingMessage(msg stanza.Message, t xmlstream.Toke
content := strings.TrimSpace(payload.Body) content := strings.TrimSpace(payload.Body)
if content == "" { if content == "" {
if payload.Request != nil && msg.ID != "" {
go c.sendDeliveryReceipt(msg)
}
return nil return nil
} }
if payload.Request != nil && msg.ID != "" {
go c.sendDeliveryReceipt(msg)
}
fromBare := msg.From.Bare().String() fromBare := msg.From.Bare().String()
chatID := msg.From.String() chatID := msg.From.String()
@ -342,6 +367,71 @@ func (c *XMPPChannel) handleIncomingMessage(msg stanza.Message, t xmlstream.Toke
return nil return nil
} }
func (c *XMPPChannel) sendChatState(to jid.JID, state string) error {
if !c.IsRunning() || c.session == nil {
return nil
}
elem := xmlstream.Wrap(
nil,
xml.StartElement{
Name: xml.Name{Local: state},
Attr: []xml.Attr{
{Name: xml.Name{Local: "xmlns"}, Value: chatStatesNS},
},
},
)
msg := stanza.Message{
Type: stanza.ChatMessage,
To: to,
}
return c.sendStanza(msg.Wrap(elem))
}
func (c *XMPPChannel) SendTyping(ctx context.Context, chatID string, composing bool) error {
if !c.IsRunning() || c.session == nil {
return nil
}
to, err := jid.Parse(chatID)
if err != nil {
return err
}
state := "active"
if composing {
state = "composing"
}
return c.sendChatState(to, state)
}
func (c *XMPPChannel) sendDeliveryReceipt(orig stanza.Message) error {
if !c.IsRunning() || c.session == nil || orig.ID == "" {
return nil
}
payload := xmlstream.Wrap(
nil,
xml.StartElement{
Name: xml.Name{Local: "received"},
Attr: []xml.Attr{
{Name: xml.Name{Local: "xmlns"}, Value: receiptsNS},
{Name: xml.Name{Local: "id"}, Value: orig.ID},
},
},
)
reply := stanza.Message{
Type: orig.Type,
To: orig.From,
}
return c.sendStanza(reply.Wrap(payload))
}
func (c *XMPPChannel) discoverUploadService(ctx context.Context) (jid.JID, error) { func (c *XMPPChannel) discoverUploadService(ctx context.Context) (jid.JID, error) {
c.uploadMu.Lock() c.uploadMu.Lock()
defer c.uploadMu.Unlock() defer c.uploadMu.Unlock()