From fbbf4b83754779fce0d8cad8de7abf5ed1629d66 Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Mon, 6 Apr 2026 07:00:49 +0800 Subject: [PATCH] Fix desktop GoAgentCore endpoint routing --- .../go_agent_core_desktop_transport.dart | 18 ++- ...go_agent_core_desktop_transport_suite.dart | 153 ++++++++++++++++++ 2 files changed, 167 insertions(+), 4 deletions(-) create mode 100644 test/runtime/go_agent_core_desktop_transport_suite.dart diff --git a/lib/runtime/go_agent_core_desktop_transport.dart b/lib/runtime/go_agent_core_desktop_transport.dart index d88fa152..f4d14f62 100644 --- a/lib/runtime/go_agent_core_desktop_transport.dart +++ b/lib/runtime/go_agent_core_desktop_transport.dart @@ -22,6 +22,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { GoCoreLocator? goCoreLocator, GoAgentCoreProcessStarter? processStarter, }) : _acpClient = acpClient, + _endpointResolver = endpointResolver, _goCoreLocator = goCoreLocator ?? GoCoreLocator(), _processStarter = processStarter ?? @@ -35,6 +36,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { }); final GatewayAcpClient _acpClient; + final Uri? Function(AssistantExecutionTarget target) _endpointResolver; final GoCoreLocator _goCoreLocator; final GoAgentCoreProcessStarter _processStarter; @@ -62,7 +64,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { required AssistantExecutionTarget target, bool forceRefresh = false, }) async { - final endpoint = await _ensureLocalEndpoint(); + final endpoint = await _resolveEndpoint(target); if (endpoint == null) { return const GoAgentCoreCapabilities.empty(); } @@ -83,7 +85,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { GoAgentCoreSessionRequest request, { required void Function(GoAgentCoreSessionUpdate update) onUpdate, }) async { - final endpoint = await _ensureLocalEndpoint(); + final endpoint = await _resolveEndpoint(request.target); if (endpoint == null) { throw const GatewayAcpException( 'Missing Go Agent-core endpoint', @@ -123,7 +125,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { required String sessionId, required String threadId, }) async { - final endpoint = await _ensureLocalEndpoint(); + final endpoint = await _resolveEndpoint(target); if (endpoint == null) { return; } @@ -140,7 +142,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { required String sessionId, required String threadId, }) async { - final endpoint = await _ensureLocalEndpoint(); + final endpoint = await _resolveEndpoint(target); if (endpoint == null) { return; } @@ -166,6 +168,14 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { } } + Future _resolveEndpoint(AssistantExecutionTarget target) async { + if (target == AssistantExecutionTarget.singleAgent || + target == AssistantExecutionTarget.auto) { + return _ensureLocalEndpoint(); + } + return _endpointResolver(target); + } + Future _ensureLocalEndpoint() async { if (_localEndpoint != null) { return _localEndpoint; diff --git a/test/runtime/go_agent_core_desktop_transport_suite.dart b/test/runtime/go_agent_core_desktop_transport_suite.dart new file mode 100644 index 00000000..43322beb --- /dev/null +++ b/test/runtime/go_agent_core_desktop_transport_suite.dart @@ -0,0 +1,153 @@ +@TestOn('vm') +library; + +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:xworkmate/runtime/gateway_acp_client.dart'; +import 'package:xworkmate/runtime/go_agent_core_client.dart'; +import 'package:xworkmate/runtime/go_agent_core_desktop_transport.dart'; +import 'package:xworkmate/runtime/runtime_models.dart'; + +void main() { + group('GoAgentCoreDesktopTransport', () { + test('uses resolved gateway endpoint for local gateway sessions', () async { + final server = await _AcpFakeServer.start(); + addTearDown(server.close); + + final transport = GoAgentCoreDesktopTransport( + acpClient: GatewayAcpClient(endpointResolver: () => null), + endpointResolver: (target) => switch (target) { + AssistantExecutionTarget.local => server.baseHttpUri, + _ => null, + }, + ); + + final result = await transport.executeSession( + const GoAgentCoreSessionRequest( + sessionId: 'session-local', + threadId: 'thread-local', + target: AssistantExecutionTarget.local, + prompt: 'ping local gateway', + workingDirectory: '/tmp', + model: '', + thinking: '', + selectedSkills: [], + inlineAttachments: [], + localAttachments: [], + aiGatewayBaseUrl: '', + aiGatewayApiKey: '', + agentId: '', + metadata: {}, + ), + onUpdate: (_) {}, + ); + + expect(result.success, isTrue); + expect(result.message, 'gateway-ok'); + expect(server.lastHttpRequestPath, '/acp/rpc'); + expect(server.rpcMethods, contains('session.start')); + expect(server.lastSessionMode, 'gateway-chat'); + }); + + test('reports missing endpoint when gateway target cannot resolve', () async { + final transport = GoAgentCoreDesktopTransport( + acpClient: GatewayAcpClient(endpointResolver: () => null), + endpointResolver: (_) => null, + ); + + await expectLater( + () => transport.executeSession( + const GoAgentCoreSessionRequest( + sessionId: 'session-local', + threadId: 'thread-local', + target: AssistantExecutionTarget.local, + prompt: 'ping local gateway', + workingDirectory: '/tmp', + model: '', + thinking: '', + selectedSkills: [], + inlineAttachments: [], + localAttachments: [], + aiGatewayBaseUrl: '', + aiGatewayApiKey: '', + agentId: '', + metadata: {}, + ), + onUpdate: (_) {}, + ), + throwsA( + isA().having( + (error) => error.code, + 'code', + 'GO_AGENT_CORE_ENDPOINT_MISSING', + ), + ), + ); + }); + }); +} + +class _AcpFakeServer { + _AcpFakeServer._(this._server); + + final HttpServer _server; + final List rpcMethods = []; + String? lastHttpRequestPath; + String? lastSessionMode; + + Uri get baseHttpUri => Uri.parse('http://127.0.0.1:${_server.port}'); + + static Future<_AcpFakeServer> start() async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + final fake = _AcpFakeServer._(server); + unawaited(fake._listen()); + return fake; + } + + Future close() async { + await _server.close(force: true); + } + + Future _listen() async { + await for (final request in _server) { + if (request.uri.path == '/acp/rpc' && request.method == 'POST') { + lastHttpRequestPath = request.uri.path; + await _handleHttpRpc(request); + continue; + } + request.response.statusCode = HttpStatus.notFound; + await request.response.close(); + } + } + + Future _handleHttpRpc(HttpRequest request) async { + final body = await utf8.decodeStream(request); + final envelope = (jsonDecode(body) as Map).cast(); + final id = envelope['id']; + final method = envelope['method']?.toString() ?? ''; + final params = + (envelope['params'] as Map?)?.cast() ?? + const {}; + rpcMethods.add(method); + + request.response.headers.set( + HttpHeaders.contentTypeHeader, + 'text/event-stream; charset=utf-8', + ); + if (method == 'session.start' || method == 'session.message') { + lastSessionMode = params['mode']?.toString(); + request.response.write( + 'data: ${jsonEncode({'jsonrpc': '2.0', 'id': id, 'result': {'success': true, 'message': 'gateway-ok', 'summary': 'gateway-ok', 'turnId': 'turn-1'}})}\n\n', + ); + await request.response.close(); + return; + } + request.response.write( + 'data: ${jsonEncode({'jsonrpc': '2.0', 'id': id, 'result': {'singleAgent': true, 'multiAgent': true, 'providers': ['codex'], 'capabilities': {'single_agent': true, 'multi_agent': true, 'providers': ['codex']}}})}\n\n', + ); + await request.response.close(); + } +}