fix: extend openclaw result recovery

This commit is contained in:
Haitao Pan 2026-05-22 18:45:23 +08:00
parent bd6e21e265
commit 0f8afd3f4f
2 changed files with 253 additions and 8 deletions

View File

@ -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<void>.delayed(const Duration(seconds: 2));
await Future<void>.delayed(_recoveryPollDelay);
}
Map<String, dynamic> 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(
<String, dynamic>{
'jsonrpc': '2.0',
@ -237,6 +248,18 @@ class ExternalCodeAgentAcpDesktopTransport
completedMessage: completedMessage,
);
}
if (status == 'failed' || status == 'cancelled' || status == 'canceled') {
return goTaskServiceResultFromAcpResponse(
<String, dynamic>{
'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<String, dynamic> _failureResultFromSessionSnapshot(
Map<String, dynamic> snapshot,
String status,
) {
final task = _castMap(snapshot['task']);
final error = _castMap(snapshot['error']);
final message = _firstNonEmptyDisplayText(
<String, dynamic>{...error, ...snapshot, 'taskMessage': task['message']},
const <String>[
'message',
'error',
'errorMessage',
'reason',
'taskMessage',
'code',
],
);
final code = _firstNonEmptyDisplayText(
<String, dynamic>{...error, ...snapshot, 'taskCode': task['code']},
const <String>['code', 'errorCode', 'taskCode'],
);
final result = <String, dynamic>{
'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<String, dynamic> result) {
for (final key in const <String>['artifacts', 'files', 'attachments']) {
if (_listValue(result[key]).isNotEmpty) {

View File

@ -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<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': '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(<String, dynamic>{
'jsonrpc': '2.0',
'id': id,
'result': <String, dynamic>{
'status': completed ? 'completed' : 'running',
'sessionId': 'unit-fixture-task-b',
'threadId': 'unit-fixture-task-b',
'task': <String, dynamic>{
'state': completed ? 'completed' : 'running',
'turnId': 'turn-recovered-running',
},
if (completed)
'result': <String, dynamic>{
'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: <String>[],
inlineAttachments: <GatewayChatAttachmentPayload>[],
localAttachments: <CollaborationAttachment>[],
agentId: '',
metadata: <String, dynamic>{},
),
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<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': '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(<String, dynamic>{
'jsonrpc': '2.0',
'id': id,
'result': <String, dynamic>{
'status': 'failed',
'sessionId': 'unit-fixture-task-c',
'threadId': 'unit-fixture-task-c',
'task': <String, dynamic>{
'state': 'failed',
'turnId': 'turn-failed',
},
'error': <String, dynamic>{
'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: <String>[],
inlineAttachments: <GatewayChatAttachmentPayload>[],
localAttachments: <CollaborationAttachment>[],
agentId: '',
metadata: <String, dynamic>{},
),
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 {