From 0f8afd3f4fa0e07842fb74e01ae90e143e4fd37f Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Fri, 22 May 2026 18:45:23 +0800 Subject: [PATCH] fix: extend openclaw result recovery --- ...rnal_code_agent_acp_desktop_transport.dart | 73 ++++++- .../runtime/gateway_acp_client_auth_test.dart | 188 ++++++++++++++++++ 2 files changed, 253 insertions(+), 8 deletions(-) diff --git a/lib/runtime/external_code_agent_acp_desktop_transport.dart b/lib/runtime/external_code_agent_acp_desktop_transport.dart index 10041813..3cc8fce8 100644 --- a/lib/runtime/external_code_agent_acp_desktop_transport.dart +++ b/lib/runtime/external_code_agent_acp_desktop_transport.dart @@ -14,13 +14,19 @@ class ExternalCodeAgentAcpDesktopTransport required GatewayAcpClient client, required Uri? Function(AssistantExecutionTarget target) endpointResolver, Uri? Function(GoTaskServiceRequest request)? taskEndpointResolver, + Duration recoveryPollDelay = const Duration(seconds: 2), + int recoveryMaxAttempts = 300, }) : _client = client, _endpointResolver = endpointResolver, - _taskEndpointResolver = taskEndpointResolver; + _taskEndpointResolver = taskEndpointResolver, + _recoveryPollDelay = recoveryPollDelay, + _recoveryMaxAttempts = recoveryMaxAttempts; final GatewayAcpClient _client; final Uri? Function(AssistantExecutionTarget target) _endpointResolver; final Uri? Function(GoTaskServiceRequest request)? _taskEndpointResolver; + final Duration _recoveryPollDelay; + final int _recoveryMaxAttempts; @visibleForTesting GatewayAcpClient get clientForTest => _client; @@ -195,9 +201,10 @@ class ExternalCodeAgentAcpDesktopTransport if (endpoint == null) { return null; } - for (var attempt = 0; attempt < 60; attempt += 1) { + final attempts = _recoveryMaxAttempts <= 0 ? 1 : _recoveryMaxAttempts; + for (var attempt = 0; attempt < attempts; attempt += 1) { if (attempt > 0) { - await Future.delayed(const Duration(seconds: 2)); + await Future.delayed(_recoveryPollDelay); } Map response; try { @@ -220,12 +227,16 @@ class ExternalCodeAgentAcpDesktopTransport .toString() .trim() .toLowerCase(); + final terminal = + status == 'completed' || + status == 'failed' || + status == 'cancelled' || + status == 'canceled'; + if (!terminal) { + continue; + } final result = _recoveredResultFromSessionSnapshot(snapshot); - if (result.isNotEmpty && - (status == 'completed' || - status == 'failed' || - status == 'cancelled' || - status == 'canceled')) { + if (result.isNotEmpty) { return goTaskServiceResultFromAcpResponse( { 'jsonrpc': '2.0', @@ -237,6 +248,18 @@ class ExternalCodeAgentAcpDesktopTransport completedMessage: completedMessage, ); } + if (status == 'failed' || status == 'cancelled' || status == 'canceled') { + return goTaskServiceResultFromAcpResponse( + { + 'jsonrpc': '2.0', + 'id': 'recovered-from-terminal-session-snapshot', + 'result': _failureResultFromSessionSnapshot(snapshot, status), + }, + route: request.route, + streamedText: streamedText, + completedMessage: completedMessage, + ); + } } return null; } @@ -345,6 +368,40 @@ class ExternalCodeAgentAcpDesktopTransport return result; } + Map _failureResultFromSessionSnapshot( + Map snapshot, + String status, + ) { + final task = _castMap(snapshot['task']); + final error = _castMap(snapshot['error']); + final message = _firstNonEmptyDisplayText( + {...error, ...snapshot, 'taskMessage': task['message']}, + const [ + 'message', + 'error', + 'errorMessage', + 'reason', + 'taskMessage', + 'code', + ], + ); + final code = _firstNonEmptyDisplayText( + {...error, ...snapshot, 'taskCode': task['code']}, + const ['code', 'errorCode', 'taskCode'], + ); + final result = { + 'success': false, + 'status': status, + 'turnId': task['turnId']?.toString().trim() ?? '', + 'error': message.isNotEmpty ? message : 'Bridge session ended: $status', + 'message': message.isNotEmpty ? message : 'Bridge session ended: $status', + }; + if (code.isNotEmpty) { + result['code'] = code; + } + return result; + } + bool _hasArtifactList(Map result) { for (final key in const ['artifacts', 'files', 'attachments']) { if (_listValue(result[key]).isNotEmpty) { diff --git a/test/runtime/gateway_acp_client_auth_test.dart b/test/runtime/gateway_acp_client_auth_test.dart index 60ce909d..1edc4b03 100644 --- a/test/runtime/gateway_acp_client_auth_test.dart +++ b/test/runtime/gateway_acp_client_auth_test.dart @@ -784,6 +784,194 @@ void main() { }, ); + test( + 'keeps polling running OpenClaw snapshot after SSE connection close', + () async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + addTearDown(() => server.close(force: true)); + var snapshotPolls = 0; + server.listen((request) async { + final body = await utf8.decoder.bind(request).join(); + 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': 'unit-fixture-task-b'}, + }); + 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') { + snapshotPolls += 1; + final completed = snapshotPolls >= 3; + request.response.headers.contentType = ContentType.json; + request.response.write( + jsonEncode({ + 'jsonrpc': '2.0', + 'id': id, + 'result': { + 'status': completed ? 'completed' : 'running', + 'sessionId': 'unit-fixture-task-b', + 'threadId': 'unit-fixture-task-b', + 'task': { + 'state': completed ? 'completed' : 'running', + 'turnId': 'turn-recovered-running', + }, + if (completed) + 'result': { + 'success': true, + 'output': 'recovered after running snapshot', + 'turnId': 'turn-recovered-running', + }, + }, + }), + ); + 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'), + recoveryPollDelay: Duration.zero, + recoveryMaxAttempts: 4, + ); + addTearDown(transport.dispose); + + final result = await transport.executeTask( + const GoTaskServiceRequest( + sessionId: 'unit-fixture-task-b', + threadId: 'unit-fixture-task-b', + target: AssistantExecutionTarget.gateway, + provider: SingleAgentProvider.openclaw, + prompt: 'wait for result', + workingDirectory: '/tmp/workspace', + model: '', + thinking: 'off', + selectedSkills: [], + inlineAttachments: [], + localAttachments: [], + agentId: '', + metadata: {}, + ), + onUpdate: (_) {}, + ); + + expect(snapshotPolls, 3); + expect(result.success, isTrue); + expect(result.message, 'recovered after running snapshot'); + }, + ); + + test( + 'recovers terminal failed OpenClaw snapshot without displayable result', + () async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + addTearDown(() => server.close(force: true)); + server.listen((request) async { + final body = await utf8.decoder.bind(request).join(); + 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': 'unit-fixture-task-c'}, + }); + 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': 'failed', + 'sessionId': 'unit-fixture-task-c', + 'threadId': 'unit-fixture-task-c', + 'task': { + 'state': 'failed', + 'turnId': 'turn-failed', + }, + 'error': { + 'code': 'OPENCLAW_WAIT_FAILED', + 'message': 'openclaw wait failed', + }, + }, + }), + ); + 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'), + recoveryPollDelay: Duration.zero, + recoveryMaxAttempts: 1, + ); + addTearDown(transport.dispose); + + final result = await transport.executeTask( + const GoTaskServiceRequest( + sessionId: 'unit-fixture-task-c', + threadId: 'unit-fixture-task-c', + target: AssistantExecutionTarget.gateway, + provider: SingleAgentProvider.openclaw, + prompt: 'fail task', + workingDirectory: '/tmp/workspace', + model: '', + thinking: 'off', + selectedSkills: [], + inlineAttachments: [], + localAttachments: [], + agentId: '', + metadata: {}, + ), + onUpdate: (_) {}, + ); + + expect(result.success, isFalse); + expect(result.status, 'failed'); + expect(result.code, 'OPENCLAW_WAIT_FAILED'); + expect(result.errorMessage, contains('openclaw wait failed')); + }, + ); + test( 'retries interrupted TLS handshakes before surfacing ACP diagnostics', () async {