fix openclaw task recovery after interrupted sse

This commit is contained in:
Haitao Pan 2026-05-15 13:56:28 +08:00
parent bc468d8f2e
commit 0db5e6b467
2 changed files with 140 additions and 1 deletions

View File

@ -95,6 +95,7 @@ class ExternalCodeAgentAcpDesktopTransport
}) async {
var streamedText = '';
String? completedMessage;
Map<String, dynamic>? completedResultSnapshot;
try {
final endpointOverride = _taskEndpointResolver == null
? _endpointResolver(request.target)
@ -123,6 +124,9 @@ class ExternalCodeAgentAcpDesktopTransport
}
if (update.isDone && update.message.trim().isNotEmpty) {
completedMessage = update.message.trim();
completedResultSnapshot = _completedResultSnapshotFromUpdate(
update,
);
}
onUpdate(update);
},
@ -133,7 +137,20 @@ class ExternalCodeAgentAcpDesktopTransport
streamedText: streamedText,
completedMessage: completedMessage,
);
} on GatewayAcpException {
} on GatewayAcpException catch (error) {
if (error.code == 'ACP_HTTP_CONNECTION_CLOSED' &&
completedResultSnapshot != null) {
return goTaskServiceResultFromAcpResponse(
<String, dynamic>{
'jsonrpc': '2.0',
'id': 'recovered-from-completed-session-update',
'result': completedResultSnapshot,
},
route: request.route,
streamedText: streamedText,
completedMessage: completedMessage,
);
}
rethrow;
} on SocketException catch (error) {
final timeout = _socketExceptionLooksLikeConnectTimeout(error);
@ -183,6 +200,48 @@ class ExternalCodeAgentAcpDesktopTransport
@override
Future<void> dispose() => _client.dispose();
Map<String, dynamic>? _completedResultSnapshotFromUpdate(
GoTaskServiceUpdate update,
) {
if (!update.isDone) {
return null;
}
final payload = update.payload;
final embeddedResult = _castMap(payload['result']);
final snapshot = <String, dynamic>{...embeddedResult, ...payload};
snapshot.remove('sessionId');
snapshot.remove('threadId');
snapshot.remove('type');
snapshot.remove('event');
snapshot.remove('pending');
snapshot.remove('result');
snapshot['turnId'] = update.turnId;
snapshot['success'] = !update.error;
final text = _firstNonEmptyString(snapshot, const <String>[
'output',
'message',
'summary',
'text',
'delta',
]);
if (text.isNotEmpty) {
snapshot['output'] = text;
snapshot['message'] = text;
snapshot['summary'] = text;
}
return snapshot;
}
String _firstNonEmptyString(Map<String, dynamic> values, List<String> keys) {
for (final key in keys) {
final value = values[key]?.toString().trim() ?? '';
if (value.isNotEmpty) {
return value;
}
}
return '';
}
Map<String, dynamic> _castMap(Object? value) {
if (value is Map<String, dynamic>) {
return value;

View File

@ -532,6 +532,86 @@ void main() {
},
);
test(
'recovers OpenClaw task result from completed session update when final SSE envelope is lost',
() async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
addTearDown(() => server.close(force: true));
server.listen((request) async {
await utf8.decoder.bind(request).join();
final event = jsonEncode(<String, dynamic>{
'jsonrpc': '2.0',
'method': 'session.update',
'params': <String, dynamic>{
'sessionId': 'draft:test-task-a',
'threadId': 'draft:test-task-a',
'turnId': 'turn-1',
'type': 'status',
'event': 'completed',
'pending': false,
'error': false,
'message': 'stable completed output',
'result': <String, dynamic>{
'success': true,
'output': 'stable completed output',
'turnId': 'turn-1',
'artifacts': <Map<String, dynamic>>[
<String, dynamic>{
'relativePath': 'exports/final.md',
'downloadUrl':
'https://xworkmate-bridge.svc.plus/artifacts/openclaw/download'
'?sessionKey=draft:test-task-a&runId=turn-1&relativePath=exports%2Ffinal.md',
'contentType': 'text/markdown',
'sizeBytes': 42,
},
],
},
},
});
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();
});
final endpoint = Uri.parse('http://127.0.0.1:${server.port}');
final transport = ExternalCodeAgentAcpDesktopTransport(
client: GatewayAcpClient(endpointResolver: () => endpoint),
endpointResolver: (_) => endpoint,
);
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, 'stable completed output');
expect(result.artifacts, hasLength(1));
expect(result.artifacts.single.relativePath, 'exports/final.md');
},
);
test(
'retries interrupted TLS handshakes before surfacing ACP diagnostics',
() async {