fix: harden artifact download sync

This commit is contained in:
Haitao Pan 2026-05-06 21:05:59 +08:00
parent 0c6f1538b9
commit d902ddc8a8
2 changed files with 521 additions and 23 deletions

View File

@ -3,6 +3,7 @@
import 'dart:async';
import 'dart:convert';
import 'dart:io';
import 'package:crypto/crypto.dart' as crypto;
import 'package:flutter/material.dart';
import 'app_metadata.dart';
import 'app_capabilities.dart';
@ -608,44 +609,73 @@ extension AppControllerDesktopRuntimeHelpers on AppController {
await root.create(recursive: true);
var wroteArtifact = false;
var failedArtifact = false;
var skippedArtifact = false;
for (final artifact in artifacts) {
final relativePath = _sanitizeArtifactRelativePathInternal(
artifact.relativePath,
);
if (relativePath.isEmpty) {
skippedArtifact = true;
continue;
}
final bytes = await artifactBytesInternal(artifact);
final bytesResult = await _artifactBytesResultInternal(artifact);
if (bytesResult.failed) {
failedArtifact = true;
}
final bytes = bytesResult.bytes;
if (bytes == null) {
skippedArtifact = true;
continue;
}
final target = await _nextArtifactTargetFileInternal(root, relativePath);
await target.parent.create(recursive: true);
await target.writeAsBytes(bytes, flush: true);
final verified = await _writeVerifiedArtifactBytesInternal(
target,
bytes,
artifact,
);
if (!verified) {
failedArtifact = true;
continue;
}
wroteArtifact = true;
}
final syncStatus = wroteArtifact
? (failedArtifact || skippedArtifact ? 'partial' : 'synced')
: failedArtifact
? 'download-failed'
: 'no-artifacts';
upsertTaskThreadInternal(
normalizedSessionKey,
lastArtifactSyncAtMs: syncedAtMs,
lastArtifactSyncStatus: wroteArtifact ? 'synced' : 'no-inline-content',
lastArtifactSyncStatus: syncStatus,
updatedAtMs: syncedAtMs,
);
}
Future<List<int>?> artifactBytesInternal(
GoTaskServiceArtifact artifact,
) async {
return (await _artifactBytesResultInternal(artifact)).bytes;
}
Future<_ArtifactBytesResult> _artifactBytesResultInternal(
GoTaskServiceArtifact artifact,
) async {
if (artifact.hasInlineContent) {
return _decodeArtifactContentInternal(artifact);
return _ArtifactBytesResult.bytes(
_decodeArtifactContentInternal(artifact),
);
}
final rawDownloadUrl = artifact.downloadUrl.trim();
if (rawDownloadUrl.isEmpty) {
return null;
return const _ArtifactBytesResult.skipped();
}
final uri = Uri.tryParse(rawDownloadUrl);
if (uri == null || (uri.scheme != 'http' && uri.scheme != 'https')) {
return null;
return const _ArtifactBytesResult.skipped();
}
final bridgeEndpoint = resolveBridgeAcpEndpointInternal();
final sameBridgeHost =
@ -653,30 +683,144 @@ extension AppControllerDesktopRuntimeHelpers on AppController {
uri.host.trim().toLowerCase() ==
bridgeEndpoint.host.trim().toLowerCase();
if (!sameBridgeHost) {
return null;
return const _ArtifactBytesResult.skipped();
}
final authorization =
await resolveBridgeArtifactAuthorizationHeaderInternal(uri);
if (authorization == null || authorization.trim().isEmpty) {
return null;
return const _ArtifactBytesResult.skipped();
}
final client = HttpClient();
final bytes = await _downloadBridgeArtifactBytesInternal(
uri,
authorization,
);
if (bytes == null) {
return const _ArtifactBytesResult.failed();
}
return _ArtifactBytesResult.bytes(bytes);
}
Future<List<int>?> _downloadBridgeArtifactBytesInternal(
Uri uri,
String authorization,
) async {
var bytes = <int>[];
for (var attempt = 1; attempt <= 3; attempt++) {
final result = await _downloadBridgeArtifactBytesOnceInternal(
uri,
authorization,
rangeStart: bytes.length,
);
if (result.reset) {
bytes = <int>[];
}
if (result.bytes.isNotEmpty) {
bytes.addAll(result.bytes);
}
if (result.completed) {
return bytes;
}
if (attempt < 3) {
await Future<void>.delayed(Duration(milliseconds: attempt * 250));
}
}
return null;
}
Future<_ArtifactDownloadAttemptResult>
_downloadBridgeArtifactBytesOnceInternal(
Uri uri,
String authorization, {
required int rangeStart,
}) async {
final client = HttpClient()
..connectionTimeout = const Duration(seconds: 12);
var reset = false;
final bytes = <int>[];
try {
final request = await client.getUrl(uri);
request.headers.set(HttpHeaders.authorizationHeader, authorization);
final response = await request.close();
if (response.statusCode != HttpStatus.ok) {
return null;
if (rangeStart > 0) {
request.headers.set(HttpHeaders.rangeHeader, 'bytes=$rangeStart-');
}
return response.fold<List<int>>(
<int>[],
(buffer, chunk) => buffer..addAll(chunk),
final response = await request.close();
if (response.statusCode == HttpStatus.ok) {
reset = rangeStart > 0;
} else if (response.statusCode == HttpStatus.partialContent) {
reset = false;
} else {
return const _ArtifactDownloadAttemptResult.retry();
}
await for (final chunk in response) {
bytes.addAll(chunk);
}
return _ArtifactDownloadAttemptResult(
bytes: bytes,
completed: true,
reset: reset,
);
} on HttpException {
return _ArtifactDownloadAttemptResult(
bytes: bytes,
completed: false,
reset: reset,
);
} on SocketException {
return _ArtifactDownloadAttemptResult(
bytes: bytes,
completed: false,
reset: reset,
);
} on TimeoutException {
return _ArtifactDownloadAttemptResult(
bytes: bytes,
completed: false,
reset: reset,
);
} on StateError {
return _ArtifactDownloadAttemptResult(
bytes: bytes,
completed: false,
reset: reset,
);
} finally {
client.close(force: true);
}
}
Future<bool> _writeVerifiedArtifactBytesInternal(
File target,
List<int> bytes,
GoTaskServiceArtifact artifact,
) async {
final expectedSize = artifact.sizeBytes;
if (expectedSize != null && expectedSize != bytes.length) {
return false;
}
final expectedSha256 = artifact.sha256.trim().toLowerCase();
if (expectedSha256.isNotEmpty &&
expectedSha256.length == 64 &&
crypto.sha256.convert(bytes).toString() != expectedSha256) {
return false;
}
final temp = File(
'${target.path}.xworkmate-sync-${DateTime.now().microsecondsSinceEpoch}.tmp',
);
try {
await temp.writeAsBytes(bytes, flush: true);
if (await target.exists()) {
await target.delete();
}
await temp.rename(target.path);
return true;
} catch (_) {
if (await temp.exists()) {
await temp.delete();
}
return false;
}
}
Uri? resolveGatewayAcpEndpointInternal() {
return resolveBridgeAcpEndpointInternal();
}
@ -867,6 +1011,35 @@ String _sanitizeArtifactRelativePathInternal(String raw) {
.join('/');
}
class _ArtifactBytesResult {
const _ArtifactBytesResult._({this.bytes, required this.failed});
const _ArtifactBytesResult.skipped() : this._(failed: false);
const _ArtifactBytesResult.failed() : this._(failed: true);
const _ArtifactBytesResult.bytes(List<int> bytes)
: this._(bytes: bytes, failed: false);
final List<int>? bytes;
final bool failed;
}
class _ArtifactDownloadAttemptResult {
const _ArtifactDownloadAttemptResult({
required this.bytes,
required this.completed,
required this.reset,
});
const _ArtifactDownloadAttemptResult.retry()
: this(bytes: const <int>[], completed: false, reset: false);
final List<int> bytes;
final bool completed;
final bool reset;
}
List<int> _decodeArtifactContentInternal(GoTaskServiceArtifact artifact) {
final encoding = artifact.encoding.trim().toLowerCase();
if (encoding == 'base64') {

View File

@ -1,5 +1,6 @@
import 'dart:io';
import 'package:crypto/crypto.dart' as crypto;
import 'package:flutter_test/flutter_test.dart';
import 'package:xworkmate/app/app_controller.dart';
import 'package:xworkmate/app/app_controller_desktop_runtime_coordination_impl.dart';
@ -201,14 +202,13 @@ void main() {
route: GoTaskServiceRoute.externalAcpSingle,
);
final proxyClient = HttpClient()
..findProxy = (_) => 'PROXY 127.0.0.1:${server.port}';
final clientFactory = _proxiedClientFactory(server.port);
await HttpOverrides.runZoned(() async {
await controller.persistGoTaskArtifactsForSessionInternal(
'session-1',
result,
);
}, createHttpClient: (_) => proxyClient);
}, createHttpClient: clientFactory);
final artifact = File('${localWorkspace.path}/reports/download.txt');
expect(await artifact.readAsString(), 'downloaded artifact body');
@ -295,7 +295,7 @@ void main() {
'contentType': 'application/octet-stream',
'sizeBytes': 8,
'sha256':
'59f56f3c87334ee2eb47024a0748f725c1ad2be2954a85bc680db4b012d0b02e',
'7fbd7ef36fdd97293aa5b3bcd597146101d3ea9a12b271ed0c88bdca25b63d12',
},
],
},
@ -304,14 +304,13 @@ void main() {
route: GoTaskServiceRoute.externalAcpSingle,
);
final proxyClient = HttpClient()
..findProxy = (_) => 'PROXY 127.0.0.1:${server.port}';
final clientFactory = _proxiedClientFactory(server.port);
await HttpOverrides.runZoned(() async {
await controller.persistGoTaskArtifactsForSessionInternal(
sessionKey,
result,
);
}, createHttpClient: (_) => proxyClient);
}, createHttpClient: clientFactory);
final artifact = File('${taskWorkspace.path}/exports/openclaw.bin');
expect(await artifact.readAsBytes(), <int>[
@ -342,6 +341,323 @@ void main() {
},
);
test(
'resumes bridge artifact downloads after a weak network disconnect',
() async {
final body = <int>[0x41, 0x52, 0x54, 0x49, 0x46, 0x41, 0x43, 0x54];
final observedRanges = <String>[];
var requestCount = 0;
final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0);
addTearDown(() => server.close());
server.listen((socket) async {
requestCount += 1;
final requestBytes = <int>[];
await for (final chunk in socket) {
requestBytes.addAll(chunk);
if (String.fromCharCodes(requestBytes).contains('\r\n\r\n')) {
break;
}
}
final rawRequest = String.fromCharCodes(requestBytes);
final rangeLine = rawRequest
.split('\r\n')
.firstWhere(
(line) => line.toLowerCase().startsWith('range:'),
orElse: () => '',
);
observedRanges.add(
rangeLine.replaceFirst(RegExp('^[Rr]ange:\\s*'), ''),
);
if (requestCount == 1) {
socket.add(
'HTTP/1.1 200 OK\r\n'
'Content-Type: application/octet-stream\r\n'
'Content-Length: 8\r\n'
'\r\n'
.codeUnits,
);
socket.add(body.take(4).toList());
await socket.flush();
socket.destroy();
return;
}
expect(rangeLine.toLowerCase(), 'range: bytes=4-');
socket.add(
'HTTP/1.1 206 Partial Content\r\n'
'Content-Type: application/octet-stream\r\n'
'Content-Range: bytes 4-7/8\r\n'
'Content-Length: 4\r\n'
'\r\n'
.codeUnits,
);
socket.add(body.skip(4).toList());
await socket.flush();
await socket.close();
});
final controller = AppController(
environmentOverride: const <String, String>{
'BRIDGE_AUTH_TOKEN': 'bridge-token',
},
);
addTearDown(controller.dispose);
final localWorkspace = await Directory.systemTemp.createTemp(
'xworkmate-resume-artifact-workspace-',
);
addTearDown(() async {
if (await localWorkspace.exists()) {
await localWorkspace.delete(recursive: true);
}
});
controller.upsertTaskThreadInternal(
'session-1',
workspaceBinding: WorkspaceBinding(
workspaceId: 'session-1',
workspaceKind: WorkspaceKind.localFs,
workspacePath: localWorkspace.path,
displayPath: localWorkspace.path,
writable: true,
),
);
final result = GoTaskServiceResult(
success: true,
message: 'hello',
turnId: 'turn-1',
raw: <String, dynamic>{
'artifacts': <Map<String, dynamic>>[
<String, dynamic>{
'relativePath': 'reports/resume.bin',
'downloadUrl':
'http://xworkmate-bridge.svc.plus:${server.port}/artifacts/openclaw/download'
'?sessionKey=session-1&runId=run-1&relativePath=reports%2Fresume.bin'
'&expires=9999999999&sig=test-signature',
'contentType': 'application/octet-stream',
'sizeBytes': body.length,
'sha256': crypto.sha256.convert(body).toString(),
},
],
},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
);
final clientFactory = _proxiedClientFactory(server.port);
await HttpOverrides.runZoned(() async {
await controller.persistGoTaskArtifactsForSessionInternal(
'session-1',
result,
);
}, createHttpClient: clientFactory);
expect(requestCount, 2);
expect(observedRanges, <String>['', 'bytes=4-']);
expect(
await File('${localWorkspace.path}/reports/resume.bin').readAsBytes(),
body,
);
expect(
controller
.requireTaskThreadForSessionInternal('session-1')
.lastArtifactSyncStatus,
'synced',
);
},
);
test('keeps syncing later artifacts when one download fails', () async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
addTearDown(() => server.close(force: true));
server.listen((request) async {
if (request.uri.path.endsWith('/failed.txt')) {
request.response.statusCode = HttpStatus.badGateway;
await request.response.close();
return;
}
request.response
..statusCode = HttpStatus.ok
..headers.contentType = ContentType.text
..write('download ok');
await request.response.close();
});
final controller = AppController(
environmentOverride: const <String, String>{
'BRIDGE_AUTH_TOKEN': 'bridge-token',
},
);
addTearDown(controller.dispose);
final localWorkspace = await Directory.systemTemp.createTemp(
'xworkmate-partial-artifact-workspace-',
);
addTearDown(() async {
if (await localWorkspace.exists()) {
await localWorkspace.delete(recursive: true);
}
});
controller.upsertTaskThreadInternal(
'session-1',
workspaceBinding: WorkspaceBinding(
workspaceId: 'session-1',
workspaceKind: WorkspaceKind.localFs,
workspacePath: localWorkspace.path,
displayPath: localWorkspace.path,
writable: true,
),
);
final result = GoTaskServiceResult(
success: true,
message: 'hello',
turnId: 'turn-1',
raw: <String, dynamic>{
'artifacts': <Map<String, dynamic>>[
<String, dynamic>{
'relativePath': 'reports/inline.txt',
'content': 'inline ok',
'contentType': 'text/plain',
},
<String, dynamic>{
'relativePath': 'reports/failed.txt',
'downloadUrl':
'http://xworkmate-bridge.svc.plus:${server.port}/failed.txt',
'contentType': 'text/plain',
},
<String, dynamic>{
'relativePath': 'reports/download.txt',
'downloadUrl':
'http://xworkmate-bridge.svc.plus:${server.port}/download.txt',
'contentType': 'text/plain',
},
],
},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
);
final clientFactory = _proxiedClientFactory(server.port);
await HttpOverrides.runZoned(() async {
await controller.persistGoTaskArtifactsForSessionInternal(
'session-1',
result,
);
}, createHttpClient: clientFactory);
expect(
await File('${localWorkspace.path}/reports/inline.txt').readAsString(),
'inline ok',
);
expect(
await File('${localWorkspace.path}/reports/download.txt').readAsString(),
'download ok',
);
expect(
await File('${localWorkspace.path}/reports/failed.txt').exists(),
isFalse,
);
final snapshot = await controller.loadAssistantArtifactSnapshot(
sessionKey: 'session-1',
);
expect(
snapshot.fileEntries.map((entry) => entry.relativePath),
containsAll(<String>['reports/inline.txt', 'reports/download.txt']),
);
expect(
controller
.requireTaskThreadForSessionInternal('session-1')
.lastArtifactSyncStatus,
'partial',
);
});
test('drops artifacts when size or sha256 validation fails', () async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
addTearDown(() => server.close(force: true));
server.listen((request) async {
request.response
..statusCode = HttpStatus.ok
..headers.contentType = ContentType.text
..write('bad body');
await request.response.close();
});
final controller = AppController(
environmentOverride: const <String, String>{
'BRIDGE_AUTH_TOKEN': 'bridge-token',
},
);
addTearDown(controller.dispose);
final localWorkspace = await Directory.systemTemp.createTemp(
'xworkmate-invalid-artifact-workspace-',
);
addTearDown(() async {
if (await localWorkspace.exists()) {
await localWorkspace.delete(recursive: true);
}
});
controller.upsertTaskThreadInternal(
'session-1',
workspaceBinding: WorkspaceBinding(
workspaceId: 'session-1',
workspaceKind: WorkspaceKind.localFs,
workspacePath: localWorkspace.path,
displayPath: localWorkspace.path,
writable: true,
),
);
final result = GoTaskServiceResult(
success: true,
message: 'hello',
turnId: 'turn-1',
raw: <String, dynamic>{
'artifacts': <Map<String, dynamic>>[
<String, dynamic>{
'relativePath': 'reports/invalid.txt',
'downloadUrl':
'http://xworkmate-bridge.svc.plus:${server.port}/invalid.txt',
'contentType': 'text/plain',
'sizeBytes': 8,
'sha256':
'0000000000000000000000000000000000000000000000000000000000000000',
},
],
},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
);
final clientFactory = _proxiedClientFactory(server.port);
await HttpOverrides.runZoned(() async {
await controller.persistGoTaskArtifactsForSessionInternal(
'session-1',
result,
);
}, createHttpClient: clientFactory);
expect(
await File('${localWorkspace.path}/reports/invalid.txt').exists(),
isFalse,
);
final leftovers = await localWorkspace
.list(recursive: true)
.where((entity) => entity.path.contains('.xworkmate-sync-'))
.toList();
expect(leftovers, isEmpty);
expect(
controller
.requireTaskThreadForSessionInternal('session-1')
.lastArtifactSyncStatus,
'download-failed',
);
});
test('skips download URL artifacts outside the bridge host', () async {
final controller = AppController(
environmentOverride: const <String, String>{
@ -401,7 +717,16 @@ void main() {
controller
.requireTaskThreadForSessionInternal('session-1')
.lastArtifactSyncStatus,
'no-inline-content',
'no-artifacts',
);
});
}
HttpClient Function(SecurityContext?) _proxiedClientFactory(int port) {
final clients = List<HttpClient>.generate(
16,
(_) => HttpClient()..findProxy = (_) => 'PROXY 127.0.0.1:$port',
);
var index = 0;
return (_) => clients[index++];
}