From 20194f1bb254db0c62448bdb2de28b321dae032f Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Mon, 18 May 2026 18:33:09 +0800 Subject: [PATCH] fix: recover interrupted bridge task results --- ...rnal_code_agent_acp_desktop_transport.dart | 90 +++++++++++++++ .../runtime/gateway_acp_client_auth_test.dart | 106 ++++++++++++++++++ 2 files changed, 196 insertions(+) diff --git a/lib/runtime/external_code_agent_acp_desktop_transport.dart b/lib/runtime/external_code_agent_acp_desktop_transport.dart index cf087de6..019ca26c 100644 --- a/lib/runtime/external_code_agent_acp_desktop_transport.dart +++ b/lib/runtime/external_code_agent_acp_desktop_transport.dart @@ -3,6 +3,7 @@ import 'dart:io'; import 'package:flutter/foundation.dart'; +import 'acp_endpoint_paths.dart'; import 'gateway_acp_client.dart'; import 'go_task_service_client.dart'; import 'runtime_models.dart'; @@ -151,6 +152,19 @@ class ExternalCodeAgentAcpDesktopTransport completedMessage: completedMessage, ); } + if (error.code == 'ACP_HTTP_CONNECTION_CLOSED') { + final recovered = await _recoverTaskResultAfterConnectionClosed( + request, + taskEndpoint: _taskEndpointResolver == null + ? _endpointResolver(request.target) + : _taskEndpointResolver.call(request), + streamedText: streamedText, + completedMessage: completedMessage, + ); + if (recovered != null) { + return recovered; + } + } rethrow; } on SocketException catch (error) { final timeout = _socketExceptionLooksLikeConnectTimeout(error); @@ -171,6 +185,82 @@ class ExternalCodeAgentAcpDesktopTransport } } + Future _recoverTaskResultAfterConnectionClosed( + GoTaskServiceRequest request, { + required Uri? taskEndpoint, + required String streamedText, + required String? completedMessage, + }) async { + final endpoint = _sessionSnapshotEndpoint(taskEndpoint); + if (endpoint == null) { + return null; + } + for (var attempt = 0; attempt < 60; attempt += 1) { + if (attempt > 0) { + await Future.delayed(const Duration(seconds: 2)); + } + Map response; + try { + response = await _client.request( + method: 'xworkmate.sessions.get', + params: { + 'sessionId': request.sessionId, + 'threadId': request.threadId, + }, + endpointOverride: endpoint, + ); + } on GatewayAcpException { + continue; + } on SocketException { + continue; + } + final snapshot = _castMap(response['result']); + final task = _castMap(snapshot['task']); + final status = (task['state'] ?? snapshot['status'] ?? '') + .toString() + .trim() + .toLowerCase(); + final result = _castMap(snapshot['result']); + if (result.isNotEmpty && + (status == 'completed' || + status == 'failed' || + status == 'cancelled' || + status == 'canceled')) { + return goTaskServiceResultFromAcpResponse( + { + 'jsonrpc': '2.0', + 'id': 'recovered-from-session-snapshot', + 'result': result, + }, + route: request.route, + streamedText: streamedText, + completedMessage: completedMessage, + ); + } + } + return null; + } + + Uri? _sessionSnapshotEndpoint(Uri? taskEndpoint) { + final controlEndpoint = resolveAcpHttpRpcEndpoint( + _endpointResolver(AssistantExecutionTarget.gateway), + ); + if (controlEndpoint != null) { + return controlEndpoint; + } + final taskPath = taskEndpoint?.path.trim() ?? ''; + if (taskEndpoint != null && + (taskPath == '/gateway/openclaw' || + taskPath.endsWith('/gateway/openclaw'))) { + return taskEndpoint.replace( + path: '/acp/rpc', + query: null, + fragment: null, + ); + } + return resolveAcpHttpRpcEndpoint(taskEndpoint); + } + @override Future cancelTask({ required AssistantExecutionTarget target, diff --git a/test/runtime/gateway_acp_client_auth_test.dart b/test/runtime/gateway_acp_client_auth_test.dart index 294b262e..e1f51677 100644 --- a/test/runtime/gateway_acp_client_auth_test.dart +++ b/test/runtime/gateway_acp_client_auth_test.dart @@ -644,6 +644,112 @@ void main() { }, ); + test( + 'recovers OpenClaw task result from bridge session snapshot after SSE connection close', + () async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + addTearDown(() => server.close(force: true)); + final requestPaths = []; + server.listen((request) async { + final body = await utf8.decoder.bind(request).join(); + requestPaths.add(request.uri.path); + final decoded = jsonDecode(body) as Map; + final method = decoded['method']?.toString() ?? ''; + final id = decoded['id']?.toString() ?? 'request-id'; + if (method == 'session.start') { + final event = jsonEncode({ + 'jsonrpc': '2.0', + 'method': 'xworkmate.bridge.accepted', + 'params': {'sessionId': 'draft:test-task-a'}, + }); + final eventBytes = utf8.encode('data: $event\n\n'); + request.response.headers.set( + HttpHeaders.contentTypeHeader, + 'text/event-stream', + ); + request.response.contentLength = eventBytes.length + 128; + final socket = await request.response.detachSocket(); + socket.add(eventBytes); + await socket.flush(); + socket.destroy(); + return; + } + if (method == 'xworkmate.sessions.get') { + request.response.headers.contentType = ContentType.json; + request.response.write( + jsonEncode({ + 'jsonrpc': '2.0', + 'id': id, + 'result': { + 'status': 'completed', + 'sessionId': 'draft:test-task-a', + 'threadId': 'draft:test-task-a', + 'task': { + 'state': 'completed', + 'turnId': 'turn-recovered', + }, + 'result': { + 'success': true, + 'output': 'recovered from bridge session snapshot', + 'turnId': 'turn-recovered', + 'artifacts': >[ + { + 'relativePath': 'exports/snapshot.md', + 'downloadUrl': + 'https://xworkmate-bridge.svc.plus/artifacts/openclaw/download' + '?sessionKey=draft:test-task-a&runId=turn-recovered&relativePath=exports%2Fsnapshot.md', + 'contentType': 'text/markdown', + 'sizeBytes': 64, + }, + ], + }, + }, + }), + ); + await request.response.close(); + return; + } + request.response.statusCode = HttpStatus.badRequest; + await request.response.close(); + }); + final endpoint = Uri.parse('http://127.0.0.1:${server.port}'); + final transport = ExternalCodeAgentAcpDesktopTransport( + client: GatewayAcpClient(endpointResolver: () => endpoint), + endpointResolver: (_) => endpoint, + taskEndpointResolver: (_) => + endpoint.replace(path: '/gateway/openclaw'), + ); + addTearDown(transport.dispose); + + final result = await transport.executeTask( + const GoTaskServiceRequest( + sessionId: 'draft:test-task-a', + threadId: 'draft:test-task-a', + target: AssistantExecutionTarget.gateway, + provider: SingleAgentProvider.openclaw, + prompt: 'create files', + workingDirectory: '/tmp/workspace', + model: '', + thinking: 'off', + selectedSkills: [], + inlineAttachments: [], + localAttachments: [], + agentId: '', + metadata: {}, + ), + onUpdate: (_) {}, + ); + + expect(result.success, isTrue); + expect(result.message, 'recovered from bridge session snapshot'); + expect(result.artifacts.single.relativePath, 'exports/snapshot.md'); + expect( + requestPaths, + containsAll(['/gateway/openclaw', '/acp/rpc']), + ); + }, + ); + test( 'retries interrupted TLS handshakes before surfacing ACP diagnostics', () async {