diff --git a/README.md b/README.md index 65a781a..0720f7e 100644 --- a/README.md +++ b/README.md @@ -107,6 +107,10 @@ adb install -r app/build/outputs/apk/debug/app-debug.apk ## 连接异常排查 +若会话被其他 Codex 客户端占用,Android 会以 external-writer 模式查看并自动更新原会话历史。可选的 [Desktop companion](companion/README.md) 支持将文字交给持有原会话的电脑客户端,避免另建对话;当前为实验功能,需要独立部署伴随代理。 + +同步时直接更新最新回合的完整处理过程;历史回合的详情按需加载并保留缓存。同一服务端摘要不会因为详情中多了中间回复而反复重置加载状态。 + 生产版本通过 Logcat 的 `CodexConnection` 标签记录连接失败、非主动断开、请求超时和会话恢复失败,正常连接、消息收发和主动取消不记录 ```bash diff --git a/app/build.gradle.kts b/app/build.gradle.kts index 4e197b7..047d8a1 100644 --- a/app/build.gradle.kts +++ b/app/build.gradle.kts @@ -1,3 +1,8 @@ +import org.gradle.api.file.ConfigurableFileCollection +import org.gradle.api.tasks.Classpath +import org.gradle.api.tasks.testing.Test +import org.gradle.process.CommandLineArgumentProvider + plugins { alias(libs.plugins.android.application) alias(libs.plugins.kotlin.compose) @@ -6,6 +11,21 @@ plugins { alias(libs.plugins.room) } +val mockitoAgent = configurations.create("mockitoAgent") + +abstract class MockitoAgentProvider : CommandLineArgumentProvider { + @get:Classpath + abstract val agentJar: ConfigurableFileCollection + + override fun asArguments(): List = listOf("-javaagent:${agentJar.asPath}") +} + +tasks.withType().configureEach { + jvmArgumentProviders.add(objects.newInstance().apply { + agentJar.from(mockitoAgent) + }) +} + val releaseSigningEnabled = providers.gradleProperty("releaseSigning").orNull == "true" fun requiredSigningEnvironmentVariable(name: String): String = @@ -102,6 +122,8 @@ dependencies { debugImplementation(libs.androidx.compose.ui.tooling) testImplementation(libs.junit) + testImplementation(libs.mockito.core) + mockitoAgent(libs.mockito.core) { isTransitive = false } testImplementation(libs.kotlinx.coroutines.test) ksp(libs.androidx.room.compiler) diff --git a/app/src/main/java/com/hebo/codex/data/remote/AppServerApi.kt b/app/src/main/java/com/hebo/codex/data/remote/AppServerApi.kt index 3c98b48..998a9c3 100644 --- a/app/src/main/java/com/hebo/codex/data/remote/AppServerApi.kt +++ b/app/src/main/java/com/hebo/codex/data/remote/AppServerApi.kt @@ -160,6 +160,7 @@ class AppServerApi( StableRpcMethod.THREAD_READ, buildJsonObject { put("threadId", threadId) + put("includeTurns", false) }, ).jsonObject @@ -213,6 +214,39 @@ class AppServerApi( timeoutMillis = input.submissionTimeoutMillis(), ).jsonObject + // Optional codex-android companion extension, separate from Codex CLI RPCs. + suspend fun discoverDesktopControl(threadId: String): Boolean = try { + val result = connection.request( + "codexAndroid/desktop/discover", + buildJsonObject { put("threadId", threadId) }, + ).jsonObject + if (result.requiredString("threadId") != threadId || result.requiredInt("protocolVersion") != 1) { + throw ProtocolException("电脑伴随服务的会话或协议不匹配") + } + result.requiredBoolean("available") + } catch (error: RpcException) { + val unsupported = error.code == -32601 || (error.code == -32600 && + error.message.startsWith("Invalid request: unknown variant `codexAndroid/desktop/discover`")) + if (!unsupported) throw error + false + } + + suspend fun sendDesktopMessage(threadId: String, clientMessageId: String, text: String): JsonObject = + connection.request( + "codexAndroid/desktop/turn/start", + buildJsonObject { + put("threadId", threadId) + put("clientUserMessageId", clientMessageId) + put("text", text) + }, + timeoutMillis = 120_000, + ).jsonObject.also { result -> + if (result.requiredString("threadId") != threadId || + result.requiredString("clientUserMessageId") != clientMessageId + ) throw ProtocolException("电脑端的发送确认与当前消息不一致") + result.requiredObject("turn").requiredString("id") + } + suspend fun steerTurn( threadId: String, turnId: String, diff --git a/app/src/main/java/com/hebo/codex/domain/CodexRepository.kt b/app/src/main/java/com/hebo/codex/domain/CodexRepository.kt index 693f92e..08a6f95 100644 --- a/app/src/main/java/com/hebo/codex/domain/CodexRepository.kt +++ b/app/src/main/java/com/hebo/codex/domain/CodexRepository.kt @@ -501,6 +501,25 @@ class CodexRepository( return conversation } + suspend fun refreshConversation(serverId: String, threadId: String) { + conversationMutex.withLock { + val key = ConversationKey(serverId, threadId) + refreshConversation(key, requireConversation(serverId, threadId)) + } + requestRuntimeQueueRefresh(ConversationKey(serverId, threadId)) + } + + internal suspend fun syncExternalConversation(serverId: String, threadId: String) { + val key = ConversationKey(serverId, threadId) + conversationMutex.withLock { + val state = requireConversation(serverId, threadId) + // Manual refresh/reconnect may have already restored the live stream. + if (!state.value.isExternalWriter) return + refreshConversation(key, state, failOnError = true, reportErrors = false) + } + requestRuntimeQueueRefresh(key) + } + suspend fun loadOlderConversationTurns(serverId: String, threadId: String): Boolean { val key = ConversationKey(serverId, threadId) val state = requireConversation(serverId, threadId) @@ -573,6 +592,8 @@ class CodexRepository( val text = input.filterIsInstance().joinToString("\n") { it.text }.trim() require(text.isNotEmpty() || input.any { it !is UserInputPart.Text }) { "消息内容不能为空" } val state = requireConversation(serverId, threadId) + state.value.requireWritable() + val targetApi = writableApi(serverId, threadId) val activeTurnId = state.value.activeTurnId val clientId = UUID.randomUUID().toString() val messageTitle = input.conversationTitle() @@ -587,9 +608,9 @@ class CodexRepository( try { val result = if (activeTurnId == null) { - api(serverId).startTurn(threadId, clientId, input, options) + targetApi.startTurn(threadId, clientId, input, options) } else { - api(serverId).steerTurn(threadId, activeTurnId, clientId, input) + targetApi.steerTurn(threadId, activeTurnId, clientId, input) } val turnId = if (activeTurnId == null) { result.requiredObject("turn").requiredString("id") @@ -611,12 +632,30 @@ class CodexRepository( } } + suspend fun sendDesktopMessage(serverId: String, threadId: String, text: String) { + require(text.isNotBlank()) { "消息内容不能为空" } + val state = requireConversation(serverId, threadId) + check(state.value.canSendToDesktop) { state.value.viewOnlyMessage ?: "电脑端连接已失效,请刷新会话" } + val targetApi = api(serverId) + // Reconnection may restore local ownership or lose the desktop owner. + check(state.value.canSendToDesktop) { "会话连接已变化,请刷新后重试" } + check(!state.value.isRunning) { "电脑正在处理此会话,请完成后再发送" } + val result = targetApi.sendDesktopMessage(threadId, UUID.randomUUID().toString(), text) + state.updateState { current -> + current.reduceConversationEvent( + ConversationEvent.TurnStarted(result.requiredObject("turn").requiredString("id")), + ) + } + // History is read from the original thread; the companion owns no writer. + // Do not retry a failed delivery or clear drafts on an unconfirmed result. + } + suspend fun queueInput(serverId: String, threadId: String, input: List) { require(input.isNotEmpty()) { "消息内容不能为空" } val conversation = requireConversation(serverId, threadId).value - check(!conversation.isViewOnly) { "此会话当前仅可查看" } + conversation.requireWritable() try { - api(serverId).queueInput(threadId, UUID.randomUUID().toString(), input) + writableApi(serverId, threadId).queueInput(threadId, UUID.randomUUID().toString(), input) } finally { requestRuntimeQueueRefresh(ConversationKey(serverId, threadId)) } @@ -624,9 +663,9 @@ class CodexRepository( suspend fun deleteQueuedInput(serverId: String, threadId: String, entryId: String): Boolean { val state = requireConversation(serverId, threadId) - check(!state.value.isViewOnly) { "此会话当前仅可查看" } + state.value.requireWritable() try { - val deleted = api(serverId).deleteQueuedInput(threadId, entryId) + val deleted = writableApi(serverId, threadId).deleteQueuedInput(threadId, entryId) if (deleted) state.updateState { current -> current.copy(runtimeStatus = current.runtimeStatus.copy( queue = current.runtimeStatus.queue?.let { queue -> @@ -647,11 +686,13 @@ class CodexRepository( val normalized = command.trim() require(normalized.isNotEmpty()) { "Shell 命令不能为空" } val state = requireConversation(serverId, threadId) + state.value.requireWritable() + val targetApi = writableApi(serverId, threadId) require(!state.value.isRunning) { "当前任务结束后才能执行 Shell 命令" } state.updateState { current -> current.copy(isRunning = true, errorMessage = null) } setThreadAttention(serverId, threadId, ThreadAttention.RUNNING) try { - api(serverId).runShellCommand(threadId, normalized) + targetApi.runShellCommand(threadId, normalized) } catch (error: Throwable) { state.updateState { current -> if (current.activeTurnId == null) { @@ -669,8 +710,9 @@ class CodexRepository( } suspend fun terminateBackgroundTerminal(serverId: String, threadId: String, processId: String) { + requireConversation(serverId, threadId).value.requireWritable() require(processId.isNotBlank()) { "后台命令进程 ID 不能为空" } - check(api(serverId).terminateBackgroundTerminal(threadId, processId)) { + check(writableApi(serverId, threadId).terminateBackgroundTerminal(threadId, processId)) { "命令已经结束或无法终止" } val key = ConversationKey(serverId, threadId) @@ -680,8 +722,10 @@ class CodexRepository( suspend fun interrupt(serverId: String, threadId: String) { val state = requireConversation(serverId, threadId) + state.value.requireWritable() + val targetApi = writableApi(serverId, threadId) val turnId = state.value.activeTurnId ?: return - api(serverId).interruptTurn(threadId, turnId) + targetApi.interruptTurn(threadId, turnId) } suspend fun forkConversation(serverId: String, threadId: String, lastTurnId: String? = null): String { @@ -693,9 +737,10 @@ class CodexRepository( } suspend fun renameConversation(serverId: String, threadId: String, name: String) { + checkConversationWritableIfOpen(serverId, threadId) val normalized = name.trim() require(normalized.isNotEmpty()) { "会话名称不能为空" } - api(serverId).request( + writableApi(serverId, threadId).request( StableRpcMethod.THREAD_NAME_SET, buildJsonObject { put("threadId", threadId) @@ -707,7 +752,8 @@ class CodexRepository( } suspend fun archiveConversation(serverId: String, threadId: String) { - api(serverId).request( + checkConversationWritableIfOpen(serverId, threadId) + writableApi(serverId, threadId).request( StableRpcMethod.THREAD_ARCHIVE, buildJsonObject { put("threadId", threadId) }, ) @@ -721,7 +767,8 @@ class CodexRepository( } suspend fun unarchiveConversation(serverId: String, threadId: String) { - api(serverId).request( + checkConversationWritableIfOpen(serverId, threadId) + writableApi(serverId, threadId).request( StableRpcMethod.THREAD_UNARCHIVE, buildJsonObject { put("threadId", threadId) }, ) @@ -729,7 +776,8 @@ class CodexRepository( } suspend fun deleteConversation(serverId: String, threadId: String) { - api(serverId).request( + checkConversationWritableIfOpen(serverId, threadId) + writableApi(serverId, threadId).request( StableRpcMethod.THREAD_DELETE, buildJsonObject { put("threadId", threadId) }, ) @@ -744,7 +792,8 @@ class CodexRepository( } suspend fun compactConversation(serverId: String, threadId: String) { - api(serverId).request( + requireConversation(serverId, threadId).value.requireWritable() + writableApi(serverId, threadId).request( StableRpcMethod.THREAD_COMPACT, buildJsonObject { put("threadId", threadId) }, ) @@ -765,7 +814,8 @@ class CodexRepository( target: ReviewTarget, delivery: ReviewDelivery, ): String { - val response = api(serverId).review(threadId, target, delivery) + requireConversation(serverId, threadId).value.requireWritable() + val response = writableApi(serverId, threadId).review(threadId, target, delivery) return response.requiredString("reviewThreadId") } @@ -791,7 +841,8 @@ class CodexRepository( tokenBudget: Long? = null, ): ThreadGoal { require(objective.isNotBlank()) { "目标不能为空" } - val targetApi = api(serverId) + requireConversation(serverId, threadId).value.requireWritable() + val targetApi = writableApi(serverId, threadId) val key = ConversationKey(serverId, threadId) val revisions = goalRevisions.computeIfAbsent(key) { AtomicLong() } val revision = revisions.incrementAndGet() @@ -810,7 +861,8 @@ class CodexRepository( } suspend fun clearGoal(serverId: String, threadId: String) { - val targetApi = api(serverId) + requireConversation(serverId, threadId).value.requireWritable() + val targetApi = writableApi(serverId, threadId) val key = ConversationKey(serverId, threadId) val revisions = goalRevisions.computeIfAbsent(key) { AtomicLong() } val revision = revisions.incrementAndGet() @@ -853,7 +905,7 @@ class CodexRepository( if (conversationLeases.isHeld(key)) return val conversation = conversations[key] ?: return if (conversation.value.requiresRuntimeRetention) return - runSuspendCatching { + if (!conversation.value.isExternalWriter) runSuspendCatching { api(key.serverId).request( StableRpcMethod.THREAD_UNSUBSCRIBE, buildJsonObject { put("threadId", key.threadId) }, @@ -1044,6 +1096,14 @@ class CodexRepository( retryFailed: Boolean = false, ): AppServerApi = AppServerApi(connections.connection(serverId, retryFailed).connection) + private suspend fun writableApi(serverId: String, threadId: String): AppServerApi { + checkConversationWritableIfOpen(serverId, threadId) + val targetApi = api(serverId) + // Connecting can restore this conversation in external-writer mode. + checkConversationWritableIfOpen(serverId, threadId) + return targetApi + } + private suspend fun establishConnection(serverId: String): ManagedConnection { val server = awaitServer(serverId) val connection = clientFactory.create( @@ -1131,11 +1191,14 @@ class CodexRepository( current.withRefreshedConversation(snapshot.toConversationState(server)) } removeRestoredItemUpdates(key, state.value) + clearExternalWriterRuntime(key, state.value) completeBackgroundTerminalRefresh( key = key, state = state, buffer = backgroundRefresh, - result = runSuspendCatching { loadAllBackgroundTerminals(api, key.threadId) } + result = (if (snapshot.accessMode == ConversationAccessMode.EXTERNAL_WRITER) { + Result.success(emptyList()) + } else runSuspendCatching { loadAllBackgroundTerminals(api, key.threadId) }) .onFailure { error -> if (error !is SocketTimeoutException) { diagnostics.warning("restore_background_terminals_failed", error) @@ -1475,6 +1538,16 @@ class CodexRepository( private fun requireConversation(serverId: String, threadId: String): MutableStateFlow = conversations[ConversationKey(serverId, threadId)] ?: error("会话尚未打开") + private fun checkConversationWritableIfOpen(serverId: String, threadId: String) { + conversations[ConversationKey(serverId, threadId)]?.value?.requireWritable() + } + + private fun clearExternalWriterRuntime(key: ConversationKey, state: ConversationState) { + if (!state.isExternalWriter) return + runtimeQueueRefreshes[key]?.revision?.incrementAndGet() + removePendingRequests(key.serverId, key.threadId) + } + suspend fun loadConversationTurnDetails(serverId: String, threadId: String, turnId: String) { requireConversation(serverId, threadId).loadTurnDetails(turnId) { loadAllTurnItems(api(serverId), threadId, turnId) @@ -1484,15 +1557,16 @@ class CodexRepository( private suspend fun resumeConversation( key: ConversationKey, api: AppServerApi, - ): ThreadSnapshot { + ): ConversationResumeResult { pendingResumeGoalSnapshots.add(key) return try { - val snapshot = ProtocolMapper.resumedThreadSnapshot(api.resumeThread(key.threadId)) - snapshot.copy( - turns = loadResumedTurns(snapshot.turns) { turnId -> - loadAllTurnItems(api, key.threadId, turnId) - }, - ) + resumeOrReadConversation(api, key.threadId) { turnId -> + loadAllTurnItems(api, key.threadId, turnId) + }.also { result -> + if (result.accessMode == ConversationAccessMode.EXTERNAL_WRITER) { + pendingResumeGoalSnapshots.remove(key) + } + } } catch (error: Throwable) { pendingResumeGoalSnapshots.remove(key) throw error @@ -1503,20 +1577,22 @@ class CodexRepository( key: ConversationKey, state: MutableStateFlow, failOnError: Boolean = false, + reportErrors: Boolean = true, ) { val buffer = beginConversationEventBuffer(key) val backgroundRefresh = beginBackgroundTerminalRefresh(key) try { val targetApi = api(key.serverId) val snapshot = resumeConversation(key, targetApi) - val backgroundTerminals = runSuspendCatching { - loadAllBackgroundTerminals(targetApi, key.threadId) - } + val backgroundTerminals = if (snapshot.accessMode == ConversationAccessMode.EXTERNAL_WRITER) { + Result.success(emptyList()) + } else runSuspendCatching { loadAllBackgroundTerminals(targetApi, key.threadId) } val server = awaitServer(key.serverId) completeConversationRefresh(key, state, buffer) { current -> current.withRefreshedConversation(snapshot.toConversationState(server)) } removeRestoredItemUpdates(key, state.value) + clearExternalWriterRuntime(key, state.value) completeBackgroundTerminalRefresh( key = key, state = state, @@ -1530,7 +1606,7 @@ class CodexRepository( } catch (error: Throwable) { cancelBackgroundTerminalRefresh(key, backgroundRefresh) completeConversationRefresh(key, state, buffer) { current -> - current.copy(errorMessage = error.toUserFacingError("刷新会话失败")) + if (reportErrors) current.copy(errorMessage = error.toUserFacingError("刷新会话失败")) else current } if (failOnError) throw error } @@ -1558,7 +1634,7 @@ class CodexRepository( } private fun requestRuntimeQueueRefresh(key: ConversationKey) { - if (!conversations.containsKey(key)) return + if (conversations[key]?.value?.isExternalWriter != false) return val refresh = runtimeQueueRefreshes.computeIfAbsent(key) { RuntimeQueueRefresh() } val revision = refresh.revision.incrementAndGet() applicationScope.launch { @@ -1584,7 +1660,7 @@ class CodexRepository( // 通知触发了更新读取或连接已断开时,丢弃此请求的旧结果 if (revision != refresh.revision.get()) return@withLock conversations[key]?.updateState { current -> - if (revision != refresh.revision.get()) return@updateState current + if (revision != refresh.revision.get() || current.isExternalWriter) return@updateState current val queue = result.fold( onSuccess = { RuntimeInputQueue(entries = it) }, onFailure = { @@ -1734,7 +1810,19 @@ class CodexRepository( .first { it != null } } ?: error("服务器不存在") - private fun com.hebo.codex.data.remote.ThreadSnapshot.toConversationState( + private fun ConversationResumeResult.toConversationState(profile: ServerProfile): ConversationState = + snapshot.toConversationState(profile).let { state -> + state.copy( + accessMode = accessMode, + desktopControlAvailable = desktopControlAvailable, + // thread/read reports notLoaded for another app-server's thread; + // the paginated history still identifies its in-progress turn. + isRunning = state.isRunning || + (accessMode == ConversationAccessMode.EXTERNAL_WRITER && state.activeTurnId != null), + ) + } + + private fun ThreadSnapshot.toConversationState( profile: ServerProfile, ): ConversationState { return ConversationState( @@ -1825,6 +1913,7 @@ class CodexRepository( key: String, result: JsonObject, ) { + requireConversation(serverId, threadId).value.requireWritable() val route = pendingRequestRoutes[key] ?: error("服务端请求已失效") check(route.threadId == threadId) { "服务端请求不属于当前会话" } val connection = connections.current(serverId)?.connection ?: error("服务器未连接") diff --git a/app/src/main/java/com/hebo/codex/domain/ConversationItems.kt b/app/src/main/java/com/hebo/codex/domain/ConversationItems.kt index 3433b24..0316e21 100644 --- a/app/src/main/java/com/hebo/codex/domain/ConversationItems.kt +++ b/app/src/main/java/com/hebo/codex/domain/ConversationItems.kt @@ -62,7 +62,7 @@ internal fun ConversationState.withRestoredConversation( if ( current != null && current.status != ConversationTurnStatus.IN_PROGRESS && current.status == turn.status && turn.detailsState == TurnDetailsState.Unloaded && - turn.items == current.items.summaryItems() + turn.items == (current.detailsSummary ?: current.items.summaryItems()) ) { current.copy(error = turn.error) } else { @@ -127,10 +127,21 @@ internal fun ConversationState.withRefreshedConversation( agentNickname = merged.agentNickname, agentRole = merged.agentRole, canAcceptDirectInput = merged.canAcceptDirectInput, + accessMode = merged.accessMode, + desktopControlAvailable = merged.desktopControlAvailable, history = merged.history, activeItemKeys = merged.activeItemKeys, activeTurnId = merged.activeTurnId, isRunning = merged.isRunning, + model = merged.model ?: model, + serviceTier = merged.serviceTier ?: serviceTier, + reasoningEffort = merged.reasoningEffort ?: reasoningEffort, + approvalPolicy = merged.approvalPolicy ?: approvalPolicy, + approvalsReviewer = merged.approvalsReviewer ?: approvalsReviewer, + sandbox = merged.sandbox ?: sandbox, + errorMessage = null, + pendingRequests = if (merged.isExternalWriter) emptyList() else pendingRequests, + runtimeStatus = if (merged.isExternalWriter) RuntimeStatus() else runtimeStatus, plan = plan?.takeIf { it.turnId == merged.activeTurnId }, ) } diff --git a/app/src/main/java/com/hebo/codex/domain/ConversationResume.kt b/app/src/main/java/com/hebo/codex/domain/ConversationResume.kt new file mode 100644 index 0000000..3bdd56f --- /dev/null +++ b/app/src/main/java/com/hebo/codex/domain/ConversationResume.kt @@ -0,0 +1,52 @@ +package com.hebo.codex.domain + +import com.hebo.codex.data.remote.AppServerApi +import com.hebo.codex.data.remote.ProtocolMapper +import com.hebo.codex.data.remote.RpcException +import com.hebo.codex.data.remote.ThreadSnapshot + +internal data class ConversationResumeResult( + val snapshot: ThreadSnapshot, + val accessMode: ConversationAccessMode, + val desktopControlAvailable: Boolean = false, +) + +internal suspend fun resumeOrReadConversation( + api: AppServerApi, + threadId: String, + loadItems: suspend (String) -> List, +): ConversationResumeResult { + // Only a rejected resume may select the external-writer path. Read, paging and + // mapping failures must still surface as errors rather than changing ownership. + val resumed = try { + api.resumeThread(threadId) + } catch (error: RpcException) { + if (error.code != -32600 || !error.message.contains("already has an active writer")) throw error + null + } + val accessMode: ConversationAccessMode + val snapshot: ThreadSnapshot + if (resumed != null) { + accessMode = ConversationAccessMode.WRITABLE + snapshot = ProtocolMapper.resumedThreadSnapshot(resumed) + } else { + accessMode = ConversationAccessMode.EXTERNAL_WRITER + val metadata = ProtocolMapper.threadMetadata(api.readThread(threadId)) + val page = ProtocolMapper.threadTurnPage(api.listThreadTurns(threadId)) + snapshot = metadata.copy(turns = page.data, olderTurnsCursor = page.nextCursor) + } + return ConversationResumeResult( + snapshot = snapshot.copy(turns = loadResumedTurns( + snapshot.turns, + loadNewestDetails = accessMode == ConversationAccessMode.EXTERNAL_WRITER, + loadItems = loadItems, + )), + accessMode = accessMode, + desktopControlAvailable = accessMode == ConversationAccessMode.EXTERNAL_WRITER && + snapshot.summary.canAcceptDirectInput != false && api.discoverDesktopControl(threadId), + ) +} + +internal fun ConversationState.requireWritable() { + check(!isViewOnly) { viewOnlyMessage ?: "此会话当前仅可查看" } +} diff --git a/app/src/main/java/com/hebo/codex/domain/ConversationRetention.kt b/app/src/main/java/com/hebo/codex/domain/ConversationRetention.kt index 299b3c1..a54da71 100644 --- a/app/src/main/java/com/hebo/codex/domain/ConversationRetention.kt +++ b/app/src/main/java/com/hebo/codex/domain/ConversationRetention.kt @@ -17,10 +17,10 @@ internal class ConversationLeaseRegistry { } internal val ConversationState.requiresRuntimeRetention: Boolean - get() = isRunning || + get() = !isExternalWriter && (isRunning || activeTurnId != null || pendingRequests.isNotEmpty() || - backgroundTerminals.isNotEmpty() + backgroundTerminals.isNotEmpty()) internal fun ConversationState.withActiveTurnPlan(plan: ActiveTurnPlan): ConversationState = if (plan.turnId == activeTurnId) copy(plan = plan) else this diff --git a/app/src/main/java/com/hebo/codex/domain/ConversationTurnDetails.kt b/app/src/main/java/com/hebo/codex/domain/ConversationTurnDetails.kt index 140a192..dbe442e 100644 --- a/app/src/main/java/com/hebo/codex/domain/ConversationTurnDetails.kt +++ b/app/src/main/java/com/hebo/codex/domain/ConversationTurnDetails.kt @@ -7,12 +7,18 @@ import kotlinx.coroutines.flow.update internal suspend fun loadResumedTurns( turnsNewestFirst: List, + loadNewestDetails: Boolean = false, loadItems: suspend (String) -> List, ): List { + val newestTurnId = turnsNewestFirst.firstOrNull()?.id return turnsNewestFirst.map { turn -> - // 运行中的回合恢复增量通知所需的基础条目;已结束的回合详情按需加载 - if (turn.status == ConversationTurnStatus.IN_PROGRESS) { - turn.copy(items = loadItems(turn.id), detailsState = TurnDetailsState.Loaded) + // 外部 writer 的磁盘快照可能把尚在运行的回合标为 interrupted。 + // 同步最新回合的完整条目,直接更新详情,不先重置成加载状态。 + if (turn.status == ConversationTurnStatus.IN_PROGRESS || (loadNewestDetails && turn.id == newestTurnId)) { + turn.copy( + items = loadItems(turn.id), detailsState = TurnDetailsState.Loaded, + detailsSummary = turn.items, + ) } else { turn } @@ -29,7 +35,9 @@ internal suspend fun MutableStateFlow.loadTurnDetails( val current = value val turn = current.history.turns.singleOrNull { it.id == turnId } ?: return if (turn.detailsState == TurnDetailsState.Loaded || turn.detailsState is TurnDetailsState.Loading) return - val loading = current.updateTurn(turnId) { it.copy(detailsState = loadingState) } + val loading = current.updateTurn(turnId) { + it.copy(detailsState = loadingState, detailsSummary = it.detailsSummary ?: it.items) + } if (compareAndSet(current, loading)) { initialTurn = turn break diff --git a/app/src/main/java/com/hebo/codex/domain/ExternalConversationSync.kt b/app/src/main/java/com/hebo/codex/domain/ExternalConversationSync.kt new file mode 100644 index 0000000..dccfde8 --- /dev/null +++ b/app/src/main/java/com/hebo/codex/domain/ExternalConversationSync.kt @@ -0,0 +1,35 @@ +package com.hebo.codex.domain + +import com.hebo.codex.util.runSuspendCatching +import kotlinx.coroutines.currentCoroutineContext +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.Flow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.combine +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.isActive + +/** Runs only while the chat screen is resumed. Snapshots converge on the same + * thread; a successful resume switches back to its normal notification stream. + * This never forks, sends input, or forces an existing writer to detach. */ +internal suspend fun monitorExternalConversation( + conversation: StateFlow, + connected: Flow, + refresh: suspend () -> Unit, + onError: (Throwable?) -> Unit, +) { + val targets = combine(conversation, connected) { state, online -> state to online } + var retryDelayMillis = 0L + while (currentCoroutineContext().isActive) { + val (state) = targets.first { (state, online) -> state.isExternalWriter && online } + delay(retryDelayMillis.takeIf { it > 0 } ?: if (state.isRunning) 2_000L else 5_000L) + val (latest, online) = targets.first() + if (!latest.isExternalWriter || !online) { + retryDelayMillis = 0 + continue + } + val error = runSuspendCatching { refresh() }.exceptionOrNull() + onError(error) + retryDelayMillis = if (error == null) 0 else (retryDelayMillis * 2).coerceIn(5_000L, 30_000L) + } +} diff --git a/app/src/main/java/com/hebo/codex/domain/Models.kt b/app/src/main/java/com/hebo/codex/domain/Models.kt index 397dfeb..508fda4 100644 --- a/app/src/main/java/com/hebo/codex/domain/Models.kt +++ b/app/src/main/java/com/hebo/codex/domain/Models.kt @@ -219,6 +219,11 @@ data class RecentProject( val threads: List, ) +enum class ConversationAccessMode { + WRITABLE, + EXTERNAL_WRITER, +} + data class ConversationState( val serverId: String, val serverName: String, @@ -250,12 +255,31 @@ data class ConversationState( val agentNickname: String? = null, val agentRole: String? = null, val canAcceptDirectInput: Boolean? = null, + val accessMode: ConversationAccessMode = ConversationAccessMode.WRITABLE, + val desktopControlAvailable: Boolean = false, ) { val items: List get() = history.items val isViewOnly: Boolean - get() = canAcceptDirectInput == false + get() = isExternalWriter || canAcceptDirectInput == false + + val isExternalWriter: Boolean + get() = accessMode == ConversationAccessMode.EXTERNAL_WRITER + + // Sending through the owner does not grant local Shell, approval or settings access. + val canSendToDesktop: Boolean + get() = isExternalWriter && desktopControlAvailable && canAcceptDirectInput != false + + val viewOnlyMessage: String? + get() = when { + isExternalWriter && canAcceptDirectInput == false -> + "此子代理正在其他 Codex 客户端中使用,自动更新历史;输入仍由父会话控制" + canSendToDesktop -> "已连接此会话的电脑客户端,可在原对话中发送文字。Shell、审批和设置请在电脑端处理" + isExternalWriter -> "此会话正在其他 Codex 客户端中使用,当前只读并自动更新历史;释放占用后可在此会话继续发送" + canAcceptDirectInput == false -> "此会话当前仅可查看" + else -> null + } } data class BackgroundTerminal( @@ -281,6 +305,9 @@ data class ConversationTurn( val items: List = emptyList(), val detailsState: TurnDetailsState = TurnDetailsState.Loaded, val error: UserFacingError? = null, + // The server summary used to load these details. Full items may contain + // commentary/async messages that the server intentionally omits from it. + val detailsSummary: List? = null, ) sealed interface TurnDetailsState { diff --git a/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatScreen.kt b/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatScreen.kt index 6aae94c..af3d598 100644 --- a/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatScreen.kt +++ b/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatScreen.kt @@ -329,6 +329,7 @@ fun ChatScreen( val scrollScope = rememberCoroutineScope() var bottomControlsHeightPx by remember { mutableStateOf(0) } var dialog by remember { mutableStateOf(null) } + var forkFromTurnId by remember { mutableStateOf(null) } var showTurnOptions by remember { mutableStateOf(false) } var activeRequestKey by rememberSaveable(conversation?.threadId) { mutableStateOf(null) } var shownNoticeId by rememberSaveable(conversation?.threadId) { mutableStateOf(null) } @@ -344,7 +345,8 @@ fun ChatScreen( focusManager.clearFocus() imagePreviewSource = source } - val activeRequest = conversation?.pendingRequests?.firstOrNull { it.requestKey == activeRequestKey } + val activeRequest = conversation?.takeUnless { it.isViewOnly } + ?.pendingRequests?.firstOrNull { it.requestKey == activeRequestKey } val density = LocalDensity.current val showLinkError: (String) -> Unit = { message -> noticeState.showError(UserFacingError(message)) @@ -502,7 +504,10 @@ fun ChatScreen( actions = ConversationHeaderActions( onNewTask = { conversation?.let { onNewTask(it.cwd) } }, onRename = { dialog = ChatDialog.Rename }, - onFork = { dialog = ChatDialog.Fork }, + onFork = { + forkFromTurnId = null + dialog = ChatDialog.Fork + }, onOpenGoal = { dialog = ChatDialog.Goal }, ), onRefreshRateLimits = onRefreshRateLimits, @@ -581,7 +586,10 @@ fun ChatScreen( ForkFromTurnAction( turn = turn, forkRequest = state.forkRequest, - onFork = viewModel::fork, + onFork = { + forkFromTurnId = it + dialog = ChatDialog.Fork + }, ) } } @@ -631,7 +639,7 @@ fun ChatScreen( .fillMaxWidth() .onSizeChanged { bottomControlsHeightPx = it.height }, ) { - conversation.pendingRequests.firstOrNull()?.let { request -> + conversation.pendingRequests.firstOrNull()?.takeUnless { conversation.isViewOnly }?.let { request -> PendingRequestBanner( request = request, requestCount = conversation.pendingRequests.size, @@ -644,7 +652,7 @@ fun ChatScreen( verticalAlignment = Alignment.CenterVertically, ) { Text( - if (conversation.isViewOnly) "子代理由父会话控制" else "子代理会话", + if (conversation.canAcceptDirectInput == false) "子代理由父会话控制" else "子代理会话", modifier = Modifier.weight(1f), style = MaterialTheme.typography.bodySmall, color = MaterialTheme.colorScheme.onSurfaceVariant, @@ -654,13 +662,72 @@ fun ChatScreen( } } } - if (conversation.isViewOnly && conversation.parentThreadId == null) { - Text( - "此会话当前仅可查看", - modifier = Modifier.padding(16.dp), - style = MaterialTheme.typography.bodySmall, - color = MaterialTheme.colorScheme.onSurfaceVariant, - ) + if (conversation.isExternalWriter || + (conversation.isViewOnly && conversation.parentThreadId == null) + ) { + Row( + modifier = Modifier.fillMaxWidth().padding(horizontal = 16.dp, vertical = 8.dp), + verticalAlignment = Alignment.CenterVertically, + ) { + Text( + requireNotNull(conversation.viewOnlyMessage), + modifier = Modifier.weight(1f), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + if (conversation.isExternalWriter) { + TextButton( + onClick = viewModel::refreshConversation, + enabled = state.connectionStatus is ConnectionStatus.Connected, + ) { Text(if (conversation.canSendToDesktop) "刷新" else "尝试继续") } + } + } + } + if (conversation.isExternalWriter) { + state.externalSyncError?.let { error -> + Text( + error.summary, + modifier = Modifier.fillMaxWidth().padding(horizontal = 16.dp), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.error, + ) + } + if (conversation.canAcceptDirectInput != false) { + OutlinedTextField( + value = state.draft.message, + onValueChange = viewModel::updateDraftMessage, + enabled = state.queueOperationId == null && !state.isSendingDraft, + modifier = Modifier.fillMaxWidth().padding(horizontal = 16.dp, vertical = 8.dp), + label = { Text(if (conversation.canSendToDesktop) "在原对话中继续" else "此会话的草稿") }, + supportingText = { Text( + if (conversation.canSendToDesktop) { + if (conversation.isRunning) "电脑正在处理,完成后可发送" else "发送文字给电脑端,沿用原会话设置" + } else "仅保存草稿;恢复发送后需手动发送,不会创建新对话", + ) }, + minLines = 1, + maxLines = 3, + ) + if (conversation.canSendToDesktop) { + if (state.draft.attachments.isNotEmpty()) { + TextButton( + onClick = { state.draft.attachments.forEach { viewModel.removeDraftAttachment(it.id) } }, + enabled = !state.isSendingDraft && state.queueOperationId == null, + modifier = Modifier.padding(horizontal = 16.dp), + ) { Text("移除草稿附件(${state.draft.attachments.size}),仅发送文字") } + } + TextButton( + onClick = { + pagingController.onSubmission() + viewModel.sendDesktopDraft() + }, + enabled = !conversation.isRunning && !state.isSendingDraft && + state.draft.message.isNotBlank() && + state.draft.attachments.isEmpty() && + state.connectionStatus is ConnectionStatus.Connected, + modifier = Modifier.padding(horizontal = 16.dp), + ) { Text(if (state.isSendingDraft) "正在发送" else "发送到原对话") } + } + } } ConversationAccessoryPanels( conversation = if (conversation.goalLoaded) conversation else conversation.copy(goal = state.goal), @@ -780,6 +847,7 @@ fun ChatScreen( ChatDialog.Fork -> ForkConversationSheet( turns = forkTurns, + initialLastTurnId = forkFromTurnId, isRunning = conversation?.isRunning == true, forkRequest = state.forkRequest, onDismiss = { dialog = null }, @@ -833,8 +901,8 @@ fun ChatScreen( } else { ActivityDetailsActions( isRunning = conversation.isRunning && !conversation.isViewOnly, - pendingRequest = conversation.pendingRequests.firstOrNull(), - pendingRequestCount = conversation.pendingRequests.size, + pendingRequest = conversation.pendingRequests.firstOrNull()?.takeUnless { conversation.isViewOnly }, + pendingRequestCount = if (conversation.isViewOnly) 0 else conversation.pendingRequests.size, onStop = viewModel::interrupt, onOpenPendingRequest = { activeRequestKey = it }, ) @@ -1130,7 +1198,7 @@ internal fun ConversationFloatingHeader( }, ) } - ChatActionMenuItem( + if (conversation?.isViewOnly != true) ChatActionMenuItem( text = "重命名", icon = { Icon(Icons.Rounded.Edit, contentDescription = null) }, onClick = { @@ -1278,7 +1346,7 @@ private fun ForkFromTurnAction( ) } Spacer(Modifier.width(6.dp)) - Text(if (isCreating) "正在创建…" else "从这里继续") + Text(if (isCreating) "正在创建…" else "分叉为新对话") } } } @@ -3034,14 +3102,15 @@ private fun RenameConversationDialog( @Composable private fun ForkConversationSheet( turns: List, + initialLastTurnId: String?, isRunning: Boolean, forkRequest: ForkRequest?, onDismiss: () -> Unit, onFork: (String?) -> Unit, ) { val canForkLatest = !isRunning - var lastTurnId by rememberSaveable(turns, canForkLatest) { - mutableStateOf(if (canForkLatest) null else turns.lastOrNull()?.turnId) + var lastTurnId by rememberSaveable(initialLastTurnId) { + mutableStateOf(initialLastTurnId ?: if (canForkLatest) null else turns.lastOrNull()?.turnId) } ModalBottomSheet( onDismissRequest = onDismiss, @@ -3057,13 +3126,13 @@ private fun ForkConversationSheet( ) { Text("分叉会话", style = MaterialTheme.typography.titleLarge) Text( - "选择新会话包含的最后一个回合", + "这会创建独立的新对话,后续消息不会追加到原会话。选择新对话包含的最后一个回合。", style = MaterialTheme.typography.bodyMedium, color = MaterialTheme.colorScheme.onSurfaceVariant, ) if (canForkLatest) { ForkSelectionCard( - title = "从最新状态继续", + title = "复制最新状态", supportingText = "包含当前会话的全部已完成回合", selected = lastTurnId == null, onClick = { lastTurnId = null }, @@ -3097,10 +3166,11 @@ private fun ForkConversationSheet( } Button( onClick = { onFork(lastTurnId) }, - enabled = forkRequest == null && (lastTurnId != null || canForkLatest), + enabled = forkRequest == null && + (lastTurnId?.let { id -> turns.any { it.turnId == id } } ?: canForkLatest), modifier = Modifier.fillMaxWidth(), ) { - Text("创建并打开") + Text("创建并打开新对话") } TextButton( onClick = onDismiss, diff --git a/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatViewModel.kt b/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatViewModel.kt index 65f1434..391e848 100644 --- a/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatViewModel.kt +++ b/app/src/main/java/com/hebo/codex/ui/screens/chat/ChatViewModel.kt @@ -31,6 +31,7 @@ import com.hebo.codex.domain.UserFacingError import com.hebo.codex.domain.newDraftRevision import com.hebo.codex.domain.quickPhraseInput import com.hebo.codex.domain.loadInputResources +import com.hebo.codex.domain.monitorExternalConversation import com.hebo.codex.util.runSuspendCatching import com.hebo.codex.util.toUserFacingError import kotlinx.coroutines.CancellationException @@ -42,6 +43,8 @@ import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.collectLatest +import kotlinx.coroutines.flow.distinctUntilChanged +import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.receiveAsFlow import kotlinx.coroutines.flow.update import kotlinx.coroutines.launch @@ -85,6 +88,7 @@ data class ChatUiState( val queueFeedback: String? = null, val terminatingBackgroundProcessIds: Set = emptySet(), val forkRequest: ForkRequest? = null, + val externalSyncError: UserFacingError? = null, val screenError: UserFacingError? = null, ) @@ -115,6 +119,7 @@ class ChatViewModel( val events = mutableEvents.receiveAsFlow() private var draftInitialized = false private var conversationJob: Job? = null + private var conversationRefreshJob: Job? = null private var conversationReady: CompletableDeferred? = null private var subAgentsJob: Job? = null private var modelListJob: Job? = null @@ -168,10 +173,25 @@ class ChatViewModel( runSuspendCatching { queueWithdrawalCoordinator.loadDraft(draftKey) } } var conversationOpened = false + var externalSyncJob: Job? = null try { val conversation = repository.openConversation(serverId, threadId) conversationOpened = true ready.complete(Unit) + externalSyncJob = launch { + monitorExternalConversation( + conversation = conversation, + connected = repository.connectionStatuses + .map { it[serverId] is ConnectionStatus.Connected } + .distinctUntilChanged(), + refresh = { repository.syncExternalConversation(serverId, threadId) }, + onError = { error -> + mutableState.update { + it.copy(externalSyncError = error?.toUserFacingError("暂时无法更新历史,稍后自动重试")) + } + }, + ) + } conversation.collectLatest { value -> val initialDraft = if (draftInitialized) { null @@ -186,6 +206,7 @@ class ChatViewModel( goal = if (value.goalLoaded) value.goal else it.goal, draft = initialDraft ?: it.draft, isLoading = false, + externalSyncError = it.externalSyncError.takeIf { value.isExternalWriter }, screenError = restoredResult.exceptionOrNull() ?.toUserFacingError("无法恢复草稿") ?: it.screenError, @@ -202,6 +223,7 @@ class ChatViewModel( ) } } finally { + externalSyncJob?.cancel() ready.cancel() if (conversationOpened) repository.releaseConversation(serverId, threadId) } @@ -209,6 +231,8 @@ class ChatViewModel( } fun pauseConversation() { + conversationRefreshJob?.cancel() + conversationRefreshJob = null conversationReady?.cancel() conversationReady = null conversationJob?.cancel() @@ -222,6 +246,15 @@ class ChatViewModel( fun refreshSubAgents() = repository.refreshSubAgents(serverId, threadId) + fun refreshConversation() { + if (conversationRefreshJob?.isActive == true) return + conversationRefreshJob = viewModelScope.launch { + runSuspendCatching { repository.refreshConversation(serverId, threadId) } + .onSuccess { mutableState.update { it.copy(externalSyncError = null) } } + .onFailure(::showError) + } + } + fun refreshRuntimeQueue() = repository.refreshRuntimeQueue(serverId, threadId) fun loadOlderTurns() = requestOlderTurns(retry = false) @@ -333,6 +366,22 @@ class ChatViewModel( } } + fun sendDesktopDraft() { + if (queueBusy) return + val current = mutableState.value + if (current.conversation?.canSendToDesktop != true || current.conversation.isRunning || + current.isSendingDraft || current.connectionStatus !is ConnectionStatus.Connected) return + val draft = current.draft + if (draft.message.isBlank()) return + if (draft.attachments.isNotEmpty()) { + reportInputError(IllegalArgumentException("电脑端续聊目前仅支持文字,请先移除草稿附件")) + return + } + submitDraft(draft, failureFeedback = "发送未确认,请检查原对话后再重试") { + repository.sendDesktopMessage(serverId, threadId, draft.message) + } + } + private fun sendShellDraft(current: ChatUiState, draft: ChatDraft) { val blockReason = shellCommandBlockReason( message = draft.message, @@ -452,6 +501,7 @@ class ChatViewModel( } fun interrupt() { + if (mutableState.value.conversation?.isViewOnly == true) return viewModelScope.launch { runSuspendCatching { repository.interrupt(serverId, threadId) } .onFailure(::showError) @@ -460,6 +510,7 @@ class ChatViewModel( fun terminateBackgroundCommand(processId: String) { val current = mutableState.value + if (current.conversation?.isViewOnly == true) return if (processId in current.terminatingBackgroundProcessIds) return mutableState.update { it.copy( @@ -518,6 +569,7 @@ class ChatViewModel( } fun rename(name: String) { + if (mutableState.value.conversation?.isViewOnly == true) return viewModelScope.launch { runSuspendCatching { repository.renameConversation(serverId, threadId, name) diff --git a/app/src/test/java/com/hebo/codex/data/remote/AppServerHistoryApiTest.kt b/app/src/test/java/com/hebo/codex/data/remote/AppServerHistoryApiTest.kt index 4e31964..36b15af 100644 --- a/app/src/test/java/com/hebo/codex/data/remote/AppServerHistoryApiTest.kt +++ b/app/src/test/java/com/hebo/codex/data/remote/AppServerHistoryApiTest.kt @@ -54,11 +54,12 @@ class AppServerHistoryApiTest { val api = AppServerApi(connection) api.resumeThread("thread-1") + api.readThread("thread-1") api.listThreadTurns("thread-1", cursor = "older-turns") api.listThreadItems("thread-1", "turn-1", cursor = "next-items") assertEquals( - listOf("thread/resume", "thread/turns/list", "thread/items/list"), + listOf("thread/resume", "thread/read", "thread/turns/list", "thread/items/list"), requests.map { it.requiredString("method") }, ) assertEquals( @@ -66,13 +67,17 @@ class AppServerHistoryApiTest { requests[0].requiredObject("params"), ) assertEquals( - json("""{"threadId":"thread-1","cursor":"older-turns","limit":25,"sortDirection":"desc","itemsView":"summary"}"""), + json("""{"threadId":"thread-1","includeTurns":false}"""), requests[1].requiredObject("params"), ) assertEquals( - json("""{"threadId":"thread-1","turnId":"turn-1","limit":100,"sortDirection":"asc","cursor":"next-items"}"""), + json("""{"threadId":"thread-1","cursor":"older-turns","limit":25,"sortDirection":"desc","itemsView":"summary"}"""), requests[2].requiredObject("params"), ) + assertEquals( + json("""{"threadId":"thread-1","turnId":"turn-1","limit":100,"sortDirection":"asc","cursor":"next-items"}"""), + requests[3].requiredObject("params"), + ) } finally { connection.close() } diff --git a/app/src/test/java/com/hebo/codex/domain/ConversationTurnDetailsTest.kt b/app/src/test/java/com/hebo/codex/domain/ConversationTurnDetailsTest.kt index 68308e9..2d7b49e 100644 --- a/app/src/test/java/com/hebo/codex/domain/ConversationTurnDetailsTest.kt +++ b/app/src/test/java/com/hebo/codex/domain/ConversationTurnDetailsTest.kt @@ -33,7 +33,9 @@ class ConversationTurnDetailsTest { assertEquals(listOf("running"), reads) assertEquals( - listOf(older, running.copy(items = runningItems, detailsState = TurnDetailsState.Loaded), newest), + listOf(older, running.copy( + items = runningItems, detailsState = TurnDetailsState.Loaded, detailsSummary = running.items, + ), newest), turns, ) } @@ -171,6 +173,92 @@ class ConversationTurnDetailsTest { assertEquals(loaded.history, restored.history) } + @Test + fun `相同服务端摘要不会反复加载包含中间回复的详情`() = runTest { + val user = ConversationItem.UserMessage("user", "turn", listOf(UserMessagePart.Text("问题"))) + val commentary = answer.copy(phase = AgentMessagePhase.COMMENTARY, text = "正在处理") + val state = state() + val summary = state.value.copy( + accessMode = ConversationAccessMode.EXTERNAL_WRITER, + history = ConversationHistory(listOf(ConversationTurn( + "turn", ConversationTurnStatus.INTERRUPTED, listOf(user), TurnDetailsState.Unloaded, + ))), + ) + state.value = summary + var reads = 0 + val details = listOf(user, commentary, command) + repeat(3) { + state.loadTurnDetails("turn") { reads++; details } + state.value = state.value.withRefreshedConversation(summary) + assertEquals(TurnDetailsState.Loaded, state.turn.detailsState) + assertEquals(details, state.turn.items) + } + assertEquals(1, reads) + } + + @Test + fun `空摘要加载的异步回复不会让相同摘要使缓存失效`() = runTest { + val state = state() + state.value = state.value.copy(history = ConversationHistory(listOf( + state.turn.copy(items = emptyList()), + ))) + val summary = state.value + val question = answer.copy(delivery = AgentMessageDelivery.ASYNC) + state.loadTurnDetails("turn") { listOf(question, command) } + val loaded = state.value + repeat(3) { state.value = state.value.withRestoredConversation(summary) } + assertEquals(loaded.history, state.value.history) + } + + @Test + fun `真实摘要变化仍丢弃包含中间回复的缓存`() = runTest { + val state = state() + val summary = state.value + state.loadTurnDetails("turn") { listOf(answer, command, answer.copy(id = "commentary", phase = AgentMessagePhase.COMMENTARY)) } + val changed = summary.withConversationItem(answer.copy(text = "新的最终结果"), isActive = false) + state.value = state.value.withRestoredConversation(changed) + assertEquals(changed.history, state.value.history) + assertEquals(TurnDetailsState.Unloaded, state.turn.detailsState) + } + + @Test + fun `请求期间新增中间回复不会让相同摘要重置加载或失败状态`() = runTest { + val state = state() + val summary = state.value + val result = CompletableDeferred>() + val job = launch { state.loadTurnDetails("turn") { result.await() } } + runCurrent() + state.value = state.value.withConversationItem( + answer.copy(id = "commentary", phase = AgentMessagePhase.COMMENTARY), isActive = false, + ) + val loading = state.turn.detailsState + state.value = state.value.withRestoredConversation(summary) + assertTrue(state.turn.detailsState === loading) + result.completeExceptionally(IllegalStateException("读取失败")) + job.join() + val failed = state.turn.detailsState + repeat(3) { state.value = state.value.withRestoredConversation(summary) } + assertTrue(failed is TurnDetailsState.Failed) + assertTrue(state.turn.detailsState === failed) + } + + @Test + fun `失败后手动重试仍使用原服务端摘要校验缓存`() = runTest { + val state = state() + val summary = state.value + val result = CompletableDeferred>() + val job = launch { state.loadTurnDetails("turn") { result.await() } } + runCurrent() + val commentary = answer.copy(id = "commentary", phase = AgentMessagePhase.COMMENTARY) + state.value = state.value.withConversationItem(commentary, isActive = false) + result.completeExceptionally(IllegalStateException("读取失败")) + job.join() + state.loadTurnDetails("turn") { listOf(answer, command, commentary) } + state.value = state.value.withRestoredConversation(summary) + assertEquals(TurnDetailsState.Loaded, state.turn.detailsState) + assertEquals(listOf(answer, command, commentary), state.turn.items) + } + @Test fun `摘要内容变化使旧缓存失效`() { val summary = state().value diff --git a/app/src/test/java/com/hebo/codex/domain/ExternalConversationSyncTest.kt b/app/src/test/java/com/hebo/codex/domain/ExternalConversationSyncTest.kt new file mode 100644 index 0000000..b0d91bf --- /dev/null +++ b/app/src/test/java/com/hebo/codex/domain/ExternalConversationSyncTest.kt @@ -0,0 +1,137 @@ +package com.hebo.codex.domain + +import kotlinx.coroutines.awaitCancellation +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.launch +import kotlinx.coroutines.test.advanceTimeBy +import kotlinx.coroutines.test.runCurrent +import kotlinx.coroutines.test.runTest +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Test + +@OptIn(kotlinx.coroutines.ExperimentalCoroutinesApi::class) +class ExternalConversationSyncTest { + @Test + fun externalChatUsesFastSnapshotsUntilWriterIsReleasedThenStopsPolling() = runTest { + val conversation = MutableStateFlow(state(running = true)) + val connected = MutableStateFlow(true) + var count = 0 + val job = backgroundScope.launch { + monitorExternalConversation(conversation, connected, { + count++ + if (count == 2) conversation.value = conversation.value.copy(accessMode = ConversationAccessMode.WRITABLE) + }, {}) + } + runCurrent() + advanceTimeBy(1_999); runCurrent() + assertEquals(0, count) + advanceTimeBy(1); runCurrent() + assertEquals(1, count) + advanceTimeBy(2_000); runCurrent() + assertEquals(2, count) + assertEquals("original", conversation.value.threadId) + assertFalse(conversation.value.isViewOnly) + advanceTimeBy(30_000); runCurrent() + assertEquals(2, count) + job.cancel() + } + + @Test + fun idleExternalChatUsesSlowerSnapshotsAndDoesNotPollWhileDisconnected() = runTest { + val conversation = MutableStateFlow(state()) + val connected = MutableStateFlow(false) + var count = 0 + val job = backgroundScope.launch { monitorExternalConversation(conversation, connected, { count++ }, {}) } + runCurrent() + advanceTimeBy(30_000); runCurrent() + assertEquals(0, count) + connected.value = true + runCurrent() + advanceTimeBy(5_000); runCurrent() + assertEquals(1, count) + connected.value = false + advanceTimeBy(30_000); runCurrent() + assertEquals(1, count) + job.cancel() + } + + @Test + fun writableAndSubAgentChatsDoNotPoll() = runTest { + val conversation = MutableStateFlow(state().copy(accessMode = ConversationAccessMode.WRITABLE)) + val connected = MutableStateFlow(true) + var count = 0 + val job = backgroundScope.launch { monitorExternalConversation(conversation, connected, { count++ }, {}) } + runCurrent() + advanceTimeBy(30_000); runCurrent() + conversation.value = conversation.value.copy(canAcceptDirectInput = false) + advanceTimeBy(30_000); runCurrent() + assertEquals(0, count) + job.cancel() + } + + @Test + fun errorsAreReportedAndBackOffWhileSuccessfulSyncClearsError() = runTest { + val conversation = MutableStateFlow(state(running = true)) + val connected = MutableStateFlow(true) + val attempts = mutableListOf() + val errors = mutableListOf() + val error = IllegalStateException("invalid params") + val job = backgroundScope.launch { + monitorExternalConversation(conversation, connected, { + attempts += testScheduler.currentTime + if (attempts.size <= 2) throw error + }, { errors += it }) + } + runCurrent() + for (interval in listOf(2_000L, 5_000L, 10_000L, 2_000L)) { + advanceTimeBy(interval); runCurrent() + } + assertEquals(listOf(2_000L, 7_000L, 17_000L, 19_000L), attempts) + assertEquals(error, errors[0]) + assertEquals(error, errors[1]) + assertNull(errors[2]) + job.cancel() + } + + @Test + fun manualRecoveryDuringDelayPreventsAnExtraResume() = runTest { + val conversation = MutableStateFlow(state()) + var count = 0 + val job = backgroundScope.launch { + monitorExternalConversation(conversation, MutableStateFlow(true), { count++ }, {}) + } + runCurrent() + conversation.value = conversation.value.copy(accessMode = ConversationAccessMode.WRITABLE) + advanceTimeBy(5_000); runCurrent() + assertEquals(0, count) + job.cancel() + } + + @Test + fun leavingScreenCancelsInFlightReadAndDoesNotReportCancellationAsAnError() = runTest { + var started = false + var cancelled = false + val errors = mutableListOf() + val job = backgroundScope.launch { + monitorExternalConversation(MutableStateFlow(state()), MutableStateFlow(true), { + started = true + try { awaitCancellation() } finally { cancelled = true } + }, { errors += it }) + } + runCurrent() + advanceTimeBy(5_000); runCurrent() + assertTrue(started) + job.cancel() + runCurrent() + assertTrue(cancelled) + assertTrue(errors.isEmpty()) + } + + private fun state(running: Boolean = false) = ConversationState( + serverId = "server", serverName = "Server", threadId = "original", title = "Original", cwd = "/workspace", + isRunning = running, accessMode = ConversationAccessMode.EXTERNAL_WRITER, + ) +} diff --git a/app/src/test/java/com/hebo/codex/domain/ExternalWriterConversationTest.kt b/app/src/test/java/com/hebo/codex/domain/ExternalWriterConversationTest.kt new file mode 100644 index 0000000..c4de2a9 --- /dev/null +++ b/app/src/test/java/com/hebo/codex/domain/ExternalWriterConversationTest.kt @@ -0,0 +1,607 @@ +package com.hebo.codex.domain + +import com.hebo.codex.data.local.ServerSettings +import com.hebo.codex.data.local.ServerStore +import com.hebo.codex.data.remote.AppServerClientFactory +import com.hebo.codex.data.remote.AppServerConnection +import com.hebo.codex.data.remote.AppServerWebSocket +import com.hebo.codex.data.remote.AppServerWebSocketFactory +import com.hebo.codex.data.remote.RpcException +import com.hebo.codex.data.remote.requiredObject +import com.hebo.codex.data.remote.requiredString +import kotlinx.coroutines.CancellationException +import kotlinx.coroutines.flow.flowOf +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.launch +import kotlinx.coroutines.test.TestScope +import kotlinx.coroutines.test.runCurrent +import kotlinx.coroutines.test.advanceTimeBy +import kotlinx.coroutines.test.runTest +import kotlinx.serialization.json.Json +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.jsonObject +import kotlinx.serialization.json.put +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue +import org.junit.Assert.fail +import org.junit.Test +import org.mockito.Mockito + +@OptIn(kotlinx.coroutines.ExperimentalCoroutinesApi::class) +class ExternalWriterConversationTest { + @Test + fun normalResumeRemainsWritableAndDoesNotReadThread() = runTest { + withFixture { fixture -> + val state = fixture.repository.openConversation("server", "thread").value + assertEquals(ConversationAccessMode.WRITABLE, state.accessMode) + assertFalse(state.isViewOnly) + assertEquals(listOf("old", "running"), state.history.turns.map { it.id }) + assertEquals("older", state.history.olderTurnsCursor) + assertFalse("thread/read" in fixture.methods) + assertFalse("thread/turns/list" in fixture.methods) + } + } + + @Test + fun conflictReadsMetadataAndPagesHistoryWithoutClaimingWriter() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val flow = fixture.repository.openConversation("server", "thread") + val state = flow.value + assertTrue(state.isExternalWriter) + assertTrue(state.isViewOnly) + assertEquals(true, state.canAcceptDirectInput) + assertEquals(null, state.parentThreadId) + assertEquals("running", state.activeTurnId) + assertEquals(listOf("old", "running"), state.history.turns.map { it.id }) + assertEquals(listOf("answer-1", "answer-2"), state.history.turns.last().items.map { it.id }) + assertEquals(TurnDetailsState.Loaded, state.history.turns.last().detailsState) + assertEquals(TurnDetailsState.Unloaded, state.history.turns.first().detailsState) + assertEquals("older", state.history.olderTurnsCursor) + assertEquals( + listOf("thread/resume", "thread/read", "thread/turns/list", "thread/items/list", "thread/items/list", "codexAndroid/desktop/discover"), + fixture.methods, + ) + assertEquals(json("""{"threadId":"thread","includeTurns":false}"""), + fixture.requests[1].requiredObject("params")) + assertEquals(json("""{"threadId":"thread","limit":25,"sortDirection":"desc","itemsView":"summary"}"""), + fixture.requests[2].requiredObject("params")) + + assertTrue(fixture.repository.loadOlderConversationTurns("server", "thread")) + assertEquals(listOf("oldest", "old", "running"), flow.value.history.turns.map { it.id }) + assertEquals(null, flow.value.history.olderTurnsCursor) + fixture.repository.loadConversationTurnDetails("server", "thread", "oldest") + assertEquals(2, flow.value.history.turns.first().items.size) + assertEquals(1, fixture.methods.count { it == "thread/resume" }) + assertFalse(flow.value.requiresRuntimeRetention) + } + } + + @Test + fun otherRpcErrorsAreNotSwallowed() = runTest { + for (error in listOf(RpcException(-32600, "invalid thread", null), + RpcException(-32001, "already has an active writer", null))) { + withFixture { fixture -> + fixture.resumeError = error + try { + fixture.repository.openConversation("server", "thread") + fail("Expected resume failure") + } catch (actual: RpcException) { + assertEquals(error.code, actual.code) + assertEquals(error.message, actual.message) + } + // The connection's existing overload retries remain intact. + assertEquals(List(if (error.code == -32001) 3 else 1) { "thread/resume" }, fixture.methods) + } + } + } + + @Test + fun failedReadAfterConflictStillFailsOpening() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.readError = RpcException(-32600, "thread not found", null) + try { + fixture.repository.openConversation("server", "thread") + fail("Expected read failure") + } catch (error: RpcException) { + assertEquals("thread not found", error.message) + } + assertEquals(listOf("thread/resume", "thread/read"), fixture.methods) + } + } + + @Test + fun failedRefreshKeepsReadOnlyStateAndReportsTheRpcError() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val flow = fixture.repository.openConversation("server", "thread") + val readCount = fixture.methods.count { it == "thread/read" } + fixture.resumeError = RpcException(-32600, "invalid resume params", null) + fixture.repository.refreshConversation("server", "thread") + assertTrue(flow.value.isExternalWriter) + assertEquals(readCount, fixture.methods.count { it == "thread/read" }) + assertTrue(requireNotNull(flow.value.errorMessage).technicalDetails!!.contains("invalid resume params")) + } + } + + @Test + fun readOnlyRepositoryRejectsWritesBeforeSendingRpcOrChangingState() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val flow = fixture.repository.openConversation("server", "thread") + val before = flow.value + val requestCount = fixture.requests.size + val mutations: List Unit> = listOf( + { fixture.repository.sendMessage("server", "thread", "hello") }, + { fixture.repository.queueInput("server", "thread", listOf(UserInputPart.Text("hello"))) }, + { fixture.repository.deleteQueuedInput("server", "thread", "entry"); Unit }, + { fixture.repository.runShellCommand("server", "thread", "pwd") }, + { fixture.repository.interrupt("server", "thread") }, + { fixture.repository.terminateBackgroundTerminal("server", "thread", "process") }, + { fixture.repository.compactConversation("server", "thread") }, + { fixture.repository.reviewConversation("server", "thread", ReviewTarget.UncommittedChanges, ReviewDelivery.INLINE); Unit }, + { fixture.repository.renameConversation("server", "thread", "name") }, + { fixture.repository.archiveConversation("server", "thread") }, + { fixture.repository.unarchiveConversation("server", "thread") }, + { fixture.repository.deleteConversation("server", "thread") }, + { fixture.repository.setGoal("server", "thread", "goal"); Unit }, + { fixture.repository.clearGoal("server", "thread") }, + { fixture.repository.respondToApproval("server", "thread", "request", ApprovalDecision.ACCEPT) }, + { fixture.repository.respondToApprovalWithExecPolicy("server", "thread", "request", json("{}")) }, + { fixture.repository.respondToApprovalWithNetworkPolicy("server", "thread", "request", json("{}")) }, + { fixture.repository.respondToQuestions("server", "thread", "request", emptyMap()) }, + { fixture.repository.respondToMcpElicitation("server", "thread", "request", "accept") }, + { fixture.repository.respondToPermissions("server", "thread", "request", json("{}"), false, false) }, + ) + for ((index, mutation) in mutations.withIndex()) { + try { + mutation() + fail("Mutation $index should be blocked") + } catch (error: IllegalStateException) { + assertEquals(before.viewOnlyMessage, error.message) + } + } + assertEquals(requestCount, fixture.requests.size) + assertEquals(before, flow.value) + assertTrue(fixture.responses.isEmpty()) + } + } + + @Test + fun forkIsAvailableWhileSourceRemainsReadOnly() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val source = fixture.repository.openConversation("server", "thread") + val forkId = fixture.repository.forkConversation("server", "thread", "old") + assertEquals("fork", forkId) + assertEquals("thread/fork", fixture.methods.last()) + assertEquals(json("""{"threadId":"thread","lastTurnId":"old","excludeTurns":true}"""), + fixture.requests.last().requiredObject("params")) + assertTrue(source.value.isExternalWriter) + fixture.resumeError = null + assertFalse(fixture.repository.openConversation("server", forkId).value.isViewOnly) + } + } + + @Test + fun refreshRetriesResumeAndRestoresWritableStateAndOptions() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val flow = fixture.repository.openConversation("server", "thread") + fixture.repository.refreshConversation("server", "thread") + assertTrue(flow.value.isExternalWriter) + fixture.resumeError = null + fixture.repository.refreshConversation("server", "thread") + assertFalse(flow.value.isViewOnly) + assertEquals(ConversationAccessMode.WRITABLE, flow.value.accessMode) + assertEquals("gpt-test", flow.value.model) + assertEquals("high", flow.value.reasoningEffort) + assertEquals(null, flow.value.errorMessage) + assertEquals(3, fixture.methods.count { it == "thread/resume" }) + fixture.repository.sendMessage("server", "thread", "continue") + assertTrue("turn/steer" in fixture.methods) + } + } + + @Test + fun automaticSyncReadsDesktopUpdatesThenSendsOnOriginalThreadWithoutForking() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val flow = fixture.repository.openConversation("server", "thread") + fixture.repository.loadOlderConversationTurns("server", "thread") + val monitor = backgroundScope.launch { + monitorExternalConversation(flow, MutableStateFlow(true), { + fixture.repository.syncExternalConversation("server", "thread") + }, {}) + } + runCurrent() + fixture.historyPage = """{"data":[${turn("next", "inProgress")},${turn("running", "completed")},${turn("old", "completed")}],"nextCursor":"older"}""" + advanceTimeBy(2_000); runCurrent() + assertTrue(flow.value.isExternalWriter) + assertEquals("next", flow.value.activeTurnId) + assertEquals(listOf("oldest", "old", "running", "next"), flow.value.history.turns.map { it.id }) + assertEquals(2, flow.value.history.turns.last().items.size) + fixture.resumeError = null + advanceTimeBy(2_000); runCurrent() + assertFalse(flow.value.isViewOnly) + val count = fixture.requests.size + advanceTimeBy(30_000); runCurrent() + assertEquals(count, fixture.requests.size) + fixture.repository.sendMessage("server", "thread", "continue in original") + assertEquals("thread", fixture.requests.last().requiredObject("params").requiredString("threadId")) + assertFalse("thread/fork" in fixture.methods) + monitor.cancel() + } + } + + @Test + fun automaticSyncPropagatesOtherErrorsWithoutReplacingTheSnapshotOrClaimingWriter() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + val flow = fixture.repository.openConversation("server", "thread") + val before = flow.value + val reads = fixture.methods.count { it == "thread/read" } + fixture.resumeError = RpcException(-32600, "invalid resume params", null) + try { + fixture.repository.syncExternalConversation("server", "thread") + fail("Expected sync error") + } catch (error: RpcException) { + assertEquals("invalid resume params", error.message) + } + assertEquals(before, flow.value) + assertEquals(reads, fixture.methods.count { it == "thread/read" }) + assertFalse("thread/fork" in fixture.methods) + } + } + + @Test + fun externalNewestTurnDetailsStayLoadedAndUpdateWhenDiskReportsInterrupted() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.isRunning = false + fixture.historyPage = """{"data":[${turn("latest", "interrupted")},${turn("old", "completed")}],"nextCursor":null}""" + fixture.itemPage = """{"data":[{"turnId":"latest","item":{"id":"progress","type":"agentMessage","phase":"commentary","text":"处理中"}}],"nextCursor":null}""" + val flow = fixture.repository.openConversation("server", "thread") + assertEquals(TurnDetailsState.Loaded, flow.value.history.turns.last().detailsState) + assertEquals("处理中", (flow.value.history.turns.last().items.single() as ConversationItem.AgentMessage).text) + assertEquals(TurnDetailsState.Unloaded, flow.value.history.turns.first().detailsState) + fixture.itemPage = fixture.itemPage!!.replace("处理中", "处理进展已更新") + fixture.repository.syncExternalConversation("server", "thread") + assertEquals(TurnDetailsState.Loaded, flow.value.history.turns.last().detailsState) + assertEquals("处理进展已更新", (flow.value.history.turns.last().items.single() as ConversationItem.AgentMessage).text) + assertEquals(2, fixture.methods.count { it == "thread/items/list" }) + assertFalse("thread/fork" in fixture.methods) + } + } + + @Test + fun reconnectRechecksOwnershipBeforeSendingFromPreviouslyWritableState() = runTest { + withFixture { fixture -> + val flow = fixture.repository.openConversation("server", "thread") + assertFalse(flow.value.isViewOnly) + fixture.prepareReconnect() + fixture.resumeError = activeWriterError() + try { + fixture.repository.sendMessage("server", "thread", "continue") + fail("Expected read-only rejection after reconnect") + } catch (error: IllegalStateException) { + assertEquals(flow.value.viewOnlyMessage, error.message) + } + assertTrue(flow.value.isExternalWriter) + assertFalse("turn/start" in fixture.methods) + assertFalse("turn/steer" in fixture.methods) + fixture.resumeError = null + fixture.prepareReconnect() + fixture.repository.connect("server") + assertFalse(flow.value.isViewOnly) + } + } + + @Test + fun reopeningRetriesResumeAndDoesNotUnsubscribeExternalWriter() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.repository.openConversation("server", "thread") + fixture.repository.releaseConversation("server", "thread") + runCurrent() + assertFalse("thread/unsubscribe" in fixture.methods) + fixture.resumeError = null + assertFalse(fixture.repository.openConversation("server", "thread").value.isViewOnly) + assertEquals(2, fixture.methods.count { it == "thread/resume" }) + } + } + + @Test + fun cancelledResumeDoesNotReadOrAcquireWriter() = runTest { + withFixture { fixture -> + fixture.holdResume = true + val job = launch { fixture.repository.openConversation("server", "thread") } + runCurrent() + job.cancel(CancellationException("test cancellation")) + job.join() + assertEquals(listOf("thread/resume"), fixture.methods) + } + } + + @Test + fun childInputRestrictionsRemainIndependentOfWriterOwnership() { + val child = ConversationState("server", "Server", "child", "Child", "/workspace", + canAcceptDirectInput = false, parentThreadId = "parent") + val external = child.copy(accessMode = ConversationAccessMode.EXTERNAL_WRITER) + assertTrue(external.isExternalWriter) + assertTrue(requireNotNull(external.viewOnlyMessage).contains("其他 Codex 客户端")) + val restored = external.withRefreshedConversation(child) + assertFalse(restored.isExternalWriter) + assertTrue(restored.isViewOnly) + } + + @Test + fun companionSendsTextThroughOwnerInSameThreadAndKeepsWriteGuards() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopAvailable = true + fixture.historyPage = idlePage + fixture.isRunning = false + val flow = fixture.repository.openConversation("server", "thread") + assertTrue(flow.value.canSendToDesktop) + assertTrue(flow.value.isViewOnly) + fixture.requests.clear() + fixture.repository.sendDesktopMessage("server", "thread", "continue here") + assertEquals(listOf("codexAndroid/desktop/turn/start"), fixture.methods) + val params = fixture.requests.single().requiredObject("params") + assertEquals("thread", params.requiredString("threadId")) + assertEquals("continue here", params.requiredString("text")) + assertEquals(setOf("threadId", "clientUserMessageId", "text"), params.keys) + assertEquals("desktop-turn", flow.value.activeTurnId) + assertEquals("thread", flow.value.threadId) + assertTrue(flow.value.isExternalWriter) + try { + fixture.repository.runShellCommand("server", "thread", "pwd") + fail("Desktop message capability must not grant Shell access") + } catch (_: IllegalStateException) { } + assertEquals(1, fixture.requests.size) + } + } + + @Test + fun companionLossDoesNotRetryForkOrChangeConversationOnFailedSend() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopAvailable = true + fixture.historyPage = idlePage + fixture.isRunning = false + val flow = fixture.repository.openConversation("server", "thread") + val before = flow.value + fixture.desktopSendError = RpcException(-32070, "client-disconnected", null) + fixture.requests.clear() + try { + fixture.repository.sendDesktopMessage("server", "thread", "hello") + fail("Delivery failure must surface") + } catch (error: RpcException) { assertEquals(-32070, error.code) } + assertEquals(before, flow.value) + assertEquals(listOf("codexAndroid/desktop/turn/start"), fixture.methods) + fixture.desktopAvailable = false + fixture.repository.refreshConversation("server", "thread") + assertFalse(flow.value.canSendToDesktop) + assertTrue(flow.value.isViewOnly) + fixture.resumeError = null + fixture.repository.refreshConversation("server", "thread") + assertFalse(flow.value.isViewOnly) + assertFalse(flow.value.desktopControlAvailable) + } + } + + @Test + fun runningDesktopConversationDoesNotAcceptAnotherStart() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopAvailable = true + val flow = fixture.repository.openConversation("server", "thread") + assertTrue(flow.value.canSendToDesktop) + val before = fixture.requests.size + try { + fixture.repository.sendDesktopMessage("server", "thread", "hello") + fail("Wait for the Desktop turn to finish") + } catch (_: IllegalStateException) { } + assertEquals(before, fixture.requests.size) + } + } + + @Test + fun externalInProgressHistoryDisablesSendWhenReaderDoesNotReportActive() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopAvailable = true + fixture.isRunning = false + val state = fixture.repository.openConversation("server", "thread").value + assertEquals("running", state.activeTurnId) + assertTrue(state.isRunning) + val before = fixture.requests.size + try { + fixture.repository.sendDesktopMessage("server", "thread", "hello") + fail("External history must gate a second start") + } catch (_: IllegalStateException) { } + assertEquals(before, fixture.requests.size) + } + } + + @Test + fun onlyMethodNotFoundErrorsDisableOptionalCompanion() = runTest { + for (error in listOf(RpcException(-32601, "Method not found", null), + RpcException(-32600, "Invalid request: unknown variant `codexAndroid/desktop/discover`, expected one of ...", null))) { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopDiscoverError = error + val state = fixture.repository.openConversation("server", "thread").value + assertTrue(state.isExternalWriter) + assertFalse(state.canSendToDesktop) + } + } + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopDiscoverError = RpcException(-32600, "Permission denied", null) + try { + fixture.repository.openConversation("server", "thread") + fail("Other errors must surface") + } catch (error: RpcException) { assertEquals("Permission denied", error.message) } + } + } + + @Test + fun mismatchedDesktopConfirmationIsNotAccepted() = runTest { + withFixture { fixture -> + fixture.resumeError = activeWriterError() + fixture.desktopAvailable = true + fixture.historyPage = idlePage + fixture.isRunning = false + val flow = fixture.repository.openConversation("server", "thread") + val before = flow.value + fixture.desktopConfirmationThread = "different-thread" + try { + fixture.repository.sendDesktopMessage("server", "thread", "hello") + fail("Confirmation must identify the original thread") + } catch (_: com.hebo.codex.data.remote.ProtocolException) { } + assertEquals(before, flow.value) + } + } + + @Test + fun companionDoesNotOverrideSubAgentInputPolicy() { + val state = ConversationState("server", "Server", "child", "Child", "", + accessMode = ConversationAccessMode.EXTERNAL_WRITER, desktopControlAvailable = true, + canAcceptDirectInput = false) + assertFalse(state.canSendToDesktop) + assertTrue(state.isViewOnly) + } + + private suspend fun TestScope.withFixture(block: suspend (Fixture) -> Unit) { + val fixture = Fixture(this) + try { + fixture.connect() + block(fixture) + } finally { + fixture.connection.close() + } + } + + private class Fixture(private val scope: TestScope) { + val requests = mutableListOf() + val responses = mutableListOf() + val methods get() = requests.map { it.requiredString("method") } + var resumeError: RpcException? = null + var readError: RpcException? = null + var desktopAvailable = false + var desktopDiscoverError: RpcException? = null + var desktopSendError: RpcException? = null + var desktopConfirmationThread = "thread" + var isRunning = true + var holdResume = false + var historyPage = initialPage + var itemPage: String? = null + var connection = createConnection() + + private fun createConnection() = AppServerConnection( + webSocketFactory = AppServerWebSocketFactory { listener -> + object : AppServerWebSocket { + override fun start() = listener.onOpen(this) + override fun setResponseTimeout(timeoutMillis: Long) = Unit + override fun send(text: String): Boolean { + val request = json(text) + if (request["method"] == null) { + responses += request + return true + } + val method = request.requiredString("method") + if (method == "initialized") return true + if (method != "initialize") requests += request + if (method == "thread/resume" && holdResume) return true + val error = when (method) { + "thread/resume" -> resumeError + "thread/read" -> readError + "codexAndroid/desktop/discover" -> desktopDiscoverError + "codexAndroid/desktop/turn/start" -> desktopSendError + else -> null + } + val response = buildJsonObject { + put("id", requireNotNull(request["id"])) + if (error != null) put("error", buildJsonObject { + put("code", error.code) + put("message", error.message) + }) else put("result", result(method, request["params"] as? JsonObject)) + } + listener.onMessage(this, response.toString()) + return true + } + override fun close(code: Int, reason: String): Boolean = true + override fun cancel() = Unit + } + }, + json = Json, + connectTimeoutMillis = 1_000, + parentContext = scope.backgroundScope.coroutineContext, + ) + private val server = ServerProfile("server", "Server", "ws://test:1234", ServerAuthentication.CAPABILITY_TOKEN) + private val store = Mockito.mock(ServerStore::class.java) + private val factory = Mockito.mock(AppServerClientFactory::class.java) + val repository: CodexRepository + + init { + Mockito.`when`(store.settings).thenReturn(flowOf(ServerSettings(listOf(server), emptyList(), "server"))) + repository = CodexRepository(store, factory, scope.backgroundScope) + } + + suspend fun prepareReconnect() { + scope.runCurrent() + connection.close() + scope.runCurrent() + connection = createConnection() + stubConnection() + } + + private suspend fun stubConnection() { + Mockito.`when`(store.readToken("server")).thenReturn("test-token") + Mockito.`when`(factory.create(server.webSocketUrl, "test-token", null, scope.backgroundScope.coroutineContext)) + .thenReturn(connection) + } + + suspend fun connect() { + stubConnection() + repository.connect("server") + scope.runCurrent() + requests.clear() + } + + private fun result(method: String, params: JsonObject?): JsonObject = when (method) { + "initialize" -> json("""{"userAgent":"test","codexHome":"/tmp/codex","platformFamily":"unix","platformOs":"linux"}""") + "thread/read" -> json("""{"thread":${fixtureThread("thread")}}""") + "thread/resume" -> json("""{"thread":${fixtureThread(params!!.requiredString("threadId"))},"model":"gpt-test","reasoningEffort":"high","initialTurnsPage":$historyPage}""") + "thread/fork" -> json("""{"thread":${thread("fork")}}""") + "codexAndroid/desktop/discover" -> json("""{"threadId":"thread","protocolVersion":1,"available":$desktopAvailable}""") + "codexAndroid/desktop/turn/start" -> json("""{"threadId":"$desktopConfirmationThread","clientUserMessageId":"${params!!.requiredString("clientUserMessageId")}","turn":{"id":"desktop-turn","status":"inProgress","items":[]}}""") + "thread/turns/list" -> if (params?.get("cursor") != null) json("""{"data":[${turn("oldest", "completed")}],"nextCursor":null}""") else json(historyPage) + "thread/items/list" -> itemPage?.let(::json) ?: run { + val turnId = params!!.requiredString("turnId") + val isFirst = params["cursor"] == null + json("""{"data":[{"turnId":"$turnId","item":{"id":"answer-${if (isFirst) 1 else 2}","type":"agentMessage","text":"history"}}],"nextCursor":${if (isFirst) "\"next-items\"" else "null"}}""") + } + "turn/steer" -> json("""{"turnId":"running"}""") + else -> json("""{"data":[],"nextCursor":null}""") + } + + private fun fixtureThread(id: String) = thread(id).let { + if (isRunning) it else it.replace("\"type\":\"active\",\"activeFlags\":[]", "\"type\":\"idle\"") + } + } + + companion object { + private fun activeWriterError() = RpcException(-32600, "thread-store conflict: thread thread already has an active writer", null) + private fun json(value: String) = Json.parseToJsonElement(value).jsonObject + private fun thread(id: String) = """{"id":"$id","sessionId":"session","name":"Test","preview":"history","cwd":"/workspace","createdAt":1,"updatedAt":2,"status":{"type":"active","activeFlags":[]},"historyMode":"paginated","canAcceptDirectInput":true}""" + private fun turn(id: String, status: String) = """{"id":"$id","status":"$status","itemsView":"summary","items":[]}""" + private val initialPage = """{"data":[${turn("running", "inProgress")},${turn("old", "completed")}],"nextCursor":"older"}""" + private val idlePage = """{"data":[${turn("old", "completed")}],"nextCursor":null}""" + } +} diff --git a/companion/.gitignore b/companion/.gitignore new file mode 100644 index 0000000..670a936 --- /dev/null +++ b/companion/.gitignore @@ -0,0 +1,2 @@ +__pycache__/ +.venv/ diff --git a/companion/README.md b/companion/README.md new file mode 100644 index 0000000..9370ad5 --- /dev/null +++ b/companion/README.md @@ -0,0 +1,63 @@ +# Desktop companion(实验性) + +让手机向正在其他 Codex Desktop 客户端中打开的原会话发送文字,保留原 thread ID。普通请求仍交给现有 app-server;发送文字时,通过本机 Desktop IPC 找到 writer,再调用该 writer 的 follower 接收入口。伴随服务不持有原会话的 writer,也不写入 rollout 文件。 + +这是 codex-android 的可选扩展,**不是 Codex CLI 的官方 RPC,也不是稳定的 Desktop 公共接口**。目前适配本机桌面包中的 Unix IPC:4 字节小端长度 + UTF-8 JSON,`thread-owner-discovery` v1、`thread-follower-start-turn` v2。当前 CLI 协议按 0.159.2 导出的 JSON Schema 核对;Desktop IPC 另行核对已安装桌面包。桌面版本、平台或 IPC 协议变化后需重新验证。 + +## 使用 + +需要 Python 3.11+、本机运行的 Codex Desktop,以及已有的 capability-token 认证 app-server。Desktop 和 app-server 必须使用同一个 Codex 历史目录。当前支持 Linux/Unix;不实现 Windows named pipe。 + +在仓库根目录安装独立依赖环境: + +```bash +python3 -m venv companion/.venv +companion/.venv/bin/pip install -r companion/requirements.txt +``` + +启动代理(将示例 Tailscale 地址替换为电脑的实际地址): + +```bash +companion/.venv/bin/python companion/desktop_companion.py \ + --backend ws://100.100.100.100:8390 \ + --bind 100.100.100.100 --port 8391 \ + --token-file ~/.codex/app-server-token +``` + +Android 的服务器地址改为 `ws://100.100.100.100:8391`,沿用已有 capability token。原 app-server 的端口保持不变;伴随代理是单独的入口。只允许监听 loopback 或具体 Tailscale IPv4 地址。令牌文件须由当前用户持有、权限为 600,不接受符号链接。Desktop IPC socket 必须是当前用户私有的 socket,Linux 还检查连接对端的用户身份。令牌用于认证代理入口和原 app-server,不出现在日志或参数中。 + +默认 IPC 为 `${CODEX_HOME:-$HOME/.codex}/ipc/ipc.sock`;需要时可通过 `--ipc` 指定已有 socket。服务不会创建、替换或删除 Desktop socket、writer lock 或历史文件。关闭代理后可以继续使用原 app-server 地址。 + +## 手机行为 + +- `thread/resume` 只有遇到 active-writer 冲突才进入 external-writer 模式,通过 `thread/read` 和分页历史查看原对话。 +- 检测到 Desktop owner 后显示“发送到原对话”。当前仅支持文字,沿用 Desktop 的模型、权限、工作目录等设置;发送端无法覆盖这些设置。 +- 正在处理的会话先等待当前回合结束。此版本不提供跨 Desktop 的排队、打断或 steer;不会伪装成已排队。 +- Shell、审批、工具操作和会话设置仍由 Desktop 端处理。子代理 `canAcceptDirectInput=false` 不因伴随服务改变。 +- 未安装伴随服务、Desktop 不可用、writer 是独立 CLI/TUI、owner 已断开或协议不兼容时,保留只读查看和草稿。不会自动 Fork。 +- 每次发送重新寻找 owner,随后只发一次。发送未确认时保留草稿,请先检查原会话,避免手动重试造成重复消息。 +- 历史在 Android 前台轮询更新;writer 释放后正常 `thread/resume` 恢复本地可写状态。“分叉为新对话”仍是独立且明确的选项。 + +## 扩展协议 v1 + +使用既有 WS 连接及 JSON-RPC 请求格式。两个方法不属于 CLI 的 `ClientRequest` Schema,代理拦截后处理;其他请求和服务器通知按原样转发。 + +`codexAndroid/desktop/discover` 参数:`{ "threadId": "" }`。 +结果:`{ "threadId": "", "protocolVersion": 1, "available": true/false }`。 + +`codexAndroid/desktop/turn/start` 参数:`{ "threadId": "", "clientUserMessageId": "", "text": "..." }`。 +结果:`{ "threadId": "", "clientUserMessageId": "", "turn": }`。 + +发送只接受上述参数及最多 1 MiB 的非空文字。无法用它指定任意 RPC、运行 Shell、代答审批或调整权限。Android 校验确认中的 thread ID、message ID 和 turn ID;发送异常不自动重试。发现操作不可用返回 `available=false`;发送失败返回 `-32070`,非法参数为 `-32602`。 + +原生 CLI 对未知方法可能返回 `-32601`,0.159.2 实测返回 `-32600` 和精确的 unknown-variant 提示;Android 仅将这一扩展方法的“未知方法”视为未安装服务,其他错误仍上报。 + +## 验证 + +```bash +companion/.venv/bin/python -m unittest discover -s companion -p 'test_*.py' -v +``` + +单元测试验证 Unix IPC framing、owner 路由、认证、RPC 原样转发、不可用状态、参数限制和发送不重试。集成测试需要 PATH 中的 `codex`:在 `/tmp` 下创建隔离的 `CODEX_HOME`,启动两个独立 app-server,使用不可达的本机模型地址,禁用插件发现;验证 active-writer 冲突后经 owner 发送的文字存入原 thread,且手机侧不发送 `turn/start` 或 `thread/fork`。Desktop IPC router 在此集成测试中为模拟实现,不等同于真实手机/Desktop UI 测试。测试结束仅清理自身进程组和临时目录。 + +本机真实 Desktop 已验证 owner 发现、代理认证和原会话读取,并在专门的测试对话中经 follower v2 提交有效文字、收到回复,确认 thread ID 未改变。手机用户也已确认文字能够发送到原对话。处理过程缓存修复通过了回归测试;包含该修复的最新 APK 尚未进行手机实机验收。当前验证仅覆盖本机 Unix IPC 环境,启用此实验功能后应先用专门的测试对话做设备验收。 diff --git a/companion/desktop_companion.py b/companion/desktop_companion.py new file mode 100644 index 0000000..1e271ba --- /dev/null +++ b/companion/desktop_companion.py @@ -0,0 +1,285 @@ +"""Optional authenticated WS proxy to a Codex app-server and its Desktop owner. + +Only the two codexAndroid/desktop extensions use Desktop's private local IPC. +Normal CLI requests, notifications and approval replies pass through unchanged. +No writer is acquired, released or replaced by the IPC adapter. +""" + +import argparse +import asyncio +from contextlib import asynccontextmanager +import hmac +from http import HTTPStatus +import ipaddress +import json +import os +from pathlib import Path +import socket +import stat +import struct +import uuid + +from websockets.asyncio.client import connect +from websockets.asyncio.server import serve +from websockets.exceptions import ConnectionClosed + +DISCOVER = "codexAndroid/desktop/discover" +START = "codexAndroid/desktop/turn/start" +PROTOCOL_VERSION = 1 +MAX_FRAME = 64 * 1024 * 1024 + + +class CompanionError(Exception): + def __init__(self, message, code=-32070): + super().__init__(message) + self.code = code + + +def validate_uuid(value): + if not isinstance(value, str): + raise CompanionError("Expected a UUID", -32602) + try: + if str(uuid.UUID(value)) != value: + raise ValueError() + except ValueError: + raise CompanionError("Expected a canonical UUID", -32602) from None + return value + + +def private_file(path): + info = path.lstat() + if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() or info.st_mode & 0o077: + raise ValueError("Token file must be a private regular file owned by this user") + token = path.read_text().strip() + if not token or "\n" in token or "\r" in token: + raise ValueError("Invalid token file") + return token + + +def validate_bind(host): + address = ipaddress.ip_address(host) + tailscale = ipaddress.ip_network("100.64.0.0/10") + if not address.is_loopback and address not in tailscale: + raise ValueError("Bind only to loopback or a specific Tailscale IPv4 address") + return host + + +class DesktopIpc: + """Versioned Unix IPC adapter. Never create the Desktop socket or router.""" + + def __init__(self, path): + self.path = Path(path) + self.client_id = "initializing-client" + self.reader = self.writer = None + + async def __aenter__(self): + info, parent = self.path.lstat(), self.path.parent.lstat() + if (not stat.S_ISSOCK(info.st_mode) or info.st_uid != os.getuid() or + info.st_mode & 0o077 or not stat.S_ISDIR(parent.st_mode) or + parent.st_uid != os.getuid() or parent.st_mode & 0o022): + raise CompanionError("Desktop IPC endpoint is not private to this user") + try: + self.reader, self.writer = await asyncio.wait_for( + asyncio.open_unix_connection(str(self.path)), 3) + # Linux additionally verifies the connected peer, after pathname checks. + if hasattr(socket, "SO_PEERCRED"): + peer = self.writer.get_extra_info("socket") + _, uid, _ = struct.unpack("3i", peer.getsockopt(socket.SOL_SOCKET, socket.SO_PEERCRED, 12)) + if uid != os.getuid(): + raise CompanionError("Desktop IPC peer belongs to another user") + result = await self.request("initialize", {"clientType": "codex-android-companion"}, 0) + self.client_id = result["result"]["clientId"] + return self + except BaseException: + await self.close() + raise + + async def __aexit__(self, *_): + await self.close() + + async def close(self): + if self.writer: + self.writer.close() + try: + await self.writer.wait_closed() + except (OSError, ConnectionError): + pass + + async def write(self, message): + body = json.dumps(message, separators=(",", ":"), ensure_ascii=False).encode() + if len(body) > MAX_FRAME: + raise CompanionError("Desktop IPC frame too large") + self.writer.write(struct.pack(" 1024 * 1024: + raise CompanionError("Expected nonempty text up to 1 MiB", -32602) + try: + async with self.ipc_factory(self.ipc_path) as ipc: + if method == DISCOVER: + await ipc.owner(thread_id) + return {"threadId": thread_id, "protocolVersion": PROTOCOL_VERSION, "available": True} + return await ipc.start(thread_id, message_id, text) + except (CompanionError, OSError, TimeoutError) as error: + # Read-only capability discovery is allowed to fail closed. A send + # failure is always visible and is never replayed by this proxy. + if method == DISCOVER: + return {"threadId": thread_id, "protocolVersion": PROTOCOL_VERSION, "available": False} + if isinstance(error, CompanionError): + raise + raise CompanionError("Desktop unavailable; message was not confirmed") from error + + async def handle(self, client): + try: + async with connect(self.backend_url, additional_headers={"Authorization": "Bearer " + self.token}, + max_size=MAX_FRAME, compression=None, proxy=None) as backend: + async def from_client(): + async for raw in client: + if not isinstance(raw, str): + await client.close(1003, "Text frames required") + return + try: + message = json.loads(raw) + except (ValueError, TypeError): + await client.close(1007, "Invalid JSON") + return + if not isinstance(message, dict): + await client.close(1007, "Object required") + return + method = message.get("method", "") + if not isinstance(method, str) or not method.startswith("codexAndroid/desktop/"): + await backend.send(raw) + continue + # The extension is a client request, never a notification + # or an approval response with a reused request id. + if "id" not in message or "result" in message or "error" in message: + await client.close(1007, "Extension request id required") + return + response = {"id": message["id"]} + try: + response["result"] = await self.extension(method, message.get("params")) + except CompanionError as error: + response["error"] = {"code": error.code, "message": str(error)} + except (ValueError, TypeError, KeyError): + response["error"] = {"code": -32070, "message": "Unsupported Desktop IPC response; delivery unconfirmed"} + await client.send(json.dumps(response, ensure_ascii=False)) + + async def from_backend(): + async for raw in backend: + await client.send(raw) + + tasks = [asyncio.create_task(from_client()), asyncio.create_task(from_backend())] + try: + done, _ = await asyncio.wait(tasks, return_when=asyncio.FIRST_COMPLETED) + for task in done: + task.result() + finally: + for task in tasks: + task.cancel() + await asyncio.gather(*tasks, return_exceptions=True) + except (ConnectionClosed, OSError): + await client.close(1011, "Companion connection closed") + + @asynccontextmanager + async def listen(self, host, port): + validate_bind(host) + async with serve(self.handle, host, port, process_request=self.authenticate, + max_size=MAX_FRAME, compression=None) as server: + yield server + + +async def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--bind", default="127.0.0.1") + parser.add_argument("--port", type=int, default=8391) + parser.add_argument("--backend", required=True, help="Existing authenticated Codex app-server WS URL") + parser.add_argument("--token-file", type=Path, required=True) + parser.add_argument("--ipc", type=Path, default=Path(os.environ.get("CODEX_HOME", str(Path.home()/".codex")))/"ipc/ipc.sock") + args = parser.parse_args() + token = private_file(args.token_file) + companion = DesktopCompanion(args.backend, token, args.ipc) + async with companion.listen(args.bind, args.port): + print("Codex Android companion ready", flush=True) + await asyncio.Future() + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except KeyboardInterrupt: + pass diff --git a/companion/requirements.txt b/companion/requirements.txt new file mode 100644 index 0000000..7bc03d7 --- /dev/null +++ b/companion/requirements.txt @@ -0,0 +1 @@ +websockets==15.0.1 diff --git a/companion/test_desktop_companion.py b/companion/test_desktop_companion.py new file mode 100644 index 0000000..e9f58af --- /dev/null +++ b/companion/test_desktop_companion.py @@ -0,0 +1,206 @@ +import asyncio +from contextlib import asynccontextmanager +import json +from pathlib import Path +import struct +import tempfile +import unittest +import uuid + +from websockets.asyncio.client import connect +from websockets.asyncio.server import serve +from websockets.exceptions import InvalidStatus + +from desktop_companion import (CompanionError, DesktopCompanion, DesktopIpc, + DISCOVER, START, private_file, validate_bind) + +THREAD = "aaaaaaaa-aaaa-4aaa-aaaa-aaaaaaaaaaaa" +MESSAGE = "bbbbbbbb-bbbb-4bbb-bbbb-bbbbbbbbbbbb" +OWNER = "cccccccc-cccc-4ccc-cccc-cccccccccccc" +TOKEN = "a-test-capability-token" + + +class IpcFixture: + """Desktop router framing and follower responses, without user conversations.""" + def __init__(self): + self.messages = [] + self.owner_available = True + self.delivery_error = None + self.confirmation_owner = OWNER + self.deliver = None + self.tasks = set() + + async def handle(self, reader, writer): + self.tasks.add(asyncio.current_task()) + try: + while True: + size = struct.unpack("