fix: improve Signal channel after code review and testing
- Replace DMsEnabled/GroupsEnabled with GroupTrigger (mention_only, prefixes) - Add @mention detection for group chats (isBotMentioned, stripMention) - Fix sendReaction/sendTyping recipient type (string → []string) - Fix Send() error wrapping to preserve root cause - Add io.LimitReader guard on SSE error body read - Remove dead voice transcription code (deferred) - Remove compound senderID (redundant with SenderInfo) - Add isGroupChat safety comment - Update README config docs Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
d1c557a728
commit
b3a19f474e
3 changed files with 102 additions and 57 deletions
|
|
@ -492,8 +492,9 @@ Register or link a phone number following the [signal-cli docs](https://github.c
|
||||||
"account": "+1234567890",
|
"account": "+1234567890",
|
||||||
"signal_cli_url": "http://localhost:8080",
|
"signal_cli_url": "http://localhost:8080",
|
||||||
"allow_from": ["+1987654321"],
|
"allow_from": ["+1987654321"],
|
||||||
"dms_enabled": true,
|
"group_trigger": {
|
||||||
"groups_enabled": false
|
"mention_only": true
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -502,7 +503,7 @@ Register or link a phone number following the [signal-cli docs](https://github.c
|
||||||
- `account`: The phone number registered with signal-cli
|
- `account`: The phone number registered with signal-cli
|
||||||
- `signal_cli_url`: URL of the signal-cli REST API (default: `http://localhost:8080`)
|
- `signal_cli_url`: URL of the signal-cli REST API (default: `http://localhost:8080`)
|
||||||
- `allow_from`: Phone numbers allowed to interact (empty = allow all)
|
- `allow_from`: Phone numbers allowed to interact (empty = allow all)
|
||||||
- `dms_enabled` / `groups_enabled`: Toggle DM and group message handling
|
- `group_trigger`: Group chat trigger config — `mention_only` requires @mention, `prefixes` triggers on message prefixes (omit for respond-to-all)
|
||||||
|
|
||||||
**3. Run**
|
**3. Run**
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -22,7 +22,6 @@ import (
|
||||||
"github.com/sipeed/picoclaw/pkg/identity"
|
"github.com/sipeed/picoclaw/pkg/identity"
|
||||||
"github.com/sipeed/picoclaw/pkg/logger"
|
"github.com/sipeed/picoclaw/pkg/logger"
|
||||||
"github.com/sipeed/picoclaw/pkg/utils"
|
"github.com/sipeed/picoclaw/pkg/utils"
|
||||||
"github.com/sipeed/picoclaw/pkg/voice"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
|
|
@ -40,10 +39,9 @@ const (
|
||||||
// Implements: channels.Channel, channels.TypingCapable, channels.ReactionCapable
|
// Implements: channels.Channel, channels.TypingCapable, channels.ReactionCapable
|
||||||
type SignalChannel struct {
|
type SignalChannel struct {
|
||||||
*channels.BaseChannel
|
*channels.BaseChannel
|
||||||
config config.SignalConfig
|
config config.SignalConfig
|
||||||
httpClient *http.Client
|
httpClient *http.Client
|
||||||
transcriber *voice.GroqTranscriber
|
ctx context.Context
|
||||||
ctx context.Context
|
|
||||||
cancel context.CancelFunc
|
cancel context.CancelFunc
|
||||||
wg sync.WaitGroup
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
@ -72,6 +70,14 @@ type signalDataMessage struct {
|
||||||
ViewOnce bool `json:"viewOnce"`
|
ViewOnce bool `json:"viewOnce"`
|
||||||
GroupInfo *signalGroupInfo `json:"groupInfo"`
|
GroupInfo *signalGroupInfo `json:"groupInfo"`
|
||||||
Attachments []signalAttachment `json:"attachments"`
|
Attachments []signalAttachment `json:"attachments"`
|
||||||
|
Mentions []signalMention `json:"mentions"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type signalMention struct {
|
||||||
|
Start int `json:"start"`
|
||||||
|
Length int `json:"length"`
|
||||||
|
UUID string `json:"uuid"`
|
||||||
|
Number string `json:"number"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type signalGroupInfo struct {
|
type signalGroupInfo struct {
|
||||||
|
|
@ -117,6 +123,7 @@ func NewSignalChannel(cfg *config.Config, b *bus.MessageBus) (channels.Channel,
|
||||||
|
|
||||||
opts := []channels.BaseChannelOption{
|
opts := []channels.BaseChannelOption{
|
||||||
channels.WithMaxMessageLength(signalMaxMessageLength),
|
channels.WithMaxMessageLength(signalMaxMessageLength),
|
||||||
|
channels.WithGroupTrigger(signalCfg.GroupTrigger),
|
||||||
}
|
}
|
||||||
if signalCfg.ReasoningChannelID != "" {
|
if signalCfg.ReasoningChannelID != "" {
|
||||||
opts = append(opts, channels.WithReasoningChannelID(signalCfg.ReasoningChannelID))
|
opts = append(opts, channels.WithReasoningChannelID(signalCfg.ReasoningChannelID))
|
||||||
|
|
@ -131,10 +138,6 @@ func NewSignalChannel(cfg *config.Config, b *bus.MessageBus) (channels.Channel,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *SignalChannel) SetTranscriber(transcriber *voice.GroqTranscriber) {
|
|
||||||
c.transcriber = transcriber
|
|
||||||
}
|
|
||||||
|
|
||||||
func (c *SignalChannel) Start(ctx context.Context) error {
|
func (c *SignalChannel) Start(ctx context.Context) error {
|
||||||
logger.InfoCF("signal", "Starting Signal channel", map[string]any{
|
logger.InfoCF("signal", "Starting Signal channel", map[string]any{
|
||||||
"signal_cli_url": c.config.SignalCLIURL,
|
"signal_cli_url": c.config.SignalCLIURL,
|
||||||
|
|
@ -190,7 +193,7 @@ func (c *SignalChannel) Send(ctx context.Context, msg bus.OutboundMessage) error
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := c.sendMessage(ctx, msg.ChatID, msg.Content); err != nil {
|
if err := c.sendMessage(ctx, msg.ChatID, msg.Content); err != nil {
|
||||||
return fmt.Errorf("signal send: %w", channels.ErrTemporary)
|
return fmt.Errorf("signal send: %w: %v", channels.ErrTemporary, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -296,7 +299,7 @@ func (c *SignalChannel) connectSSE() error {
|
||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
if resp.StatusCode != http.StatusOK {
|
if resp.StatusCode != http.StatusOK {
|
||||||
body, _ := io.ReadAll(resp.Body)
|
body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
|
||||||
return fmt.Errorf("SSE returned status %d: %s", resp.StatusCode, string(body))
|
return fmt.Errorf("SSE returned status %d: %s", resp.StatusCode, string(body))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -382,22 +385,25 @@ func (c *SignalChannel) handleEvent(event signalEvent) {
|
||||||
peerID := senderPhone
|
peerID := senderPhone
|
||||||
|
|
||||||
if isGroup {
|
if isGroup {
|
||||||
if c.config.GroupsEnabled != nil && !*c.config.GroupsEnabled {
|
|
||||||
logger.DebugCF("signal", "Group message ignored (groups_enabled=false)", map[string]any{
|
|
||||||
"group_id": dm.GroupInfo.GroupID,
|
|
||||||
})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
chatID = dm.GroupInfo.GroupID
|
chatID = dm.GroupInfo.GroupID
|
||||||
peerKind = "group"
|
peerKind = "group"
|
||||||
peerID = dm.GroupInfo.GroupID
|
peerID = dm.GroupInfo.GroupID
|
||||||
} else {
|
|
||||||
if c.config.DMsEnabled != nil && !*c.config.DMsEnabled {
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
content := dm.Message
|
content := dm.Message
|
||||||
|
|
||||||
|
// In group chats, apply unified group trigger filtering
|
||||||
|
if isGroup {
|
||||||
|
isMentioned := c.isBotMentioned(dm.Mentions)
|
||||||
|
if isMentioned {
|
||||||
|
content = c.stripMention(content, dm.Mentions)
|
||||||
|
}
|
||||||
|
respond, cleaned := c.ShouldRespondInGroup(isMentioned, content)
|
||||||
|
if !respond {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
content = cleaned
|
||||||
|
}
|
||||||
mediaPaths := []string{}
|
mediaPaths := []string{}
|
||||||
localFiles := []string{}
|
localFiles := []string{}
|
||||||
|
|
||||||
|
|
@ -423,8 +429,7 @@ func (c *SignalChannel) handleEvent(event signalEvent) {
|
||||||
if strings.HasPrefix(att.ContentType, "image/") {
|
if strings.HasPrefix(att.ContentType, "image/") {
|
||||||
content = appendContent(content, "[image: photo]")
|
content = appendContent(content, "[image: photo]")
|
||||||
} else if utils.IsAudioFile(att.Filename, att.ContentType) {
|
} else if utils.IsAudioFile(att.Filename, att.ContentType) {
|
||||||
transcribedText := c.transcribeAudio(localPath)
|
content = appendContent(content, "[voice message]")
|
||||||
content = appendContent(content, transcribedText)
|
|
||||||
} else {
|
} else {
|
||||||
name := att.Filename
|
name := att.Filename
|
||||||
if name == "" {
|
if name == "" {
|
||||||
|
|
@ -441,12 +446,6 @@ func (c *SignalChannel) handleEvent(event signalEvent) {
|
||||||
content = "[media only]"
|
content = "[media only]"
|
||||||
}
|
}
|
||||||
|
|
||||||
// Build compound senderID for backward compat
|
|
||||||
senderID := senderPhone
|
|
||||||
if envelope.SourceName != "" {
|
|
||||||
senderID = fmt.Sprintf("%s|%s", senderPhone, envelope.SourceName)
|
|
||||||
}
|
|
||||||
|
|
||||||
peer := bus.Peer{Kind: peerKind, ID: peerID}
|
peer := bus.Peer{Kind: peerKind, ID: peerID}
|
||||||
|
|
||||||
// Encode messageID as "timestamp:senderPhone" so ReactToMessage can extract both
|
// Encode messageID as "timestamp:senderPhone" so ReactToMessage can extract both
|
||||||
|
|
@ -473,7 +472,71 @@ func (c *SignalChannel) handleEvent(event signalEvent) {
|
||||||
"preview": utils.Truncate(content, 50),
|
"preview": utils.Truncate(content, 50),
|
||||||
})
|
})
|
||||||
|
|
||||||
c.HandleMessage(c.ctx, peer, messageID, senderID, chatID, content, mediaPaths, metadata, sender)
|
c.HandleMessage(c.ctx, peer, messageID, senderPhone, chatID, content, mediaPaths, metadata, sender)
|
||||||
|
}
|
||||||
|
|
||||||
|
// isBotMentioned checks whether the bot was @mentioned in a group message
|
||||||
|
// by looking for its account number or UUID in the structured mentions array.
|
||||||
|
//
|
||||||
|
// Note: signal-cli v0.13.24 has a bug (https://github.com/AsamK/signal-cli/issues/1940)
|
||||||
|
// where the mentions array is empty due to binary ACI parsing issues. This will
|
||||||
|
// work correctly once the fix (PR #1944) is released.
|
||||||
|
func (c *SignalChannel) isBotMentioned(mentions []signalMention) bool {
|
||||||
|
for _, m := range mentions {
|
||||||
|
if m.Number == c.config.Account || m.UUID == c.config.Account {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// stripMention removes the bot's @mention from the message content using
|
||||||
|
// the precise UTF-16 offsets from the structured mention data.
|
||||||
|
// Signal represents mentions as U+FFFC (object replacement character) in the text.
|
||||||
|
func (c *SignalChannel) stripMention(content string, mentions []signalMention) string {
|
||||||
|
for _, m := range mentions {
|
||||||
|
if m.Number != c.config.Account && m.UUID != c.config.Account {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
runes := []rune(content)
|
||||||
|
runeStart, runeLen := utf16PosToRunePos(runes, m.Start, m.Length)
|
||||||
|
if runeStart >= 0 && runeStart+runeLen <= len(runes) {
|
||||||
|
before := strings.TrimRight(string(runes[:runeStart]), " ")
|
||||||
|
after := strings.TrimLeft(string(runes[runeStart+runeLen:]), " ")
|
||||||
|
if before == "" {
|
||||||
|
return after
|
||||||
|
}
|
||||||
|
if after == "" {
|
||||||
|
return before
|
||||||
|
}
|
||||||
|
return before + " " + after
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return content
|
||||||
|
}
|
||||||
|
|
||||||
|
// utf16PosToRunePos converts a UTF-16 code unit position and length to rune position and length.
|
||||||
|
func utf16PosToRunePos(runes []rune, utf16Start, utf16Len int) (int, int) {
|
||||||
|
pos := 0
|
||||||
|
runeStart := -1
|
||||||
|
runeLen := 0
|
||||||
|
for i, r := range runes {
|
||||||
|
if pos == utf16Start {
|
||||||
|
runeStart = i
|
||||||
|
}
|
||||||
|
units := 1
|
||||||
|
if r >= 0x10000 {
|
||||||
|
units = 2 // surrogate pair
|
||||||
|
}
|
||||||
|
if runeStart >= 0 {
|
||||||
|
runeLen++
|
||||||
|
if pos+units >= utf16Start+utf16Len {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
pos += units
|
||||||
|
}
|
||||||
|
return runeStart, runeLen
|
||||||
}
|
}
|
||||||
|
|
||||||
// Media handling
|
// Media handling
|
||||||
|
|
@ -494,26 +557,6 @@ func (c *SignalChannel) downloadAttachment(att signalAttachment) string {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *SignalChannel) transcribeAudio(localPath string) string {
|
|
||||||
if c.transcriber != nil && c.transcriber.IsAvailable() {
|
|
||||||
tCtx, cancel := context.WithTimeout(c.ctx, 30*time.Second)
|
|
||||||
result, err := c.transcriber.Transcribe(tCtx, localPath)
|
|
||||||
cancel()
|
|
||||||
|
|
||||||
if err != nil {
|
|
||||||
logger.ErrorCF("signal", "Voice transcription failed", map[string]any{
|
|
||||||
"error": err.Error(),
|
|
||||||
})
|
|
||||||
return "[voice (transcription failed)]"
|
|
||||||
}
|
|
||||||
logger.InfoCF("signal", "Voice transcribed successfully", map[string]any{
|
|
||||||
"text": result.Text,
|
|
||||||
})
|
|
||||||
return fmt.Sprintf("[voice transcription: %s]", result.Text)
|
|
||||||
}
|
|
||||||
return "[voice]"
|
|
||||||
}
|
|
||||||
|
|
||||||
func extensionFromMIME(mime string) string {
|
func extensionFromMIME(mime string) string {
|
||||||
switch {
|
switch {
|
||||||
case strings.HasPrefix(mime, "image/jpeg"):
|
case strings.HasPrefix(mime, "image/jpeg"):
|
||||||
|
|
@ -579,7 +622,7 @@ func (c *SignalChannel) sendReaction(ctx context.Context, chatID, targetAuthor s
|
||||||
if isGroupChat(chatID) {
|
if isGroupChat(chatID) {
|
||||||
params["groupId"] = chatID
|
params["groupId"] = chatID
|
||||||
} else {
|
} else {
|
||||||
params["recipient"] = chatID
|
params["recipient"] = []string{chatID}
|
||||||
}
|
}
|
||||||
|
|
||||||
if _, err := c.rpcCall(ctx, "sendReaction", params); err != nil {
|
if _, err := c.rpcCall(ctx, "sendReaction", params); err != nil {
|
||||||
|
|
@ -643,7 +686,7 @@ func (c *SignalChannel) sendTyping(chatID string) {
|
||||||
if isGroupChat(chatID) {
|
if isGroupChat(chatID) {
|
||||||
params["groupId"] = chatID
|
params["groupId"] = chatID
|
||||||
} else {
|
} else {
|
||||||
params["recipient"] = chatID
|
params["recipient"] = []string{chatID}
|
||||||
}
|
}
|
||||||
|
|
||||||
ctx, cancel := context.WithTimeout(c.ctx, 5*time.Second)
|
ctx, cancel := context.WithTimeout(c.ctx, 5*time.Second)
|
||||||
|
|
@ -658,6 +701,8 @@ func (c *SignalChannel) sendTyping(chatID string) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// isGroupChat determines if a chatID is a Signal group (base64-encoded) or a phone number.
|
// isGroupChat determines if a chatID is a Signal group (base64-encoded) or a phone number.
|
||||||
|
// This is safe because chatID is always set by handleEvent: either the sender's E.164 phone
|
||||||
|
// number (starts with "+") for DMs, or the base64-encoded GroupInfo.GroupID for groups.
|
||||||
func isGroupChat(chatID string) bool {
|
func isGroupChat(chatID string) bool {
|
||||||
return chatID != "" && !strings.HasPrefix(chatID, "+")
|
return chatID != "" && !strings.HasPrefix(chatID, "+")
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -455,8 +455,7 @@ type SignalConfig struct {
|
||||||
Account string `json:"account" env:"PICOCLAW_CHANNELS_SIGNAL_ACCOUNT"`
|
Account string `json:"account" env:"PICOCLAW_CHANNELS_SIGNAL_ACCOUNT"`
|
||||||
SignalCLIURL string `json:"signal_cli_url" env:"PICOCLAW_CHANNELS_SIGNAL_CLI_URL"`
|
SignalCLIURL string `json:"signal_cli_url" env:"PICOCLAW_CHANNELS_SIGNAL_CLI_URL"`
|
||||||
AllowFrom FlexibleStringSlice `json:"allow_from" env:"PICOCLAW_CHANNELS_SIGNAL_ALLOW_FROM"`
|
AllowFrom FlexibleStringSlice `json:"allow_from" env:"PICOCLAW_CHANNELS_SIGNAL_ALLOW_FROM"`
|
||||||
DMsEnabled *bool `json:"dms_enabled,omitempty"`
|
GroupTrigger GroupTriggerConfig `json:"group_trigger,omitempty"`
|
||||||
GroupsEnabled *bool `json:"groups_enabled,omitempty"`
|
|
||||||
ReasoningChannelID string `json:"reasoning_channel_id" env:"PICOCLAW_CHANNELS_SIGNAL_REASONING_CHANNEL_ID"`
|
ReasoningChannelID string `json:"reasoning_channel_id" env:"PICOCLAW_CHANNELS_SIGNAL_REASONING_CHANNEL_ID"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue