fix: 作废投递走统一终态函数并写回执收尾

This commit is contained in:
Nixevol
2026-09-30 16:21:05 +08:00
parent 91e887ba46
commit f139b9ed9b
10 changed files with 365 additions and 166 deletions
+2 -2
View File
@@ -48,7 +48,7 @@ WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`,
if e := insertReceiptTx(tx, req.From, seq, endpointID, DeliveryAccepted, "", nowMs); e != nil {
return e
}
return tryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays)
return TryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays)
}
var state string
err = tx.QueryRow(`
@@ -175,7 +175,7 @@ SELECT COUNT(*) FROM deliveries WHERE seq = ? AND state IN ('expired','dropped',
default:
data.Result = "failed"
}
return tryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays)
return TryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays)
})
if err != nil {
return data, err
+49 -7
View File
@@ -8,6 +8,9 @@ import (
"git.asio.asia/nixevol/NixMsg/internal/protocol"
)
// DefaultVoidRetentionDays 是 group/identity 作废路径未接线配置时的记录保留天数。
const DefaultVoidRetentionDays = 7
// 投递状态(DEVELOPMENT 7.1)。
const (
DeliveryPending = "pending"
@@ -116,7 +119,7 @@ WHERE gm.group_id = ? AND gm.endpoint_id != ?`, destID, senderID)
}
if completeEarly {
if err := finalizeMessageTx(tx, seq, wantReceipt, senderID, "", msgReason, nowMs, a.lim.RecordRetentionDays); err != nil {
if err := FinalizeMessageTx(tx, seq, wantReceipt, senderID, "", msgReason, nowMs, a.lim.RecordRetentionDays); err != nil {
return "", true, err
}
return StateCompleted, true, nil
@@ -192,7 +195,7 @@ VALUES(?,?,?,?,?,?,?,NULL,NULL,0,?)`,
}
return StateDispatched, true, nil
}
if err := finalizeMessageTx(tx, seq, wantReceipt, senderID, "", "", nowMs, a.lim.RecordRetentionDays); err != nil {
if err := FinalizeMessageTx(tx, seq, wantReceipt, senderID, "", "", nowMs, a.lim.RecordRetentionDays); err != nil {
return "", true, err
}
return StateCompleted, true, nil
@@ -212,9 +215,9 @@ func (a *App) lookupConn(endpointID string) (LiveConn, bool) {
return a.conns.Current(endpointID)
}
// finalizeMessageTx 无 pending 时收尾:completed、删正文;记录天数 0 则删消息与投递。
// FinalizeMessageTx 无 pending 时收尾:completed、删正文;记录天数 0 则删消息与投递。
// msgReason 非空时写入消息 reason(发送前结束);endpointID 为空表示消息级回执。
func finalizeMessageTx(tx *sql.Tx, seq int64, wantReceipt bool, senderID, endpointID, msgReason string, nowMs int64, recordDays int) error {
func FinalizeMessageTx(tx *sql.Tx, seq int64, wantReceipt bool, senderID, endpointID, msgReason string, nowMs int64, recordDays int) error {
var msgID string
var receipt int
if err := tx.QueryRow(`SELECT id, receipt FROM messages WHERE seq = ?`, seq).Scan(&msgID, &receipt); err != nil {
@@ -244,8 +247,8 @@ func finalizeMessageTx(tx *sql.Tx, seq int64, wantReceipt bool, senderID, endpoi
return nil
}
// tryFinalizeTx 若无 pending 则收尾。
func tryFinalizeTx(tx *sql.Tx, seq int64, nowMs int64, recordDays int) error {
// TryFinalizeTx 若无 pending 则收尾。
func TryFinalizeTx(tx *sql.Tx, seq int64, nowMs int64, recordDays int) error {
var n int
if err := tx.QueryRow(`SELECT COUNT(*) FROM deliveries WHERE seq = ? AND state = 'pending'`, seq).Scan(&n); err != nil {
return err
@@ -261,7 +264,46 @@ func tryFinalizeTx(tx *sql.Tx, seq int64, nowMs int64, recordDays int) error {
}
return err
}
return finalizeMessageTx(tx, seq, receipt != 0, senderID, "", "", nowMs, recordDays)
return FinalizeMessageTx(tx, seq, receipt != 0, senderID, "", "", nowMs, recordDays)
}
func skipVoidReceipt(reason string) bool {
return reason == "sender_disabled" || reason == "sender_deleted"
}
// RejectPendingTx 把一条 pending 投递改为 rejected;消息要求回执且发送方存在时写回执。
// 停用/删除发送方(sender_disabled / sender_deleted)不写回执(DEVELOPMENT 7.6)。
// 返回该投递是否曾推送,供调用方发 revoked。
func RejectPendingTx(tx *sql.Tx, seq int64, endpointID, reason string, nowMs int64) (pushed bool, err error) {
var pushedAt sql.NullInt64
err = tx.QueryRow(`SELECT pushed_at FROM deliveries WHERE seq = ? AND endpoint_id = ?`, seq, endpointID).Scan(&pushedAt)
if err == sql.ErrNoRows {
return false, nil
}
if err != nil {
return false, err
}
res, err := tx.Exec(`
UPDATE deliveries SET state = ?, reason = ?, updated_at = ?
WHERE seq = ? AND endpoint_id = ? AND state = ?`,
DeliveryRejected, reason, nowMs, seq, endpointID, DeliveryPending)
if err != nil {
return false, err
}
aff, _ := res.RowsAffected()
if aff == 0 {
return false, nil
}
if !skipVoidReceipt(reason) {
var senderID string
if err := tx.QueryRow(`SELECT sender_id FROM messages WHERE seq = ?`, seq).Scan(&senderID); err != nil {
return false, err
}
if err := insertReceiptTx(tx, senderID, seq, endpointID, DeliveryRejected, reason, nowMs); err != nil {
return false, err
}
}
return pushedAt.Valid, nil
}
func insertReceiptTx(tx *sql.Tx, senderID string, seq int64, endpointID, state, reason string, nowMs int64) error {
+2 -2
View File
@@ -259,7 +259,7 @@ WHERE seq = ? AND endpoint_id = ? AND state = 'pending' AND pushed_conn IS NULL`
if err := insertReceiptTx(tx, senderID, seq, endpointID, DeliveryRejected, ReasonTooLarge, nowMs); err != nil {
return err
}
return tryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays)
return TryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays)
})
}
@@ -361,7 +361,7 @@ WHERE seq = ? AND endpoint_id = ? AND state = 'pending'`,
return err
}
}
if err := tryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays); err != nil {
if err := TryFinalizeTx(tx, seq, nowMs, a.lim.RecordRetentionDays); err != nil {
return err
}
if sendRevoked {
+37
View File
@@ -78,6 +78,10 @@ WHERE d.state = 'pending' AND d.pushed_conn IS NULL
}
}
if err := finalizeStuckDispatchedTx(tx, nowMs, a.lim.RecordRetentionDays); err != nil {
return err
}
if a.lim.RecordRetentionDays > 0 {
cutoff := nowMs - int64(a.lim.RecordRetentionDays)*24*3600*1000
if _, err := tx.Exec(`
@@ -120,3 +124,36 @@ DELETE FROM send_keys WHERE rowid IN (
a.flushRevokes(ctx)
return nil
}
// finalizeStuckDispatchedTx 收尾「dispatched 且已无 pending 投递」的消息(C-04 兜底,修复已卡住的数据)。
func finalizeStuckDispatchedTx(tx *sql.Tx, nowMs int64, recordDays int) error {
rows, err := tx.Query(`
SELECT seq FROM messages
WHERE state = ?
AND NOT EXISTS (
SELECT 1 FROM deliveries d WHERE d.seq = messages.seq AND d.state = 'pending'
)
LIMIT 500`, StateDispatched)
if err != nil {
return err
}
var seqs []int64
for rows.Next() {
var seq int64
if err := rows.Scan(&seq); err != nil {
_ = rows.Close()
return err
}
seqs = append(seqs, seq)
}
_ = rows.Close()
if err := rows.Err(); err != nil {
return err
}
for _, seq := range seqs {
if err := TryFinalizeTx(tx, seq, nowMs, recordDays); err != nil {
return err
}
}
return nil
}
+102
View File
@@ -0,0 +1,102 @@
package message
import (
"context"
"database/sql"
"testing"
"git.asio.asia/nixevol/NixMsg/internal/app/port"
)
func TestRejectPendingAndFinalizeRetentionZero(t *testing.T) {
t.Parallel()
lim := defaultTestLimits()
lim.RecordRetentionDays = 0
app, db := openTestApp(t, lim)
insertEndpoint(t, db, "alice", "", 1, 0)
insertEndpoint(t, db, "bob", "", 1, 0)
ctx := context.Background()
req := baseSend("z1", "bob")
req.Offline = keepTrue()
if _, err := app.Submit(ctx, "alice", port.ConnInfo{}, req); err != nil {
t.Fatal(err)
}
err := db.Queue.Do(ctx, func(tx *sql.Tx) error {
var seq int64
if e := tx.QueryRow(`SELECT seq FROM messages WHERE id='z1'`).Scan(&seq); e != nil {
return e
}
if _, e := RejectPendingTx(tx, seq, "bob", ReasonEndpointDisabled, 1_700_000_000_000); e != nil {
return e
}
return TryFinalizeTx(tx, seq, 1_700_000_000_000, 0)
})
if err != nil {
t.Fatal(err)
}
var n int
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM messages WHERE id='z1'`).Scan(&n); err != nil {
t.Fatal(err)
}
if n != 0 {
t.Fatalf("message row should be deleted when retention=0, n=%d", n)
}
var receipts int
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM receipts WHERE msg_id='z1' AND state='rejected'`).Scan(&receipts); err != nil {
t.Fatal(err)
}
if receipts != 1 {
t.Fatalf("receipts=%d", receipts)
}
}
func TestCleanupStuckDispatched(t *testing.T) {
t.Parallel()
lim := defaultTestLimits()
app, db := openTestApp(t, lim)
insertEndpoint(t, db, "alice", "", 1, 0)
insertEndpoint(t, db, "bob", "", 1, 0)
ctx := context.Background()
nowMs := int64(1_700_000_000_000)
err := db.Queue.Do(ctx, func(tx *sql.Tx) error {
res, e := tx.Exec(`
INSERT INTO messages(id, sender_id, dest_kind, dest_id, meta, content_type, body_enc,
send_at, keep, ttl_seconds, receipt, state, reason, created_at)
VALUES('stuck','alice','endpoint','bob','{}','text/plain; charset=utf-8','utf8',
?,0,0,1,'dispatched','',?)`, nowMs, nowMs)
if e != nil {
return e
}
seq, e := res.LastInsertId()
if e != nil {
return e
}
if _, e = tx.Exec(`INSERT INTO message_bodies(seq, body) VALUES(?, ?)`, seq, []byte("x")); e != nil {
return e
}
_, e = tx.Exec(`
INSERT INTO deliveries(seq, endpoint_id, send_at, keep, state, reason, updated_at)
VALUES(?,?,?,0,'rejected','left_group',?)`, seq, "bob", nowMs, nowMs)
return e
})
if err != nil {
t.Fatal(err)
}
if err := app.CleanupOnce(ctx, nowMs); err != nil {
t.Fatal(err)
}
var state string
if err := db.Read.QueryRow(`SELECT state FROM messages WHERE id='stuck'`).Scan(&state); err != nil {
t.Fatal(err)
}
if state != StateCompleted {
t.Fatalf("state=%s want completed", state)
}
var bodies int
if err := db.Read.QueryRow(`SELECT COUNT(*) FROM message_bodies b JOIN messages m ON m.seq=b.seq WHERE m.id='stuck'`).Scan(&bodies); err != nil {
t.Fatal(err)
}
if bodies != 0 {
t.Fatalf("body still present: %d", bodies)
}
}