package message import ( "context" "database/sql" ) // RecoverOnStart 启动恢复(DEVELOPMENT 7.8)。 func (a *App) RecoverOnStart(ctx context.Context) error { nowMs := a.now().UnixMilli() graceMs := a.lim.GraceSeconds * 1000 minExpire := nowMs + graceMs err := a.db.Queue.Do(ctx, func(tx *sql.Tx) error { if _, err := tx.Exec(` UPDATE deliveries SET pushed_conn = NULL, expire_at = CASE WHEN expire_at IS NULL OR expire_at < ? THEN ? ELSE expire_at END, updated_at = ? WHERE state = 'pending'`, minExpire, minExpire, nowMs); err != nil { return err } // 停机前在线:online_since 晚于 offline_since,或 offline_since 空而 online_since 非空 _, err := tx.Exec(` UPDATE endpoints SET offline_since = ? WHERE online_since IS NOT NULL AND (offline_since IS NULL OR online_since > offline_since)`, nowMs) return err }) if err != nil { return err } // 停机期间到点的 scheduled 立即分发 _, err = a.DispatchDue(ctx, nowMs, 1000) return err } // CleanupOnce 处理未推送且到期的 pending,并做记录/回执/防重清理。 func (a *App) CleanupOnce(ctx context.Context, nowMs int64) error { err := a.db.Queue.Do(ctx, func(tx *sql.Tx) error { rows, err := tx.Query(` SELECT d.seq, d.endpoint_id, d.keep, m.sender_id, m.id FROM deliveries d JOIN messages m ON m.seq = d.seq WHERE d.state = 'pending' AND d.pushed_conn IS NULL AND d.expire_at IS NOT NULL AND d.expire_at <= ?`, nowMs) if err != nil { return err } type item struct { seq, keep int64 endpointID, senderID string msgID string } var list []item for rows.Next() { var it item if err := rows.Scan(&it.seq, &it.endpointID, &it.keep, &it.senderID, &it.msgID); err != nil { _ = rows.Close() return err } list = append(list, it) } _ = rows.Close() for _, it := range list { state := DeliveryDropped reason := ReasonOffline if it.keep != 0 { state = DeliveryExpired reason = ReasonTTL } if err := a.finishDeliveryTx(tx, it.seq, it.endpointID, it.senderID, it.msgID, state, reason, false, nowMs); err != nil { return err } } if a.lim.RecordRetentionDays > 0 { cutoff := nowMs - int64(a.lim.RecordRetentionDays)*24*3600*1000 if _, err := tx.Exec(` DELETE FROM messages WHERE seq IN ( SELECT seq FROM messages WHERE state = 'completed' AND created_at < ? LIMIT 5000 )`, cutoff); err != nil { return err } } if a.lim.ReceiptRetentionDays > 0 { cutoff := nowMs - int64(a.lim.ReceiptRetentionDays)*24*3600*1000 if _, err := tx.Exec(` DELETE FROM receipts WHERE receipt_id IN ( SELECT receipt_id FROM receipts WHERE created_at < ? LIMIT 5000 )`, cutoff); err != nil { return err } } if a.lim.IdempotencyHours > 0 { cutoff := nowMs - int64(a.lim.IdempotencyHours)*3600*1000 if _, err := tx.Exec(` DELETE FROM send_keys WHERE rowid IN ( SELECT sk.rowid FROM send_keys sk WHERE sk.created_at < ? AND NOT EXISTS ( SELECT 1 FROM messages m WHERE m.sender_id = sk.sender_id AND m.id = sk.msg_id ) LIMIT 5000 )`, cutoff); err != nil { return err } } return nil }) if err != nil { return err } a.flushRevokes(ctx) return nil }