Files
2026-07-20 19:01:03 +08:00

23 KiB
Raw Permalink Blame History

步骤 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>>>

任务

  1. 成员实体表改造为 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
        }
    }
    
  2. OnClientEnter 归并逻辑

    // ChannelViewModel.kt
    fun handleClientEnter(data: String) {
        val client = Json.decodeFromString<ClientInfo>(data)
        // 按 ID 覆盖,重复事件不会重复计数
        repository.upsertClient(client)
    }
    

    关键约束

    • 按 ClientInfo.ID 覆盖,禁止使用 +1 增量累加频道人数
    • 频道人数由 channelClients[channelId].size 实时派生
    • 已有基线时,进入事件覆盖旧数据;无基线时,插入新条目
  3. 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
    • 未知成员不得静默忽略,必须触发补偿同步
  4. 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-opMap.remove 天然满足)
    • 不使用 -1 减量维护频道人数
    • IsSelf 为 true 时走踢出/断开流程,不从成员表删除
  5. 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

任务

  1. 补偿同步触发器

    // 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)
            // 补偿同步失败不阻塞业务,等待下次触发
        }
    }
    
  2. 频道基线补偿同步

    // 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)
            }
        }
    }
    
  3. 补偿同步与增量事件的协调

    // 在 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 重连流程

目标:断开重连后,清空旧会话状态,重新执行完整首次同步,确保数据与服务器完全一致。

触发条件

  • 网络恢复后自动重连
  • 用户手动触发重连
  • 被踢出后重新连接

任务

  1. 重连状态机

    // 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()
    }
    
  2. 重连流程实现

    // 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
    }
    
  3. 断开事件处理中的重连触发

    // 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
            }
        }
    }
    
  4. 清理会话数据的完整性

    // 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 请求失败时,正确流转同步状态,允许重试,防止在不完整数据上执行业务操作。

任务

  1. 同步状态机完善

    // 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(补偿同步或重连)
    
  2. 首次同步失败处理

    // 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)
            // 不进入业务就绪,允许重试
        }
    }
    
  3. 重试机制

    // 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 显示重试按钮
    }
    
  4. 手动重试入口

    // ServerViewModel.kt
    fun retrySync() {
        viewModelScope.launch {
            performSyncWithRetry()
        }
    }
    
  5. 同步失败时的 UI 保护

    // 在业务操作前检查同步状态
    fun switchChannel(channelId: Long, password: String? = null) {
        if (_syncState.value !is SyncState.Synchronized) {
            _error.value = "数据未同步,请等待同步完成或点击重试"
            return
        }
        // 执行频道切换...
    }
    

关键约束

  • 任一核心请求(ListChannels / ListClients / ClientID)失败则不标记业务就绪
  • 允许重试,避免在不完整数据上执行业务操作
  • 同步失败期间,业务操作(切换频道、发消息等)应被阻止
  • 补偿同步失败不进入 SyncFailed,仅记录日志等待下次触发

10.5 数据一致性

目标:在长期运行中,检测并修复可能的数据不一致(如频道列表过期、成员数据漂移)。

任务

  1. 频道列表过期检测

    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 分钟
    }
    
  2. 被动刷新策略

    在用户执行关键操作时,检查并刷新过期数据:

    // 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()
    }
    
  3. 成员数据校验

    当成员引用了本地不存在的频道时,触发频道基线刷新:

    // 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()
        }
    }
    
  4. 后台一致性检查(可选)

    // 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 API
  • docs/implementation/04_连接与首次同步.md - 首次同步实现、SyncState 状态机
  • docs/implementation/09_断开连接.md - 断开连接清理逻辑