diff --git a/lib/app/app_controller_desktop_runtime_helpers.dart b/lib/app/app_controller_desktop_runtime_helpers.dart index 36fa3cac..e349f517 100644 --- a/lib/app/app_controller_desktop_runtime_helpers.dart +++ b/lib/app/app_controller_desktop_runtime_helpers.dart @@ -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?> 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?> _downloadBridgeArtifactBytesInternal( + Uri uri, + String authorization, + ) async { + var bytes = []; + for (var attempt = 1; attempt <= 3; attempt++) { + final result = await _downloadBridgeArtifactBytesOnceInternal( + uri, + authorization, + rangeStart: bytes.length, + ); + if (result.reset) { + bytes = []; + } + if (result.bytes.isNotEmpty) { + bytes.addAll(result.bytes); + } + if (result.completed) { + return bytes; + } + if (attempt < 3) { + await Future.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 = []; 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>( - [], - (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 _writeVerifiedArtifactBytesInternal( + File target, + List 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 bytes) + : this._(bytes: bytes, failed: false); + + final List? bytes; + final bool failed; +} + +class _ArtifactDownloadAttemptResult { + const _ArtifactDownloadAttemptResult({ + required this.bytes, + required this.completed, + required this.reset, + }); + + const _ArtifactDownloadAttemptResult.retry() + : this(bytes: const [], completed: false, reset: false); + + final List bytes; + final bool completed; + final bool reset; +} + List _decodeArtifactContentInternal(GoTaskServiceArtifact artifact) { final encoding = artifact.encoding.trim().toLowerCase(); if (encoding == 'base64') { diff --git a/test/runtime/app_controller_thread_workspace_binding_test.dart b/test/runtime/app_controller_thread_workspace_binding_test.dart index 8d1c74a3..235a2958 100644 --- a/test/runtime/app_controller_thread_workspace_binding_test.dart +++ b/test/runtime/app_controller_thread_workspace_binding_test.dart @@ -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(), [ @@ -342,6 +341,323 @@ void main() { }, ); + test( + 'resumes bridge artifact downloads after a weak network disconnect', + () async { + final body = [0x41, 0x52, 0x54, 0x49, 0x46, 0x41, 0x43, 0x54]; + final observedRanges = []; + var requestCount = 0; + final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0); + addTearDown(() => server.close()); + server.listen((socket) async { + requestCount += 1; + final requestBytes = []; + 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 { + '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: { + 'artifacts': >[ + { + '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, ['', '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 { + '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: { + 'artifacts': >[ + { + 'relativePath': 'reports/inline.txt', + 'content': 'inline ok', + 'contentType': 'text/plain', + }, + { + 'relativePath': 'reports/failed.txt', + 'downloadUrl': + 'http://xworkmate-bridge.svc.plus:${server.port}/failed.txt', + 'contentType': 'text/plain', + }, + { + '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(['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 { + '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: { + 'artifacts': >[ + { + '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 { @@ -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.generate( + 16, + (_) => HttpClient()..findProxy = (_) => 'PROXY 127.0.0.1:$port', + ); + var index = 0; + return (_) => clients[index++]; +}