fix: recover interrupted bridge task results

This commit is contained in:
Haitao Pan 2026-05-18 18:33:09 +08:00
parent ab1a6be90f
commit 20194f1bb2
2 changed files with 196 additions and 0 deletions

View File

@ -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<GoTaskServiceResult?> _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<void>.delayed(const Duration(seconds: 2));
}
Map<String, dynamic> response;
try {
response = await _client.request(
method: 'xworkmate.sessions.get',
params: <String, dynamic>{
'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(
<String, dynamic>{
'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<void> cancelTask({
required AssistantExecutionTarget target,

View File

@ -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 = <String>[];
server.listen((request) async {
final body = await utf8.decoder.bind(request).join();
requestPaths.add(request.uri.path);
final decoded = jsonDecode(body) as Map<String, dynamic>;
final method = decoded['method']?.toString() ?? '';
final id = decoded['id']?.toString() ?? 'request-id';
if (method == 'session.start') {
final event = jsonEncode(<String, dynamic>{
'jsonrpc': '2.0',
'method': 'xworkmate.bridge.accepted',
'params': <String, dynamic>{'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(<String, dynamic>{
'jsonrpc': '2.0',
'id': id,
'result': <String, dynamic>{
'status': 'completed',
'sessionId': 'draft:test-task-a',
'threadId': 'draft:test-task-a',
'task': <String, dynamic>{
'state': 'completed',
'turnId': 'turn-recovered',
},
'result': <String, dynamic>{
'success': true,
'output': 'recovered from bridge session snapshot',
'turnId': 'turn-recovered',
'artifacts': <Map<String, dynamic>>[
<String, dynamic>{
'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: <String>[],
inlineAttachments: <GatewayChatAttachmentPayload>[],
localAttachments: <CollaborationAttachment>[],
agentId: '',
metadata: <String, dynamic>{},
),
onUpdate: (_) {},
);
expect(result.success, isTrue);
expect(result.message, 'recovered from bridge session snapshot');
expect(result.artifacts.single.relativePath, 'exports/snapshot.md');
expect(
requestPaths,
containsAll(<String>['/gateway/openclaw', '/acp/rpc']),
);
},
);
test(
'retries interrupted TLS handshakes before surfacing ACP diagnostics',
() async {