From 5fa7c86f9c265bbc4f3e93aeb19bad3c5db3e20a Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Fri, 8 May 2026 10:54:23 +0800 Subject: [PATCH] fix: stabilize ACP closed flow handling --- ...app_controller_desktop_thread_binding.dart | 36 +++++++--- lib/runtime/gateway_acp_client.dart | 15 ++-- ...troller_thread_workspace_binding_test.dart | 72 +++++++++++++++++++ .../runtime/gateway_acp_client_auth_test.dart | 61 ++++++++++++++++ 4 files changed, 166 insertions(+), 18 deletions(-) diff --git a/lib/app/app_controller_desktop_thread_binding.dart b/lib/app/app_controller_desktop_thread_binding.dart index c359a668..d6583416 100644 --- a/lib/app/app_controller_desktop_thread_binding.dart +++ b/lib/app/app_controller_desktop_thread_binding.dart @@ -95,19 +95,28 @@ extension AppControllerDesktopThreadBinding on AppController { final normalizedSessionKey = normalizedAssistantSessionKeyInternal( sessionKey, ); - final baseWorkspace = settings.workspacePath.trim().isNotEmpty - ? settings.workspacePath.trim() - : resolvedUserHomeDirectoryInternal.trim(); - if (baseWorkspace.isEmpty) { + final homeDirectory = resolvedUserHomeDirectoryInternal.trim(); + if (homeDirectory.isEmpty) { return ''; } final threadWorkspace = - '${trimTrailingPathSeparatorInternal(baseWorkspace)}/.xworkmate/threads/${threadWorkspaceDirectoryNameInternal(normalizedSessionKey)}'; + '${trimTrailingPathSeparatorInternal(homeDirectory)}/.xworkmate/threads/${threadWorkspaceDirectoryNameInternal(normalizedSessionKey)}'; return ensureLocalWorkspaceDirectoryInternal(threadWorkspace) ? threadWorkspace : ''; } + String localThreadWorkspaceDisplayPathInternal(String sessionKey) { + final homeDirectory = resolvedUserHomeDirectoryInternal.trim(); + if (homeDirectory.isEmpty) { + return ''; + } + final normalizedSessionKey = normalizedAssistantSessionKeyInternal( + sessionKey, + ); + return '\$HOME/.xworkmate/threads/${threadWorkspaceDirectoryNameInternal(normalizedSessionKey)}'; + } + String remoteThreadWorkspacePathInternal( String sessionKey, ThreadOwnerScope ownerScope, @@ -190,11 +199,14 @@ extension AppControllerDesktopThreadBinding on AppController { WorkspaceBinding? existingBinding, }) { final localPath = localThreadWorkspacePathInternal(sessionKey); + final displayPath = localPath.isEmpty + ? '' + : localThreadWorkspaceDisplayPathInternal(sessionKey); return WorkspaceBinding( workspaceId: normalizedAssistantSessionKeyInternal(sessionKey), workspaceKind: WorkspaceKind.localFs, workspacePath: localPath, - displayPath: localPath, + displayPath: displayPath, writable: existingBinding?.writable ?? true, ); } @@ -277,11 +289,13 @@ extension AppControllerDesktopThreadBinding on AppController { normalizedSessionKey, ownerScope: ownerScope, workspaceBinding: workspaceBinding, - lastRemoteWorkingDirectory: remoteThreadWorkspacePathInternal( - normalizedSessionKey, - ownerScope, - ), - lastRemoteWorkspaceRefKind: WorkspaceRefKind.remotePath, + lastRemoteWorkingDirectory: + snapshot.record?.lastRemoteWorkingDirectory?.trim().isNotEmpty == true + ? snapshot.record?.lastRemoteWorkingDirectory + : remoteThreadWorkspacePathInternal(normalizedSessionKey, ownerScope), + lastRemoteWorkspaceRefKind: + snapshot.record?.lastRemoteWorkspaceRefKind ?? + WorkspaceRefKind.remotePath, executionBinding: buildDesktopExecutionBindingInternal( executionTarget: snapshot.executionTarget, existingBinding: snapshot.record?.executionBinding, diff --git a/lib/runtime/gateway_acp_client.dart b/lib/runtime/gateway_acp_client.dart index 79dfcc2c..83d2327c 100644 --- a/lib/runtime/gateway_acp_client.dart +++ b/lib/runtime/gateway_acp_client.dart @@ -863,7 +863,7 @@ class GatewayAcpClient { required String requestId, required void Function(Map) onNotification, }) async { - final completer = Completer>(); + Map? resolvedResponse; final eventLines = []; void consumeEventPayload(String payload) { @@ -874,9 +874,7 @@ class GatewayAcpClient { final json = _decodeMap(trimmed); if (stringValue(json['id']) == requestId && (json.containsKey('result') || json.containsKey('error'))) { - if (!completer.isCompleted) { - completer.complete(json); - } + resolvedResponse ??= json; return; } if ((stringValue(json['method']) ?? '').isNotEmpty) { @@ -890,6 +888,9 @@ class GatewayAcpClient { if (eventLines.isNotEmpty) { consumeEventPayload(eventLines.join('\n')); eventLines.clear(); + if (resolvedResponse != null) { + break; + } } continue; } @@ -898,16 +899,16 @@ class GatewayAcpClient { } } - if (eventLines.isNotEmpty) { + if (eventLines.isNotEmpty && resolvedResponse == null) { consumeEventPayload(eventLines.join('\n')); } - if (!completer.isCompleted) { + final resolved = resolvedResponse; + if (resolved == null) { throw GatewayAcpException( 'ACP SSE ended without JSON-RPC response for request $requestId', code: 'ACP_SSE_NO_RESULT', ); } - final resolved = await completer.future; _throwIfJsonRpcError(resolved); return resolved; } diff --git a/test/runtime/app_controller_thread_workspace_binding_test.dart b/test/runtime/app_controller_thread_workspace_binding_test.dart index 235a2958..1f569b79 100644 --- a/test/runtime/app_controller_thread_workspace_binding_test.dart +++ b/test/runtime/app_controller_thread_workspace_binding_test.dart @@ -4,10 +4,82 @@ import 'package:crypto/crypto.dart' as crypto; import 'package:flutter_test/flutter_test.dart'; import 'package:xworkmate/app/app_controller.dart'; import 'package:xworkmate/app/app_controller_desktop_runtime_coordination_impl.dart'; +import 'package:xworkmate/app/app_controller_desktop_thread_binding.dart'; import 'package:xworkmate/runtime/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; void main() { + test( + 'converges managed local thread workspaces to the user home root', + () async { + final controller = AppController( + environmentOverride: const {}, + ); + addTearDown(controller.dispose); + + final home = await Directory.systemTemp.createTemp( + 'xworkmate-home-thread-root-', + ); + final oldWorkspace = await Directory.systemTemp.createTemp( + 'xworkmate-app-worktree-thread-root-', + ); + addTearDown(() async { + if (await home.exists()) { + await home.delete(recursive: true); + } + if (await oldWorkspace.exists()) { + await oldWorkspace.delete(recursive: true); + } + }); + controller.resolvedUserHomeDirectoryInternal = home.path; + + const sessionKey = 'draft-1778207741322'; + final oldThreadWorkspace = Directory( + '${oldWorkspace.path}/.xworkmate/threads/$sessionKey', + ); + await oldThreadWorkspace.create(recursive: true); + + controller.upsertTaskThreadInternal( + sessionKey, + workspaceBinding: WorkspaceBinding( + workspaceId: sessionKey, + workspaceKind: WorkspaceKind.localFs, + workspacePath: oldThreadWorkspace.path, + displayPath: oldThreadWorkspace.path, + writable: true, + ), + messages: const [ + GatewayChatMessage( + id: 'assistant-1', + role: 'assistant', + text: 'kept message', + timestampMs: 1, + toolCallId: null, + toolName: null, + stopReason: null, + pending: false, + error: false, + ), + ], + lastRemoteWorkingDirectory: '/remote/thread/workspace', + lastRemoteWorkspaceRefKind: WorkspaceRefKind.remotePath, + ); + + await controller.ensureDesktopTaskThreadBindingInternal(sessionKey); + + final expectedWorkspace = '${home.path}/.xworkmate/threads/$sessionKey'; + final thread = controller.requireTaskThreadForSessionInternal(sessionKey); + expect(thread.workspaceBinding.workspacePath, expectedWorkspace); + expect( + thread.workspaceBinding.displayPath, + '\$HOME/.xworkmate/threads/$sessionKey', + ); + expect(Directory(expectedWorkspace).existsSync(), isTrue); + expect(thread.lastRemoteWorkingDirectory, '/remote/thread/workspace'); + expect(thread.messages.single.text, 'kept message'); + }, + ); + test( 'keeps local workspace binding separate from remote execution workspace', () { diff --git a/test/runtime/gateway_acp_client_auth_test.dart b/test/runtime/gateway_acp_client_auth_test.dart index b3e2e66e..29a8eb10 100644 --- a/test/runtime/gateway_acp_client_auth_test.dart +++ b/test/runtime/gateway_acp_client_auth_test.dart @@ -249,6 +249,67 @@ void main() { expect((response['result'] as Map)['ok'], true); }); + test( + 'returns SSE final response before a truncated chunked close is reported', + () async { + final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0); + addTearDown(() => server.close()); + server.listen((socket) async { + final requestBytes = []; + var headerEnd = -1; + await for (final chunk in socket) { + requestBytes.addAll(chunk); + final raw = utf8.decode(requestBytes, allowMalformed: true); + headerEnd = raw.indexOf('\r\n\r\n'); + if (headerEnd < 0) { + continue; + } + if (raw.contains('"id"') && raw.contains('"method"')) { + break; + } + } + final rawRequest = utf8.decode(requestBytes, allowMalformed: true); + final id = + RegExp( + r'"id"\s*:\s*"([^"]+)"', + ).firstMatch(rawRequest)?.group(1) ?? + 'request-id'; + final event = utf8.encode( + 'data: ${jsonEncode({ + 'jsonrpc': '2.0', + 'id': id, + 'result': {'ok': true}, + })}\n\n', + ); + socket + ..add( + ascii.encode( + 'HTTP/1.1 200 OK\r\n' + 'Content-Type: text/event-stream\r\n' + 'Transfer-Encoding: chunked\r\n' + 'Connection: keep-alive\r\n' + '\r\n' + '${event.length.toRadixString(16)}\r\n', + ), + ) + ..add(event) + ..add(ascii.encode('\r\n')); + await socket.flush(); + socket.destroy(); + }); + + final endpoint = Uri.parse('http://127.0.0.1:${server.port}'); + final client = GatewayAcpClient(endpointResolver: () => endpoint); + + final response = await client.request( + method: 'acp.capabilities', + params: const {}, + ); + + expect((response['result'] as Map)['ok'], true); + }, + ); + test( 'normalizes raw authorization override into bearer header for HTTP', () async {