feat(memory): P1 长期记忆升级 —— 异步攒批 Consolidate + 软删 + importance/last_seen #1

Merged
Blizzard merged 181 commits from feat/wails3 into main 2026-07-17 01:12:32 +00:00
5 changed files with 46 additions and 19 deletions
Showing only changes of commit ee1e9cfdca - Show all commits
+2 -1
View File
@@ -77,7 +77,8 @@ func (h *Handler) GuardrailEvents(c *gin.Context) {
// 区别于 stats/overview(桌面端个人工作台):这里一律系统级——全平台用户/任务/评测/ // 区别于 stats/overview(桌面端个人工作台):这里一律系统级——全平台用户/任务/评测/
// 模型配置态/提示词控制面态/服务健康。Task/Eval 表无 owner 即全量;用户/KB/Doc 走全局计数。 // 模型配置态/提示词控制面态/服务健康。Task/Eval 表无 owner 即全量;用户/KB/Doc 走全局计数。
func (h *Handler) AdminOverview(c *gin.Context) { func (h *Handler) AdminOverview(c *gin.Context) {
ctx := c.Request.Context() // 系统级口径:显式跨租户(否则受租户表 KB/Doc/Task/Eval 会被插件按 admin 自己的租户过滤)。
ctx := store.WithoutTenant(c.Request.Context())
ov := h.db.StatsOverview(ctx, "") // owner="" → 跳过个人 KB 口径,仅取全局任务/评测 ov := h.db.StatsOverview(ctx, "") // owner="" → 跳过个人 KB 口径,仅取全局任务/评测
users, kbs, docs := h.db.SystemCounts(ctx) users, kbs, docs := h.db.SystemCounts(ctx)
@@ -65,7 +65,7 @@ func (h *Handler) SubmitTask(c *gin.Context) {
task.Meta[contract.MetaSafetyCheck] = true task.Meta[contract.MetaSafetyCheck] = true
} }
// 持久化任务提交(best-effort:降级模式下静默跳过,不阻断发布)。 // 持久化任务提交(best-effort:降级模式下静默跳过,不阻断发布)。
if err := h.db.SaveTask(c.Request.Context(), task.ID, string(task.Graph)); err != nil { if err := h.db.SaveTask(c.Request.Context(), userID(c), task.ID, string(task.Graph)); err != nil {
log.Printf("[gateway] save task %s failed: %v", task.ID, err) log.Printf("[gateway] save task %s failed: %v", task.ID, err)
} }
if err := h.bus.PublishTask(c.Request.Context(), task); err != nil { if err := h.bus.PublishTask(c.Request.Context(), task); err != nil {
@@ -419,7 +419,7 @@ func (h *Handler) StatsOverview(c *gin.Context) {
// 近期运行 feed。 // 近期运行 feed。
recent := make([]gin.H, 0, 8) recent := make([]gin.H, 0, 8)
for _, t := range h.db.RecentTasks(ctx, 8) { for _, t := range h.db.RecentTasks(ctx, uid, 8) {
recent = append(recent, gin.H{ recent = append(recent, gin.H{
"task_id": t.TaskID, "status": t.Status, "detail": t.Detail, "at": t.CreatedAt, "task_id": t.TaskID, "status": t.Status, "detail": t.Detail, "at": t.CreatedAt,
}) })
@@ -464,7 +464,7 @@ func (h *Handler) Runs(c *gin.Context) {
limit = n limit = n
} }
} }
c.JSON(http.StatusOK, gin.H{"runs": h.db.RecentRuns(c.Request.Context(), limit)}) c.JSON(http.StatusOK, gin.H{"runs": h.db.RecentRuns(c.Request.Context(), userID(c), limit)})
} }
// ListMemory: GET /api/v1/memory —— 列出当前用户的全部偏好(结构化,供记忆面板)。 // ListMemory: GET /api/v1/memory —— 列出当前用户的全部偏好(结构化,供记忆面板)。
+8 -4
View File
@@ -17,15 +17,19 @@ type User struct {
// 业务 idtask_xxx,用于 NATS subject/stream)单列 TaskID,主键统一雪花。 // 业务 idtask_xxx,用于 NATS subject/stream)单列 TaskID,主键统一雪花。
type Task struct { type Task struct {
BaseModel BaseModel
TaskID string `gorm:"uniqueIndex;size:64"` // task_xxx TenantID string `gorm:"size:64;index"` // 多租户作用域(tenant 插件按 ctx 自动填/过滤)
Graph string `gorm:"type:jsonb"` // React Flow 导出的 DSL 原文 Owner string `gorm:"size:64;index"` // 提交者 user.id(个人工作台按此过滤"我的运行")
Status string `gorm:"size:32"` // submitted / running / done / failed / timeout TaskID string `gorm:"uniqueIndex;size:64"` // task_xxx
Detail string `gorm:"type:text"` // 失败/超时原因等(状态机回写) Graph string `gorm:"type:jsonb"` // React Flow 导出的 DSL 原文
Status string `gorm:"size:32"` // submitted / running / done / failed / timeout
Detail string `gorm:"type:text"` // 失败/超时原因等(状态机回写)
// 收尾持久化:供「运行历史复盘」永久回放(Redis 流仅 10min TTL,过期后历史任务靠这两列)。 // 收尾持久化:供「运行历史复盘」永久回放(Redis 流仅 10min TTL,过期后历史任务靠这两列)。
Output string `gorm:"type:text"` // 最终模型输出(收尾时由网关从流快照落库) Output string `gorm:"type:text"` // 最终模型输出(收尾时由网关从流快照落库)
Trace string `gorm:"type:text"` // 执行轨迹事件 JSON 数组(存为文本,容忍空串;不在库内查它) Trace string `gorm:"type:text"` // 执行轨迹事件 JSON 数组(存为文本,容忍空串;不在库内查它)
} }
func (Task) isTenantScoped() {}
// Eval 是一次任务的自动化评测结果(dispatcher 评完经 NATS 回写,每任务一条,按 task_id upsert)。 // Eval 是一次任务的自动化评测结果(dispatcher 评完经 NATS 回写,每任务一条,按 task_id upsert)。
type Eval struct { type Eval struct {
BaseModel BaseModel
+16 -9
View File
@@ -117,11 +117,12 @@ func migrateDocLinkToID(db *gorm.DB) {
func (p *Postgres) Enabled() bool { return p.db != nil } func (p *Postgres) Enabled() bool { return p.db != nil }
// SaveTask 持久化一次任务提交(best-effort:降级模式下静默跳过)。 // SaveTask 持久化一次任务提交(best-effort:降级模式下静默跳过)。
func (p *Postgres) SaveTask(ctx context.Context, id, graph string) error { func (p *Postgres) SaveTask(ctx context.Context, owner, id, graph string) error {
if p.db == nil { if p.db == nil {
return nil return nil
} }
return p.db.WithContext(ctx).Create(&Task{TaskID: id, Graph: graph, Status: contract.TaskSubmitted}).Error // TenantID 由 tenant 插件按请求 ctx 自动填;Owner 显式记录提交者(供个人工作台过滤)。
return p.db.WithContext(ctx).Create(&Task{Owner: owner, TaskID: id, Graph: graph, Status: contract.TaskSubmitted}).Error
} }
// UpdateTaskStatus 流转任务状态(running/done/failed/timeout),由 dispatcher 经 NATS 回写驱动。 // UpdateTaskStatus 流转任务状态(running/done/failed/timeout),由 dispatcher 经 NATS 回写驱动。
@@ -270,12 +271,14 @@ func (p *Postgres) StatsOverview(ctx context.Context, owner string) *Overview {
} }
// RecentTasks 返回最近 n 条任务(工作台「近期运行」feed)。 // RecentTasks 返回最近 n 条任务(工作台「近期运行」feed)。
func (p *Postgres) RecentTasks(ctx context.Context, n int) []Task { // RecentTasks 返回某用户最近 n 条任务(个人工作台「近期运行」feed)。
// owner 过滤"我的运行";tenant 由插件自动叠加(双保险:跨用户/跨租户都隔离)。
func (p *Postgres) RecentTasks(ctx context.Context, owner string, n int) []Task {
if p.db == nil { if p.db == nil {
return nil return nil
} }
var out []Task var out []Task
p.db.WithContext(ctx).Order("created_at desc").Limit(n).Find(&out) p.db.WithContext(ctx).Where("owner = ?", owner).Order("created_at desc").Limit(n).Find(&out)
return out return out
} }
@@ -289,18 +292,22 @@ type RunRow struct {
EvalOverall float64 `json:"eval_overall"` EvalOverall float64 `json:"eval_overall"`
} }
// RecentRuns 返回最近 n 条运行(含评测分级,供「运行历史」列表)。 // RecentRuns 返回某用户最近 n 条运行(含评测分级,供「运行历史」列表)。
func (p *Postgres) RecentRuns(ctx context.Context, n int) []RunRow { // 注:raw Table 查询绕过 gorm 模型回调 → 租户插件不生效,故此处**手动**按 owner(+ctx 租户) 过滤。
func (p *Postgres) RecentRuns(ctx context.Context, owner string, n int) []RunRow {
if p.db == nil { if p.db == nil {
return nil return nil
} }
var out []RunRow var out []RunRow
p.db.WithContext(ctx).Table("sundynix_task as t"). q := p.db.WithContext(ctx).Table("sundynix_task as t").
Select("t.task_id, t.status, t.detail, t.created_at as at, " + Select("t.task_id, t.status, t.detail, t.created_at as at, " +
"coalesce(e.level,'') as eval_level, coalesce(e.overall,0) as eval_overall"). "coalesce(e.level,'') as eval_level, coalesce(e.overall,0) as eval_overall").
Joins("left join sundynix_eval e on e.task_id = t.task_id"). Joins("left join sundynix_eval e on e.task_id = t.task_id").
Where("t.deleted_at is null"). Where("t.deleted_at is null AND t.owner = ?", owner)
Order("t.created_at desc").Limit(n).Scan(&out) if tid := tenantFromCtx(ctx); tid != "" && !isSystemCtx(ctx) {
q = q.Where("t.tenant_id = ?", tid)
}
q.Order("t.created_at desc").Limit(n).Scan(&out)
return out return out
} }
@@ -16,6 +16,7 @@ type tenantScopedMarker interface{ isTenantScoped() }
// ---- 请求上下文携带 tenant_id(中间件注入 → 传到 store → 插件读取)---- // ---- 请求上下文携带 tenant_id(中间件注入 → 传到 store → 插件读取)----
type ctxKeyTenant struct{} type ctxKeyTenant struct{}
type ctxKeySystem struct{}
// WithTenant 把 tenant_id 放进 context(空则原样返回,避免污染系统/回填查询)。 // WithTenant 把 tenant_id 放进 context(空则原样返回,避免污染系统/回填查询)。
func WithTenant(ctx context.Context, tenantID string) context.Context { func WithTenant(ctx context.Context, tenantID string) context.Context {
@@ -25,6 +26,20 @@ func WithTenant(ctx context.Context, tenantID string) context.Context {
return context.WithValue(ctx, ctxKeyTenant{}, tenantID) return context.WithValue(ctx, ctxKeyTenant{}, tenantID)
} }
// WithoutTenant 标记本次操作为「系统/跨租户」——即使 ctx 里带着某租户(如 admin 自己的),
// 插件也**不**加 tenant 过滤。用于 admin 系统级聚合、跨租户后台任务等需全平台可见的路径。
func WithoutTenant(ctx context.Context) context.Context {
return context.WithValue(ctx, ctxKeySystem{}, true)
}
func isSystemCtx(ctx context.Context) bool {
if ctx == nil {
return false
}
b, _ := ctx.Value(ctxKeySystem{}).(bool)
return b
}
func tenantFromCtx(ctx context.Context) string { func tenantFromCtx(ctx context.Context) string {
if ctx == nil { if ctx == nil {
return "" return ""
@@ -53,7 +68,7 @@ func isTenantScopedStmt(db *gorm.DB) bool {
// addTenantWhere:受租户模型 + ctx 有 tenant → 追加 tenant_id 过滤。 // addTenantWhere:受租户模型 + ctx 有 tenant → 追加 tenant_id 过滤。
// ctx 无 tenant(系统/回填/未登录)→ 不过滤(这些路径本就需跨租户;用户面由中间件保证有 tenant)。 // ctx 无 tenant(系统/回填/未登录)→ 不过滤(这些路径本就需跨租户;用户面由中间件保证有 tenant)。
func addTenantWhere(db *gorm.DB) { func addTenantWhere(db *gorm.DB) {
if !isTenantScopedStmt(db) { if !isTenantScopedStmt(db) || isSystemCtx(db.Statement.Context) {
return return
} }
if tid := tenantFromCtx(db.Statement.Context); tid != "" { if tid := tenantFromCtx(db.Statement.Context); tid != "" {
@@ -65,7 +80,7 @@ func addTenantWhere(db *gorm.DB) {
// setTenantOnCreate:受租户模型 + ctx 有 tenant → 强制把 tenant_id 设为 ctx 租户(防越权写他租)。 // setTenantOnCreate:受租户模型 + ctx 有 tenant → 强制把 tenant_id 设为 ctx 租户(防越权写他租)。
func setTenantOnCreate(db *gorm.DB) { func setTenantOnCreate(db *gorm.DB) {
if !isTenantScopedStmt(db) { if !isTenantScopedStmt(db) || isSystemCtx(db.Statement.Context) {
return return
} }
if tid := tenantFromCtx(db.Statement.Context); tid != "" { if tid := tenantFromCtx(db.Statement.Context); tid != "" {