From ac2be3c1d58b7851596f4513318bd0de01a0d043 Mon Sep 17 00:00:00 2001 From: Miyamura80 Date: Mon, 7 Sep 2026 11:18:15 +0100 Subject: [PATCH 1/3] Fix long-running Mobile Bash computer actions --- .../stdiod/mcp/AndroidComputerSource.kt | 21 ++++- .../stdiod/mcp/MobileCommandRouter.kt | 32 ++++++- .../java/ai/sealgate/stdiod/mcp/UiSettler.kt | 3 +- .../ai/sealgate/stdiod/tunnel/TunnelClient.kt | 21 ++++- .../sealgate/stdiod/mcp/ComputerModuleTest.kt | 8 +- .../stdiod/mcp/MobileBashRuntimeTest.kt | 83 +++++++++++++++++++ .../ai/sealgate/stdiod/mcp/UiSettlerTest.kt | 21 +++++ 7 files changed, 177 insertions(+), 12 deletions(-) diff --git a/app/src/main/java/ai/sealgate/stdiod/mcp/AndroidComputerSource.kt b/app/src/main/java/ai/sealgate/stdiod/mcp/AndroidComputerSource.kt index f9ef8bc..a56e1db 100644 --- a/app/src/main/java/ai/sealgate/stdiod/mcp/AndroidComputerSource.kt +++ b/app/src/main/java/ai/sealgate/stdiod/mcp/AndroidComputerSource.kt @@ -166,7 +166,14 @@ class AndroidComputerSource(context: Context) : ComputerSource { val baseline = uiStateMarker(service) val path = Path().apply { moveTo(x.toFloat(), y.toFloat()) } val performed = dispatchGesture(service, path, durationMillis) - finishAction(service, "tap", performed, if (performed) null else "Android rejected the tap gesture", baseline) + finishAction( + service, + "tap", + performed, + if (performed) null else "Android rejected the tap gesture", + baseline, + requireTransition = false, + ) } override fun swipe( @@ -185,7 +192,14 @@ class AndroidComputerSource(context: Context) : ComputerSource { lineTo(endX.toFloat(), endY.toFloat()) } val performed = dispatchGesture(service, path, durationMillis) - finishAction(service, "swipe", performed, if (performed) null else "Android rejected the swipe gesture", baseline) + finishAction( + service, + "swipe", + performed, + if (performed) null else "Android rejected the swipe gesture", + baseline, + requireTransition = false, + ) } override fun globalAction(action: String): ComputerOperationResult = withService { service -> @@ -263,9 +277,10 @@ class AndroidComputerSource(context: Context) : ComputerSource { performed: Boolean, error: String?, baseline: UiStateMarker, + requireTransition: Boolean = true, ): ComputerOperationResult { if (!performed) return actionFailure(action, error ?: "action was not performed") - val settle = uiSettler.awaitPostAction(baseline) { uiStateMarker(service) } + val settle = uiSettler.awaitPostAction(baseline, requireTransition) { uiStateMarker(service) } val observation = captureObservation(service) val payload = buildJsonObject { put("action", buildJsonObject { diff --git a/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt b/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt index a07dcba..916b060 100644 --- a/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt +++ b/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt @@ -52,6 +52,7 @@ class MobileCommandRouter(modules: List) { private val supplementLock = Any() private val supplements = LinkedHashMap() private var pendingSupplementBytes = 0L + private var latestComputerSupplementToken: String? = null init { val mappedTools = SPECS.map(CommandSpec::tool).toSet() @@ -123,7 +124,10 @@ class MobileCommandRouter(modules: List) { val typedContent = content.filter { it["type"]?.jsonPrimitive?.content != "text" } val structuredContent = result["structuredContent"] as? JsonObject val supplementToken = if (typedContent.isNotEmpty() || structuredContent != null) { - retainSupplement(MobileCommandSupplement(typedContent, structuredContent)) + retainSupplement( + MobileCommandSupplement(typedContent, structuredContent), + replacePreviousComputerObservation = spec.module == ComputerModule.NAME, + ) } else { null } @@ -138,6 +142,7 @@ class MobileCommandRouter(modules: List) { fun clearSupplements() = synchronized(supplementLock) { supplements.clear() pendingSupplementBytes = 0L + latestComputerSupplementToken = null } fun availableNamespacesJson(): String = buildJsonArray { @@ -152,18 +157,39 @@ class MobileCommandRouter(modules: List) { tokens.mapNotNull(supplements::remove).also { supplements.clear() pendingSupplementBytes = 0L + latestComputerSupplementToken = null } } - private fun retainSupplement(supplement: MobileCommandSupplement): String = synchronized(supplementLock) { - check(supplements.size < MAX_PENDING_SUPPLEMENTS) { "too many pending mobile command attachments" } + private fun retainSupplement( + supplement: MobileCommandSupplement, + replacePreviousComputerObservation: Boolean, + ): String = synchronized(supplementLock) { val supplementBytes = supplement.serializedBytes() + check(supplementBytes <= MAX_PENDING_SUPPLEMENT_BYTES) { + "mobile command attachment exceeds 4 MiB" + } + // A shell script can perform many observation-producing actions while + // redirecting their textual output. Typed MCP attachments do not flow + // through Bash file descriptors, so retaining every one would grow an + // invisible side channel until the script failed. Keep only the most + // recent computer observation: it represents the device state after + // the latest action and bounds loops independently of their length. + // Other typed results (for example multiple camera snapshots) retain + // their existing multi-attachment behavior. + if (replacePreviousComputerObservation) { + latestComputerSupplementToken?.let(supplements::remove)?.let { previous -> + pendingSupplementBytes -= previous.serializedBytes() + } + } + check(supplements.size < MAX_PENDING_SUPPLEMENTS) { "too many pending mobile command attachments" } check(supplementBytes <= MAX_PENDING_SUPPLEMENT_BYTES - pendingSupplementBytes) { "pending mobile command attachments exceed 4 MiB" } val token = nextSupplementId.incrementAndGet().toString() supplements[token] = supplement pendingSupplementBytes += supplementBytes + if (replacePreviousComputerObservation) latestComputerSupplementToken = token token } diff --git a/app/src/main/java/ai/sealgate/stdiod/mcp/UiSettler.kt b/app/src/main/java/ai/sealgate/stdiod/mcp/UiSettler.kt index bd91013..625cb89 100644 --- a/app/src/main/java/ai/sealgate/stdiod/mcp/UiSettler.kt +++ b/app/src/main/java/ai/sealgate/stdiod/mcp/UiSettler.kt @@ -33,6 +33,7 @@ internal class UiSettler( ) { fun awaitPostAction( baseline: UiStateMarker, + requireTransition: Boolean = true, sample: () -> UiStateMarker, ): UiSettleResult { val startedAt = uptimeMillis() @@ -61,7 +62,7 @@ internal class UiSettler( val transitionObserved = eventObserved || windowChanged val eventStreamIsQuiet = now - current.lastEventUptimeMillis >= quietMillis val activeWindowIsStable = current.hasActiveWindow && now - stableSince >= quietMillis - if (transitionObserved && eventStreamIsQuiet && activeWindowIsStable) { + if ((!requireTransition || transitionObserved) && eventStreamIsQuiet && activeWindowIsStable) { return UiSettleResult( settled = true, postActionEventObserved = eventObserved, diff --git a/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt b/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt index df28224..00a295b 100644 --- a/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt +++ b/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt @@ -16,6 +16,9 @@ import okhttp3.Response import okhttp3.WebSocket import okhttp3.WebSocketListener import java.util.concurrent.TimeUnit +import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.Executors +import java.util.concurrent.RejectedExecutionException import java.util.concurrent.atomic.AtomicBoolean import kotlin.coroutines.resume import kotlin.concurrent.thread @@ -60,7 +63,7 @@ class TunnelClient( private val modulesByName: Map = modules.associateBy { it.name } /** server_id (backend's key for `mcp_frame`s) → built-in module. */ - private val modulesByServerId = HashMap() + private val modulesByServerId = ConcurrentHashMap() private val _state = MutableStateFlow(TunnelState.Disconnected) val state: StateFlow = _state @@ -71,6 +74,9 @@ class TunnelClient( private val modulesClosed = AtomicBoolean(false) private val modulesCloseFinished = CompletableDeferred() private val moduleLock = Any() + private val moduleExecutor = Executors.newSingleThreadExecutor { runnable -> + Thread(runnable, "mobile-mcp-requests").apply { isDaemon = true } + } fun start() { if (stopped.get()) return @@ -84,6 +90,7 @@ class TunnelClient( loopJob = null webSocket?.close(NORMAL_CLOSURE, "client stopping") webSocket = null + moduleExecutor.shutdownNow() if (modulesClosed.compareAndSet(false, true)) { // A QuickJS evaluation may hold its runtime lock until the 60-second // execution limit. Never make the service/main thread wait for it. @@ -179,7 +186,17 @@ class TunnelClient( bindServers(webSocket, frame.added + frame.updated) frame.removed.forEach(modulesByServerId::remove) } - is McpFrame -> routeMcpFrame(webSocket, frame) + is McpFrame -> { + // Local modules may perform gestures, screenshots, or a + // full Bash script. Never run them on OkHttp's reader + // callback: doing so prevents WebSocket control frames + // and unrelated tunnel messages from being processed. + try { + moduleExecutor.execute { routeMcpFrame(webSocket, frame) } + } catch (_: RejectedExecutionException) { + // The client was stopped between parsing and queueing. + } + } is Ping -> send(webSocket, Pong) is Pong -> Unit // Built-in modules have no spawn-time env/spec to store. diff --git a/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt b/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt index 002a584..31ed511 100644 --- a/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt @@ -66,7 +66,7 @@ class ComputerModuleTest { } @Test - fun routerBoundsPendingSupplementsBySerializedBytes() { + fun routerKeepsOnlyTheLatestPendingSupplement() { val source = FakeComputerSource().apply { screenshotData = "a".repeat((MobileCommandRouter.MAX_PENDING_SUPPLEMENT_BYTES / 2 + 1024).toInt()) } @@ -77,8 +77,10 @@ class ComputerModuleTest { val second = Json.parseToJsonElement(router.executeJson(request)).jsonObject assertEquals(0, first["exitCode"]!!.jsonPrimitive.content.toInt()) - assertEquals(1, second["exitCode"]!!.jsonPrimitive.content.toInt()) - assertTrue(second["stderr"]!!.jsonPrimitive.content.contains("attachments exceed 4 MiB")) + assertEquals(0, second["exitCode"]!!.jsonPrimitive.content.toInt()) + val firstToken = first["supplementToken"]!!.jsonPrimitive.content + val secondToken = second["supplementToken"]!!.jsonPrimitive.content + assertTrue(router.takeSupplements(listOf(firstToken, secondToken)).size == 1) router.clearSupplements() } diff --git a/app/src/test/java/ai/sealgate/stdiod/mcp/MobileBashRuntimeTest.kt b/app/src/test/java/ai/sealgate/stdiod/mcp/MobileBashRuntimeTest.kt index 55225b3..068940b 100644 --- a/app/src/test/java/ai/sealgate/stdiod/mcp/MobileBashRuntimeTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/mcp/MobileBashRuntimeTest.kt @@ -6,6 +6,7 @@ import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonPrimitive import kotlinx.serialization.json.buildJsonArray import kotlinx.serialization.json.buildJsonObject +import kotlinx.serialization.json.jsonPrimitive import org.junit.Assert.assertEquals import org.junit.Assert.assertTrue import org.junit.Test @@ -133,4 +134,86 @@ class MobileBashRuntimeTest { runtime.close() } } + + @Test + fun repeatedRedirectedObservationsKeepOnlyTheLatestAttachment() { + var observation = 0 + val probe = object : BaseMcpModule() { + override val name = "computer" + override fun toolDescriptors(): JsonElement = buildJsonArray { + add(buildJsonObject { + put("name", JsonPrimitive("computer_observe")) + put("description", JsonPrimitive("probe")) + put("inputSchema", buildJsonObject { + put("type", JsonPrimitive("object")) + put("properties", buildJsonObject {}) + }) + }) + } + override fun callTool(id: JsonElement, toolName: String, arguments: JsonObject): JsonObject { + observation++ + return JsonRpc.toolResult( + id, + listOf( + JsonRpc.textContent("observation $observation"), + JsonRpc.imageContent("a".repeat(1024 * 1024), "image/jpeg"), + ), + buildJsonObject { put("observationId", JsonPrimitive("obs_$observation")) }, + ) + } + } + val source = File("src/main/assets/mobile-bash-runtime.js").readText() + val runtime = QuickJsMobileBashRuntime({ source }, MobileCommandRouter(listOf(probe))) + try { + val result = runtime.execute( + "for i in ${'$'}(seq 1 17); do computer observe > /tmp/observation; done; computer observe", + ) + + assertEquals(0, result.exitCode) + assertEquals(18, observation) + assertEquals(1, result.supplements.size) + assertEquals( + "obs_18", + result.supplements.single().structuredContent!!["observationId"]!!.jsonPrimitive.content, + ) + } finally { + runtime.close() + } + } + + @Test + fun nonComputerAttachmentsStillAccumulateWithinTheExistingBudget() { + var snapshot = 0 + val probe = object : BaseMcpModule() { + override val name = "camera" + override fun toolDescriptors(): JsonElement = buildJsonArray { + add(buildJsonObject { + put("name", JsonPrimitive("camera_snap")) + put("description", JsonPrimitive("probe")) + put("inputSchema", buildJsonObject { + put("type", JsonPrimitive("object")) + put("properties", buildJsonObject {}) + }) + }) + } + override fun callTool(id: JsonElement, toolName: String, arguments: JsonObject): JsonObject { + snapshot++ + return JsonRpc.toolResult( + id, + listOf(JsonRpc.imageContent("aA==", "image/jpeg")), + buildJsonObject { put("snapshot", JsonPrimitive(snapshot)) }, + ) + } + } + val source = File("src/main/assets/mobile-bash-runtime.js").readText() + val runtime = QuickJsMobileBashRuntime({ source }, MobileCommandRouter(listOf(probe))) + try { + val result = runtime.execute("camera snap; camera snap") + + assertEquals(0, result.exitCode) + assertEquals(2, result.supplements.size) + } finally { + runtime.close() + } + } } diff --git a/app/src/test/java/ai/sealgate/stdiod/mcp/UiSettlerTest.kt b/app/src/test/java/ai/sealgate/stdiod/mcp/UiSettlerTest.kt index 015573a..4f63752 100644 --- a/app/src/test/java/ai/sealgate/stdiod/mcp/UiSettlerTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/mcp/UiSettlerTest.kt @@ -51,6 +51,27 @@ class UiSettlerTest { assertTrue(now >= 1_500L) } + @Test + fun completedGestureCanSettleWithoutAUiTransition() { + var now = 1_000L + val baseline = marker(sequence = 7, eventAt = 100, packageName = "same", windowId = 1) + val settler = UiSettler( + uptimeMillis = { now }, + sleepMillis = { now += it }, + quietMillis = 300L, + pollMillis = 50L, + timeoutMillis = 2_000L, + ) + + val result = settler.awaitPostAction(baseline, requireTransition = false) { baseline } + + assertTrue(result.settled) + assertFalse(result.postActionEventObserved) + assertFalse(result.activeWindowChanged) + assertTrue(now >= 1_300L) + assertTrue(now < 3_000L) + } + @Test fun sameWindowContentEventCanSettleWithoutWindowIdentityChange() { var now = 1_000L From bb48f695a508bec75e91bb4a3693439b8c5f2a1a Mon Sep 17 00:00:00 2001 From: Miyamura80 Date: Mon, 7 Sep 2026 11:34:32 +0100 Subject: [PATCH 2/3] Address review feedback for long-running tasks --- .../stdiod/mcp/MobileCommandRouter.kt | 21 +++-- .../ai/sealgate/stdiod/tunnel/TunnelClient.kt | 80 +++++++++++++++---- .../sealgate/stdiod/mcp/ComputerModuleTest.kt | 54 ++++++++++++- .../tunnel/TunnelClientLifecycleTest.kt | 41 ++++++++++ 4 files changed, 169 insertions(+), 27 deletions(-) diff --git a/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt b/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt index 916b060..bc2af19 100644 --- a/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt +++ b/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt @@ -126,7 +126,8 @@ class MobileCommandRouter(modules: List) { val supplementToken = if (typedContent.isNotEmpty() || structuredContent != null) { retainSupplement( MobileCommandSupplement(typedContent, structuredContent), - replacePreviousComputerObservation = spec.module == ComputerModule.NAME, + replacePreviousComputerObservation = + spec.module == ComputerModule.NAME && spec.tool != "computer_status", ) } else { null @@ -177,15 +178,19 @@ class MobileCommandRouter(modules: List) { // the latest action and bounds loops independently of their length. // Other typed results (for example multiple camera snapshots) retain // their existing multi-attachment behavior. - if (replacePreviousComputerObservation) { - latestComputerSupplementToken?.let(supplements::remove)?.let { previous -> - pendingSupplementBytes -= previous.serializedBytes() - } - } - check(supplements.size < MAX_PENDING_SUPPLEMENTS) { "too many pending mobile command attachments" } - check(supplementBytes <= MAX_PENDING_SUPPLEMENT_BYTES - pendingSupplementBytes) { + val previousToken = latestComputerSupplementToken.takeIf { replacePreviousComputerObservation } + val previous = previousToken?.let(supplements::get) + val previousBytes = previous?.serializedBytes() ?: 0L + val projectedCount = supplements.size - if (previous == null) 0 else 1 + val projectedBytes = pendingSupplementBytes - previousBytes + supplementBytes + check(projectedCount < MAX_PENDING_SUPPLEMENTS) { "too many pending mobile command attachments" } + check(projectedBytes <= MAX_PENDING_SUPPLEMENT_BYTES) { "pending mobile command attachments exceed 4 MiB" } + if (previousToken != null && previous != null) { + supplements.remove(previousToken) + pendingSupplementBytes -= previousBytes + } val token = nextSupplementId.incrementAndGet().toString() supplements[token] = supplement pendingSupplementBytes += supplementBytes diff --git a/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt b/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt index 00a295b..e72d78d 100644 --- a/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt +++ b/app/src/main/java/ai/sealgate/stdiod/tunnel/TunnelClient.kt @@ -16,10 +16,12 @@ import okhttp3.Response import okhttp3.WebSocket import okhttp3.WebSocketListener import java.util.concurrent.TimeUnit +import java.util.concurrent.ArrayBlockingQueue import java.util.concurrent.ConcurrentHashMap -import java.util.concurrent.Executors import java.util.concurrent.RejectedExecutionException +import java.util.concurrent.ThreadPoolExecutor import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicReference import kotlin.coroutines.resume import kotlin.concurrent.thread import kotlin.random.Random @@ -74,9 +76,7 @@ class TunnelClient( private val modulesClosed = AtomicBoolean(false) private val modulesCloseFinished = CompletableDeferred() private val moduleLock = Any() - private val moduleExecutor = Executors.newSingleThreadExecutor { runnable -> - Thread(runnable, "mobile-mcp-requests").apply { isDaemon = true } - } + private val activeDispatcher = AtomicReference() fun start() { if (stopped.get()) return @@ -90,7 +90,7 @@ class TunnelClient( loopJob = null webSocket?.close(NORMAL_CLOSURE, "client stopping") webSocket = null - moduleExecutor.shutdownNow() + activeDispatcher.getAndSet(null)?.close() if (modulesClosed.compareAndSet(false, true)) { // A QuickJS evaluation may hold its runtime lock until the 60-second // execution limit. Never make the service/main thread wait for it. @@ -138,6 +138,9 @@ class TunnelClient( /** Runs one WebSocket session to completion. Returns true if `server_hello` arrived. */ private suspend fun runOneConnection(): Boolean = suspendCancellableCoroutine { cont -> + val sessionActive = AtomicBoolean(true) + val dispatcher = McpRequestDispatcher() + activeDispatcher.getAndSet(dispatcher)?.close() val request = Request.Builder() .url(gatewayUrl) .header("Authorization", "Bearer $authToken") @@ -191,10 +194,17 @@ class TunnelClient( // full Bash script. Never run them on OkHttp's reader // callback: doing so prevents WebSocket control frames // and unrelated tunnel messages from being processed. - try { - moduleExecutor.execute { routeMcpFrame(webSocket, frame) } - } catch (_: RejectedExecutionException) { - // The client was stopped between parsing and queueing. + // Resolve the binding at receipt time. A later desired-state + // update must not retroactively change an earlier request. + val module = modulesByServerId[frame.serverId] ?: modulesByName[frame.serverId] + if (!dispatcher.submit { + routeMcpFrame(webSocket, frame, module, sessionActive) + } + ) { + // Backpressure is explicit: retaining an unbounded number + // of long-running requests would eventually exhaust memory. + webSocket.close(TRY_AGAIN_LATER, "MCP request queue full") + finish() } } is Ping -> send(webSocket, Pong) @@ -222,6 +232,9 @@ class TunnelClient( } fun finish() { + if (!sessionActive.compareAndSet(true, false)) return + dispatcher.close() + activeDispatcher.compareAndSet(dispatcher, null) this@TunnelClient.webSocket = null modulesByServerId.clear() if (cont.isActive) cont.resume(sawServerHello) @@ -229,7 +242,12 @@ class TunnelClient( } val socket = httpClient.newWebSocket(request, listener) - cont.invokeOnCancellation { socket.cancel() } + cont.invokeOnCancellation { + sessionActive.set(false) + dispatcher.close() + activeDispatcher.compareAndSet(dispatcher, null) + socket.cancel() + } } /** @@ -265,10 +283,13 @@ class TunnelClient( } } - private fun routeMcpFrame(webSocket: WebSocket, frame: McpFrame) { - // Accept a module addressed by bare name too, so a backend that keys - // built-ins by name (and tests) can skip the desired-state handshake. - val module = modulesByServerId[frame.serverId] ?: modulesByName[frame.serverId] + private fun routeMcpFrame( + webSocket: WebSocket, + frame: McpFrame, + module: LocalMcpModule?, + sessionActive: AtomicBoolean, + ) { + if (!sessionActive.get() || stopped.get()) return if (module == null) { send( webSocket, @@ -281,10 +302,10 @@ class TunnelClient( return } val response = synchronized(moduleLock) { - if (stopped.get()) return + if (!sessionActive.get() || stopped.get()) return module.handle(frame.frame) } ?: return - if (stopped.get()) return + if (!sessionActive.get() || stopped.get()) return send(webSocket, McpFrame(serverId = frame.serverId, frame = response)) } @@ -295,6 +316,7 @@ class TunnelClient( companion object { private const val TAG = "TunnelClient" private const val NORMAL_CLOSURE = 1000 + private const val TRY_AGAIN_LATER = 1013 private const val INITIAL_BACKOFF_MILLIS = 1_000L private const val MAX_BACKOFF_MILLIS = 60_000L @@ -307,3 +329,29 @@ class TunnelClient( .build() } } + +internal const val MCP_REQUEST_QUEUE_CAPACITY = 16 + +/** A bounded, session-scoped serial dispatcher for potentially slow module calls. */ +internal class McpRequestDispatcher { + private val executor = ThreadPoolExecutor( + 1, + 1, + 0L, + TimeUnit.MILLISECONDS, + ArrayBlockingQueue(MCP_REQUEST_QUEUE_CAPACITY), + { runnable -> Thread(runnable, "mobile-mcp-requests").apply { isDaemon = true } }, + ThreadPoolExecutor.AbortPolicy(), + ) + + fun submit(task: () -> Unit): Boolean = try { + executor.execute(task) + true + } catch (_: RejectedExecutionException) { + false + } + + fun close() { + executor.shutdownNow() + } +} diff --git a/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt b/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt index 31ed511..77e9987 100644 --- a/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt @@ -80,25 +80,73 @@ class ComputerModuleTest { assertEquals(0, second["exitCode"]!!.jsonPrimitive.content.toInt()) val firstToken = first["supplementToken"]!!.jsonPrimitive.content val secondToken = second["supplementToken"]!!.jsonPrimitive.content - assertTrue(router.takeSupplements(listOf(firstToken, secondToken)).size == 1) + val retained = router.takeSupplements(listOf(firstToken, secondToken)).single() + assertEquals("obs_2", retained.structuredContent!!["observationId"]!!.jsonPrimitive.content) router.clearSupplements() } + @Test + fun computerStatusDoesNotReplaceTheLatestObservation() { + val router = MobileCommandRouter(listOf(ComputerModule(FakeComputerSource()))) + + val observation = router.execute("computer", listOf("observe")) + val status = router.execute("computer", listOf("status")) + + assertEquals(0, observation.exitCode) + assertEquals(0, status.exitCode) + val supplements = router.takeSupplements(listOf(observation.supplementToken!!, status.supplementToken!!)) + assertEquals(2, supplements.size) + assertEquals("image", supplements.first().content.single()["type"]!!.jsonPrimitive.content) + assertTrue(supplements.last().content.isEmpty()) + } + + @Test + fun failedComputerReplacementPreservesThePreviousObservation() { + val source = FakeComputerSource().apply { screenshotData = "a".repeat(256 * 1024) } + val camera = object : CameraSource { + private val result = CameraOperationResult( + payload = buildJsonObject { put("lens", JsonPrimitive("back")) }, + photo = CameraPhoto("a".repeat(3 * 1024 * 1024), "image/jpeg"), + ) + override fun status() = result + override fun list() = result + override fun snap(options: CameraSnapOptions) = result + } + val router = MobileCommandRouter(listOf(CameraModule(camera), ComputerModule(source))) + val cameraResult = router.execute("camera", listOf("snap")) + val firstObservation = router.execute("computer", listOf("observe")) + source.screenshotData = "b".repeat(2 * 1024 * 1024) + + val failedReplacement = Json.parseToJsonElement( + router.executeJson("""{"namespace":"computer","args":["observe"]}"""), + ).jsonObject + + assertEquals(1, failedReplacement["exitCode"]!!.jsonPrimitive.content.toInt()) + val retained = router.takeSupplements( + listOf(cameraResult.supplementToken!!, firstObservation.supplementToken!!), + ) + assertEquals(2, retained.size) + assertEquals("obs_1", retained.last().structuredContent!!["observationId"]!!.jsonPrimitive.content) + } + private class FakeComputerSource : ComputerSource { var nodeId = "" var text = "" var tapDurationMillis = 0 var screenshotData = "aGVsbG8=" + var observationNumber = 0 private fun result() = ComputerOperationResult( payload = buildJsonObject { - put("observationId", JsonPrimitive("obs_1")) + put("observationId", JsonPrimitive("obs_${++observationNumber}")) put("accessibilityTree", buildJsonObject { put("nodes", kotlinx.serialization.json.buildJsonArray {}) }) }, screenshot = ComputerScreenshot(screenshotData, "image/jpeg"), ) - override fun status() = result() + override fun status() = ComputerOperationResult( + payload = buildJsonObject { put("enabled", JsonPrimitive(true)) }, + ) override fun observe() = result() override fun click(nodeId: String): ComputerOperationResult = result().also { this.nodeId = nodeId } override fun setText(nodeId: String, text: String): ComputerOperationResult = result().also { diff --git a/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt b/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt index 6ef9bdf..26e3c0d 100644 --- a/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt @@ -9,9 +9,50 @@ import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.runBlocking import kotlinx.serialization.json.JsonObject import org.junit.Assert.assertTrue +import org.junit.Assert.assertFalse import org.junit.Test class TunnelClientLifecycleTest { + @Test + fun requestDispatcherRejectsWorkWhenItsBoundedQueueIsFull() { + val dispatcher = McpRequestDispatcher() + val running = CountDownLatch(1) + val release = CountDownLatch(1) + assertTrue(dispatcher.submit { + running.countDown() + release.await() + }) + assertTrue(running.await(1, TimeUnit.SECONDS)) + repeat(MCP_REQUEST_QUEUE_CAPACITY) { assertTrue(dispatcher.submit {}) } + + assertFalse(dispatcher.submit {}) + release.countDown() + dispatcher.close() + } + + @Test + fun closingRequestDispatcherDropsQueuedSessionWork() { + val dispatcher = McpRequestDispatcher() + val running = CountDownLatch(1) + val release = CountDownLatch(1) + val queuedRan = CountDownLatch(1) + assertTrue(dispatcher.submit { + running.countDown() + try { + release.await() + } catch (_: InterruptedException) { + // The simulated in-flight module is released below. + } + }) + assertTrue(running.await(1, TimeUnit.SECONDS)) + assertTrue(dispatcher.submit { queuedRan.countDown() }) + + dispatcher.close() + release.countDown() + + assertFalse(queuedRan.await(200, TimeUnit.MILLISECONDS)) + } + @Test fun stopReturnsImmediatelyAndStopAndAwaitObservesModuleTeardown() { val closed = CountDownLatch(1) From daa569ed603ddc24b0f08c9a901b9f3bde349522 Mon Sep 17 00:00:00 2001 From: Miyamura80 Date: Mon, 7 Sep 2026 11:50:48 +0100 Subject: [PATCH 3/3] Bound repeated computer status polling --- .../ai/sealgate/stdiod/mcp/MobileCommandRouter.kt | 8 +++++--- .../ai/sealgate/stdiod/mcp/ComputerModuleTest.kt | 14 +++++++------- .../stdiod/tunnel/TunnelClientLifecycleTest.kt | 6 +++++- 3 files changed, 17 insertions(+), 11 deletions(-) diff --git a/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt b/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt index bc2af19..444abf9 100644 --- a/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt +++ b/app/src/main/java/ai/sealgate/stdiod/mcp/MobileCommandRouter.kt @@ -123,11 +123,13 @@ class MobileCommandRouter(modules: List) { .joinToString("\n") val typedContent = content.filter { it["type"]?.jsonPrimitive?.content != "text" } val structuredContent = result["structuredContent"] as? JsonObject - val supplementToken = if (typedContent.isNotEmpty() || structuredContent != null) { + // Status is already emitted as complete JSON text. Retaining its duplicate + // structured payload would make polling grow an invisible side channel. + val isComputerStatus = spec.module == ComputerModule.NAME && spec.tool == "computer_status" + val supplementToken = if (typedContent.isNotEmpty() || structuredContent != null && !isComputerStatus) { retainSupplement( MobileCommandSupplement(typedContent, structuredContent), - replacePreviousComputerObservation = - spec.module == ComputerModule.NAME && spec.tool != "computer_status", + replacePreviousComputerObservation = spec.module == ComputerModule.NAME, ) } else { null diff --git a/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt b/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt index 77e9987..0206287 100644 --- a/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/mcp/ComputerModuleTest.kt @@ -86,18 +86,18 @@ class ComputerModuleTest { } @Test - fun computerStatusDoesNotReplaceTheLatestObservation() { + fun repeatedComputerStatusDoesNotAccumulateOrReplaceTheLatestObservation() { val router = MobileCommandRouter(listOf(ComputerModule(FakeComputerSource()))) val observation = router.execute("computer", listOf("observe")) - val status = router.execute("computer", listOf("status")) + val statuses = List(65) { router.execute("computer", listOf("status")) } assertEquals(0, observation.exitCode) - assertEquals(0, status.exitCode) - val supplements = router.takeSupplements(listOf(observation.supplementToken!!, status.supplementToken!!)) - assertEquals(2, supplements.size) - assertEquals("image", supplements.first().content.single()["type"]!!.jsonPrimitive.content) - assertTrue(supplements.last().content.isEmpty()) + assertTrue(statuses.all { it.exitCode == 0 }) + assertTrue(statuses.all { it.supplementToken == null }) + val supplement = router.takeSupplements(listOf(observation.supplementToken!!)).single() + assertEquals("image", supplement.content.single()["type"]!!.jsonPrimitive.content) + assertEquals("obs_1", supplement.structuredContent!!["observationId"]!!.jsonPrimitive.content) } @Test diff --git a/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt b/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt index 26e3c0d..5ba0ab1 100644 --- a/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt +++ b/app/src/test/java/ai/sealgate/stdiod/tunnel/TunnelClientLifecycleTest.kt @@ -20,7 +20,11 @@ class TunnelClientLifecycleTest { val release = CountDownLatch(1) assertTrue(dispatcher.submit { running.countDown() - release.await() + try { + release.await() + } catch (_: InterruptedException) { + // Dispatcher shutdown interrupts the simulated in-flight work. + } }) assertTrue(running.await(1, TimeUnit.SECONDS)) repeat(MCP_REQUEST_QUEUE_CAPACITY) { assertTrue(dispatcher.submit {}) }