Skip to content

Commit b82668d

Browse files
committed
fix(sync): resolve low-risk logic issues SYNC2-L1~L5
- SYNC2-L1: oplog pull cursor advances only to last successfully applied version, preventing permanent skip of failed operations - SYNC2-L2: autoSync no longer counts lock/offline states as failures, avoiding spurious exponential backoff after network recovery - SYNC2-L3: getPendingCRDTChanges supports per-table filtering, fixing push starvation when multiple tables have pending changes - SYNC2-L4: offline queue max version read by version instead of createdAt to avoid duplicate versions within the same millisecond - SYNC2-L5: reset per-table CRDT in-memory doc on persistence failure, preventing baseline drift and MissingDependencyError
1 parent 7ef8ad8 commit b82668d

7 files changed

Lines changed: 93 additions & 24 deletions

File tree

‎client/src/lib/storage/writeWithLog.ts‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -169,8 +169,12 @@ async function applyCRDTChange(
169169
await enqueueCRDTChange(changeset);
170170
await crdtEngine.persistDoc(tableName);
171171
} catch (err) {
172-
// CRDT 路径失败不应阻塞主写入流程,仅记录警告
173-
console.warn('[writeWithLog] CRDT change failed (non-blocking):', err);
172+
// SYNC2-L5: CRDT 路径失败不阻塞主写入,但必须重置该表内存文档——
173+
// applyLocalChange 已前滚内存 doc,若 enqueue/persistDoc 失败,
174+
// 内存与队列/快照基线漂移,下次变更会生成非法 changeset
175+
//(MissingDependencyError),该表 CRDT 路径从此永久失效
176+
await crdtEngine.resetTable(tableName);
177+
console.warn('[writeWithLog] CRDT change failed, table reset (non-blocking):', err);
174178
}
175179
}
176180

‎client/src/lib/sync/OfflineQueue.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,8 +42,10 @@ export class OfflineQueue {
4242
const deviceId = getDeviceId();
4343

4444
// 在 Dexie 事务内完成"读 max version + 1 + 写入",避免版本号竞态
45+
// SYNC2-L4: 取最大值必须按 version(单调自增主键)而非 createdAt——
46+
// Date 毫秒精度下同一毫秒的并发写入会取到错误基线,产生重复版本号
4547
await db.transaction('rw', db.offlineQueue, async () => {
46-
const items = await db.offlineQueue.orderBy('createdAt').reverse().limit(1).toArray();
48+
const items = await db.offlineQueue.orderBy('version').reverse().limit(1).toArray();
4749
const version = items.length > 0 ? (items[0].version || 0) + 1 : 1;
4850

4951
await db.offlineQueue.add({

‎client/src/lib/sync/SyncEngine.ts‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,13 @@ export class SyncEngine {
7070
if (now - this.lastAutoSyncAt < gapMs) return;
7171
this.lastAutoSyncAt = now;
7272
const result = await this.sync();
73-
if (result.errors.length > 0) {
73+
// SYNC2-L2: 并发锁/离线是正常状态而非故障——原实现把这两个错误计入
74+
// autoSyncFailures,网络恢复瞬间的并发 sync 会误触发指数退避,
75+
// 导致自动同步长时间被抑制
76+
const actionableErrors = result.errors.filter(
77+
(e) => e !== 'Sync already in progress' && e !== 'Device is offline',
78+
);
79+
if (actionableErrors.length > 0) {
7480
this.autoSyncFailures += 1;
7581
} else {
7682
this.autoSyncFailures = 0;

‎client/src/lib/sync/crdtEngine.ts‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -285,6 +285,18 @@ export class CRDTEngine {
285285
this.docs.clear();
286286
this.initialized = false;
287287
}
288+
289+
/**
290+
* 重置单表 CRDT 文档(SYNC2-L5: 本地变更持久化失败后调用)
291+
* 从 crdt_docs 最新快照重建该表内存状态,丢弃本次未落盘的变更——
292+
* 避免内存 doc 前滚而队列/快照未跟上导致的基线漂移
293+
*(下次 applyLocalChange 基于漂移基线生成非法 changeset)。
294+
* 该变更最终仍会随下次业务写入重新生成,数据不丢。
295+
*/
296+
async resetTable(tableName: string): Promise<void> {
297+
this.docs.delete(tableName);
298+
await this.initDoc(tableName);
299+
}
288300
}
289301

290302
// ─── 单例 ────────────────────────────────────────────────────────────────────

‎client/src/lib/sync/crdtPersistence.ts‎

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,9 +23,22 @@ export async function enqueueCRDTChange(changeset: CRDTChangeset): Promise<void>
2323

2424
/**
2525
* 获取待上传的 changesets(按 seq 排序)
26+
* SYNC2-L3: 支持按表过滤——原实现仅支持全局 limit,多表启用 CRDT 时
27+
* crdtSyncChannel 每表取前 50 条再 filter,靠后表(数据量大)的变更
28+
* 可能被全局前 50 条挤掉,形成推送饥饿。
2629
*/
27-
export async function getPendingCRDTChanges(limit: number = 50): Promise<CRDTChangeRecord[]> {
30+
export async function getPendingCRDTChanges(
31+
limit: number = 50,
32+
tableName?: string,
33+
): Promise<CRDTChangeRecord[]> {
2834
const { db } = await import('@/lib/storage/database');
35+
if (tableName) {
36+
return db.crdtChanges
37+
.orderBy('seq')
38+
.filter(c => c.tableName === tableName)
39+
.limit(limit)
40+
.toArray();
41+
}
2942
return db.crdtChanges
3043
.orderBy('seq')
3144
.limit(limit)

‎client/src/lib/sync/crdtSyncChannel.ts‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,10 +25,12 @@ export async function crdtPush(
2525

2626
try {
2727
// 收集所有启用 CRDT 的表的待上传变更
28+
// SYNC2-L3: 按表分别取(每表各自 limit 50)——原实现每表取全局前 50
29+
// 条再 filter,靠后表的高 seq 变更会被靠前表挤掉,形成推送饥饿
2830
const allPending: CRDTChangeRecord[] = [];
2931
for (const tableName of CRDT_ENABLED_TABLES) {
30-
const pending = await getPendingCRDTChanges(50);
31-
allPending.push(...pending.filter(p => p.tableName === tableName));
32+
const pending = await getPendingCRDTChanges(50, tableName);
33+
allPending.push(...pending);
3234
}
3335

3436
if (allPending.length === 0) return { pushed: 0, errors };

‎client/src/lib/sync/oplogSyncChannel.ts‎

Lines changed: 47 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -122,39 +122,69 @@ export async function oplogPull(
122122
}>(`${basePath}/pull?deviceId=${encodeURIComponent(deviceId)}&sinceVersion=${lastVersion}`);
123123

124124
if (response.operations.length > 0) {
125-
await applyRemoteOperations(response.operations);
126-
setLastSyncVersion(response.latestVersion);
125+
// SYNC2-L1: 游标只推进到最后成功应用的版本——原实现直接用
126+
// response.latestVersion 一次性跳过全部,若中间某条 apply 失败
127+
//(表缺失/数据异常),该操作被永久跳过(sinceVersion 已越过它)
128+
const { applied, lastVersion: appliedVersion } = await applyRemoteOperations(response.operations);
129+
if (appliedVersion > lastVersion) {
130+
setLastSyncVersion(appliedVersion);
131+
}
132+
return { pulled: applied, errors };
127133
}
128134

129-
return { pulled: response.operations.length, errors };
135+
return { pulled: 0, errors };
130136
} catch (error: unknown) {
131137
const message = error instanceof Error ? error.message : String(error);
132138
errors.push(`Pull failed: ${message}`);
133139
return { pulled: 0, errors };
134140
}
135141
}
136142

137-
/** 将远端操作应用到本地 Dexie(create/update 用 put 幂等覆盖) */
143+
/**
144+
* 将远端操作应用到本地 Dexie(create/update 用 put 幂等覆盖)
145+
* 逐条 try/catch:单条失败即停止,返回已成功应用的计数与最后成功版本。
146+
* 未知 entityType 的操作跳过但不阻断游标推进(与应用成功等价,
147+
* 否则 sinceVersion 不推进会反复拉取同一批数据)。
148+
*/
138149
async function applyRemoteOperations(
139-
operations: Array<{ entityType: string; entityId: string; operation: string; data: unknown }>,
140-
): Promise<void> {
150+
operations: Array<{ entityType: string; entityId: string; operation: string; data: unknown; version: number }>,
151+
): Promise<{ applied: number; lastVersion: number }> {
141152
const { db } = await import('../storage/database');
153+
let applied = 0;
154+
let lastVersion = -1;
142155

143156
for (const op of operations) {
144157
const tableName = getEntityTableName(op.entityType);
145-
if (!tableName) continue;
146-
147-
const table = db.table(tableName);
148-
switch (op.operation) {
149-
case 'create':
150-
case 'update':
151-
await table.put(op.data);
152-
break;
153-
case 'delete':
154-
await table.delete(op.entityId);
155-
break;
158+
if (!tableName) {
159+
// 未知类型:跳过但继续推进游标,避免永久重复拉取
160+
lastVersion = Math.max(lastVersion, op.version);
161+
continue;
162+
}
163+
164+
try {
165+
const table = db.table(tableName);
166+
switch (op.operation) {
167+
case 'create':
168+
case 'update':
169+
await table.put(op.data);
170+
break;
171+
case 'delete':
172+
await table.delete(op.entityId);
173+
break;
174+
}
175+
applied += 1;
176+
lastVersion = Math.max(lastVersion, op.version);
177+
} catch (err) {
178+
// 单条失败:停止应用,游标停在失败前的最后成功版本,
179+
// 下轮拉取重试该操作(避免数据永久丢失)
180+
const message = err instanceof Error ? err.message : String(err);
181+
// eslint-disable-next-line no-console -- 远端操作应用失败需记录以便排查
182+
console.warn(`[oplogSyncChannel] 远端操作应用失败 [${op.entityType}:${op.entityId}] v${op.version}: ${message}`);
183+
break;
156184
}
157185
}
186+
187+
return { applied, lastVersion };
158188
}
159189

160190
/** entityType(服务端操作记录中的键)→ Dexie 表名映射 */

0 commit comments

Comments
 (0)