From 20cf61cd6b8da9b28614c052da992a32dc0a9ad7 Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Sun, 29 Mar 2026 14:56:27 +0800 Subject: [PATCH] refactor(runtime): retire legacy direct single-agent path --- ...direct_single_agent_app_server_client.dart | 8 +- ...t_single_agent_app_server_client_core.dart | 224 ----- ...ingle_agent_app_server_client_helpers.dart | 179 ---- ...ngle_agent_app_server_client_protocol.dart | 133 --- ...gle_agent_app_server_client_transport.dart | 885 ----------------- lib/runtime/single_agent_runner.dart | 231 +---- .../app_controller_ai_gateway_chat_suite.dart | 1 - ...controller_ai_gateway_chat_suite_core.dart | 1 - ...ontroller_ai_gateway_chat_suite_fakes.dart | 69 +- .../direct_single_agent_app_server_suite.dart | 937 ------------------ test/runtime/go_agent_core_client_suite.dart | 123 +++ ...st.dart => go_agent_core_client_test.dart} | 2 +- .../no_direct_cli_execution_guard_suite.dart | 24 +- 13 files changed, 158 insertions(+), 2659 deletions(-) delete mode 100644 lib/runtime/direct_single_agent_app_server_client_core.dart delete mode 100644 lib/runtime/direct_single_agent_app_server_client_helpers.dart delete mode 100644 lib/runtime/direct_single_agent_app_server_client_transport.dart delete mode 100644 test/runtime/direct_single_agent_app_server_suite.dart create mode 100644 test/runtime/go_agent_core_client_suite.dart rename test/runtime/{direct_single_agent_app_server_test.dart => go_agent_core_client_test.dart} (53%) diff --git a/lib/runtime/direct_single_agent_app_server_client.dart b/lib/runtime/direct_single_agent_app_server_client.dart index 03c06182..67cad407 100644 --- a/lib/runtime/direct_single_agent_app_server_client.dart +++ b/lib/runtime/direct_single_agent_app_server_client.dart @@ -1,4 +1,6 @@ +// Legacy compatibility surface retained while the app imports are cleaned up. +// +// The direct single-agent app-server runtime has been retired in favor of the +// GoAgentCore ACP path. This library intentionally exports only the capability +// DTOs still consumed by the UI-facing state layer. export 'direct_single_agent_app_server_client_protocol.dart'; -export 'direct_single_agent_app_server_client_transport.dart'; -export 'direct_single_agent_app_server_client_helpers.dart'; -export 'direct_single_agent_app_server_client_core.dart'; diff --git a/lib/runtime/direct_single_agent_app_server_client_core.dart b/lib/runtime/direct_single_agent_app_server_client_core.dart deleted file mode 100644 index 7b117d69..00000000 --- a/lib/runtime/direct_single_agent_app_server_client_core.dart +++ /dev/null @@ -1,224 +0,0 @@ -// ignore_for_file: unused_import, unnecessary_import - -import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; -import 'runtime_models.dart'; -import 'direct_single_agent_app_server_client_protocol.dart'; -import 'direct_single_agent_app_server_client_transport.dart'; -import 'direct_single_agent_app_server_client_helpers.dart'; - -class DirectSingleAgentAppServerClient { - DirectSingleAgentAppServerClient({required this.endpointResolver}); - - final Uri? Function(SingleAgentProvider provider) endpointResolver; - final DirectSingleAgentWebSocketTransportInternal webSocketTransportInternal = - DirectSingleAgentWebSocketTransportInternal(); - final DirectSingleAgentRestTransportInternal restTransportInternal = - DirectSingleAgentRestTransportInternal(); - - final Map - cachedCapabilitiesInternal = - {}; - final Map capabilitiesRefreshedAtInternal = - {}; - final Map - transportKindsInternal = - {}; - - Future loadCapabilities({ - required SingleAgentProvider provider, - bool forceRefresh = false, - String gatewayToken = '', - }) async { - final cached = cachedCapabilitiesInternal[provider]; - final refreshedAt = capabilitiesRefreshedAtInternal[provider]; - if (!forceRefresh && - cached != null && - refreshedAt != null && - DateTime.now().difference(refreshedAt) < const Duration(seconds: 15)) { - return cached; - } - - final descriptor = describeEndpointInternal(provider); - if (!descriptor.isSupported || descriptor.baseUri == null) { - final unavailable = const DirectSingleAgentCapabilities.unavailable( - endpoint: '', - errorMessage: 'Single-agent app-server endpoint is not configured.', - ); - cachedCapabilitiesInternal[provider] = unavailable; - capabilitiesRefreshedAtInternal[provider] = DateTime.now(); - return unavailable; - } - - try { - final transport = await resolveTransportInternal( - provider, - descriptor: descriptor, - gatewayToken: gatewayToken, - ); - transportKindsInternal[provider] = transport.kind; - cachedCapabilitiesInternal[provider] = DirectSingleAgentCapabilities( - available: true, - supportedProviders: [provider], - endpoint: transport.endpoint.toString(), - ); - } catch (error) { - cachedCapabilitiesInternal[provider] = - DirectSingleAgentCapabilities.unavailable( - endpoint: descriptor.baseUri.toString(), - errorMessage: error.toString(), - ); - transportKindsInternal.remove(provider); - } finally { - capabilitiesRefreshedAtInternal[provider] = DateTime.now(); - } - - return cachedCapabilitiesInternal[provider]!; - } - - Future run( - DirectSingleAgentRunRequest request, - ) async { - final descriptor = describeEndpointInternal(request.provider); - if (!descriptor.isSupported || descriptor.baseUri == null) { - return const DirectSingleAgentRunResult( - success: false, - output: '', - errorMessage: 'Single-agent app-server endpoint is missing.', - ); - } - late final ResolvedSingleAgentTransportInternal transport; - try { - transport = await resolveTransportInternal( - request.provider, - descriptor: descriptor, - gatewayToken: request.gatewayToken, - ); - } catch (error) { - return DirectSingleAgentRunResult( - success: false, - output: '', - errorMessage: error.toString(), - ); - } - if (transport.kind == - DirectSingleAgentTransportKindInternal.restSessionApi) { - return transport.rest!.run( - request, - base: transport.endpoint, - ); - } - return transport.websocket!.run( - request, - endpoint: transport.endpoint, - ); - } - - Future abort(String sessionId) async { - await restTransportInternal.abort( - sessionId, - candidateBases: [ - for (final entry in transportKindsInternal.entries) - if (entry.value == - DirectSingleAgentTransportKindInternal.restSessionApi) ...[ - if (describeEndpointInternal(entry.key).baseUri != null) - describeEndpointInternal(entry.key).baseUri!, - ], - ], - ); - await webSocketTransportInternal.abort(sessionId); - } - - Future dispose() async { - await webSocketTransportInternal.dispose(); - } - - DirectSingleAgentEndpointDescriptor describeEndpointInternal( - SingleAgentProvider provider, - ) { - return DirectSingleAgentEndpointDescriptor.describe( - endpointResolver(provider), - ); - } - - Future resolveTransportInternal( - SingleAgentProvider provider, { - required DirectSingleAgentEndpointDescriptor descriptor, - required String gatewayToken, - }) async { - final cachedKind = transportKindsInternal[provider]; - if (cachedKind != null) { - final cachedEndpoint = - cachedKind == - DirectSingleAgentTransportKindInternal.websocketAppServer - ? descriptor.websocketUri - : descriptor.baseUri; - if (cachedEndpoint != null) { - return ResolvedSingleAgentTransportInternal( - kind: cachedKind, - endpoint: cachedEndpoint, - websocket: - cachedKind == - DirectSingleAgentTransportKindInternal.websocketAppServer - ? webSocketTransportInternal - : null, - rest: - cachedKind == - DirectSingleAgentTransportKindInternal.restSessionApi - ? restTransportInternal - : null, - ); - } - } - - if (descriptor.prefersWebSocket) { - final endpoint = descriptor.websocketUri; - if (endpoint == null) { - throw StateError('Single-agent websocket endpoint is not configured.'); - } - await webSocketTransportInternal.probe( - endpoint, - gatewayToken: gatewayToken, - ); - return ResolvedSingleAgentTransportInternal( - kind: DirectSingleAgentTransportKindInternal.websocketAppServer, - endpoint: endpoint, - websocket: webSocketTransportInternal, - ); - } - - if (descriptor.allowsRest) { - final base = descriptor.baseUri; - if (base == null) { - throw StateError('Single-agent endpoint is not configured.'); - } - try { - await restTransportInternal.probe(base, gatewayToken: gatewayToken); - return ResolvedSingleAgentTransportInternal( - kind: DirectSingleAgentTransportKindInternal.restSessionApi, - endpoint: base, - rest: restTransportInternal, - ); - } catch (_) { - final websocket = descriptor.websocketUri; - if (websocket == null) { - rethrow; - } - await webSocketTransportInternal.probe( - websocket, - gatewayToken: gatewayToken, - ); - return ResolvedSingleAgentTransportInternal( - kind: DirectSingleAgentTransportKindInternal.websocketAppServer, - endpoint: websocket, - websocket: webSocketTransportInternal, - ); - } - } - - throw StateError( - 'Single-agent endpoint mode ${descriptor.mode.name} is not supported.', - ); - } -} diff --git a/lib/runtime/direct_single_agent_app_server_client_helpers.dart b/lib/runtime/direct_single_agent_app_server_client_helpers.dart deleted file mode 100644 index c73531db..00000000 --- a/lib/runtime/direct_single_agent_app_server_client_helpers.dart +++ /dev/null @@ -1,179 +0,0 @@ -// ignore_for_file: unused_import, unnecessary_import - -import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; -import 'runtime_models.dart'; -import 'direct_single_agent_app_server_client_protocol.dart'; -import 'direct_single_agent_app_server_client_transport.dart'; -import 'direct_single_agent_app_server_client_core.dart'; - -Uri buildRestUriInternal( - Uri base, - String path, { - Map? queryParameters, -}) { - final normalizedPath = path.startsWith('/') ? path : '/$path'; - return base.replace( - path: normalizedPath, - queryParameters: queryParameters, - fragment: null, - ); -} - -Future> fetchJsonInternal( - Uri uri, { - required String gatewayToken, -}) async { - final client = HttpClient()..connectionTimeout = const Duration(seconds: 8); - try { - final request = await client.getUrl(uri); - final normalizedToken = gatewayToken.trim(); - if (normalizedToken.isNotEmpty) { - request.headers.set( - HttpHeaders.authorizationHeader, - 'Bearer $normalizedToken', - ); - } - final response = await request.close(); - final body = await response.transform(utf8.decoder).join(); - return decodeMapInternal(body); - } finally { - client.close(force: true); - } -} - -Future> postJsonInternal( - Uri uri, { - required Object? body, - required String gatewayToken, -}) async { - final client = HttpClient()..connectionTimeout = const Duration(seconds: 8); - try { - final request = await client.postUrl(uri); - request.headers.set( - HttpHeaders.contentTypeHeader, - 'application/json; charset=utf-8', - ); - final normalizedToken = gatewayToken.trim(); - if (normalizedToken.isNotEmpty) { - request.headers.set( - HttpHeaders.authorizationHeader, - 'Bearer $normalizedToken', - ); - } - if (body != null) { - request.add(utf8.encode(jsonEncode(body))); - } - final response = await request.close(); - final text = await response.transform(utf8.decoder).join(); - if (text.trim().isEmpty) { - return const {}; - } - return decodeMapInternal(text); - } finally { - client.close(force: true); - } -} - -Future> fetchJsonListInternal( - Uri uri, { - required String gatewayToken, -}) async { - final client = HttpClient()..connectionTimeout = const Duration(seconds: 8); - try { - final request = await client.getUrl(uri); - final normalizedToken = gatewayToken.trim(); - if (normalizedToken.isNotEmpty) { - request.headers.set( - HttpHeaders.authorizationHeader, - 'Bearer $normalizedToken', - ); - } - final response = await request.close(); - final body = await response.transform(utf8.decoder).join(); - final decoded = jsonDecode(body); - if (decoded is List) { - return decoded; - } - if (decoded is List) { - return decoded.cast(); - } - return const []; - } finally { - client.close(force: true); - } -} - -String? extractThreadIdInternal(Map payload) { - final topLevelId = payload['id']?.toString().trim() ?? ''; - if (topLevelId.isNotEmpty) { - return topLevelId; - } - final thread = asMapInternal(payload['thread']); - final nestedId = thread['id']?.toString().trim() ?? ''; - if (nestedId.isNotEmpty) { - return nestedId; - } - return null; -} - -String? extractModelInternal(Map payload) { - final model = payload['model']?.toString().trim() ?? ''; - if (model.isNotEmpty) { - return model; - } - return null; -} - -String? extractThreadPathInternal(Map payload) { - final directPath = payload['path']?.toString().trim() ?? ''; - if (directPath.isNotEmpty) { - return directPath; - } - final thread = asMapInternal(payload['thread']); - final nestedPath = thread['path']?.toString().trim() ?? ''; - if (nestedPath.isNotEmpty) { - return nestedPath; - } - return null; -} - -Map decodeMapInternal(Object raw) { - if (raw is Map) { - return raw; - } - if (raw is Map) { - return raw.cast(); - } - final decoded = jsonDecode(raw.toString()); - if (decoded is Map) { - return decoded; - } - if (decoded is Map) { - return decoded.cast(); - } - return const {}; -} - -Map asMapInternal(Object? value) { - if (value is Map) { - return value; - } - if (value is Map) { - return value.cast(); - } - return const {}; -} - -bool isLocalHostInternal(String host) { - final normalized = host.trim().toLowerCase(); - if (normalized.isEmpty || - normalized == 'localhost' || - normalized == '127.0.0.1' || - normalized == '::1') { - return true; - } - final address = InternetAddress.tryParse(normalized); - return address?.isLoopback ?? false; -} diff --git a/lib/runtime/direct_single_agent_app_server_client_protocol.dart b/lib/runtime/direct_single_agent_app_server_client_protocol.dart index 114ecc39..7a29e1b3 100644 --- a/lib/runtime/direct_single_agent_app_server_client_protocol.dart +++ b/lib/runtime/direct_single_agent_app_server_client_protocol.dart @@ -1,12 +1,4 @@ -// ignore_for_file: unused_import, unnecessary_import - -import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; import 'runtime_models.dart'; -import 'direct_single_agent_app_server_client_transport.dart'; -import 'direct_single_agent_app_server_client_helpers.dart'; -import 'direct_single_agent_app_server_client_core.dart'; class DirectSingleAgentCapabilities { const DirectSingleAgentCapabilities({ @@ -32,128 +24,3 @@ class DirectSingleAgentCapabilities { bool supportsProvider(SingleAgentProvider provider) => supportedProviders.contains(provider); } - -class DirectSingleAgentRunResult { - const DirectSingleAgentRunResult({ - required this.success, - required this.output, - required this.errorMessage, - this.aborted = false, - this.resolvedModel = '', - this.resolvedWorkingDirectory = '', - this.resolvedWorkspaceRefKind = WorkspaceRefKind.localPath, - }); - - final bool success; - final String output; - final String errorMessage; - final bool aborted; - final String resolvedModel; - final String resolvedWorkingDirectory; - final WorkspaceRefKind resolvedWorkspaceRefKind; -} - -class DirectSingleAgentRunRequest { - const DirectSingleAgentRunRequest({ - required this.sessionId, - required this.provider, - required this.prompt, - required this.model, - required this.workingDirectory, - required this.gatewayToken, - this.selectedSkills = const [], - this.onOutput, - }); - - final String sessionId; - final SingleAgentProvider provider; - final String prompt; - final String model; - final String workingDirectory; - final String gatewayToken; - final List selectedSkills; - final void Function(String text)? onOutput; -} - -enum DirectSingleAgentEndpointMode { - wsLocal, - wss, - httpLocal, - https, - unsupported, -} - -enum DirectSingleAgentTransportKindInternal { - websocketAppServer, - restSessionApi, -} - -class DirectSingleAgentEndpointDescriptor { - const DirectSingleAgentEndpointDescriptor({ - required this.mode, - required this.baseUri, - this.websocketUri, - }); - - final DirectSingleAgentEndpointMode mode; - final Uri? baseUri; - final Uri? websocketUri; - - bool get isSupported => mode != DirectSingleAgentEndpointMode.unsupported; - - bool get prefersWebSocket => - mode == DirectSingleAgentEndpointMode.wsLocal || - mode == DirectSingleAgentEndpointMode.wss; - - bool get allowsRest => - mode == DirectSingleAgentEndpointMode.httpLocal || - mode == DirectSingleAgentEndpointMode.https; - - static DirectSingleAgentEndpointDescriptor describe(Uri? endpoint) { - if (endpoint == null) { - return const DirectSingleAgentEndpointDescriptor( - mode: DirectSingleAgentEndpointMode.unsupported, - baseUri: null, - ); - } - final scheme = endpoint.scheme.toLowerCase(); - final normalizedBase = endpoint.replace( - path: '', - query: null, - fragment: null, - ); - final isLocal = isLocalHostInternal(endpoint.host); - if (scheme == 'ws' && isLocal) { - return DirectSingleAgentEndpointDescriptor( - mode: DirectSingleAgentEndpointMode.wsLocal, - baseUri: normalizedBase, - websocketUri: normalizedBase, - ); - } - if (scheme == 'wss') { - return DirectSingleAgentEndpointDescriptor( - mode: DirectSingleAgentEndpointMode.wss, - baseUri: normalizedBase, - websocketUri: normalizedBase, - ); - } - if (scheme == 'http' && isLocal) { - return DirectSingleAgentEndpointDescriptor( - mode: DirectSingleAgentEndpointMode.httpLocal, - baseUri: normalizedBase, - websocketUri: normalizedBase.replace(scheme: 'ws'), - ); - } - if (scheme == 'https') { - return DirectSingleAgentEndpointDescriptor( - mode: DirectSingleAgentEndpointMode.https, - baseUri: normalizedBase, - websocketUri: normalizedBase.replace(scheme: 'wss'), - ); - } - return DirectSingleAgentEndpointDescriptor( - mode: DirectSingleAgentEndpointMode.unsupported, - baseUri: normalizedBase, - ); - } -} diff --git a/lib/runtime/direct_single_agent_app_server_client_transport.dart b/lib/runtime/direct_single_agent_app_server_client_transport.dart deleted file mode 100644 index 7fa82e4c..00000000 --- a/lib/runtime/direct_single_agent_app_server_client_transport.dart +++ /dev/null @@ -1,885 +0,0 @@ -// ignore_for_file: unused_import, unnecessary_import - -import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; -import 'runtime_models.dart'; -import 'direct_single_agent_app_server_client_protocol.dart'; -import 'direct_single_agent_app_server_client_helpers.dart'; -import 'direct_single_agent_app_server_client_core.dart'; - -class ResolvedSingleAgentTransportInternal { - const ResolvedSingleAgentTransportInternal({ - required this.kind, - required this.endpoint, - this.websocket, - this.rest, - }); - - final DirectSingleAgentTransportKindInternal kind; - final Uri endpoint; - final DirectSingleAgentWebSocketTransportInternal? websocket; - final DirectSingleAgentRestTransportInternal? rest; -} - -class ResolvedDirectThreadInternal { - const ResolvedDirectThreadInternal({ - required this.threadId, - this.workingDirectory = '', - }); - - final String threadId; - final String workingDirectory; -} - -class DirectSingleAgentWebSocketTransportInternal { - final Map - activeConnectionsInternal = {}; - final Map threadIdsInternal = {}; - final Map threadWorkingDirectoriesInternal = - {}; - final Set abortedSessionsInternal = {}; - - Future probe(Uri endpoint, {required String gatewayToken}) async { - DirectAppServerConnectionInternal? connection; - try { - connection = await DirectAppServerConnectionInternal.connect( - endpoint, - gatewayToken: gatewayToken, - ); - await connection.initialize(); - } finally { - await connection?.close(); - } - } - - Future run( - DirectSingleAgentRunRequest request, { - required Uri endpoint, - }) async { - final normalizedSessionId = request.sessionId.trim(); - if (normalizedSessionId.isEmpty) { - return const DirectSingleAgentRunResult( - success: false, - output: '', - errorMessage: 'Single-agent session id is missing.', - ); - } - - abortedSessionsInternal.remove(normalizedSessionId); - final connection = await DirectAppServerConnectionInternal.connect( - endpoint, - gatewayToken: request.gatewayToken, - ); - activeConnectionsInternal[normalizedSessionId] = connection; - - try { - await connection.initialize(); - final resolvedThread = await ensureThreadInternal( - connection, - sessionId: normalizedSessionId, - workingDirectory: request.workingDirectory, - model: request.model, - ); - final threadId = resolvedThread.threadId; - final resolvedWorkingDirectory = resolvedThread.workingDirectory.trim(); - - final output = StringBuffer(); - String resolvedModel = ''; - final completion = Completer(); - late final StreamSubscription> subscription; - subscription = connection.notifications.listen( - (notification) { - final method = notification['method']?.toString().trim() ?? ''; - final params = asMapInternal(notification['params']); - if (params['threadId']?.toString() != threadId) { - return; - } - if (method == 'item/agentMessage/delta') { - final delta = params['delta']?.toString() ?? ''; - if (delta.isNotEmpty) { - output.write(delta); - request.onOutput?.call(delta); - } - return; - } - if (method == 'turn/completed' && !completion.isCompleted) { - completion.complete( - DirectSingleAgentRunResult( - success: true, - output: output.toString(), - errorMessage: '', - resolvedModel: resolvedModel, - resolvedWorkingDirectory: resolvedWorkingDirectory, - ), - ); - return; - } - if ((method == 'turn/failed' || method == 'turn/error') && - !completion.isCompleted) { - final aborted = - abortedSessionsInternal.contains(normalizedSessionId) || - (params['message']?.toString().toLowerCase().contains( - 'abort', - ) ?? - false); - completion.complete( - DirectSingleAgentRunResult( - success: false, - output: output.toString(), - aborted: aborted, - resolvedModel: resolvedModel, - resolvedWorkingDirectory: resolvedWorkingDirectory, - errorMessage: - params['message']?.toString() ?? - params['error']?.toString() ?? - 'Single-agent app-server turn failed.', - ), - ); - } - }, - onError: (Object error, StackTrace stackTrace) { - if (!completion.isCompleted) { - completion.complete( - DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: error.toString(), - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: resolvedModel, - resolvedWorkingDirectory: resolvedWorkingDirectory, - ), - ); - } - }, - onDone: () { - if (!completion.isCompleted) { - completion.complete( - DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: - abortedSessionsInternal.contains(normalizedSessionId) - ? 'Single-agent app-server run aborted.' - : 'Single-agent app-server connection closed before completion.', - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: resolvedModel, - resolvedWorkingDirectory: resolvedWorkingDirectory, - ), - ); - } - }, - ); - - try { - final input = >[ - {'type': 'text', 'text': request.prompt}, - for (final skill in request.selectedSkills) - if (skill.label.trim().isNotEmpty && - skill.sourcePath.trim().isNotEmpty) - { - 'type': 'skill', - 'name': skill.label.trim(), - 'path': skill.sourcePath.trim(), - }, - ]; - final started = await connection.request( - 'turn/start', - params: {'threadId': threadId, 'input': input}, - ); - resolvedModel = extractModelInternal(started) ?? resolvedModel; - return await completion.future.timeout( - const Duration(minutes: 10), - onTimeout: () => DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: 'Single-agent app-server request timed out.', - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: resolvedModel, - resolvedWorkingDirectory: resolvedWorkingDirectory, - ), - ); - } finally { - await subscription.cancel(); - } - } catch (error) { - return DirectSingleAgentRunResult( - success: false, - output: '', - errorMessage: error.toString(), - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: '', - resolvedWorkingDirectory: request.workingDirectory, - ); - } finally { - activeConnectionsInternal.remove(normalizedSessionId); - await connection.close(); - abortedSessionsInternal.remove(normalizedSessionId); - } - } - - Future abort(String sessionId) async { - final normalizedSessionId = sessionId.trim(); - if (normalizedSessionId.isEmpty) { - return; - } - abortedSessionsInternal.add(normalizedSessionId); - final connection = activeConnectionsInternal[normalizedSessionId]; - final threadId = threadIdsInternal[normalizedSessionId]; - if (connection == null || threadId == null || threadId.isEmpty) { - return; - } - try { - await connection.request( - 'turn/interrupt', - params: {'threadId': threadId}, - ); - } catch (_) { - // Best effort only. - } - await connection.close(); - } - - Future dispose() async { - final connections = activeConnectionsInternal.values.toList( - growable: false, - ); - activeConnectionsInternal.clear(); - for (final connection in connections) { - await connection.close(); - } - } - - Future ensureThreadInternal( - DirectAppServerConnectionInternal connection, { - required String sessionId, - required String workingDirectory, - required String model, - }) async { - final normalizedWorkingDirectory = workingDirectory.trim(); - final existingThreadId = threadIdsInternal[sessionId]?.trim() ?? ''; - final existingWorkingDirectory = - threadWorkingDirectoriesInternal[sessionId]?.trim() ?? ''; - final canReuseExistingThread = - existingThreadId.isNotEmpty && - (normalizedWorkingDirectory.isEmpty || - (existingWorkingDirectory.isNotEmpty && - existingWorkingDirectory == normalizedWorkingDirectory)); - if (existingThreadId.isNotEmpty) { - if (!canReuseExistingThread) { - threadIdsInternal.remove(sessionId); - threadWorkingDirectoriesInternal.remove(sessionId); - } - } - if (canReuseExistingThread) { - try { - final resumed = await connection.request( - 'thread/resume', - params: { - 'threadId': existingThreadId, - if (normalizedWorkingDirectory.isNotEmpty) - 'cwd': normalizedWorkingDirectory, - }, - ); - final resumedId = extractThreadIdInternal(resumed) ?? existingThreadId; - final resumedWorkingDirectory = - extractThreadPathInternal(resumed)?.trim() ?? - normalizedWorkingDirectory; - threadIdsInternal[sessionId] = resumedId; - if (resumedWorkingDirectory.isNotEmpty) { - threadWorkingDirectoriesInternal[sessionId] = resumedWorkingDirectory; - } - return ResolvedDirectThreadInternal( - threadId: resumedId, - workingDirectory: resumedWorkingDirectory, - ); - } catch (_) { - threadIdsInternal.remove(sessionId); - threadWorkingDirectoriesInternal.remove(sessionId); - } - } - - final created = await connection.request( - 'thread/start', - params: { - if (normalizedWorkingDirectory.isNotEmpty) - 'cwd': normalizedWorkingDirectory, - if (model.trim().isNotEmpty) 'model': model.trim(), - }, - ); - final threadId = extractThreadIdInternal(created) ?? ''; - if (threadId.isEmpty) { - throw StateError('Single-agent app-server returned an empty thread id.'); - } - final createdWorkingDirectory = - extractThreadPathInternal(created)?.trim() ?? - normalizedWorkingDirectory; - threadIdsInternal[sessionId] = threadId; - if (createdWorkingDirectory.isNotEmpty) { - threadWorkingDirectoriesInternal[sessionId] = createdWorkingDirectory; - } - return ResolvedDirectThreadInternal( - threadId: threadId, - workingDirectory: createdWorkingDirectory, - ); - } -} - -class DirectSingleAgentRestTransportInternal { - final Map restSessionIdsInternal = {}; - final Map restSessionWorkingDirectoriesInternal = - {}; - final Set abortedSessionsInternal = {}; - - Future probe(Uri base, {required String gatewayToken}) async { - await fetchJsonInternal( - buildRestUriInternal(base, '/global/health'), - gatewayToken: gatewayToken, - ); - } - - Future run( - DirectSingleAgentRunRequest request, { - required Uri base, - }) async { - final normalizedSessionId = request.sessionId.trim(); - if (normalizedSessionId.isEmpty) { - return const DirectSingleAgentRunResult( - success: false, - output: '', - errorMessage: 'Single-agent session id is missing.', - ); - } - - abortedSessionsInternal.remove(normalizedSessionId); - final remoteSessionId = await ensureRestSessionInternal( - base, - sessionId: normalizedSessionId, - workingDirectory: request.workingDirectory, - gatewayToken: request.gatewayToken, - ); - - final output = StringBuffer(); - final completion = Completer(); - String? activeAssistantMessageId; - String? lastAssistantText; - var busySeen = false; - - void completeFailure(String message) { - if (completion.isCompleted) { - return; - } - completion.complete( - DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: message, - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: request.model, - resolvedWorkingDirectory: request.workingDirectory, - ), - ); - } - - final eventClient = HttpClient() - ..connectionTimeout = const Duration(seconds: 8); - late final HttpClientRequest eventRequest; - late final HttpClientResponse eventResponse; - StreamSubscription? lineSubscription; - - void completeSuccess() { - if (completion.isCompleted) { - return; - } - final resolvedOutput = output.toString().trim().isNotEmpty - ? output.toString() - : (lastAssistantText ?? ''); - if (resolvedOutput.trim().isEmpty) { - completeFailure( - 'OpenCode REST session completed without assistant content.', - ); - return; - } - completion.complete( - DirectSingleAgentRunResult( - success: true, - output: resolvedOutput, - errorMessage: '', - resolvedModel: request.model, - resolvedWorkingDirectory: request.workingDirectory, - ), - ); - } - - try { - final eventUri = buildRestUriInternal(base, '/global/event'); - eventRequest = await eventClient.getUrl(eventUri); - eventRequest.headers.set(HttpHeaders.acceptHeader, 'text/event-stream'); - final normalizedToken = request.gatewayToken.trim(); - if (normalizedToken.isNotEmpty) { - eventRequest.headers.set( - HttpHeaders.authorizationHeader, - 'Bearer $normalizedToken', - ); - } - eventResponse = await eventRequest.close(); - lineSubscription = eventResponse - .transform(utf8.decoder) - .transform(const LineSplitter()) - .listen( - (line) { - if (!line.startsWith('data: ')) { - return; - } - final event = decodeMapInternal(line.substring(6)); - final payload = asMapInternal(event['payload']); - final type = payload['type']?.toString().trim() ?? ''; - final properties = asMapInternal(payload['properties']); - if (properties['sessionID']?.toString().trim() != - remoteSessionId) { - return; - } - if (type == 'session.status') { - final status = asMapInternal(properties['status']); - final statusType = status['type']?.toString().trim() ?? ''; - if (statusType == 'busy') { - busySeen = true; - } - if (statusType == 'idle' && busySeen) { - completeSuccess(); - } - return; - } - if (type == 'session.idle' && busySeen) { - completeSuccess(); - return; - } - if (type == 'session.error' && !completion.isCompleted) { - final error = asMapInternal(properties['error']); - completeFailure( - error['message']?.toString() ?? - error['name']?.toString() ?? - 'OpenCode session failed.', - ); - return; - } - if (type == 'message.updated') { - final info = asMapInternal(properties['info']); - if (info['role']?.toString().trim() == 'assistant') { - activeAssistantMessageId = info['id']?.toString().trim(); - } - return; - } - if (type == 'message.part.delta') { - final part = asMapInternal(properties['part']); - if (activeAssistantMessageId != null && - part['messageID']?.toString().trim() == - activeAssistantMessageId) { - final delta = - properties['text']?.toString() ?? - properties['delta']?.toString() ?? - ''; - if (delta.isNotEmpty) { - output.write(delta); - request.onOutput?.call(delta); - } - } - return; - } - if (type == 'message.part.updated') { - final part = asMapInternal(properties['part']); - if (activeAssistantMessageId != null && - part['messageID']?.toString().trim() == - activeAssistantMessageId && - part['type']?.toString().trim() == 'text') { - lastAssistantText = part['text']?.toString(); - if ((lastAssistantText?.trim().isNotEmpty ?? false)) { - completeSuccess(); - } - } - } - }, - onError: (Object error, StackTrace stackTrace) {}, - onDone: () {}, - cancelOnError: true, - ); - - await postJsonInternal( - buildRestUriInternal( - base, - '/session/$remoteSessionId/message', - queryParameters: { - 'directory': request.workingDirectory, - }, - ), - body: { - 'agent': 'build', - 'parts': >[ - {'type': 'text', 'text': request.prompt}, - ], - }, - gatewayToken: request.gatewayToken, - ); - unawaited( - pollRestAssistantMessageInternal( - base, - remoteSessionId: remoteSessionId, - workingDirectory: request.workingDirectory, - gatewayToken: request.gatewayToken, - onResolved: (text) { - if (text.trim().isNotEmpty) { - lastAssistantText = text; - if (output.toString().trim().isEmpty) { - output.write(text); - request.onOutput?.call(text); - } - completeSuccess(); - } - }, - onError: completeFailure, - ), - ); - - return await completion.future.timeout( - const Duration(minutes: 10), - onTimeout: () => DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: 'OpenCode REST request timed out.', - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: request.model, - resolvedWorkingDirectory: request.workingDirectory, - ), - ); - } catch (error) { - return DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: error.toString(), - aborted: abortedSessionsInternal.contains(normalizedSessionId), - resolvedModel: request.model, - resolvedWorkingDirectory: request.workingDirectory, - ); - } finally { - unawaited(lineSubscription?.cancel()); - eventClient.close(force: true); - abortedSessionsInternal.remove(normalizedSessionId); - } - } - - Future abort( - String sessionId, { - required List candidateBases, - }) async { - final normalizedSessionId = sessionId.trim(); - if (normalizedSessionId.isEmpty) { - return; - } - abortedSessionsInternal.add(normalizedSessionId); - final restSessionId = - restSessionIdsInternal[normalizedSessionId]?.trim() ?? ''; - if (restSessionId.isEmpty) { - return; - } - for (final base in candidateBases) { - try { - await postJsonInternal( - buildRestUriInternal(base, '/session/$restSessionId/abort'), - body: null, - gatewayToken: '', - ); - } catch (_) { - // Best effort only. - } - break; - } - } - - Future ensureRestSessionInternal( - Uri base, { - required String sessionId, - required String workingDirectory, - required String gatewayToken, - }) async { - final normalizedWorkingDirectory = workingDirectory.trim(); - final existing = restSessionIdsInternal[sessionId]?.trim() ?? ''; - if (existing.isNotEmpty) { - final existingWorkingDirectory = - restSessionWorkingDirectoriesInternal[sessionId]?.trim() ?? ''; - final canReuseExistingSession = - normalizedWorkingDirectory.isEmpty || - (existingWorkingDirectory.isNotEmpty && - existingWorkingDirectory == normalizedWorkingDirectory); - if (canReuseExistingSession) { - return existing; - } - restSessionIdsInternal.remove(sessionId); - restSessionWorkingDirectoriesInternal.remove(sessionId); - } - final created = await postJsonInternal( - buildRestUriInternal( - base, - '/session', - queryParameters: { - 'directory': normalizedWorkingDirectory, - }, - ), - body: {'title': sessionId}, - gatewayToken: gatewayToken, - ); - final createdId = created['id']?.toString().trim() ?? ''; - if (createdId.isEmpty) { - throw StateError('OpenCode REST endpoint returned an empty session id.'); - } - restSessionIdsInternal[sessionId] = createdId; - if (normalizedWorkingDirectory.isNotEmpty) { - restSessionWorkingDirectoriesInternal[sessionId] = - normalizedWorkingDirectory; - } - return createdId; - } - - Future pollRestAssistantMessageInternal( - Uri base, { - required String remoteSessionId, - required String workingDirectory, - required String gatewayToken, - required void Function(String text) onResolved, - required void Function(String message) onError, - }) async { - String? previousText; - var stableCount = 0; - for (var attempt = 0; attempt < 100; attempt++) { - try { - final items = await fetchJsonListInternal( - buildRestUriInternal( - base, - '/session/$remoteSessionId/message', - queryParameters: { - 'directory': workingDirectory, - 'limit': '20', - }, - ), - gatewayToken: gatewayToken, - ); - final text = latestAssistantTextFromRestMessagesInternal(items); - if (text.trim().isNotEmpty) { - if (text == previousText) { - stableCount += 1; - } else { - previousText = text; - stableCount = 1; - } - if (stableCount >= 2) { - onResolved(text); - return; - } - } - } catch (error) { - onError(error.toString()); - return; - } - await Future.delayed(const Duration(milliseconds: 200)); - } - onError('OpenCode REST session completed without assistant content.'); - } - - String latestAssistantTextFromRestMessagesInternal(List items) { - for (final raw in items.reversed) { - final item = asMapInternal(raw); - final info = asMapInternal(item['info']); - if (info['role']?.toString().trim() != 'assistant') { - continue; - } - final parts = item['parts']; - if (parts is! List) { - continue; - } - for (final rawPart in parts) { - final part = asMapInternal(rawPart); - if (part['type']?.toString().trim() == 'text') { - final text = part['text']?.toString() ?? ''; - if (text.trim().isNotEmpty) { - return text; - } - } - } - } - return ''; - } -} - -class DirectAppServerConnectionInternal { - DirectAppServerConnectionInternal(this.socketInternal); - - final WebSocket socketInternal; - final StreamController> notificationsInternal = - StreamController>.broadcast(); - final Map>> pendingRequestsInternal = - >>{}; - int requestCounterInternal = 0; - bool initializedInternal = false; - StreamSubscription? subscriptionInternal; - - Stream> get notifications => - notificationsInternal.stream; - - static Future connect( - Uri endpoint, { - String gatewayToken = '', - }) async { - final headers = {}; - final normalizedToken = gatewayToken.trim(); - if (normalizedToken.isNotEmpty) { - headers[HttpHeaders.authorizationHeader] = 'Bearer $normalizedToken'; - } - final socket = - await WebSocket.connect( - endpoint.toString(), - headers: headers.isEmpty ? null : headers, - ).timeout( - const Duration(seconds: 8), - onTimeout: () => throw TimeoutException( - 'Single-agent app-server websocket connect timed out.', - ), - ); - final connection = DirectAppServerConnectionInternal(socket); - connection.attachInternal(); - return connection; - } - - Future initialize() async { - if (initializedInternal) { - return; - } - await request( - 'initialize', - params: const { - 'clientInfo': {'name': 'xworkmate', 'version': '0'}, - 'capabilities': { - 'optOutNotificationMethods': [], - }, - }, - ); - await notify('initialized', params: const {}); - initializedInternal = true; - } - - Future> request( - String method, { - Map params = const {}, - Duration timeout = const Duration(seconds: 60), - }) async { - final id = - '${DateTime.now().microsecondsSinceEpoch}-${requestCounterInternal++}'; - final completer = Completer>(); - pendingRequestsInternal[id] = completer; - socketInternal.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'method': method, - 'params': params, - }), - ); - return completer.future.timeout( - timeout, - onTimeout: () { - pendingRequestsInternal.remove(id); - throw TimeoutException( - 'Single-agent app-server request $method timed out.', - ); - }, - ); - } - - Future notify( - String method, { - required Map params, - }) async { - socketInternal.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'method': method, - 'params': params, - }), - ); - } - - void attachInternal() { - subscriptionInternal = socketInternal.listen( - (dynamic raw) { - final message = decodeMapInternal(raw); - final id = message['id']?.toString(); - if (id != null && message.containsKey('result')) { - final completer = pendingRequestsInternal.remove(id); - if (completer != null && !completer.isCompleted) { - completer.complete(asMapInternal(message['result'])); - } - return; - } - if (id != null && message.containsKey('error')) { - final completer = pendingRequestsInternal.remove(id); - if (completer != null && !completer.isCompleted) { - final error = asMapInternal(message['error']); - completer.completeError( - StateError( - error['message']?.toString() ?? - 'Single-agent app-server request failed.', - ), - ); - } - return; - } - if (message.containsKey('method')) { - notificationsInternal.add(message); - } - }, - onError: (Object error, StackTrace stackTrace) { - for (final completer in pendingRequestsInternal.values) { - if (!completer.isCompleted) { - completer.completeError(error); - } - } - pendingRequestsInternal.clear(); - notificationsInternal.addError(error, stackTrace); - }, - onDone: () { - final error = StateError( - 'Single-agent app-server websocket closed unexpectedly.', - ); - for (final completer in pendingRequestsInternal.values) { - if (!completer.isCompleted) { - completer.completeError(error); - } - } - pendingRequestsInternal.clear(); - if (!notificationsInternal.isClosed) { - unawaited(notificationsInternal.close()); - } - }, - cancelOnError: true, - ); - } - - Future close() async { - await subscriptionInternal?.cancel(); - subscriptionInternal = null; - for (final completer in pendingRequestsInternal.values) { - if (!completer.isCompleted) { - completer.completeError( - StateError('Single-agent app-server connection closed.'), - ); - } - } - pendingRequestsInternal.clear(); - if (!notificationsInternal.isClosed) { - await notificationsInternal.close(); - } - try { - await socketInternal.close(); - } catch (_) { - // Best effort only. - } - } -} diff --git a/lib/runtime/single_agent_runner.dart b/lib/runtime/single_agent_runner.dart index ddf85299..965961a9 100644 --- a/lib/runtime/single_agent_runner.dart +++ b/lib/runtime/single_agent_runner.dart @@ -1,227 +1,4 @@ -import 'direct_single_agent_app_server_client.dart'; -import 'runtime_models.dart'; - -class SingleAgentProviderResolution { - const SingleAgentProviderResolution({ - required this.selection, - required this.resolvedProvider, - required this.fallbackReason, - }); - - final SingleAgentProvider selection; - final SingleAgentProvider? resolvedProvider; - final String? fallbackReason; -} - -class SingleAgentRunRequest { - const SingleAgentRunRequest({ - required this.sessionId, - required this.provider, - required this.prompt, - required this.model, - required this.workingDirectory, - required this.gatewayToken, - required this.attachments, - required this.selectedSkills, - required this.aiGatewayBaseUrl, - required this.aiGatewayApiKey, - required this.config, - this.onOutput, - this.configuredCodexCliPath = '', - }); - - final String sessionId; - final SingleAgentProvider provider; - final String prompt; - final String model; - final String workingDirectory; - final String gatewayToken; - final List attachments; - final List selectedSkills; - final String aiGatewayBaseUrl; - final String aiGatewayApiKey; - final MultiAgentConfig config; - final void Function(String text)? onOutput; - final String configuredCodexCliPath; -} - -class SingleAgentRunResult { - const SingleAgentRunResult({ - required this.provider, - required this.output, - required this.success, - required this.errorMessage, - required this.shouldFallbackToAiChat, - this.aborted = false, - this.fallbackReason, - this.resolvedModel = '', - this.resolvedWorkingDirectory = '', - }); - - final SingleAgentProvider provider; - final String output; - final bool success; - final String errorMessage; - final bool shouldFallbackToAiChat; - final bool aborted; - final String? fallbackReason; - final String resolvedModel; - final String resolvedWorkingDirectory; -} - -abstract class SingleAgentRunner { - Future resolveProvider({ - required SingleAgentProvider selection, - required List availableProviders, - required String configuredCodexCliPath, - required String gatewayToken, - }); - - Future run(SingleAgentRunRequest request); - - Future abort(String sessionId); -} - -class DefaultSingleAgentRunner implements SingleAgentRunner { - DefaultSingleAgentRunner({ - required DirectSingleAgentAppServerClient appServerClient, - }) : _appServerClient = appServerClient; - - final DirectSingleAgentAppServerClient _appServerClient; - - @override - Future resolveProvider({ - required SingleAgentProvider selection, - required List availableProviders, - required String configuredCodexCliPath, - required String gatewayToken, - }) async { - try { - if (selection != SingleAgentProvider.auto) { - final capabilities = await _appServerClient.loadCapabilities( - provider: selection, - gatewayToken: gatewayToken, - ); - if (!capabilities.available || - !capabilities.supportsProvider(selection)) { - return SingleAgentProviderResolution( - selection: selection, - resolvedProvider: null, - fallbackReason: - capabilities.errorMessage ?? - '${selection.label} endpoint is unavailable.', - ); - } - return SingleAgentProviderResolution( - selection: selection, - resolvedProvider: selection, - fallbackReason: null, - ); - } - - String? fallbackReason; - for (final provider in availableProviders) { - if (provider == SingleAgentProvider.auto) { - continue; - } - final capabilities = await _appServerClient.loadCapabilities( - provider: provider, - gatewayToken: gatewayToken, - ); - if (capabilities.available && capabilities.supportsProvider(provider)) { - return SingleAgentProviderResolution( - selection: selection, - resolvedProvider: provider, - fallbackReason: null, - ); - } - fallbackReason ??= capabilities.errorMessage; - } - return SingleAgentProviderResolution( - selection: selection, - resolvedProvider: null, - fallbackReason: - fallbackReason ?? - 'No external ACP endpoint is currently available.', - ); - } catch (error) { - return SingleAgentProviderResolution( - selection: selection, - resolvedProvider: null, - fallbackReason: 'Single-agent app-server negotiation failed: $error', - ); - } - } - - @override - Future run(SingleAgentRunRequest request) async { - try { - final result = await _appServerClient.run( - DirectSingleAgentRunRequest( - sessionId: request.sessionId, - provider: request.provider, - prompt: _augmentPrompt(request), - model: request.model, - workingDirectory: request.workingDirectory, - gatewayToken: request.gatewayToken, - selectedSkills: request.selectedSkills, - onOutput: request.onOutput, - ), - ); - return SingleAgentRunResult( - provider: request.provider, - output: result.output, - success: result.success, - errorMessage: result.errorMessage, - shouldFallbackToAiChat: !result.success && result.output.isEmpty, - aborted: result.aborted, - resolvedModel: result.resolvedModel, - resolvedWorkingDirectory: result.resolvedWorkingDirectory, - fallbackReason: !result.success - ? 'Single-agent app-server run failed: ${result.errorMessage}' - : null, - ); - } catch (error) { - final shouldFallback = _shouldFallbackToAiChat(error.toString()); - return SingleAgentRunResult( - provider: request.provider, - output: '', - success: false, - errorMessage: error.toString(), - shouldFallbackToAiChat: shouldFallback, - resolvedModel: '', - fallbackReason: shouldFallback - ? '${request.provider.label} provider is unavailable from the direct app-server endpoint.' - : null, - ); - } - } - - @override - Future abort(String sessionId) async { - final normalized = sessionId.trim(); - if (normalized.isEmpty) { - return; - } - await _appServerClient.abort(normalized); - } - - bool _shouldFallbackToAiChat(String message) { - final normalizedMessage = message.toLowerCase(); - return normalizedMessage.contains('timeout') || - normalizedMessage.contains('unavailable') || - normalizedMessage.contains('missing') || - normalizedMessage.contains('closed') || - normalizedMessage.contains('connect'); - } - - String _augmentPrompt(SingleAgentRunRequest request) { - if (request.attachments.isEmpty) { - return request.prompt; - } - final attachmentLines = request.attachments - .map((item) => '- ${item.name}: ${item.path}') - .join('\n'); - return 'User-selected local attachments:\n$attachmentLines\n\n${request.prompt}'; - } -} +// Legacy compatibility shim retained until remaining imports are cleaned up. +// +// Single-agent execution now flows through GoAgentCoreClient and the ACP +// transport; the previous direct runner no longer owns runtime strategy. diff --git a/test/runtime/app_controller_ai_gateway_chat_suite.dart b/test/runtime/app_controller_ai_gateway_chat_suite.dart index a36a334f..68f8ebe7 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite.dart @@ -15,7 +15,6 @@ import 'package:xworkmate/runtime/gateway_runtime.dart'; import 'package:xworkmate/runtime/runtime_coordinator.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; -import 'package:xworkmate/runtime/single_agent_runner.dart'; import 'app_controller_ai_gateway_chat_suite_core.dart'; import 'app_controller_ai_gateway_chat_suite_chat.dart'; import 'app_controller_ai_gateway_chat_suite_single_agent.dart'; diff --git a/test/runtime/app_controller_ai_gateway_chat_suite_core.dart b/test/runtime/app_controller_ai_gateway_chat_suite_core.dart index d31acbfd..8ec685ed 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite_core.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite_core.dart @@ -12,7 +12,6 @@ import 'package:xworkmate/runtime/gateway_runtime.dart'; import 'package:xworkmate/runtime/runtime_coordinator.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; -import 'package:xworkmate/runtime/single_agent_runner.dart'; import 'app_controller_ai_gateway_chat_suite_chat.dart'; import 'app_controller_ai_gateway_chat_suite_single_agent.dart'; import 'app_controller_ai_gateway_chat_suite_fakes.dart'; diff --git a/test/runtime/app_controller_ai_gateway_chat_suite_fakes.dart b/test/runtime/app_controller_ai_gateway_chat_suite_fakes.dart index e15e6dd8..f64e079e 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite_fakes.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite_fakes.dart @@ -13,7 +13,6 @@ import 'package:xworkmate/runtime/go_agent_core_client.dart'; import 'package:xworkmate/runtime/runtime_coordinator.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; -import 'package:xworkmate/runtime/single_agent_runner.dart'; import 'app_controller_ai_gateway_chat_suite_core.dart'; import 'app_controller_ai_gateway_chat_suite_chat.dart'; import 'app_controller_ai_gateway_chat_suite_single_agent.dart'; @@ -113,71 +112,6 @@ class FakeCodexRuntimeInternal extends CodexRuntime { Future stop() async {} } -class FakeSingleAgentRunnerInternal implements SingleAgentRunner { - FakeSingleAgentRunnerInternal({ - required this.resolvedProvider, - this.result, - this.fallbackReason, - }); - - final SingleAgentProvider? resolvedProvider; - final SingleAgentRunResult? result; - final String? fallbackReason; - - int resolveCalls = 0; - int runCalls = 0; - int abortCalls = 0; - SingleAgentRunRequest? lastRequest; - final List requests = []; - - @override - Future resolveProvider({ - required SingleAgentProvider selection, - required List availableProviders, - required String configuredCodexCliPath, - required String gatewayToken, - }) async { - resolveCalls += 1; - return SingleAgentProviderResolution( - selection: selection, - resolvedProvider: resolvedProvider, - fallbackReason: fallbackReason, - ); - } - - @override - Future run(SingleAgentRunRequest request) async { - runCalls += 1; - lastRequest = request; - requests.add(request); - if (result?.output.isNotEmpty == true) { - request.onOutput?.call(result!.output); - } - return result ?? - SingleAgentRunResult( - provider: request.provider, - output: '', - success: false, - errorMessage: 'no result configured', - shouldFallbackToAiChat: false, - ); - } - - @override - Future abort(String sessionId) async { - abortCalls += 1; - } -} - -class FallbackOnlySingleAgentRunnerInternal - extends FakeSingleAgentRunnerInternal { - FallbackOnlySingleAgentRunnerInternal() - : super( - resolvedProvider: null, - fallbackReason: 'No supported external CLI provider is available.', - ); -} - class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { FakeGoAgentCoreClientInternal({ this.capabilities = const GoAgentCoreCapabilities.empty(), @@ -198,7 +132,8 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { int executeCalls = 0; int cancelCalls = 0; GoAgentCoreSessionRequest? lastRequest; - final List requests = []; + final List requests = + []; @override Future loadCapabilities({ diff --git a/test/runtime/direct_single_agent_app_server_suite.dart b/test/runtime/direct_single_agent_app_server_suite.dart deleted file mode 100644 index 6bd3db8b..00000000 --- a/test/runtime/direct_single_agent_app_server_suite.dart +++ /dev/null @@ -1,937 +0,0 @@ -@TestOn('vm') -library; - -import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; - -import 'package:flutter_test/flutter_test.dart'; -import 'package:xworkmate/runtime/direct_single_agent_app_server_client.dart'; -import 'package:xworkmate/runtime/runtime_models.dart'; - -void main() { - group('DirectSingleAgentAppServerClient', () { - test( - 'direct single-agent app-server core file stays split into focused parts', - () { - final lines = File( - 'lib/runtime/direct_single_agent_app_server_client_core.dart', - ).readAsLinesSync(); - - expect( - lines.length, - lessThanOrEqualTo(1000), - reason: 'The core file should stay under the target line budget.', - ); - }, - ); - - test('classifies the four endpoint modes', () { - expect( - DirectSingleAgentEndpointDescriptor.describe( - Uri.parse('ws://127.0.0.1:9001'), - ).mode, - DirectSingleAgentEndpointMode.wsLocal, - ); - expect( - DirectSingleAgentEndpointDescriptor.describe( - Uri.parse('wss://agent.example.com'), - ).mode, - DirectSingleAgentEndpointMode.wss, - ); - expect( - DirectSingleAgentEndpointDescriptor.describe( - Uri.parse('http://localhost:38992'), - ).mode, - DirectSingleAgentEndpointMode.httpLocal, - ); - expect( - DirectSingleAgentEndpointDescriptor.describe( - Uri.parse('https://agent.example.com'), - ).mode, - DirectSingleAgentEndpointMode.https, - ); - }); - - test('probes websocket endpoint and reports provider support', () async { - final server = await FakeAppServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - - final capabilities = await client.loadCapabilities( - provider: SingleAgentProvider.opencode, - ); - - expect(capabilities.available, isTrue); - expect( - capabilities.supportsProvider(SingleAgentProvider.opencode), - isTrue, - ); - expect(capabilities.endpoint, 'ws://127.0.0.1:${server.port}'); - expect(server.methods, contains('initialize')); - }); - - test('runs single-agent turns over direct websocket app-server', () async { - final server = await FakeAppServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final deltas = []; - final result = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-1', - provider: SingleAgentProvider.opencode, - prompt: 'hello world', - model: 'gpt-4.1', - workingDirectory: '/tmp', - gatewayToken: 'token-1', - ).copyWith(onOutput: deltas.add), - ); - - expect(result.success, isTrue, reason: result.errorMessage); - expect(result.output, 'hello world from app server'); - expect(result.resolvedModel, 'codex-sonnet'); - expect(result.resolvedWorkingDirectory, '/tmp'); - expect(result.resolvedWorkspaceRefKind, WorkspaceRefKind.localPath); - expect(server.lastTurnInput, [ - {'type': 'text', 'text': 'hello world'}, - ]); - expect(deltas.join(), 'hello world from app server'); - expect( - server.methods, - containsAll(['initialize', 'thread/start', 'turn/start']), - ); - expect(server.authorizationHeaders, contains('Bearer token-1')); - }); - - test( - 'starts a new websocket thread when working directory changes for a session', - () async { - final server = await FakeAppServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final first = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-cwd-change', - provider: SingleAgentProvider.opencode, - prompt: 'first turn', - model: 'gpt-4.1', - workingDirectory: '/tmp/a', - gatewayToken: '', - ), - ); - final second = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-cwd-change', - provider: SingleAgentProvider.opencode, - prompt: 'second turn', - model: 'gpt-4.1', - workingDirectory: '/tmp/b', - gatewayToken: '', - ), - ); - - expect(first.success, isTrue, reason: first.errorMessage); - expect(second.success, isTrue, reason: second.errorMessage); - expect(second.resolvedWorkingDirectory, '/tmp/b'); - expect( - server.methods.where((method) => method == 'thread/start').length, - 2, - ); - }, - ); - - test('sends selected skills as structured app-server inputs', () async { - final server = await FakeAppServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final result = await client.run( - DirectSingleAgentRunRequest( - sessionId: 'session-skills', - provider: SingleAgentProvider.opencode, - prompt: 'use the selected skills', - model: 'gpt-4.1', - workingDirectory: '/tmp', - gatewayToken: '', - selectedSkills: const [ - AssistantThreadSkillEntry( - key: '/tmp/ppt', - label: 'PPT', - description: 'Slides', - source: 'codex', - sourcePath: '/tmp/ppt/SKILL.md', - scope: 'user', - sourceLabel: 'codex · user · ppt', - ), - AssistantThreadSkillEntry( - key: '/tmp/browser', - label: 'Browser Automation', - description: 'Browser', - source: 'agents', - sourcePath: '/tmp/browser/SKILL.md', - scope: 'user', - sourceLabel: 'agents · user · browser', - ), - ], - ), - ); - - expect(result.success, isTrue); - expect(server.lastTurnInput, [ - {'type': 'text', 'text': 'use the selected skills'}, - { - 'type': 'skill', - 'name': 'PPT', - 'path': '/tmp/ppt/SKILL.md', - }, - { - 'type': 'skill', - 'name': 'Browser Automation', - 'path': '/tmp/browser/SKILL.md', - }, - ]); - }); - - test('interrupts active turns on abort', () async { - final server = await FakeAppServerInternal.start(delayCompletion: true); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final runFuture = client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-abort', - provider: SingleAgentProvider.opencode, - prompt: 'abort me', - model: 'gpt-4.1', - workingDirectory: '/tmp', - gatewayToken: '', - ), - ); - - await server.waitForMethod('turn/start'); - await client.abort('session-abort'); - final result = await runFuture; - - expect(result.aborted, isTrue); - expect(server.methods, contains('turn/interrupt')); - }); - - test( - 'accepts nested thread objects returned by codex app-server', - () async { - final server = await FakeAppServerInternal.start( - nestedThreadResult: true, - ); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final result = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-nested', - provider: SingleAgentProvider.opencode, - prompt: 'hello nested world', - model: 'qwen2.5-coder:latest', - workingDirectory: '/tmp', - gatewayToken: '', - ), - ); - - expect(result.success, isTrue); - expect(result.output, 'hello world from app server'); - expect(result.resolvedModel, 'codex-sonnet'); - expect(result.resolvedWorkingDirectory, '/tmp'); - expect(result.resolvedWorkspaceRefKind, WorkspaceRefKind.localPath); - }, - ); - - test('captures the resolved thread path returned by app-server', () async { - final server = await FakeAppServerInternal.start( - resolvedThreadPath: '/tmp/app-server-thread', - ); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final result = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-thread-path', - provider: SingleAgentProvider.opencode, - prompt: 'hello thread path', - model: 'gpt-4.1', - workingDirectory: '/tmp/requested-thread', - gatewayToken: '', - ), - ); - - expect(result.success, isTrue); - expect(result.resolvedWorkingDirectory, '/tmp/app-server-thread'); - expect(result.resolvedWorkspaceRefKind, WorkspaceRefKind.localPath); - }); - - test( - 'probes OpenCode REST endpoint and reports provider support', - () async { - final server = await FakeOpenCodeRestServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - - final capabilities = await client.loadCapabilities( - provider: SingleAgentProvider.opencode, - ); - - expect(capabilities.available, isTrue); - expect( - capabilities.supportsProvider(SingleAgentProvider.opencode), - isTrue, - ); - expect(server.healthRequested, isTrue); - }, - ); - - test('runs OpenCode turns over REST session api', () async { - final server = await FakeOpenCodeRestServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final deltas = []; - final result = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-opencode', - provider: SingleAgentProvider.opencode, - prompt: 'hello opencode', - model: '', - workingDirectory: '/tmp', - gatewayToken: '', - ).copyWith(onOutput: deltas.add), - ); - - expect(result.success, isTrue); - expect(result.output, 'hello world from opencode'); - expect(deltas.join(), 'hello world from opencode'); - expect(server.createdSessionCount, 1); - expect(server.lastPromptText, 'hello opencode'); - }); - - test( - 'creates a new REST session when working directory changes for a session', - () async { - final server = await FakeOpenCodeRestServerInternal.start(); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final first = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-opencode-cwd-change', - provider: SingleAgentProvider.opencode, - prompt: 'first', - model: '', - workingDirectory: '/tmp/a', - gatewayToken: '', - ), - ); - final second = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-opencode-cwd-change', - provider: SingleAgentProvider.opencode, - prompt: 'second', - model: '', - workingDirectory: '/tmp/b', - gatewayToken: '', - ), - ); - - expect(first.success, isTrue, reason: first.errorMessage); - expect(second.success, isTrue, reason: second.errorMessage); - expect(server.createdSessionCount, 2); - }, - ); - - test( - 'fails OpenCode REST turns that complete without assistant content', - () async { - final server = await FakeOpenCodeRestServerInternal.start( - emitAssistantContent: false, - ); - addTearDown(server.close); - - final client = DirectSingleAgentAppServerClient( - endpointResolver: (_) => server.baseHttpUri, - ); - addTearDown(client.dispose); - - final result = await client.run( - const DirectSingleAgentRunRequest( - sessionId: 'session-opencode-empty', - provider: SingleAgentProvider.opencode, - prompt: 'hello opencode', - model: '', - workingDirectory: '/tmp', - gatewayToken: '', - ), - ); - - expect(result.success, isFalse); - expect(result.output, isEmpty); - expect(result.errorMessage, contains('without assistant content')); - }, - ); - }); -} - -class FakeAppServerInternal { - FakeAppServerInternal._( - this.serverInternal, { - required this.delayCompletion, - required this.nestedThreadResult, - required this.resolvedThreadPath, - }); - - final HttpServer serverInternal; - final bool delayCompletion; - final bool nestedThreadResult; - final String? resolvedThreadPath; - final List methods = []; - final List authorizationHeaders = []; - final Map> methodWaitersInternal = - >{}; - int threadCounterInternal = 0; - List? lastTurnInput; - - int get port => serverInternal.port; - Uri get baseHttpUri => Uri.parse('http://127.0.0.1:${serverInternal.port}'); - - static Future start({ - bool delayCompletion = false, - bool nestedThreadResult = false, - String? resolvedThreadPath, - }) async { - final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); - final fake = FakeAppServerInternal._( - server, - delayCompletion: delayCompletion, - nestedThreadResult: nestedThreadResult, - resolvedThreadPath: resolvedThreadPath, - ); - unawaited(fake.listenInternal()); - return fake; - } - - Future close() async { - await serverInternal.close(force: true); - } - - Future waitForMethod(String method) async { - if (methods.contains(method)) { - return; - } - final completer = methodWaitersInternal.putIfAbsent( - method, - Completer.new, - ); - await completer.future.timeout(const Duration(seconds: 3)); - } - - Future listenInternal() async { - await for (final request in serverInternal) { - authorizationHeaders.add( - request.headers.value(HttpHeaders.authorizationHeader) ?? '', - ); - if (request.uri.path == '/' && - WebSocketTransformer.isUpgradeRequest(request)) { - final socket = await WebSocketTransformer.upgrade(request); - unawaited(handleSocketInternal(socket)); - continue; - } - request.response.statusCode = HttpStatus.notFound; - await request.response.close(); - } - } - - Future handleSocketInternal(WebSocket socket) async { - await for (final raw in socket) { - final message = decodeMapInternal(raw); - final method = message['method']?.toString() ?? ''; - final id = message['id']; - final params = asMapInternal(message['params']); - if (method.isEmpty) { - continue; - } - methods.add(method); - methodWaitersInternal.remove(method)?.complete(); - switch (method) { - case 'initialize': - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'result': { - 'serverInfo': {'name': 'fake-codex'}, - }, - }), - ); - break; - case 'initialized': - break; - case 'thread/start': - threadCounterInternal += 1; - final threadPath = resolvedThreadPath ?? params['cwd'] ?? '/tmp'; - final result = nestedThreadResult - ? { - 'thread': { - 'id': 'thread-$threadCounterInternal', - 'path': threadPath, - 'ephemeral': false, - }, - } - : { - 'id': 'thread-$threadCounterInternal', - 'path': threadPath, - 'ephemeral': false, - }; - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'result': result, - }), - ); - break; - case 'thread/resume': - final threadPath = resolvedThreadPath ?? params['cwd'] ?? '/tmp'; - final result = nestedThreadResult - ? { - 'thread': { - 'id': params['threadId'] ?? 'thread-resumed', - 'path': threadPath, - 'ephemeral': false, - }, - } - : { - 'id': params['threadId'] ?? 'thread-resumed', - 'path': threadPath, - 'ephemeral': false, - }; - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'result': result, - }), - ); - break; - case 'turn/start': - final threadId = params['threadId']?.toString() ?? 'thread-1'; - if (params.containsKey('userInput')) { - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'error': { - 'code': -32600, - 'message': 'Invalid request: missing field `input`', - }, - }), - ); - break; - } - final input = params['input']; - if (input is! List) { - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'error': { - 'code': -32600, - 'message': - 'Invalid request: invalid type: expected a sequence', - }, - }), - ); - break; - } - lastTurnInput = List.from(input); - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'result': { - 'id': 'turn-1', - 'threadId': threadId, - 'status': 'started', - 'model': 'codex-sonnet', - }, - }), - ); - unawaited(emitTurnInternal(socket, threadId)); - break; - case 'turn/interrupt': - final threadId = params['threadId']?.toString() ?? 'thread-1'; - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'result': {'ok': true}, - }), - ); - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'method': 'turn/error', - 'params': { - 'threadId': threadId, - 'message': 'aborted', - }, - }), - ); - await socket.close(); - break; - default: - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': id, - 'error': { - 'code': -32601, - 'message': 'unknown method $method', - }, - }), - ); - } - } - } - - Future emitTurnInternal(WebSocket socket, String threadId) async { - const parts = ['hello ', 'world ', 'from app server']; - for (final part in parts) { - try { - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'method': 'item/agentMessage/delta', - 'params': { - 'threadId': threadId, - 'turnId': 'turn-1', - 'delta': part, - }, - }), - ); - } catch (_) { - return; - } - await Future.delayed(const Duration(milliseconds: 5)); - } - if (delayCompletion) { - return; - } - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'method': 'turn/completed', - 'params': {'threadId': threadId, 'turnId': 'turn-1'}, - }), - ); - } -} - -class FakeOpenCodeRestServerInternal { - FakeOpenCodeRestServerInternal._( - this.serverInternal, { - required this.emitAssistantContent, - }); - - final HttpServer serverInternal; - final bool emitAssistantContent; - final List eventResponsesInternal = []; - var sessionCounterInternal = 0; - var messageCounterInternal = 0; - bool healthRequested = false; - int createdSessionCount = 0; - String lastPromptText = ''; - final Map assistantTextBySessionInternal = {}; - - Uri get baseHttpUri => Uri.parse('http://127.0.0.1:${serverInternal.port}'); - - static Future start({ - bool emitAssistantContent = true, - }) async { - final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); - final fake = FakeOpenCodeRestServerInternal._( - server, - emitAssistantContent: emitAssistantContent, - ); - unawaited(fake.listenInternal()); - return fake; - } - - Future close() async { - for (final response in eventResponsesInternal.toList(growable: false)) { - try { - await response.close(); - } catch (_) { - // Best effort. - } - } - await serverInternal.close(force: true); - } - - Future listenInternal() async { - await for (final request in serverInternal) { - if (request.uri.path == '/global/health') { - healthRequested = true; - request.response.headers.contentType = ContentType.json; - request.response.write( - jsonEncode({'healthy': true, 'version': '1.3.3'}), - ); - await request.response.close(); - continue; - } - if (request.uri.path == '/global/event') { - request.response.headers.set( - HttpHeaders.contentTypeHeader, - 'text/event-stream', - ); - request.response.headers.set( - HttpHeaders.cacheControlHeader, - 'no-cache', - ); - request.response.write( - 'data: ${jsonEncode({ - 'payload': {'type': 'server.connected', 'properties': {}}, - })}\n\n', - ); - await request.response.flush(); - eventResponsesInternal.add(request.response); - continue; - } - if (request.uri.path == '/session' && request.method == 'POST') { - createdSessionCount += 1; - final sessionId = 'ses-${sessionCounterInternal++}'; - request.response.headers.contentType = ContentType.json; - request.response.write( - jsonEncode({ - 'id': sessionId, - 'title': 'test', - 'directory': - request.uri.queryParameters['directory'] ?? - Directory.current.path, - }), - ); - await request.response.close(); - continue; - } - final sessionMatch = RegExp( - r'^/session/([^/]+)/message$', - ).firstMatch(request.uri.path); - if (sessionMatch != null && request.method == 'GET') { - final sessionId = sessionMatch.group(1)!; - final text = assistantTextBySessionInternal[sessionId] ?? ''; - request.response.headers.contentType = ContentType.json; - request.response.write( - jsonEncode(>[ - { - 'info': {'id': 'msg-user', 'role': 'user'}, - 'parts': >[ - {'type': 'text', 'text': lastPromptText}, - ], - }, - if (text.isNotEmpty) - { - 'info': { - 'id': 'msg-assistant', - 'role': 'assistant', - }, - 'parts': >[ - {'type': 'text', 'text': text}, - ], - }, - ]), - ); - await request.response.close(); - continue; - } - if (sessionMatch != null && request.method == 'POST') { - final sessionId = sessionMatch.group(1)!; - final body = jsonDecode(await utf8.decodeStream(request)); - final parts = - (body as Map)['parts'] as List? ?? - const []; - if (parts.isNotEmpty) { - lastPromptText = - (parts.first as Map)['text']?.toString() ?? ''; - } - final assistantMessageId = 'msg-assistant-${messageCounterInternal++}'; - await broadcastEventInternal({ - 'payload': { - 'type': 'session.status', - 'properties': { - 'sessionID': sessionId, - 'status': {'type': 'busy'}, - }, - }, - }); - await broadcastEventInternal({ - 'payload': { - 'type': 'message.updated', - 'properties': { - 'sessionID': sessionId, - 'info': { - 'id': assistantMessageId, - 'role': 'assistant', - }, - }, - }, - }); - if (emitAssistantContent) { - for (final delta in [ - 'hello ', - 'world ', - 'from ', - 'opencode', - ]) { - await broadcastEventInternal({ - 'payload': { - 'type': 'message.part.delta', - 'properties': { - 'sessionID': sessionId, - 'part': {'messageID': assistantMessageId}, - 'text': delta, - }, - }, - }); - } - await broadcastEventInternal({ - 'payload': { - 'type': 'message.part.updated', - 'properties': { - 'sessionID': sessionId, - 'part': { - 'messageID': assistantMessageId, - 'type': 'text', - 'text': 'hello world from opencode', - }, - }, - }, - }); - assistantTextBySessionInternal[sessionId] = - 'hello world from opencode'; - } - await broadcastEventInternal({ - 'payload': { - 'type': 'session.status', - 'properties': { - 'sessionID': sessionId, - 'status': {'type': 'idle'}, - }, - }, - }); - request.response.headers.contentType = ContentType.json; - request.response.write(''); - await request.response.close(); - continue; - } - final abortMatch = RegExp( - r'^/session/([^/]+)/abort$', - ).firstMatch(request.uri.path); - if (abortMatch != null && request.method == 'POST') { - request.response.headers.contentType = ContentType.json; - request.response.write('{}'); - await request.response.close(); - continue; - } - request.response.statusCode = HttpStatus.notFound; - await request.response.close(); - } - } - - Future broadcastEventInternal(Map event) async { - final payload = 'data: ${jsonEncode(event)}\n\n'; - for (final response in eventResponsesInternal.toList(growable: false)) { - response.write(payload); - await response.flush(); - } - } -} - -Map decodeMapInternal(Object raw) { - if (raw is Map) { - return raw; - } - if (raw is Map) { - return raw.cast(); - } - final decoded = jsonDecode(raw.toString()); - if (decoded is Map) { - return decoded; - } - if (decoded is Map) { - return decoded.cast(); - } - return const {}; -} - -Map asMapInternal(Object? value) { - if (value is Map) { - return value; - } - if (value is Map) { - return value.cast(); - } - return const {}; -} - -extension on DirectSingleAgentRunRequest { - DirectSingleAgentRunRequest copyWith({ - void Function(String text)? onOutput, - List? selectedSkills, - }) { - return DirectSingleAgentRunRequest( - sessionId: sessionId, - provider: provider, - prompt: prompt, - model: model, - workingDirectory: workingDirectory, - gatewayToken: gatewayToken, - selectedSkills: selectedSkills ?? this.selectedSkills, - onOutput: onOutput ?? this.onOutput, - ); - } -} diff --git a/test/runtime/go_agent_core_client_suite.dart b/test/runtime/go_agent_core_client_suite.dart new file mode 100644 index 00000000..2fb993f3 --- /dev/null +++ b/test/runtime/go_agent_core_client_suite.dart @@ -0,0 +1,123 @@ +@TestOn('vm') +library; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:xworkmate/runtime/go_agent_core_client.dart'; +import 'package:xworkmate/runtime/runtime_models.dart'; + +void main() { + group('GoAgentCore client mapping', () { + test('session request maps skills, attachments, and provider into ACP', () { + const request = GoAgentCoreSessionRequest( + sessionId: 'session-1', + threadId: 'thread-1', + target: AssistantExecutionTarget.singleAgent, + prompt: 'hello world', + workingDirectory: '/tmp/workspace', + model: 'codex-sonnet', + thinking: 'medium', + selectedSkills: ['PPT', 'Browser Automation'], + inlineAttachments: [ + GatewayChatAttachmentPayload( + type: 'inline', + fileName: 'note.txt', + mimeType: 'text/plain', + content: 'aGVsbG8=', + ), + ], + localAttachments: [ + CollaborationAttachment( + name: 'spec.md', + path: '/tmp/workspace/spec.md', + description: 'workspace spec', + ), + ], + aiGatewayBaseUrl: 'https://gateway.example.com', + aiGatewayApiKey: 'secret', + agentId: '', + metadata: {}, + provider: SingleAgentProvider.opencode, + ); + + final params = request.toAcpParams(); + + expect(params['sessionId'], 'session-1'); + expect(params['threadId'], 'thread-1'); + expect(params['mode'], 'single-agent'); + expect(params['workingDirectory'], '/tmp/workspace'); + expect(params['provider'], 'opencode'); + expect(params['model'], 'codex-sonnet'); + expect(params['thinking'], 'medium'); + expect(params['selectedSkills'], ['PPT', 'Browser Automation']); + expect(params['attachments'], >[ + { + 'name': 'spec.md', + 'description': 'workspace spec', + 'path': '/tmp/workspace/spec.md', + }, + { + 'name': 'note.txt', + 'description': 'text/plain', + 'path': '', + }, + ]); + expect(params['inlineAttachments'], >[ + { + 'name': 'note.txt', + 'mimeType': 'text/plain', + 'content': 'aGVsbG8=', + 'sizeBytes': 5, + }, + ]); + }); + + test( + 'run result prefers completion text and preserves resolved workspace', + () { + final result = goAgentCoreRunResultFromResponse( + { + 'result': { + 'success': true, + 'turnId': 'turn-7', + 'summary': 'summary text', + 'resolvedModel': 'codex-sonnet', + 'resolvedWorkingDirectory': '/tmp/thread', + 'resolvedWorkspaceRefKind': 'remotePath', + }, + }, + streamedText: 'partial output', + completedMessage: 'final output', + ); + + expect(result.success, isTrue); + expect(result.turnId, 'turn-7'); + expect(result.message, 'final output'); + expect(result.resolvedModel, 'codex-sonnet'); + expect(result.resolvedWorkingDirectory, '/tmp/thread'); + expect(result.resolvedWorkspaceRefKind, WorkspaceRefKind.remotePath); + }, + ); + + test('session update recognizes delta notifications', () { + final update = goAgentCoreUpdateFromNotification({ + 'method': 'session.update', + 'params': { + 'sessionId': 'session-2', + 'threadId': 'thread-2', + 'turnId': 'turn-2', + 'type': 'delta', + 'delta': 'hello', + 'pending': true, + }, + }); + + expect(update, isNotNull); + expect(update!.sessionId, 'session-2'); + expect(update.threadId, 'thread-2'); + expect(update.turnId, 'turn-2'); + expect(update.isDelta, isTrue); + expect(update.text, 'hello'); + expect(update.pending, isTrue); + }); + }); +} diff --git a/test/runtime/direct_single_agent_app_server_test.dart b/test/runtime/go_agent_core_client_test.dart similarity index 53% rename from test/runtime/direct_single_agent_app_server_test.dart rename to test/runtime/go_agent_core_client_test.dart index 672872a8..b0eef63d 100644 --- a/test/runtime/direct_single_agent_app_server_test.dart +++ b/test/runtime/go_agent_core_client_test.dart @@ -1,5 +1,5 @@ import '../test_suite_stub.dart' - if (dart.library.io) 'direct_single_agent_app_server_suite.dart' + if (dart.library.io) 'go_agent_core_client_suite.dart' as suite; void main() { diff --git a/test/runtime/no_direct_cli_execution_guard_suite.dart b/test/runtime/no_direct_cli_execution_guard_suite.dart index dfbe0f81..12a5c1fc 100644 --- a/test/runtime/no_direct_cli_execution_guard_suite.dart +++ b/test/runtime/no_direct_cli_execution_guard_suite.dart @@ -19,7 +19,7 @@ void main() { ]; const guardedFiles = [ 'lib/app/app_controller_desktop.dart', - 'lib/runtime/single_agent_runner.dart', + 'lib/runtime/go_agent_core_client.dart', 'lib/runtime/runtime_coordinator.dart', 'lib/runtime/gateway_acp_client.dart', ]; @@ -53,5 +53,27 @@ void main() { } }, ); + + test('legacy direct single-agent runtime implementation stays removed', () { + const removedFiles = [ + 'lib/runtime/direct_single_agent_app_server_client_core.dart', + 'lib/runtime/direct_single_agent_app_server_client_helpers.dart', + 'lib/runtime/direct_single_agent_app_server_client_transport.dart', + ]; + + for (final relativePath in removedFiles) { + expect( + File(relativePath).existsSync(), + isFalse, + reason: '$relativePath should stay removed after GoAgentCore cutover', + ); + } + + final runnerShim = File('lib/runtime/single_agent_runner.dart'); + expect(runnerShim.existsSync(), isTrue); + final shimContent = runnerShim.readAsStringSync(); + expect(shimContent.contains('DefaultSingleAgentRunner'), isFalse); + expect(shimContent.contains('DirectSingleAgentAppServerClient'), isFalse); + }); }); }