You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

181 lines
5.6 KiB

package memory
import (
"fmt"
"strings"
"time"
"yunyan/comm"
"yunyan/lego/sys/mysql"
"yunyan/pb"
)
/*
存量 DBTask(allhelp_task)→ memory_item 的一次性迁移。
在模块 Start 时跑一次,**幂等**:靠 (uid, client_key) 唯一索引挡重复,
client_key 用 "task:<id>",同一条任务再迁一次会被唯一索引挡掉。
⚠️ **不删 allhelp_task 表,也不改它的数据**。这是不可逆迁移的通例:
出问题时旧表是唯一的回退依据。确认新表跑稳后由人手工处置。
⚠️ 迁移失败**不阻断服务启动**。它是补数据,不是启动的前置条件;
在这里返回 error 会让 home 退出,而五个服务同镜像、entrypoint 会连坐全杀。
*/
// 每轮最多搬多少条,避免一次把全表读进内存
const migrateBatchSize = 500
// migrateTasks 把存量任务搬进记忆项。返回搬了几条。
func (this *modelComp) migrateTasks() (int, error) {
// 先看有没有活儿。表不存在(老部署没有 allhelp 模块)时这里会报错,
// 直接当作「没有存量」返回,不是故障。
var pending int64
if err := mysql.Table(comm.TableAllhelpTask).Count(&pending).Error; err != nil {
return 0, nil
}
if pending == 0 {
return 0, nil
}
migrated := 0
var lastID uint64
for {
tasks := make([]*pb.DBTask, 0, migrateBatchSize)
if err := mysql.Table(comm.TableAllhelpTask).
Where("id > ?", lastID).
Order("id asc").Limit(migrateBatchSize).
Find(&tasks).Error; err != nil {
return migrated, err
}
if len(tasks) == 0 {
break
}
for _, t := range tasks {
lastID = t.Id
item := taskToMemoryItem(t)
if item == nil {
continue
}
// 已迁过就跳过:查一次比撞唯一索引再吞错误干净,
// 而且能区分「已迁过」和「真的写失败了」。
if exist, err := this.findItemByClientKey(item.Uid, item.ClientKey); err != nil {
this.module.Warnf("迁移任务 %d 查重失败已跳过: %v", t.Id, err)
continue
} else if exist != nil {
continue
}
if err := this.addItem(item); err != nil {
this.module.Warnf("迁移任务 %d 写入失败已跳过: %v", t.Id, err)
continue
}
migrated++
}
if len(tasks) < migrateBatchSize {
break
}
}
return migrated, nil
}
// taskToMemoryItem 字段映射。无法映射出日期的条目返回 nil(跳过而不是硬塞今天)。
func taskToMemoryItem(t *pb.DBTask) *pb.DBMemoryItem {
if t == nil || t.Uid == "" {
return nil
}
title := strings.TrimSpace(t.TaskName)
if title == "" {
title = strings.TrimSpace(t.TaskDesc)
}
if title == "" {
// 标题和描述都空的任务没有迁的价值,迁过去也只是一条空白记录
return nil
}
// trigger_time 是带时区的 RFC3339(如 2025-04-08T15:00:00+08:00)。
// ⚠️ 用 time.Parse 保留原时区再取日期,**不要先转 UTC**——
// 那会让 00:30 这类凌晨任务整体退到前一天。
var (
happenDate string
happenTime string
tz string
)
if ts := strings.TrimSpace(t.TriggerTime); ts != "" {
if parsed, err := time.Parse(time.RFC3339, ts); err == nil {
happenDate = comm.FormatMemoryDate(parsed)
happenTime = parsed.Format("15:04")
// 原串里的时区偏移信息只有偏移没有 IANA 名,存不了 tz,
// 留空让服务端按容器时区处理——这是已知的精度损失。
_ = tz
}
}
if happenDate == "" {
// 没有可用的触发时间就退回创建日期:记忆项的 happen_date 永不为空,
// 日历按日期排,没有日期的项根本渲染不出来。
if t.CreateTime > 0 {
happenDate = comm.FormatMemoryDate(time.Unix(t.CreateTime, 0))
} else {
return nil
}
}
item := &pb.DBMemoryItem{
Uid: t.Uid,
Category: comm.MemoryCatTodo, // 存量任务一律按待办迁,不猜分类
Source: comm.MemorySrcAssistant,
SourceId: fmt.Sprintf("task:%d", t.Id),
// 幂等键:同一条任务重复迁移会被 (uid, client_key) 唯一索引挡掉
ClientKey: fmt.Sprintf("task:%d", t.Id),
Title: title,
Detail: strings.TrimSpace(t.TaskDesc),
HappenDate: happenDate,
HappenTime: happenTime,
DateCertain: happenTime != "",
RepeatRule: taskTypeToRepeat(t.TaskType),
Currency: "CNY",
State: taskStatusToState(t.Status),
// 存量任务在旧体系里「服务端仅记录、不自动触发」,用户从没被它提醒过。
// 迁过来就默认开提醒等于突然给所有人补发一堆历史提醒,所以置 0。
RemindAhead: 0,
// 迁移来的数据算用户资产,不该被会议重新生成之类的清理逻辑碰到
UserEdited: true,
CreateTime: t.CreateTime,
FinishTime: t.FinishTime,
}
if item.RepeatRule == comm.MemoryRepeatWeekly {
if d, ok := comm.ParseMemoryDate(happenDate); ok {
w := int32(d.Weekday())
if w == 0 {
w = 7 // 周日:Go 的 0 → 客户端口径的 7
}
item.Weekday = w
}
}
return item
}
func taskTypeToRepeat(t pb.TaskType) string {
switch t {
case pb.TaskType_TaskType_Daily:
return comm.MemoryRepeatDaily
case pb.TaskType_TaskType_Weekly:
return comm.MemoryRepeatWeekly
case pb.TaskType_TaskType_Monthly:
return comm.MemoryRepeatMonthly
default:
// Cron 表达式没有对应的 repeat_rule,按一次性迁——
// 猜错的重复规则会让用户每天被一个他没设过的闹钟吵醒。
return comm.MemoryRepeatOnce
}
}
func taskStatusToState(s pb.TaskStatus) pb.MemoryState {
switch s {
case pb.TaskStatus_TaskStatus_Completed:
return pb.MemoryState_MemoryState_Done
case pb.TaskStatus_TaskStatus_Cancelled:
return pb.MemoryState_MemoryState_Dropped
default:
return pb.MemoryState_MemoryState_Pending
}
}