Skip to content

Commit c7bdb7b

Browse files
committed
fix(sync): 同步服务修复——版本冲突判定/WS查询洪泛/Resolve安全/CRDT幂等/广播完整性
- SYNC-H1: Push 冲突判定收紧为 client version <= server version——相同版本号不同内容的静默覆盖改为返回冲突(双设备并发编辑不再丢数据;幂等重试由 op.ID 去重覆盖) - SYNC-H2: WS sync_request 每连接在飞查询上限 2(信号槽),单连接洪泛不再打满数据库连接池(50)拖垮全部用户 - SYNC-M1: Resolve 增加 FOR UPDATE 行锁 + 版本必须大于服务端版本(拒绝版本回退);成功后广播到其他在线设备(实时收敛) - SYNC-M2: Push 广播条件从 opSeqNoByID 非空改为 accepted 非空——无 ID 操作(fallback ID)落库后同样广播 - SYNC-M3: CRDTPush changeset SHA-256 哈希幂等去重((user_id, changeset_hash) 唯一索引 + 并发兜底),网络重试不再重复落库膨胀 - SYNC-L2: Pull/CRDTPull 的 deviceId 白名单校验(防回环与超长参数) 验证: go build + go vet + go test 全部通过
1 parent af1b353 commit c7bdb7b

7 files changed

Lines changed: 139 additions & 26 deletions

File tree

‎server/sync-service/handlers/crdt.go‎

Lines changed: 48 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55
package handlers
66

77
import (
8+
"crypto/sha256"
9+
"encoding/hex"
810
"net/http"
911
"strconv"
1012
"strings"
@@ -54,21 +56,45 @@ func CRDTPush(c *gin.Context) {
5456
// M7: 整个循环包裹在单个事务内
5557
txErr := models.DB.Transaction(func(tx *gorm.DB) error {
5658
for _, ch := range req.Changes {
59+
// SYNC-M3: changeset 哈希幂等去重——网络重试的重复 changeset
60+
// 不再重复分配序号与落库((user_id, changeset_hash) 唯一索引兜底)
61+
if ch.Changeset != "" {
62+
hash := sha256.Sum256([]byte(ch.Changeset))
63+
hashHex := hex.EncodeToString(hash[:])
64+
var dupCount int64
65+
if err := tx.Model(&models.CRDTChange{}).
66+
Where("user_id = ? AND changeset_hash = ?", userID, hashHex).
67+
Count(&dupCount).Error; err != nil {
68+
return err
69+
}
70+
if dupCount > 0 {
71+
// 已处理过:视为已接受(幂等),不重复落库
72+
accepted = append(accepted, ch.Seq)
73+
continue
74+
}
75+
}
76+
5777
seqNo, err := nextSeqNo(tx)
5878
if err != nil {
5979
return err
6080
}
6181

6282
if err := tx.Create(&models.CRDTChange{
63-
ServerSeqNo: seqNo,
64-
DeviceID: req.DeviceID,
65-
UserID: userID,
66-
TableName: ch.TableName,
67-
EntityID: ch.EntityID,
68-
Changeset: ch.Changeset,
69-
Operation: ch.Operation,
70-
CreatedAt: ch.CreatedAt,
83+
ServerSeqNo: seqNo,
84+
DeviceID: req.DeviceID,
85+
UserID: userID,
86+
TableName: ch.TableName,
87+
EntityID: ch.EntityID,
88+
Changeset: ch.Changeset,
89+
ChangesetHash: hashChangeSet(ch.Changeset),
90+
Operation: ch.Operation,
91+
CreatedAt: ch.CreatedAt,
7192
}).Error; err != nil {
93+
// 唯一索引兜底:并发重复推送视为已处理
94+
if isUniqueViolation(err) && ch.Changeset != "" {
95+
accepted = append(accepted, ch.Seq)
96+
continue
97+
}
7298
return err
7399
}
74100
accepted = append(accepted, ch.Seq)
@@ -88,6 +114,15 @@ func CRDTPush(c *gin.Context) {
88114
})
89115
}
90116

117+
// SYNC-M3: 计算 changeset 的 SHA-256 hex(空串返回空串,不参与幂等)
118+
func hashChangeSet(changeset string) string {
119+
if changeset == "" {
120+
return ""
121+
}
122+
hash := sha256.Sum256([]byte(changeset))
123+
return hex.EncodeToString(hash[:])
124+
}
125+
91126
// ─── CRDT Changeset Pull (GET /api/v1/sync/crdt/changes) ──────────────────────
92127

93128
// CRDTPull returns all CRDT changesets since `since` (exclusive) that were
@@ -96,6 +131,11 @@ func CRDTPush(c *gin.Context) {
96131
func CRDTPull(c *gin.Context) {
97132
userID := c.GetString("user_id")
98133
deviceID := c.Query("deviceId")
134+
// SYNC-L2: deviceId 白名单校验——空值拉回本设备变更造成回环
135+
if !isValidDeviceID(deviceID) {
136+
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid deviceId: must match [A-Za-z0-9_-]{1,64}"})
137+
return
138+
}
99139
sinceStr := c.DefaultQuery("since", "0")
100140
since, err := strconv.ParseInt(sinceStr, 10, 64)
101141
if err != nil {

‎server/sync-service/handlers/sync.go‎

Lines changed: 24 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -91,8 +91,11 @@ func Push(c *gin.Context) {
9191
}
9292
entityExists := result.RowsAffected > 0
9393

94-
// M1: 版本比较在事务锁内执行
95-
if entityExists && op.Version < ev.Version {
94+
// SYNC-H1: 冲突判定收紧为 client version <= server version。
95+
// 原实现仅拒绝 <,两台设备基于同一版本各自编辑会产生相同版本号
96+
// 不同内容,后到者静默覆盖先到者(数据丢失且无冲突提示)。
97+
// 相同版本的幂等重试由 op.ID 去重处理,不受影响。
98+
if entityExists && op.Version <= ev.Version {
9699
// Conflict: client is behind server.
97100
conflicts = append(conflicts, ConflictInfo{
98101
EntityType: op.EntityType,
@@ -190,11 +193,26 @@ func Push(c *gin.Context) {
190193

191194
// Broadcast accepted operations to the user's other online devices via WebSocket.
192195
// M11: 广播同时携带 ServerSeqNo(Pull 游标)与 Version(实体版本),语义各自独立。
193-
if len(accepted) > 0 && len(opSeqNoByID) > 0 {
196+
// SYNC-M2: 广播条件从 opSeqNoByID 非空改为 accepted 非空——
197+
// 全部操作无 ID(服务端 fallback ID)时变更已落库但不再被广播;
198+
// 空 ID 操作同样需要广播(其 ServerSeqNo 为 0,客户端按实体版本收敛)。
199+
if len(accepted) > 0 {
200+
skippedSet := make(map[string]bool, len(skipped))
201+
for _, id := range skipped {
202+
skippedSet[id] = true
203+
}
204+
conflictKeys := make(map[string]bool, len(conflicts))
205+
for _, cf := range conflicts {
206+
conflictKeys[cf.EntityType+":"+cf.EntityID] = true
207+
}
208+
194209
wsOps := make([]WSOperationPayload, 0, len(accepted))
195210
for _, op := range req.Operations {
196-
seqNo, ok := opSeqNoByID[op.ID]
197-
if !ok {
211+
// 跳过被去重跳过或冲突的操作
212+
if op.ID != "" && skippedSet[op.ID] {
213+
continue
214+
}
215+
if conflictKeys[op.EntityType+":"+op.EntityID] {
198216
continue
199217
}
200218
wsOps = append(wsOps, WSOperationPayload{
@@ -203,7 +221,7 @@ func Push(c *gin.Context) {
203221
Operation: op.Operation,
204222
Data: op.Payload,
205223
Version: op.Version,
206-
ServerSeqNo: seqNo,
224+
ServerSeqNo: opSeqNoByID[op.ID], // 空 ID 操作未登记序号时为 0
207225
DeviceID: req.DeviceID,
208226
})
209227
}

‎server/sync-service/handlers/sync_query.go‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,12 @@ import (
2020
func Pull(c *gin.Context) {
2121
userID := c.GetString("user_id")
2222
deviceID := c.Query("deviceId")
23+
// SYNC-L2: deviceId 白名单校验(与 WS 入口一致)——空值会让
24+
// device_id != '' 拉回本设备变更造成回环;超长参数增加索引扫描成本
25+
if !isValidDeviceID(deviceID) {
26+
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid deviceId: must match [A-Za-z0-9_-]{1,64}"})
27+
return
28+
}
2329
sinceVersionStr := c.DefaultQuery("sinceVersion", "0")
2430
sinceVersion, err := strconv.ParseInt(sinceVersionStr, 10, 64)
2531
if err != nil {

‎server/sync-service/handlers/sync_resolve.go‎

Lines changed: 32 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55
package handlers
66

77
import (
8+
"errors"
9+
"log"
810
"net/http"
911
"strings"
1012
"time"
@@ -13,6 +15,7 @@ import (
1315

1416
"github.com/gin-gonic/gin"
1517
"gorm.io/gorm"
18+
"gorm.io/gorm/clause"
1619
)
1720

1821
type resolveRequest struct {
@@ -56,14 +59,26 @@ func Resolve(c *gin.Context) {
5659
}
5760

5861
txErr := models.DB.Transaction(func(tx *gorm.DB) error {
62+
// SYNC-M1: 版本校验 + FOR UPDATE 行锁——与 Push 的 M1 防护一致,
63+
// 拒绝回退版本(客户端可提交任意小版本号制造数据回退);
64+
// 加锁消除 TOCTOU 竞态(另一设备并发推送时基于过期版本覆盖)
65+
var ev models.EntityVersion
66+
res := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
67+
Where("user_id = ? AND entity_type = ? AND entity_id = ?", userID, req.EntityType, req.EntityID).
68+
First(&ev)
69+
if res.Error != nil && res.Error != gorm.ErrRecordNotFound {
70+
return res.Error
71+
}
72+
if res.RowsAffected > 0 && req.Version <= ev.Version {
73+
return errors.New("resolve rejected: client version must be newer than server version")
74+
}
75+
5976
seqNo, err := nextSeqNo(tx)
6077
if err != nil {
6178
return err
6279
}
6380

6481
// Upsert EntityVersion.
65-
var ev models.EntityVersion
66-
res := tx.Where("user_id = ? AND entity_type = ? AND entity_id = ?", userID, req.EntityType, req.EntityID).First(&ev)
6782
if res.RowsAffected > 0 {
6883
if err := tx.Model(&ev).Updates(map[string]interface{}{
6984
"version": req.Version,
@@ -98,10 +113,25 @@ func Resolve(c *gin.Context) {
98113
})
99114

100115
if txErr != nil {
116+
log.Printf("Resolve failed: %v", txErr)
101117
c.JSON(http.StatusInternalServerError, gin.H{"error": txErr.Error()})
102118
return
103119
}
104120

121+
// SYNC-M1: 广播解决结果到其他在线设备(原实现不广播,
122+
// 其他设备只能等下一次 Pull 才收敛,实时性延迟)
123+
BroadcastOperation(userID, req.DeviceID, []WSOperationPayload{
124+
{
125+
EntityType: req.EntityType,
126+
EntityID: req.EntityID,
127+
Operation: "update",
128+
Data: req.Data,
129+
Version: req.Version,
130+
ServerSeqNo: 0, // Resolve 的广播无明确序号,客户端按实体版本收敛
131+
DeviceID: req.DeviceID,
132+
},
133+
})
134+
105135
case "remote":
106136
// Server wins: no changes needed; client will pull latest.
107137

‎server/sync-service/handlers/websocket.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -156,6 +156,8 @@ func HandleWebSocketWithGin(c *gin.Context, userID string, deviceID string) {
156156
DeviceID: deviceID,
157157
Conn: conn,
158158
Send: make(chan []byte, 256),
159+
// SYNC-H2: 每连接最多 2 个在飞 sync_request 查询
160+
syncReqSlots: make(chan struct{}, 2),
159161
}
160162

161163
// M12: 每用户连接数上限检查,超限则拒绝(发送 ClosePolicyViolation 帧后关闭)

‎server/sync-service/handlers/ws_connection.go‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,10 @@ type WSConnection struct {
4040
Send chan []byte
4141
closed bool
4242
mu sync.Mutex
43+
44+
// SYNC-H2: sync_request 在飞查询信号槽(容量 2)——客户端可连续发送
45+
// 任意数量 sync_request,每个都启动 goroutine 查库会打满数据库连接池
46+
syncReqSlots chan struct{}
4347
}
4448

4549
// close safely marks the connection as closed and closes the underlying socket.
@@ -121,7 +125,17 @@ func (c *WSConnection) readPump() {
121125
}
122126
// M10: 查库改为异步执行(goroutine),结果经 Send 通道回投,
123127
// 避免慢查询阻塞 readPump 导致心跳超时。
124-
go c.handleSyncRequest(payload.SinceVersion)
128+
// SYNC-H2: 在飞查询超限时丢弃本次请求(客户端轮询会重发),
129+
// 防止单连接洪泛打满数据库连接池(上限 50)拖垮全部用户
130+
select {
131+
case c.syncReqSlots <- struct{}{}:
132+
go func() {
133+
defer func() { <-c.syncReqSlots }()
134+
c.handleSyncRequest(payload.SinceVersion)
135+
}()
136+
default:
137+
log.Printf("[ws] sync_request dropped (in-flight limit reached) user=%s device=%s", c.UserID, c.DeviceID)
138+
}
125139

126140
case "operation":
127141
// Client-pushed operation over WebSocket (future enhancement).

‎server/sync-service/models/models.go‎

Lines changed: 12 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -42,14 +42,17 @@ type GlobalSeqNo struct {
4242

4343
// CRDTChange stores CRDT changesets uploaded by clients.
4444
// Each row is an immutable, append-only change record keyed by a global sequence number.
45+
// ChangesetHash 为 changeset 的 SHA-256(hex),配合 (user_id, changeset_hash) 唯一
46+
// 索引实现 CRDT 推送幂等(SYNC-M3):网络重试的重复 changeset 不再重复落库。
4547
type CRDTChange struct {
46-
ID uint `gorm:"primaryKey"`
47-
ServerSeqNo int64 `gorm:"uniqueIndex;not null" json:"serverSeqNo"`
48-
DeviceID string `gorm:"index;not null" json:"deviceId"`
49-
UserID string `gorm:"index:idx_crdt_user;not null" json:"userId"`
50-
TableName string `gorm:"index;not null" json:"tableName"`
51-
EntityID string `gorm:"not null" json:"entityId"`
52-
Changeset string `gorm:"type:text;not null" json:"changeset"` // base64-encoded Automerge changes
53-
Operation string `gorm:"not null" json:"operation"` // create|update|delete
54-
CreatedAt string `gorm:"not null" json:"createdAt"`
48+
ID uint `gorm:"primaryKey"`
49+
ServerSeqNo int64 `gorm:"uniqueIndex;not null" json:"serverSeqNo"`
50+
DeviceID string `gorm:"index;not null" json:"deviceId"`
51+
UserID string `gorm:"index:idx_crdt_user;uniqueIndex:idx_crdt_user_hash,priority:1;not null" json:"userId"`
52+
TableName string `gorm:"index;not null" json:"tableName"`
53+
EntityID string `gorm:"not null" json:"entityId"`
54+
Changeset string `gorm:"type:text;not null" json:"changeset"` // base64-encoded Automerge changes
55+
ChangesetHash string `gorm:"uniqueIndex:idx_crdt_user_hash,priority:2;not null" json:"-"` // sha256(changeset)
56+
Operation string `gorm:"not null" json:"operation"` // create|update|delete
57+
CreatedAt string `gorm:"not null" json:"createdAt"`
5558
}

0 commit comments

Comments
 (0)