Files
meshtastic_mqtt_server/internal/autoreply/queue.go
T
kevinandClaude 8e7d11a162 feat(llm): 处理完成的消息自动软删除并隐藏删除按钮
- MarkAsProcessed 在写入 reply/processed_at 的同时设置 deleted_at,
  避免 /admin/llm 队列页堆积已处理消息(仍可勾选"包含已删除"查看)
- 前端列表中已删除的消息行隐藏删除按钮,改为显示灰色"已删除"标签

Co-Authored-By: Claude <noreply@anthropic.com>
2026-06-19 17:50:56 +08:00

114 lines
3.8 KiB
Go

package autoreply
import (
"fmt"
"time"
"gorm.io/gorm"
)
const (
statusPending = "pending"
statusProcessing = "processing"
statusProcessed = "processed"
statusFailed = "error"
)
// DBMessageQueue implements MessageQueue using GORM
type DBMessageQueue struct {
db *gorm.DB
}
// NewDBMessageQueue creates a new database-backed message queue
func NewDBMessageQueue(db *gorm.DB) *DBMessageQueue {
return &DBMessageQueue{db: db}
}
// llmMessageQueueRecord is the database record for LLM messages
type llmMessageQueueRecord struct {
ID uint64 `gorm:"column:id;primaryKey;autoIncrement"`
BotID uint64 `gorm:"column:bot_id;not null;index"`
BotNodeID string `gorm:"column:bot_node_id;not null"`
BotNodeNum int64 `gorm:"column:bot_node_num;not null"`
FromNodeID string `gorm:"column:from_node_id;not null"`
FromNodeNum int64 `gorm:"column:from_node_num;not null"`
LongName *string `gorm:"column:long_name"`
ShortName *string `gorm:"column:short_name"`
Text string `gorm:"column:text;type:text;not null"`
PacketID int64 `gorm:"column:packet_id;not null"`
ChannelID *string `gorm:"column:channel_id"`
Topic string `gorm:"column:topic;not null"`
MessageType string `gorm:"column:message_type;not null;default:'direct'"` // "channel" 或 "direct"
Status string `gorm:"column:status;not null;index"`
Error string `gorm:"column:error;type:text"`
Reply string `gorm:"column:reply;type:text"`
ReceivedAt time.Time `gorm:"column:received_at;not null"`
ProcessedAt *time.Time `gorm:"column:processed_at;index"`
CreatedAt time.Time `gorm:"column:created_at;autoCreateTime;index"`
}
func (llmMessageQueueRecord) TableName() string {
return "llm_message_queue"
}
// GetPendingMessages returns pending messages from the queue
func (q *DBMessageQueue) GetPendingMessages(botID uint64, limit int) ([]QueuedMessage, error) {
var records []llmMessageQueueRecord
query := q.db.Where("status = ?", statusPending).Order("created_at ASC")
if botID > 0 {
query = query.Where("bot_id = ?", botID)
}
if limit > 0 {
query = query.Limit(limit)
}
if err := query.Find(&records).Error; err != nil {
return nil, fmt.Errorf("failed to query pending messages: %w", err)
}
messages := make([]QueuedMessage, 0, len(records))
for _, r := range records {
messages = append(messages, QueuedMessage{
ID: r.ID,
BotID: r.BotID,
BotNodeID: r.BotNodeID,
BotNodeNum: r.BotNodeNum,
FromNodeID: r.FromNodeID,
FromNodeNum: r.FromNodeNum,
LongName: r.LongName,
ShortName: r.ShortName,
Text: r.Text,
PacketID: r.PacketID,
ChannelID: r.ChannelID,
Topic: r.Topic,
MessageType: r.MessageType,
ReceivedAt: r.ReceivedAt,
})
}
return messages, nil
}
// MarkAsProcessing marks a message as being processed
func (q *DBMessageQueue) MarkAsProcessing(id uint64) error {
return q.db.Model(&llmMessageQueueRecord{}).Where("id = ?", id).Update("status", statusProcessing).Error
}
// MarkAsProcessed marks a message as successfully processed and soft-deletes it
// 处理完成的消息直接软删除,避免队列页堆积;保留 reply/processed_at 便于查询历史
func (q *DBMessageQueue) MarkAsProcessed(id uint64, reply string) error {
now := time.Now()
return q.db.Model(&llmMessageQueueRecord{}).Where("id = ?", id).Updates(map[string]any{
"status": statusProcessed,
"reply": reply,
"processed_at": &now,
"deleted_at": &now,
}).Error
}
// MarkAsFailed marks a message as failed
func (q *DBMessageQueue) MarkAsFailed(id uint64, error string) error {
return q.db.Model(&llmMessageQueueRecord{}).Where("id = ?", id).Updates(map[string]any{
"status": statusFailed,
"error": error,
}).Error
}