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

702 lines
23 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 步骤 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) 查询和幂等更新。
```kotlin
// 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 归并逻辑**
```kotlin
// ChannelViewModel.kt
fun handleClientEnter(data: String) {
val client = Json.decodeFromString<ClientInfo>(data)
// 按 ID 覆盖,重复事件不会重复计数
repository.upsertClient(client)
}
```
**关键约束**
- 按 ClientInfo.ID 覆盖,**禁止**使用 `+1` 增量累加频道人数
- 频道人数由 `channelClients[channelId].size` 实时派生
- 已有基线时,进入事件覆盖旧数据;无基线时,插入新条目
3. **OnClientMoved 归并逻辑**
```kotlin
// 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 归并逻辑**
```kotlin
// 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 数据类**
```kotlin
// 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. **补偿同步触发器**
```kotlin
// 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. **频道基线补偿同步**
```kotlin
// 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. **补偿同步与增量事件的协调**
```kotlin
// 在 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. **重连状态机**
```kotlin
// 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. **重连流程实现**
```kotlin
// 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. **断开事件处理中的重连触发**
```kotlin
// 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. **清理会话数据的完整性**
```kotlin
// 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. **同步状态机完善**
```kotlin
// 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. **首次同步失败处理**
```kotlin
// 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. **重试机制**
```kotlin
// 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. **手动重试入口**
```kotlin
// ServerViewModel.kt
fun retrySync() {
viewModelScope.launch {
performSyncWithRetry()
}
}
```
5. **同步失败时的 UI 保护**
```kotlin
// 在业务操作前检查同步状态
fun switchChannel(channelId: Long, password: String? = null) {
if (_syncState.value !is SyncState.Synchronized) {
_error.value = "数据未同步,请等待同步完成或点击重试"
return
}
// 执行频道切换...
}
```
**关键约束**
- 任一核心请求(ListChannels / ListClients / ClientID)失败则不标记业务就绪
- 允许重试,避免在不完整数据上执行业务操作
- 同步失败期间,业务操作(切换频道、发消息等)应被阻止
- 补偿同步失败不进入 SyncFailed,仅记录日志等待下次触发
### 10.5 数据一致性
**目标**:在长期运行中,检测并修复可能的数据不一致(如频道列表过期、成员数据漂移)。
**任务**
1. **频道列表过期检测**
SDK 不提供频道创建/更新/删除事件,因此频道列表可能随时间过期。
```kotlin
// 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. **被动刷新策略**
在用户执行关键操作时,检查并刷新过期数据:
```kotlin
// 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. **成员数据校验**
当成员引用了本地不存在的频道时,触发频道基线刷新:
```kotlin
// 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. **后台一致性检查(可选)**
```kotlin
// 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` - 断开连接清理逻辑