342 lines
15 KiB
Go
342 lines
15 KiB
Go
package mysql
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
|
||
"hyapp/pkg/appcode"
|
||
"hyapp/pkg/xerr"
|
||
)
|
||
|
||
const walletOutboxArchiveStateVerified = "VERIFIED"
|
||
|
||
// WalletOutboxArchiveCursor 是 delivered 时间 + 事件 ID 的稳定复合游标。
|
||
// delivered 后这两列不再变化,因此翻页不会因新写入而重复或跳行。
|
||
type WalletOutboxArchiveCursor struct {
|
||
UpdatedAtMS int64 `json:"updated_at_ms"`
|
||
EventID string `json:"event_id"`
|
||
}
|
||
|
||
// WalletOutboxArchiveRecord 保留 wallet_outbox 的完整行快照,便于恢复工具按原始事实校验。
|
||
type WalletOutboxArchiveRecord struct {
|
||
AppCode string `json:"app_code"`
|
||
EventID string `json:"event_id"`
|
||
EventType string `json:"event_type"`
|
||
TransactionID string `json:"transaction_id"`
|
||
CommandID string `json:"command_id"`
|
||
UserID int64 `json:"user_id"`
|
||
AssetType string `json:"asset_type"`
|
||
AvailableDelta int64 `json:"available_delta"`
|
||
FrozenDelta int64 `json:"frozen_delta"`
|
||
PayloadJSON string `json:"payload_json"`
|
||
Status string `json:"status"`
|
||
WorkerID string `json:"worker_id"`
|
||
LockUntilMS *int64 `json:"lock_until_ms"`
|
||
RetryCount int `json:"retry_count"`
|
||
NextRetryAtMS *int64 `json:"next_retry_at_ms"`
|
||
LastError *string `json:"last_error"`
|
||
CreatedAtMS int64 `json:"created_at_ms"`
|
||
UpdatedAtMS int64 `json:"updated_at_ms"`
|
||
}
|
||
|
||
// WalletOutboxArchiveMember 是一个已经进入同批 COS data/manifest/_SUCCESS 并完成读回验证的精确事件键。
|
||
// 它只保存后续清理所需的窄字段,不复制 payload;旧 receipt 没有 member 时天然不具备清理资格。
|
||
type WalletOutboxArchiveMember struct {
|
||
AppCode string
|
||
EventID string
|
||
UpdatedAtMS int64
|
||
}
|
||
|
||
// WalletOutboxArchiveReceipt 只在 COS data/manifest/_SUCCESS 都读回校验后写入。
|
||
// VERIFIED 只代表对象完整性,不代表已完成真实 MySQL 恢复演练。
|
||
type WalletOutboxArchiveReceipt struct {
|
||
AppCode string
|
||
BatchID string
|
||
State string
|
||
DataObjectKey string
|
||
ManifestObjectKey string
|
||
SuccessObjectKey string
|
||
SchemaVersion int
|
||
RowCount int
|
||
FirstCursor WalletOutboxArchiveCursor
|
||
LastCursor WalletOutboxArchiveCursor
|
||
MinCreatedAtMS int64
|
||
MaxCreatedAtMS int64
|
||
AvailableDeltaTotal string
|
||
FrozenDeltaTotal string
|
||
NDJSONSHA256 string
|
||
GzipSHA256 string
|
||
DataCRC64 string
|
||
ManifestCRC64 string
|
||
SuccessCRC64 string
|
||
ManifestSHA256 string
|
||
UncompressedBytes int64
|
||
CompressedBytes int64
|
||
VerifiedAtMS int64
|
||
Members []WalletOutboxArchiveMember
|
||
}
|
||
|
||
// AcquireWalletOutboxArchiveAppLock 用 MySQL connection-scoped advisory lock 保证多副本下每个 App 只有一个归档者。
|
||
// release 必须 defer 调用;专用 Conn 在释放前不会回到连接池,避免锁跟随错误会话。
|
||
func (r *Repository) AcquireWalletOutboxArchiveAppLock(ctx context.Context) (release func(), acquired bool, err error) {
|
||
if r == nil || r.db == nil {
|
||
return nil, false, xerr.New(xerr.Unavailable, "mysql repository is not configured")
|
||
}
|
||
lockName := "wallet_outbox_archive:" + appcode.FromContext(ctx)
|
||
conn, err := r.db.Conn(ctx)
|
||
if err != nil {
|
||
return nil, false, err
|
||
}
|
||
var result sql.NullInt64
|
||
if err := conn.QueryRowContext(ctx, `SELECT GET_LOCK(?, 0)`, lockName).Scan(&result); err != nil {
|
||
_ = conn.Close()
|
||
return nil, false, err
|
||
}
|
||
if !result.Valid || result.Int64 != 1 {
|
||
_ = conn.Close()
|
||
return nil, false, nil
|
||
}
|
||
release = func() {
|
||
releaseCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
|
||
defer cancel()
|
||
var released sql.NullInt64
|
||
// Close 同样会释放 connection-scoped lock;即使显式 release 超时也不能把会话留在连接池。
|
||
_ = conn.QueryRowContext(releaseCtx, `SELECT RELEASE_LOCK(?)`, lockName).Scan(&released)
|
||
_ = conn.Close()
|
||
}
|
||
return release, true, nil
|
||
}
|
||
|
||
// LatestVerifiedWalletOutboxArchiveCursor 只信任持久化的 VERIFIED 回执,上传中断的对象不会推进游标。
|
||
func (r *Repository) LatestVerifiedWalletOutboxArchiveCursor(ctx context.Context) (WalletOutboxArchiveCursor, error) {
|
||
if r == nil || r.db == nil {
|
||
return WalletOutboxArchiveCursor{}, xerr.New(xerr.Unavailable, "mysql repository is not configured")
|
||
}
|
||
var cursor WalletOutboxArchiveCursor
|
||
err := r.db.QueryRowContext(ctx, `
|
||
SELECT last_updated_at_ms, last_event_id
|
||
FROM wallet_outbox_archive_receipts FORCE INDEX (idx_wallet_outbox_archive_cursor)
|
||
WHERE app_code = ? AND state = ?
|
||
ORDER BY last_updated_at_ms DESC, last_event_id DESC
|
||
LIMIT 1`, appcode.FromContext(ctx), walletOutboxArchiveStateVerified).Scan(&cursor.UpdatedAtMS, &cursor.EventID)
|
||
if err == sql.ErrNoRows {
|
||
return WalletOutboxArchiveCursor{}, nil
|
||
}
|
||
return cursor, err
|
||
}
|
||
|
||
// ListDeliveredWalletOutboxArchiveBatch 只读取一个 App 的旧 delivered 行。
|
||
// WHERE 完整命中 (app_code,status,updated_at_ms,event_id) 索引前缀,游标下界和 cutoff 上界把每次 IO 限制在小范围内。
|
||
func (r *Repository) ListDeliveredWalletOutboxArchiveBatch(ctx context.Context, cursor WalletOutboxArchiveCursor, cutoffMS int64, limit int) ([]WalletOutboxArchiveRecord, error) {
|
||
if r == nil || r.db == nil {
|
||
return nil, xerr.New(xerr.Unavailable, "mysql repository is not configured")
|
||
}
|
||
if cutoffMS <= 0 {
|
||
return nil, xerr.New(xerr.InvalidArgument, "archive cutoff_ms must be positive")
|
||
}
|
||
if limit <= 0 {
|
||
limit = 100
|
||
}
|
||
if limit > 500 {
|
||
limit = 500
|
||
}
|
||
rows, err := r.db.QueryContext(ctx, `
|
||
SELECT app_code, event_id, event_type, transaction_id, command_id, user_id, asset_type,
|
||
available_delta, frozen_delta, CAST(payload AS CHAR), status, worker_id,
|
||
lock_until_ms, retry_count, next_retry_at_ms, last_error, created_at_ms, updated_at_ms
|
||
FROM wallet_outbox FORCE INDEX (idx_wallet_outbox_retention)
|
||
WHERE app_code = ? AND status = ?
|
||
AND updated_at_ms >= ? AND updated_at_ms < ?
|
||
AND (updated_at_ms > ? OR (updated_at_ms = ? AND event_id > ?))
|
||
ORDER BY updated_at_ms ASC, event_id ASC
|
||
LIMIT ?`,
|
||
appcode.FromContext(ctx), outboxStatusDelivered,
|
||
cursor.UpdatedAtMS, cutoffMS,
|
||
cursor.UpdatedAtMS, cursor.UpdatedAtMS, strings.TrimSpace(cursor.EventID), limit,
|
||
)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
records := make([]WalletOutboxArchiveRecord, 0, limit)
|
||
for rows.Next() {
|
||
var record WalletOutboxArchiveRecord
|
||
var lockUntilMS, nextRetryAtMS sql.NullInt64
|
||
var lastError sql.NullString
|
||
if err := rows.Scan(
|
||
&record.AppCode, &record.EventID, &record.EventType, &record.TransactionID, &record.CommandID,
|
||
&record.UserID, &record.AssetType, &record.AvailableDelta, &record.FrozenDelta, &record.PayloadJSON,
|
||
&record.Status, &record.WorkerID, &lockUntilMS, &record.RetryCount, &nextRetryAtMS, &lastError,
|
||
&record.CreatedAtMS, &record.UpdatedAtMS,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
if lockUntilMS.Valid {
|
||
value := lockUntilMS.Int64
|
||
record.LockUntilMS = &value
|
||
}
|
||
if nextRetryAtMS.Valid {
|
||
value := nextRetryAtMS.Int64
|
||
record.NextRetryAtMS = &value
|
||
}
|
||
if lastError.Valid {
|
||
value := lastError.String
|
||
record.LastError = &value
|
||
}
|
||
records = append(records, record)
|
||
}
|
||
return records, rows.Err()
|
||
}
|
||
|
||
// RecordVerifiedWalletOutboxArchive 是归档链路的最后一步;调用方必须先完成 COS HEAD + GET 读回校验。
|
||
func (r *Repository) RecordVerifiedWalletOutboxArchive(ctx context.Context, receipt WalletOutboxArchiveReceipt) error {
|
||
if r == nil || r.db == nil {
|
||
return xerr.New(xerr.Unavailable, "mysql repository is not configured")
|
||
}
|
||
receipt.AppCode = appcode.Normalize(receipt.AppCode)
|
||
if receipt.State != walletOutboxArchiveStateVerified || receipt.BatchID == "" || receipt.RowCount <= 0 ||
|
||
receipt.DataObjectKey == "" || receipt.ManifestObjectKey == "" || receipt.SuccessObjectKey == "" ||
|
||
receipt.NDJSONSHA256 == "" || receipt.GzipSHA256 == "" || receipt.ManifestSHA256 == "" ||
|
||
receipt.DataCRC64 == "" || receipt.ManifestCRC64 == "" || receipt.SuccessCRC64 == "" ||
|
||
receipt.VerifiedAtMS <= 0 || len(receipt.Members) != receipt.RowCount {
|
||
return xerr.New(xerr.InvalidArgument, "verified archive receipt is incomplete")
|
||
}
|
||
if err := validateWalletOutboxArchiveMembers(receipt); err != nil {
|
||
return err
|
||
}
|
||
tx, err := r.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer func() { _ = tx.Rollback() }()
|
||
_, err = tx.ExecContext(ctx, `
|
||
INSERT INTO wallet_outbox_archive_receipts (
|
||
app_code, batch_id, state, data_object_key, manifest_object_key, success_object_key,
|
||
schema_version, row_count, first_updated_at_ms, first_event_id, last_updated_at_ms, last_event_id,
|
||
min_created_at_ms, max_created_at_ms, available_delta_total, frozen_delta_total,
|
||
ndjson_sha256, gzip_sha256, data_crc64, manifest_crc64, success_crc64, manifest_sha256, uncompressed_bytes, compressed_bytes,
|
||
verified_at_ms, created_at_ms, updated_at_ms
|
||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||
ON DUPLICATE KEY UPDATE
|
||
state = VALUES(state), data_object_key = VALUES(data_object_key), manifest_object_key = VALUES(manifest_object_key),
|
||
success_object_key = VALUES(success_object_key), verified_at_ms = VALUES(verified_at_ms), updated_at_ms = VALUES(updated_at_ms)`,
|
||
receipt.AppCode, receipt.BatchID, receipt.State, receipt.DataObjectKey, receipt.ManifestObjectKey, receipt.SuccessObjectKey,
|
||
receipt.SchemaVersion, receipt.RowCount, receipt.FirstCursor.UpdatedAtMS, receipt.FirstCursor.EventID,
|
||
receipt.LastCursor.UpdatedAtMS, receipt.LastCursor.EventID, receipt.MinCreatedAtMS, receipt.MaxCreatedAtMS,
|
||
receipt.AvailableDeltaTotal, receipt.FrozenDeltaTotal, receipt.NDJSONSHA256, receipt.GzipSHA256,
|
||
receipt.DataCRC64, receipt.ManifestCRC64, receipt.SuccessCRC64,
|
||
receipt.ManifestSHA256, receipt.UncompressedBytes, receipt.CompressedBytes,
|
||
receipt.VerifiedAtMS, receipt.VerifiedAtMS, receipt.VerifiedAtMS,
|
||
)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if err := verifyWalletOutboxArchiveReceiptControls(ctx, tx, receipt); err != nil {
|
||
return err
|
||
}
|
||
// receipt 与 exact members 必须原子提交:只看到 receipt、看不到成员的中间态会让
|
||
// watermark 推断重新混入清理链路;只看到成员、receipt 未 VERIFIED 则无法证明 COS 对象完整。
|
||
for offset := 0; offset < len(receipt.Members); offset += 500 {
|
||
end := min(offset+500, len(receipt.Members))
|
||
if err := upsertAndVerifyWalletOutboxArchiveMembers(ctx, tx, receipt, receipt.Members[offset:end]); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return tx.Commit()
|
||
}
|
||
|
||
func verifyWalletOutboxArchiveReceiptControls(ctx context.Context, tx *sql.Tx, receipt WalletOutboxArchiveReceipt) error {
|
||
var state, manifestSHA, dataCRC64, manifestCRC64, successCRC64 string
|
||
var schemaVersion, rowCount int
|
||
var verifiedAtMS int64
|
||
if err := tx.QueryRowContext(ctx, `
|
||
SELECT state, schema_version, row_count, manifest_sha256,
|
||
data_crc64, manifest_crc64, success_crc64, verified_at_ms
|
||
FROM wallet_outbox_archive_receipts FORCE INDEX (PRIMARY)
|
||
WHERE app_code = ? AND batch_id = ?`, receipt.AppCode, receipt.BatchID).
|
||
Scan(&state, &schemaVersion, &rowCount, &manifestSHA, &dataCRC64, &manifestCRC64, &successCRC64, &verifiedAtMS); err != nil {
|
||
return err
|
||
}
|
||
// batch_id 是内容指纹,但存储层仍逐字段核对恢复授权所依赖的控制值;同 PK 下任何
|
||
// 非幂等 receipt 都必须在写 members 前失败,不能留下永远无法与 receipt JOIN 的孤立证明。
|
||
if state != receipt.State || schemaVersion != receipt.SchemaVersion || rowCount != receipt.RowCount ||
|
||
manifestSHA != receipt.ManifestSHA256 || dataCRC64 != receipt.DataCRC64 ||
|
||
manifestCRC64 != receipt.ManifestCRC64 || successCRC64 != receipt.SuccessCRC64 || verifiedAtMS != receipt.VerifiedAtMS {
|
||
return xerr.New(xerr.Conflict, "verified archive receipt controls do not match existing batch")
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func validateWalletOutboxArchiveMembers(receipt WalletOutboxArchiveReceipt) error {
|
||
seen := make(map[string]struct{}, len(receipt.Members))
|
||
for _, member := range receipt.Members {
|
||
if appcode.Normalize(member.AppCode) != receipt.AppCode || strings.TrimSpace(member.EventID) == "" || member.UpdatedAtMS <= 0 {
|
||
return xerr.New(xerr.InvalidArgument, "verified archive member is incomplete")
|
||
}
|
||
if _, exists := seen[member.EventID]; exists {
|
||
return xerr.New(xerr.InvalidArgument, "verified archive members contain duplicate event_id")
|
||
}
|
||
seen[member.EventID] = struct{}{}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func upsertAndVerifyWalletOutboxArchiveMembers(ctx context.Context, tx *sql.Tx, receipt WalletOutboxArchiveReceipt, members []WalletOutboxArchiveMember) error {
|
||
values := make([]string, 0, len(members))
|
||
args := make([]any, 0, len(members)*6)
|
||
for _, member := range members {
|
||
values = append(values, "(?, ?, ?, ?, ?, ?)")
|
||
args = append(args, receipt.AppCode, member.EventID, member.UpdatedAtMS, receipt.BatchID, receipt.ManifestSHA256, receipt.VerifiedAtMS)
|
||
}
|
||
// archive worker 的网络结果可能在 DB commit 后丢失;no-op duplicate 允许同一批重试,
|
||
// 随后的逐键读回会拒绝任何已经被另一 batch 占用或 metadata 不一致的成员。
|
||
if _, err := tx.ExecContext(ctx, `
|
||
INSERT INTO wallet_outbox_archive_members (
|
||
app_code, event_id, updated_at_ms, batch_id, manifest_sha256, receipt_verified_at_ms
|
||
) VALUES `+strings.Join(values, ",")+`
|
||
ON DUPLICATE KEY UPDATE event_id = VALUES(event_id)`, args...); err != nil {
|
||
return err
|
||
}
|
||
|
||
placeholders := make([]string, 0, len(members))
|
||
lookupArgs := make([]any, 0, len(members)+1)
|
||
lookupArgs = append(lookupArgs, receipt.AppCode)
|
||
expected := make(map[string]int64, len(members))
|
||
for _, member := range members {
|
||
placeholders = append(placeholders, "?")
|
||
lookupArgs = append(lookupArgs, member.EventID)
|
||
expected[member.EventID] = member.UpdatedAtMS
|
||
}
|
||
rows, err := tx.QueryContext(ctx, `
|
||
SELECT event_id, updated_at_ms, batch_id, manifest_sha256, receipt_verified_at_ms
|
||
FROM wallet_outbox_archive_members FORCE INDEX (PRIMARY)
|
||
WHERE app_code = ? AND event_id IN (`+strings.Join(placeholders, ",")+`)`, lookupArgs...)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer rows.Close()
|
||
matched := 0
|
||
for rows.Next() {
|
||
var eventID, batchID, manifestSHA string
|
||
var updatedAtMS, receiptVerifiedAtMS int64
|
||
if err := rows.Scan(&eventID, &updatedAtMS, &batchID, &manifestSHA, &receiptVerifiedAtMS); err != nil {
|
||
return err
|
||
}
|
||
expectedUpdatedAtMS, exists := expected[eventID]
|
||
if !exists || updatedAtMS != expectedUpdatedAtMS || batchID != receipt.BatchID || manifestSHA != receipt.ManifestSHA256 || receiptVerifiedAtMS != receipt.VerifiedAtMS {
|
||
return fmt.Errorf("wallet outbox archive member conflicts with verified receipt: event_id=%s", eventID)
|
||
}
|
||
matched++
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return err
|
||
}
|
||
if matched != len(members) {
|
||
return fmt.Errorf("wallet outbox archive member readback mismatch: expected=%d matched=%d", len(members), matched)
|
||
}
|
||
return nil
|
||
}
|