Files

123 lines
3.2 KiB
Go

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
}