23 KiB
步骤 10:状态同步进阶
实现高级状态同步:增量同步完善、补偿同步触发、重连全量同步、同步失败处理、数据一致性校验。 依赖:步骤 09(断开连接)已完成,首次同步(①)已在步骤 04 实现。
一、目标
- 增量同步(②)完善 — OnClientEnter / OnClientMoved / OnClientLeave 幂等归并
- 补偿同步(③)触发机制 — 未知实体引用时自动修复基线
- 重连全量同步(⑥) — 清理旧状态 → 重新执行首次同步
- 同步失败处理(⑦) — 状态机流转、重试、降级
- 数据一致性校验 — 周期性校验与频道列表刷新
二、任务清单
10.1 增量同步完善
目标:将步骤 04 中注册的事件处理器补全为完整的幂等归并逻辑,确保成员实体表在持续事件流中保持一致。
前置条件:
- 步骤 04 已实现首次同步(ListChannels + ListClients + ClientID)
- 步骤 04 已注册 OnClientEnter / OnClientLeave / OnClientMoved 事件处理器
ChannelRepository已有_clients: MutableStateFlow<List<ClientInfo>>和_channelClients: MutableStateFlow<Map<Long, List<ClientInfo>>>
任务:
-
成员实体表改造为 Map 结构
将成员存储从
List<ClientInfo>改为Map<Int, ClientInfo>,以 ClientID 为 key 实现 O(1) 查询和幂等更新。// data/Repository.kt class ChannelRepository { // 成员实体表:ClientID → ClientInfo private val _clientMap = MutableStateFlow<Map<Int, ClientInfo>>(emptyMap()) val clientMap: StateFlow<Map<Int, ClientInfo>> = _clientMap // 派生:按频道 ID 索引的成员列表(由 clientMap 自动计算) val channelClients: StateFlow<Map<Long, List<ClientInfo>>> = _clientMap.map { map -> map.values.groupBy { it.channelId } }.stateIn(scope, SharingStarted.WhileSubscribed(), emptyMap()) // 派生:成员列表(兼容旧接口) val clients: StateFlow<List<ClientInfo>> = _clientMap.map { it.values.toList() } .stateIn(scope, SharingStarted.WhileSubscribed(), emptyList()) /** 首次同步:原子替换整个成员表 */ fun setClientBaseline(clients: List<ClientInfo>) { _clientMap.value = clients.associateBy { it.id } } /** 增量更新:按 ID 幂等插入或覆盖 */ fun upsertClient(client: ClientInfo) { _clientMap.update { it + (client.id to client) } } /** 增量更新:按 ID 幂等删除(重复删除为 no-op) */ fun removeClient(clientId: Int) { _clientMap.update { it - clientId } } /** 增量更新:移动成员到目标频道 */ fun moveClient(clientId: Int, targetChannelId: Long) { _clientMap.update { map -> val existing = map[clientId] ?: return@update map // 不存在则 no-op if (existing.channelId == targetChannelId) return@update map // 相同频道则 no-op map + (clientId to existing.copy(channelId = targetChannelId)) } } /** 查询成员是否存在 */ fun hasClient(clientId: Int): Boolean = _clientMap.value.containsKey(clientId) /** 查询频道是否存在 */ fun hasChannel(channelId: Long): Boolean = _channels.value.any { it.id == channelId } /** 清理会话数据 */ fun clearSession() { _clientMap.value = emptyMap() _channels.value = emptyList() _selfClientId.value = null } } -
OnClientEnter 归并逻辑
// ChannelViewModel.kt fun handleClientEnter(data: String) { val client = Json.decodeFromString<ClientInfo>(data) // 按 ID 覆盖,重复事件不会重复计数 repository.upsertClient(client) }关键约束:
- 按 ClientInfo.ID 覆盖,禁止使用
+1增量累加频道人数 - 频道人数由
channelClients[channelId].size实时派生 - 已有基线时,进入事件覆盖旧数据;无基线时,插入新条目
- 按 ClientInfo.ID 覆盖,禁止使用
-
OnClientMoved 归并逻辑
// ChannelViewModel.kt fun handleClientMoved(data: String) { val event = Json.decodeFromString<ClientMovedEvent>(data) // 判断是否为自己 if (event.clientId == repository.selfClientId.value) { repository.updateSelfChannel(event.targetChannelId) } // 更新成员位置 if (repository.hasClient(event.clientId)) { repository.moveClient(event.clientId, event.targetChannelId) } else { // 成员不存在 → 触发补偿同步(见 10.2) triggerClientCompensationSync() } }关键约束:
- 目标频道 ID 是移动后的服务器事实,直接覆盖旧的 ChannelID
- 相同目标频道视为 no-op
- 未知成员不得静默忽略,必须触发补偿同步
-
OnClientLeave 归并逻辑
// ChannelViewModel.kt fun handleClientLeave(data: String) { val event = Json.decodeFromString<ClientLeftViewEvent>(data) if (event.isSelf) { // 自己被踢出或离开 → 由步骤 09 处理 handleKicked(event.reasonMessage) return } // 按 ID 删除,重复删除安全地保持 no-op repository.removeClient(event.clientId) }关键约束:
- 重复删除必须为 no-op(Map.remove 天然满足)
- 不使用
-1减量维护频道人数 IsSelf为 true 时走踢出/断开流程,不从成员表删除
-
ClientMovedEvent / ClientLeftViewEvent 数据类
// data/Models.kt data class ClientMovedEvent( val clientId: Int, val targetChannelId: Long, val reasonId: Int = 0, val invokerId: Int = 0, val invokerName: String = "", val invokerUid: String = "" ) data class ClientLeftViewEvent( val clientId: Int, val reasonId: Int = 0, // 0=正常离开, 4=频道踢, 5=服务器踢 val reasonMessage: String = "", val isSelf: Boolean = false )
10.2 补偿同步机制
目标:当增量事件引用了本地不存在的实体(ClientID 或 ChannelID),自动触发完整列表请求修复基线。
触发条件:
| 场景 | 检测方式 | 补偿动作 |
|---|---|---|
| OnClientMoved 引用未知 ClientID | !repository.hasClient(event.clientId) |
重新调用 ListClients |
| OnClientLeave 引用未知 ClientID | !repository.hasClient(event.clientId) |
重新调用 ListClients(可选,删除本身是 no-op) |
| 成员引用未知 ChannelID | !repository.hasChannel(member.channelId) |
重新调用 ListChannels |
任务:
-
补偿同步触发器
// ChannelViewModel.kt private var compensationSyncJob: Job? = null /** * 触发成员基线补偿同步。 * 使用防抖:连续多个未知实体事件只触发一次 ListClients。 */ private fun triggerClientCompensationSync() { compensationSyncJob?.cancel() compensationSyncJob = viewModelScope.launch { delay(300) // 防抖 300ms performCompensationSync() } } private suspend fun performCompensationSync() { try { Log.w(TAG, "Compensation sync: rebuilding client baseline") val clients = TSBridge.listClients() repository.setClientBaseline(clients) Log.i(TAG, "Compensation sync completed: ${clients.size} clients") } catch (e: Exception) { Log.e(TAG, "Compensation sync failed", e) // 补偿同步失败不阻塞业务,等待下次触发 } } -
频道基线补偿同步
// ChannelViewModel.kt private fun triggerChannelCompensationSync() { viewModelScope.launch { try { Log.w(TAG, "Compensation sync: rebuilding channel baseline") val channels = TSBridge.listChannels() repository.setChannelBaseline(channels) Log.i(TAG, "Compensation sync completed: ${channels.size} channels") } catch (e: Exception) { Log.e(TAG, "Channel compensation sync failed", e) } } } -
补偿同步与增量事件的协调
// 在 handleClientMoved 中集成 fun handleClientMoved(data: String) { val event = Json.decodeFromString<ClientMovedEvent>(data) if (event.clientId == repository.selfClientId.value) { repository.updateSelfChannel(event.targetChannelId) } if (repository.hasClient(event.clientId)) { repository.moveClient(event.clientId, event.targetChannelId) } else { // 检测到未知成员,触发补偿同步 Log.w(TAG, "Unknown client ${event.clientId} in move event, triggering compensation") triggerClientCompensationSync() } }
关键约束:
- 补偿同步使用完整列表替换,不是增量合并
- 使用防抖避免事件风暴时重复调用 ListClients
- 补偿同步失败不阻塞业务,等待下次事件触发重试
- 补偿同步完成后,UI 通过 StateFlow 自动刷新
时序(对应流程文档 §六):
OnClientMoved(未知 ClientID)
→ 检测到不一致
→ 调用 ListClients()
→ 完整成员列表返回
→ 替换成员基线(setClientBaseline)
→ UI 自动刷新
10.3 重连流程
目标:断开重连后,清空旧会话状态,重新执行完整首次同步,确保数据与服务器完全一致。
触发条件:
- 网络恢复后自动重连
- 用户手动触发重连
- 被踢出后重新连接
任务:
-
重连状态机
// ServerViewModel.kt sealed class ReconnectState { object Idle : ReconnectState() object Detecting : ReconnectState() // 检测到断开 object Cleaning : ReconnectState() // 清理旧会话 object Reconnecting : ReconnectState() // 重新连接中 object Syncing : ReconnectState() // 首次同步中 object Ready : ReconnectState() // 恢复就绪 data class Failed(val error: Throwable) : ReconnectState() } -
重连流程实现
// ServerViewModel.kt private suspend fun performReconnect(config: ConnectionConfig) { _reconnectState.value = ReconnectState.Cleaning // 1. 清理旧会话数据(频道、成员、Pending 全部移除) repository.clearSession() _connectionState.value = ConnectionState.Disconnected // 2. 断开旧连接(确保资源释放) try { TSBridge.disconnect() } catch (_: Exception) {} _reconnectState.value = ReconnectState.Reconnecting // 3. 重新执行连接流程(参见步骤 04) registerEventHandlers() val connectResult = TSBridge.connect(config) if (connectResult.isFailure) { _reconnectState.value = ReconnectState.Failed(connectResult.exceptionOrNull()!!) return } val waitResult = TSBridge.waitConnected() if (waitResult.isFailure) { _reconnectState.value = ReconnectState.Failed(waitResult.exceptionOrNull()!!) return } // 4. OnConnected 事件将自动触发首次同步(步骤 04 已实现) _reconnectState.value = ReconnectState.Syncing } -
断开事件处理中的重连触发
// ServerViewModel.kt fun handleDisconnected(data: String) { val error = Json.decodeFromString<DisconnectedEvent>(data) when { error.isKicked -> { // 被踢出:显示原因,不自动重连 _kickReason.value = error.reasonMessage repository.clearSession() _connectionState.value = ConnectionState.Disconnected } isAutoReconnectEnabled -> { // 网络断开:尝试自动重连 viewModelScope.launch { delay(RECONNECT_DELAY) // 等待网络恢复 performReconnect(lastConfig) } } else -> { // 手动断开或不自动重连 repository.clearSession() _connectionState.value = ConnectionState.Disconnected } } } -
清理会话数据的完整性
// data/Repository.kt fun clearSession() { _clientMap.value = emptyMap() _channels.value = emptyList() _selfClientId.value = null _currentChannelId.value = null // 清除所有待确认状态 _pendingChannelMove.value = null }
关键约束:
- 重连前必须清空旧会话状态,防止旧成员、频道和 Pending 污染新连接
- 清理操作在连接断开后执行,避免竞态
- 重连后的首次同步复用步骤 04 的
performInitialSync()逻辑 - 被踢出不自动重连,由用户决定
时序(对应流程文档 §七):
检测到断开
→ 状态改为 Cleaning
→ 清理旧会话数据(频道、成员、Pending)
→ 断开旧连接
→ 重新 Connect + WaitConnected
→ OnConnected 触发首次同步
→ ListChannels + ListClients + ClientID
→ 原子提交新基线
→ 恢复业务就绪
10.4 同步失败处理
目标:当 ListChannels 或 ListClients 请求失败时,正确流转同步状态,允许重试,防止在不完整数据上执行业务操作。
任务:
-
同步状态机完善
// data/Models.kt sealed class SyncState { object Unsynced : SyncState() // 已连接但尚无完整数据 object Syncing : SyncState() // 调用 ListChannels 和 ListClients object Synchronized : SyncState() // 列表基线可供 UI 使用 data class SyncFailed(val error: Throwable) : SyncState() // 同步失败 }状态流转:
Unsynced → Syncing → Synchronized(正常路径) Syncing → SyncFailed → Syncing(重试路径) Synchronized → Syncing(补偿同步或重连) -
首次同步失败处理
// ServerViewModel.kt private suspend fun performInitialSync() { _syncState.value = SyncState.Syncing try { // 并行请求 val channelsDeferred = async { TSBridge.listChannels() } val clientsDeferred = async { TSBridge.listClients() } val selfIdDeferred = async { TSBridge.getClientId() } val channels = channelsDeferred.await() val clients = clientsDeferred.await() val selfId = selfIdDeferred.await() // 原子提交 repository.setChannelBaseline(channels) repository.setClientBaseline(clients) repository.setSelfClientId(selfId) _syncState.value = SyncState.Synchronized _connectionState.value = ConnectionState.Ready } catch (e: Exception) { Log.e(TAG, "Initial sync failed", e) _syncState.value = SyncState.SyncFailed(e) // 不进入业务就绪,允许重试 } } -
重试机制
// ServerViewModel.kt private suspend fun performSyncWithRetry(maxRetries: Int = 3) { var retryCount = 0 while (retryCount < maxRetries) { performInitialSync() if (_syncState.value is SyncState.Synchronized) { return // 成功 } retryCount++ if (retryCount < maxRetries) { Log.w(TAG, "Sync retry $retryCount/$maxRetries") delay(1000L * retryCount) // 递增延迟 } } Log.e(TAG, "Sync failed after $maxRetries retries") // 保持 SyncFailed 状态,UI 显示重试按钮 } -
手动重试入口
// ServerViewModel.kt fun retrySync() { viewModelScope.launch { performSyncWithRetry() } } -
同步失败时的 UI 保护
// 在业务操作前检查同步状态 fun switchChannel(channelId: Long, password: String? = null) { if (_syncState.value !is SyncState.Synchronized) { _error.value = "数据未同步,请等待同步完成或点击重试" return } // 执行频道切换... }
关键约束:
- 任一核心请求(ListChannels / ListClients / ClientID)失败则不标记业务就绪
- 允许重试,避免在不完整数据上执行业务操作
- 同步失败期间,业务操作(切换频道、发消息等)应被阻止
- 补偿同步失败不进入 SyncFailed,仅记录日志等待下次触发
10.5 数据一致性
目标:在长期运行中,检测并修复可能的数据不一致(如频道列表过期、成员数据漂移)。
任务:
-
频道列表过期检测
SDK 不提供频道创建/更新/删除事件,因此频道列表可能随时间过期。
// ChannelViewModel.kt private var lastChannelRefreshTime: Long = 0 /** 检查频道列表是否需要刷新 */ private fun isChannelListStale(): Boolean { val elapsed = System.currentTimeMillis() - lastChannelRefreshTime return elapsed > CHANNEL_LIST_STALE_THRESHOLD // 建议 5 分钟 } companion object { const val CHANNEL_LIST_STALE_THRESHOLD = 5 * 60 * 1000L // 5 分钟 } -
被动刷新策略
在用户执行关键操作时,检查并刷新过期数据:
// ChannelViewModel.kt suspend fun refreshChannelsIfNeeded() { if (isChannelListStale()) { try { val channels = TSBridge.listChannels() repository.setChannelBaseline(channels) lastChannelRefreshTime = System.currentTimeMillis() } catch (e: Exception) { Log.e(TAG, "Channel refresh failed", e) // 不阻塞业务,使用旧数据 } } } /** 浏览频道时刷新 */ fun onChannelListVisible() { viewModelScope.launch { refreshChannelsIfNeeded() } } /** 切换频道前刷新 */ suspend fun beforeChannelSwitch() { refreshChannelsIfNeeded() } -
成员数据校验
当成员引用了本地不存在的频道时,触发频道基线刷新:
// ChannelViewModel.kt fun validateMemberData() { val channels = repository.channels.value.map { it.id }.toSet() val clients = repository.clientMap.value.values val unknownChannelIds = clients .map { it.channelId } .filter { it !in channels } .toSet() if (unknownChannelIds.isNotEmpty()) { Log.w(TAG, "Found clients referencing unknown channels: $unknownChannelIds") triggerChannelCompensationSync() } } -
后台一致性检查(可选)
// ServerViewModel.kt private var consistencyCheckJob: Job? = null fun startConsistencyCheck() { consistencyCheckJob = viewModelScope.launch { while (isActive) { delay(CONSISTENCY_CHECK_INTERVAL) if (_syncState.value is SyncState.Synchronized) { channelViewModel.validateMemberData() } } } } fun stopConsistencyCheck() { consistencyCheckJob?.cancel() consistencyCheckJob = null } companion object { const val CONSISTENCY_CHECK_INTERVAL = 60 * 1000L // 1 分钟 }
关键约束:
ListChannels是频道目录的唯一权威来源(SDK 无频道变更事件)- 不能假设频道列表依靠
On*事件永久保持最新 - 刷新失败时使用旧数据,不阻塞业务操作
- 一致性检查为低优先级,不影响正常事件流的实时性
三、验收标准
功能验收
-
增量同步
- OnClientEnter 事件正确添加成员到基线,重复事件不重复计数
- OnClientMoved 事件正确更新成员频道位置,相同目标频道为 no-op
- OnClientLeave 事件正确删除成员,重复删除为 no-op
- 自己的移动事件正确更新自身频道位置
- 频道人数由成员实体表实时派生,不使用独立计数器
-
补偿同步
- OnClientMoved 引用未知 ClientID 时自动触发 ListClients
- 补偿同步使用完整列表替换成员基线
- 连续多个未知实体事件只触发一次补偿同步(防抖)
- 补偿同步失败不阻塞业务
-
重连全量同步
- 重连前清空旧会话数据(频道、成员、Pending)
- 重连后自动执行首次同步
- 同步完成后恢复业务就绪
- 被踢出不自动重连
-
同步失败处理
- 首次同步失败进入 SyncFailed 状态
- 同步失败期间业务操作被阻止
- 提供手动重试入口
- 重试最多 3 次,递增延迟
-
数据一致性
- 频道列表超过 5 分钟未刷新时,关键操作前自动刷新
- 成员引用未知频道时触发频道基线补偿同步
- 刷新失败时使用旧数据,不阻塞业务
性能验收
- 增量同步单次事件处理 < 10ms
- 补偿同步(ListClients)在 3 秒内完成
- 重连全量同步在 5 秒内完成
- 防抖机制避免事件风暴时的重复请求
代码质量验收
- 成员实体表操作线程安全(StateFlow + immutable Map)
- 补偿同步防抖使用协程取消,无泄漏
- 所有网络操作在 IO 线程执行
- 日志覆盖关键状态转换和异常
测试用例
| 场景 | 操作 | 预期结果 |
|---|---|---|
| 正常增量流 | 其他用户进入/移动/离开 | 成员表实时更新,频道人数正确 |
| 重复进入事件 | 同一用户连续两次 OnClientEnter | 成员表只有一条记录,无重复计数 |
| 未知成员移动 | OnClientMoved 引用不存在的 ClientID | 自动触发 ListClients 补偿同步 |
| 补偿同步防抖 | 连续 5 个未知实体事件 | 只触发 1 次 ListClients |
| 正常重连 | 网络断开后恢复 | 清理旧数据 → 重连 → 同步 → 就绪 |
| 被踢后重连 | 被服务器踢出后手动重连 | 显示踢出原因 → 清理 → 重连 → 同步 |
| 同步失败重试 | 首次同步网络超时 | 进入 SyncFailed → 点击重试 → 成功 |
| 同步失败阻塞 | 同步失败时切换频道 | 显示"数据未同步"提示 |
| 频道列表过期 | 5 分钟后切换频道 | 自动刷新频道列表再切换 |
| 成员引用未知频道 | 成员的 ChannelID 在本地不存在 | 触发频道基线补偿同步 |
四、参考文档
docs/流程/08_状态同步.md- 完整同步机制(②③⑥⑦)docs/流程/02_浏览频道.md- 成员实体状态树、事件依赖矩阵docs/sdk文档-go.md- OnClientEnter / OnClientLeave / OnClientMoved 事件处理器、ListClients / ListChannels APIdocs/implementation/04_连接与首次同步.md- 首次同步实现、SyncState 状态机docs/implementation/09_断开连接.md- 断开连接清理逻辑