diff --git a/lib/app/app_controller_desktop_runtime_helpers.dart b/lib/app/app_controller_desktop_runtime_helpers.dart index 7c8d5407..579d4736 100644 --- a/lib/app/app_controller_desktop_runtime_helpers.dart +++ b/lib/app/app_controller_desktop_runtime_helpers.dart @@ -230,12 +230,21 @@ extension AppControllerDesktopRuntimeHelpers on AppController { final lowered = raw.toLowerCase(); final detailCode = gatewayExecutionDetailCodeInternal(error); final primaryCode = gatewayExecutionPrimaryCodeInternal(error); - if (isAcpHttpConnectionClosedErrorInternal(error)) { + final recoverableTransportCode = recoverableAcpHttpTransportCodeInternal( + error, + ); + if (recoverableTransportCode == 'ACP_HTTP_CONNECTION_CLOSED') { return appText( 'Bridge 响应读取中断;当前对话已保留,下一次发送会继续同一会话。错误码:ACP_HTTP_CONNECTION_CLOSED', 'Bridge response was interrupted; this conversation was kept, and the next send will continue the same session. Error code: ACP_HTTP_CONNECTION_CLOSED', ); } + if (recoverableTransportCode == 'ACP_HTTP_HANDSHAKE_INTERRUPTED') { + return appText( + 'Bridge 握手中断;当前对话已保留,下一次发送会继续同一会话。错误码:ACP_HTTP_HANDSHAKE_INTERRUPTED', + 'Bridge handshake was interrupted; this conversation was kept, and the next send will continue the same session. Error code: ACP_HTTP_HANDSHAKE_INTERRUPTED', + ); + } final continuationUnavailable = primaryCode == 'SESSION_CONTINUATION_UNAVAILABLE' || detailCode == 'SESSION_CONTINUATION_UNAVAILABLE' || @@ -295,12 +304,25 @@ extension AppControllerDesktopRuntimeHelpers on AppController { } bool isAcpHttpConnectionClosedErrorInternal(Object error) { + return recoverableAcpHttpTransportCodeInternal(error) == + 'ACP_HTTP_CONNECTION_CLOSED'; + } + + String? recoverableAcpHttpTransportCodeInternal(Object error) { final raw = error.toString().trim(); final primaryCode = gatewayExecutionPrimaryCodeInternal(error); final detailCode = gatewayExecutionDetailCodeInternal(error); - return primaryCode == 'ACP_HTTP_CONNECTION_CLOSED' || + if (primaryCode == 'ACP_HTTP_CONNECTION_CLOSED' || detailCode == 'ACP_HTTP_CONNECTION_CLOSED' || - raw.contains('ACP_HTTP_CONNECTION_CLOSED'); + raw.contains('ACP_HTTP_CONNECTION_CLOSED')) { + return 'ACP_HTTP_CONNECTION_CLOSED'; + } + if (primaryCode == 'ACP_HTTP_HANDSHAKE_INTERRUPTED' || + detailCode == 'ACP_HTTP_HANDSHAKE_INTERRUPTED' || + raw.contains('ACP_HTTP_HANDSHAKE_INTERRUPTED')) { + return 'ACP_HTTP_HANDSHAKE_INTERRUPTED'; + } + return null; } String formatAiGatewayHttpErrorInternal(int statusCode, String detail) { diff --git a/lib/app/app_controller_desktop_thread_actions.dart b/lib/app/app_controller_desktop_thread_actions.dart index 8a50c57c..dc587bc6 100644 --- a/lib/app/app_controller_desktop_thread_actions.dart +++ b/lib/app/app_controller_desktop_thread_actions.dart @@ -470,17 +470,20 @@ extension AppControllerDesktopThreadActions on AppController { ); } catch (error) { clearAiGatewayStreamingTextInternal(sessionKey); - final connectionClosed = isAcpHttpConnectionClosedErrorInternal( - error, - ); + final recoverableTransportCode = + recoverableAcpHttpTransportCodeInternal(error); + final recoverableTransportInterrupted = + recoverableTransportCode != null; upsertTaskThreadInternal( sessionKey, - lifecycleStatus: connectionClosed ? 'interrupted' : 'ready', + lifecycleStatus: recoverableTransportInterrupted + ? 'interrupted' + : 'ready', lastRunAtMs: DateTime.now().millisecondsSinceEpoch.toDouble(), - lastResultCode: connectionClosed - ? 'ACP_HTTP_CONNECTION_CLOSED' - : 'error', - lastArtifactSyncStatus: connectionClosed ? 'interrupted' : null, + lastResultCode: recoverableTransportCode ?? 'error', + lastArtifactSyncStatus: recoverableTransportInterrupted + ? 'interrupted' + : null, updatedAtMs: DateTime.now().millisecondsSinceEpoch.toDouble(), ); appendLocalSessionMessageInternal( diff --git a/lib/runtime/gateway_acp_client.dart b/lib/runtime/gateway_acp_client.dart index 2077f7cf..f55bf8c1 100644 --- a/lib/runtime/gateway_acp_client.dart +++ b/lib/runtime/gateway_acp_client.dart @@ -5,6 +5,10 @@ import 'dart:io'; import 'acp_endpoint_paths.dart'; import 'runtime_models.dart'; +const int gatewayAcpHttpHandshakeInterruptedRetryCount = 5; +const String gatewayAcpHttpHandshakeInterruptedCode = + 'ACP_HTTP_HANDSHAKE_INTERRUPTED'; + class GatewayAcpException implements Exception { const GatewayAcpException( this.message, { @@ -560,6 +564,43 @@ class GatewayAcpClient { ); } + GatewayAcpException? lastHandshakeError; + for ( + var attempt = 0; + attempt <= gatewayAcpHttpHandshakeInterruptedRetryCount; + attempt += 1 + ) { + try { + return await _requestViaHttpAttempt( + request, + endpoint: endpoint, + onNotification: onNotification, + authorizationOverride: authorizationOverride, + retryAttempt: attempt, + ); + } on GatewayAcpException catch (error) { + if (error.code != gatewayAcpHttpHandshakeInterruptedCode || + attempt == gatewayAcpHttpHandshakeInterruptedRetryCount) { + rethrow; + } + lastHandshakeError = error; + await Future.delayed(Duration(milliseconds: 50 * (attempt + 1))); + } + } + throw lastHandshakeError ?? + const GatewayAcpException( + 'ACP HTTP handshake was interrupted before the response started', + code: gatewayAcpHttpHandshakeInterruptedCode, + ); + } + + Future> _requestViaHttpAttempt( + _GatewayAcpRpcRequest request, { + required Uri endpoint, + required void Function(Map) onNotification, + required String authorizationOverride, + required int retryAttempt, + }) async { final client = HttpClient()..connectionTimeout = const Duration(seconds: 8); var statusCode = 0; var contentType = ''; @@ -652,6 +693,32 @@ class GatewayAcpClient { }; } on GatewayAcpException { rethrow; + } on HandshakeException catch (error) { + throw _handshakeInterruptedException( + endpoint: endpoint, + statusCode: statusCode, + contentType: contentType, + bodyRead: bodyRead, + retryAttempt: retryAttempt, + originalError: error, + ); + } on SocketException catch (error) { + if (_looksLikeHandshakeInterruptedSocketError( + error.toString(), + endpoint: endpoint, + statusCode: statusCode, + bodyRead: bodyRead, + )) { + throw _handshakeInterruptedException( + endpoint: endpoint, + statusCode: statusCode, + contentType: contentType, + bodyRead: bodyRead, + retryAttempt: retryAttempt, + originalError: error, + ); + } + rethrow; } on HttpException catch (error) { if (_looksLikeConnectionClosedBeforeResponse(error.toString())) { throw GatewayAcpException( @@ -672,6 +739,44 @@ class GatewayAcpClient { } } + GatewayAcpException _handshakeInterruptedException({ + required Uri endpoint, + required int statusCode, + required String contentType, + required bool bodyRead, + required int retryAttempt, + required Object originalError, + }) { + return GatewayAcpException( + 'ACP HTTP handshake was interrupted before the response started', + code: gatewayAcpHttpHandshakeInterruptedCode, + details: { + 'requestUrl': endpoint.toString(), + 'statusCode': statusCode, + 'contentType': contentType, + 'bodyRead': bodyRead, + 'retryAttempt': retryAttempt, + 'maxRetryAttempts': gatewayAcpHttpHandshakeInterruptedRetryCount, + 'originalError': originalError.toString(), + }, + ); + } + + bool _looksLikeHandshakeInterruptedSocketError( + String raw, { + required Uri endpoint, + required int statusCode, + required bool bodyRead, + }) { + if (endpoint.scheme != 'https' || statusCode != 0 || bodyRead) { + return false; + } + final lowered = raw.toLowerCase(); + return lowered.contains('connection reset') || + lowered.contains('read failed') || + lowered.contains('connection terminated during handshake'); + } + bool _looksLikeConnectionClosedBeforeResponse(String raw) { final lowered = raw.toLowerCase(); return lowered.contains('connection closed before full header') || diff --git a/lib/widgets/assistant_task_progress_bar.dart b/lib/widgets/assistant_task_progress_bar.dart index 9bf47bc2..2bc55322 100644 --- a/lib/widgets/assistant_task_progress_bar.dart +++ b/lib/widgets/assistant_task_progress_bar.dart @@ -125,17 +125,25 @@ AssistantTaskProgressState assistantTaskProgressState({ ); } final result = lastResultCode.trim().toUpperCase(); - if (status == 'interrupted' || - syncStatus == 'interrupted' || - result == 'ACP_HTTP_CONNECTION_CLOSED') { + if (status == 'interrupted' || syncStatus == 'interrupted') { return AssistantTaskProgressState( phase: AssistantTaskProgressPhase.interrupted, - label: appText( - 'Bridge 响应中断,等待下一次发送续写同一会话。', - 'Bridge response interrupted; the next send will continue this session.', - ), + label: _interruptedTaskProgressLabel(result), value: 0.48, ); } return const AssistantTaskProgressState.idle(); } + +String _interruptedTaskProgressLabel(String result) { + if (result == 'ACP_HTTP_HANDSHAKE_INTERRUPTED') { + return appText( + 'Bridge 握手中断,等待下一次发送续写同一会话。', + 'Bridge handshake interrupted; the next send will continue this session.', + ); + } + return appText( + 'Bridge 响应中断,等待下一次发送续写同一会话。', + 'Bridge response interrupted; the next send will continue this session.', + ); +} diff --git a/test/features/assistant/assistant_task_progress_bar_test.dart b/test/features/assistant/assistant_task_progress_bar_test.dart index 8618f8db..7e7ee3ea 100644 --- a/test/features/assistant/assistant_task_progress_bar_test.dart +++ b/test/features/assistant/assistant_task_progress_bar_test.dart @@ -114,6 +114,27 @@ void main() { expect(indicator.value, 0.48); }); + testWidgets('shows interrupted state after ACP handshake interruption', ( + tester, + ) async { + await tester.pumpWidget( + _buildTestApp( + assistantTaskProgressState( + pending: false, + lifecycleStatus: 'interrupted', + lastResultCode: 'ACP_HTTP_HANDSHAKE_INTERRUPTED', + artifactSyncStatus: 'interrupted', + ), + ), + ); + + expect(find.text('Bridge 握手中断,等待下一次发送续写同一会话。'), findsOneWidget); + final indicator = tester.widget( + find.byKey(const Key('assistant-task-progress-indicator')), + ); + expect(indicator.value, 0.48); + }); + testWidgets('hides idle progress state', (tester) async { await tester.pumpWidget( _buildTestApp(const AssistantTaskProgressState.idle()), diff --git a/test/runtime/assistant_connection_state_test.dart b/test/runtime/assistant_connection_state_test.dart index 5ad2a1a7..cc03e6e5 100644 --- a/test/runtime/assistant_connection_state_test.dart +++ b/test/runtime/assistant_connection_state_test.dart @@ -182,6 +182,32 @@ void main() { }, ); + test( + 'labels interrupted ACP HTTP handshakes as recoverable conversation state', + () async { + final controller = await _isolatedController(); + addTearDown(controller.dispose); + + final label = controller.gatewayExecutionErrorLabelInternal( + const GatewayAcpException( + 'ACP HTTP handshake was interrupted before the response started', + code: gatewayAcpHttpHandshakeInterruptedCode, + ), + target: AssistantExecutionTarget.gateway, + ); + + expect( + label, + 'Bridge 握手中断;当前对话已保留,下一次发送会继续同一会话。错误码:ACP_HTTP_HANDSHAKE_INTERRUPTED', + ); + expect( + label, + isNot(contains('Connection terminated during handshake')), + ); + expect(label, isNot(contains('handshake was interrupted'))); + }, + ); + test( 'labels unavailable session continuation without starting a new flow', () async { diff --git a/test/runtime/assistant_execution_target_test.dart b/test/runtime/assistant_execution_target_test.dart index e31326fd..73e25d6a 100644 --- a/test/runtime/assistant_execution_target_test.dart +++ b/test/runtime/assistant_execution_target_test.dart @@ -725,6 +725,99 @@ void main() { }, ); + test( + 'sendChatMessage continues the same session after ACP HTTP handshake interruption', + () async { + final localWorkspace = await Directory.systemTemp.createTemp( + 'xworkmate-acp-handshake-interrupt-artifacts-', + ); + addTearDown(() async { + if (await localWorkspace.exists()) { + await localWorkspace.delete(recursive: true); + } + }); + final fakeGoTaskService = _RecordingGoTaskServiceClient() + ..updatesBeforeNextOutcome.add( + const GoTaskServiceUpdate( + sessionId: 'session-1', + threadId: 'session-1', + turnId: 'turn-1', + type: 'delta', + text: 'handshake partial output must not persist', + message: '', + pending: true, + error: false, + route: GoTaskServiceRoute.externalAcpSingle, + payload: {}, + ), + ) + ..outcomes.add( + const GatewayAcpException( + 'ACP HTTP handshake was interrupted before the response started', + code: gatewayAcpHttpHandshakeInterruptedCode, + ), + ) + ..outcomes.add( + GoTaskServiceResult( + success: true, + message: '全部 6 个文件已生成 ✅', + turnId: 'turn-2', + raw: {'artifacts': _generatedArtifactPayloads()}, + errorMessage: '', + resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, + ), + ); + final controller = _connectedController(fakeGoTaskService); + addTearDown(controller.dispose); + controller.resolvedUserHomeDirectoryInternal = localWorkspace.path; + + await controller.sessionsController.switchSession('session-1'); + + await controller.sendChatMessage('first turn'); + + expect(fakeGoTaskService.requests, hasLength(1)); + expect(fakeGoTaskService.requests.single.resumeSession, isFalse); + final interruptedThread = controller.taskThreadForSessionInternal( + 'session-1', + ); + expect(interruptedThread?.lifecycleState.status, 'interrupted'); + expect( + interruptedThread?.lifecycleState.lastResultCode, + gatewayAcpHttpHandshakeInterruptedCode, + ); + expect( + controller.chatMessages.last.text, + 'Bridge 握手中断;当前对话已保留,下一次发送会继续同一会话。错误码:ACP_HTTP_HANDSHAKE_INTERRUPTED', + ); + expect( + controller.chatMessages.map((message) => message.text), + isNot(contains('handshake partial output must not persist')), + ); + + await controller.sendChatMessage('follow up'); + + expect(fakeGoTaskService.requests, hasLength(2)); + expect(fakeGoTaskService.requests.last.resumeSession, isTrue); + expect(controller.chatMessages.last.text, '全部 6 个文件已生成 ✅'); + final thread = controller.taskThreadForSessionInternal('session-1'); + expect(thread?.lifecycleState.status, 'ready'); + expect(thread?.lastArtifactSyncStatus, 'synced'); + expect(thread?.lastArtifactSyncAtMs, greaterThan(0)); + final workspacePath = controller.assistantWorkspacePathForSession( + 'session-1', + ); + for (final artifact in _generatedArtifactPayloads()) { + final relativePath = artifact['relativePath']! as String; + final content = artifact['content']! as String; + expect( + await File('$workspacePath/$relativePath').readAsString(), + content, + ); + } + }, + ); + test( 'chatMessages does not duplicate persisted local turn messages', () async { @@ -797,7 +890,7 @@ void main() { final controller = _connectedController(fakeGoTaskService); addTearDown(controller.dispose); - await controller.sessionsController.switchSession('task-a'); + await controller.switchSession('task-a'); final taskAFuture = controller.sendChatMessage('task A'); await fakeGoTaskService.waitForRequestCount(1); expect(fakeGoTaskService.requests.single.sessionId, 'task-a'); diff --git a/test/runtime/gateway_acp_client_auth_test.dart b/test/runtime/gateway_acp_client_auth_test.dart index 8e07a82c..695272dc 100644 --- a/test/runtime/gateway_acp_client_auth_test.dart +++ b/test/runtime/gateway_acp_client_auth_test.dart @@ -476,6 +476,60 @@ void main() { }, ); + test( + 'retries interrupted TLS handshakes before surfacing ACP diagnostics', + () async { + final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0); + var acceptedSockets = 0; + addTearDown(() => server.close()); + server.listen((socket) { + acceptedSockets += 1; + socket.destroy(); + }); + final endpoint = Uri.parse('https://127.0.0.1:${server.port}'); + final client = GatewayAcpClient(endpointResolver: () => endpoint); + + await expectLater( + client.request( + method: 'session.start', + params: const {}, + ), + throwsA( + isA() + .having( + (error) => error.code, + 'code', + gatewayAcpHttpHandshakeInterruptedCode, + ) + .having( + (error) => error.message, + 'message', + contains('handshake was interrupted'), + ) + .having( + (error) => error.details, + 'details', + allOf( + containsPair('requestUrl', '$endpoint/acp/rpc'), + containsPair( + 'maxRetryAttempts', + gatewayAcpHttpHandshakeInterruptedRetryCount, + ), + containsPair( + 'retryAttempt', + gatewayAcpHttpHandshakeInterruptedRetryCount, + ), + ), + ), + ), + ); + expect( + acceptedSockets, + gatewayAcpHttpHandshakeInterruptedRetryCount + 1, + ); + }, + ); + test('desktop bridge auth resolver skips unrelated endpoints', () async { final storeRoot = await Directory.systemTemp.createTemp( 'xworkmate-acp-auth-unrelated-',