diff --git a/go/go_core/internal/acp/stdio.go b/go/go_core/internal/acp/stdio.go new file mode 100644 index 00000000..8020aac8 --- /dev/null +++ b/go/go_core/internal/acp/stdio.go @@ -0,0 +1,95 @@ +package acp + +import ( + "bufio" + "encoding/json" + "errors" + "fmt" + "io" + "strings" + "sync" + + "xworkmate/go_core/internal/shared" +) + +func RunStdio(input io.Reader, output io.Writer) { + server := NewServer() + reader := bufio.NewReader(input) + var writeMu sync.Mutex + + writeMessage := func(message map[string]any) { + payload, _ := jsonMarshal(message) + writeMu.Lock() + defer writeMu.Unlock() + _, _ = output.Write(append(payload, '\n')) + } + + for { + payload, err := readStdioMessage(reader) + if err != nil { + if errors.Is(err, io.EOF) { + return + } + writeMessage(shared.ErrorEnvelope(nil, -32700, err.Error())) + continue + } + if len(strings.TrimSpace(string(payload))) == 0 { + continue + } + + request, err := shared.DecodeRPCRequest(payload) + if err != nil { + writeMessage(shared.ErrorEnvelope(nil, -32700, err.Error())) + continue + } + response, rpcErr := server.handleRequest(request, writeMessage) + if request.ID == nil { + continue + } + if rpcErr != nil { + writeMessage( + shared.ErrorEnvelope(request.ID, rpcErr.Code, rpcErr.Message), + ) + continue + } + writeMessage(shared.ResultEnvelope(request.ID, response)) + } +} + +func readStdioMessage(reader *bufio.Reader) ([]byte, error) { + line, err := reader.ReadString('\n') + if err != nil { + return nil, err + } + line = strings.TrimSpace(line) + if line == "" { + return nil, nil + } + if strings.HasPrefix(strings.ToLower(line), "content-length:") { + var contentLength int + if _, err := fmt.Sscanf(line, "Content-Length: %d", &contentLength); err != nil { + if _, err2 := fmt.Sscanf(line, "content-length: %d", &contentLength); err2 != nil { + return nil, fmt.Errorf("invalid content-length header") + } + } + for { + headerLine, err := reader.ReadString('\n') + if err != nil { + return nil, err + } + if strings.TrimSpace(headerLine) == "" { + break + } + } + body := make([]byte, contentLength) + if _, err := io.ReadFull(reader, body); err != nil { + return nil, err + } + return body, nil + } + return []byte(line), nil +} + +func jsonMarshal(message map[string]any) ([]byte, error) { + return json.Marshal(message) +} diff --git a/go/go_core/main.go b/go/go_core/main.go index fdf4b210..bdf71606 100644 --- a/go/go_core/main.go +++ b/go/go_core/main.go @@ -16,6 +16,10 @@ func main() { } return } + if len(os.Args) > 1 && os.Args[1] == "acp-stdio" { + acp.RunStdio(os.Stdin, os.Stdout) + return + } toolbridge.Run(os.Stdin, os.Stdout) } diff --git a/lib/app/app_controller_desktop_core.dart b/lib/app/app_controller_desktop_core.dart index 7e4e9c87..00e9470e 100644 --- a/lib/app/app_controller_desktop_core.dart +++ b/lib/app/app_controller_desktop_core.dart @@ -207,20 +207,13 @@ class AppController extends ChangeNotifier { arisBundleRepository ?? ArisBundleRepository(); goCoreLocatorInternal = GoCoreLocator(); runtimeCoordinatorInternal.attachDispatchResolver( - GoRuntimeDispatchDesktopClient( - acpClient: gatewayAcpClientInternal, - goCoreLocator: goCoreLocatorInternal, - ), + GoRuntimeDispatchDesktopClient(), ); goTaskServiceClientInternal = goTaskServiceClient ?? DesktopGoTaskService( gateway: runtimeCoordinatorInternal.gateway, - acpTransport: ExternalCodeAgentAcpDesktopTransport( - acpClient: gatewayAcpClientInternal, - endpointResolver: resolveExternalAcpEndpointForTargetInternal, - goCoreLocator: goCoreLocatorInternal, - ), + acpTransport: ExternalCodeAgentAcpDesktopTransport(), ); multiAgentOrchestratorInternal = MultiAgentOrchestrator( config: resolveMultiAgentConfigInternal( @@ -234,10 +227,7 @@ class AppController extends ChangeNotifier { MultiAgentMountManager( arisBundleRepository: arisBundleRepositoryInternal, goCoreLocator: goCoreLocatorInternal, - resolver: GoMultiAgentMountDesktopClient( - acpClient: gatewayAcpClientInternal, - goCoreLocator: goCoreLocatorInternal, - ), + resolver: GoMultiAgentMountDesktopClient(), ); attachChildListenersInternal(); diff --git a/lib/app/app_controller_web_sessions.dart b/lib/app/app_controller_web_sessions.dart index c153d63a..9e9e1b96 100644 --- a/lib/app/app_controller_web_sessions.dart +++ b/lib/app/app_controller_web_sessions.dart @@ -271,7 +271,7 @@ extension AppControllerWebSessions on AppController { assistantExecutionTargetForSession(currentSessionKeyInternal) == AssistantExecutionTarget.singleAgent && !availableSingleAgentProviders.any( - webAcpClientInternal.capabilities.providers.contains, + acpCapabilitiesInternal.providers.contains, ); List get secretReferences { diff --git a/lib/runtime/external_code_agent_acp_desktop_transport.dart b/lib/runtime/external_code_agent_acp_desktop_transport.dart index 99a69f72..43e2f84a 100644 --- a/lib/runtime/external_code_agent_acp_desktop_transport.dart +++ b/lib/runtime/external_code_agent_acp_desktop_transport.dart @@ -1,49 +1,16 @@ import 'dart:async'; -import 'dart:io'; -import 'embedded_agent_launch_policy.dart'; import 'gateway_acp_client.dart'; -import 'go_core.dart'; +import 'go_acp_stdio_bridge.dart'; import 'go_task_service_client.dart'; import 'runtime_models.dart'; -typedef ExternalCodeAgentAcpProcessStarter = - Future Function( - String executable, - List arguments, { - Map? environment, - String? workingDirectory, - }); - class ExternalCodeAgentAcpDesktopTransport implements ExternalCodeAgentAcpTransport { - ExternalCodeAgentAcpDesktopTransport({ - required GatewayAcpClient acpClient, - required Uri? Function(AssistantExecutionTarget target) endpointResolver, - GoCoreLocator? goCoreLocator, - ExternalCodeAgentAcpProcessStarter? processStarter, - }) : _acpClient = acpClient, - _endpointResolver = endpointResolver, - _goCoreLocator = goCoreLocator ?? GoCoreLocator(), - _processStarter = - processStarter ?? - ((executable, arguments, {environment, workingDirectory}) { - return Process.start( - executable, - arguments, - environment: environment, - workingDirectory: workingDirectory, - ); - }); + ExternalCodeAgentAcpDesktopTransport({GoAcpStdioBridge? bridge}) + : _bridge = bridge ?? GoAcpStdioBridge(); - final GatewayAcpClient _acpClient; - final Uri? Function(AssistantExecutionTarget target) _endpointResolver; - final GoCoreLocator _goCoreLocator; - final ExternalCodeAgentAcpProcessStarter _processStarter; - - Process? _localProcess; - Uri? _localEndpoint; - Future? _localEndpointFuture; + final GoAcpStdioBridge _bridge; List _syncedProviders = const []; @@ -54,11 +21,7 @@ class ExternalCodeAgentAcpDesktopTransport _syncedProviders = List.unmodifiable( providers, ); - final endpoint = await _ensureLocalEndpoint(); - if (endpoint == null) { - return; - } - await _syncProvidersToEndpoint(endpoint, _syncedProviders); + await _syncProviders(); } @override @@ -66,22 +29,39 @@ class ExternalCodeAgentAcpDesktopTransport required AssistantExecutionTarget target, bool forceRefresh = false, }) async { - final endpoint = await _resolveEndpoint(target); - if (endpoint == null) { - return const ExternalCodeAgentAcpCapabilities.empty(); - } - if (target == AssistantExecutionTarget.singleAgent) { - await _syncProvidersToEndpoint(endpoint, _syncedProviders); - } - final capabilities = await _acpClient.loadCapabilities( - forceRefresh: forceRefresh, - endpointOverride: endpoint, + await _syncProviders(); + final response = await _bridge.request( + method: 'acp.capabilities', + params: const {}, ); + final result = _castMap(response['result']); + final caps = _castMap(result['capabilities']); + final providers = {}; + for (final raw in [ + ..._asList(result['providers']), + ..._asList(caps['providers']), + ]) { + if (raw == null) { + continue; + } + final provider = SingleAgentProviderCopy.fromJsonValue( + raw.toString().trim().toLowerCase(), + ); + if (provider != SingleAgentProvider.auto) { + providers.add(provider); + } + } return ExternalCodeAgentAcpCapabilities( - singleAgent: capabilities.singleAgent, - multiAgent: capabilities.multiAgent, - providers: capabilities.providers, - raw: capabilities.raw, + singleAgent: + _boolValue(result['singleAgent']) ?? + _boolValue(caps['single_agent']) ?? + providers.isNotEmpty, + multiAgent: + _boolValue(result['multiAgent']) ?? + _boolValue(caps['multi_agent']) ?? + true, + providers: providers, + raw: result, ); } @@ -90,42 +70,46 @@ class ExternalCodeAgentAcpDesktopTransport GoTaskServiceRequest request, { required void Function(GoTaskServiceUpdate update) onUpdate, }) async { - final endpoint = await _resolveEndpoint(request.target); - if (endpoint == null) { - throw const GatewayAcpException( - 'Missing external ACP endpoint', - code: 'EXTERNAL_ACP_ENDPOINT_MISSING', - ); - } - if (request.target == AssistantExecutionTarget.singleAgent) { - await _syncProvidersToEndpoint(endpoint, _syncedProviders); - } + await _syncProviders(); + late final StreamSubscription> subscription; var streamedText = ''; String? completedMessage; - final response = await _acpClient.request( - method: request.resumeSession ? 'session.message' : 'session.start', - params: request.toExternalAcpParams(), - endpointOverride: endpoint, - onNotification: (notification) { - final update = goTaskServiceUpdateFromAcpNotification(notification); - if (update == null) { - return; - } - if (update.isDelta) { - streamedText += update.text; - } - if (update.isDone && update.message.trim().isNotEmpty) { - completedMessage = update.message.trim(); - } - onUpdate(update); - }, - ); - return goTaskServiceResultFromAcpResponse( - response, - route: request.route, - streamedText: streamedText, - completedMessage: completedMessage, - ); + subscription = _bridge.notifications.listen((notification) { + final update = goTaskServiceUpdateFromAcpNotification(notification); + if (update == null) { + return; + } + if (update.sessionId != request.sessionId || + update.threadId != request.threadId) { + return; + } + if (update.isDelta) { + streamedText += update.text; + } + if (update.isDone && update.message.trim().isNotEmpty) { + completedMessage = update.message.trim(); + } + onUpdate(update); + }); + try { + final response = await _bridge.request( + method: request.resumeSession ? 'session.message' : 'session.start', + params: request.toExternalAcpParams(), + ); + return goTaskServiceResultFromAcpResponse( + response, + route: request.route, + streamedText: streamedText, + completedMessage: completedMessage, + ); + } catch (error) { + throw GatewayAcpException( + error.toString(), + code: 'EXTERNAL_ACP_STDIO_ERROR', + ); + } finally { + await subscription.cancel(); + } } @override @@ -134,14 +118,9 @@ class ExternalCodeAgentAcpDesktopTransport required String sessionId, required String threadId, }) async { - final endpoint = await _resolveEndpoint(target); - if (endpoint == null) { - return; - } - await _acpClient.cancelSession( - sessionId: sessionId, - threadId: threadId, - endpointOverride: endpoint, + await _bridge.request( + method: 'session.cancel', + params: {'sessionId': sessionId, 'threadId': threadId}, ); } @@ -151,127 +130,71 @@ class ExternalCodeAgentAcpDesktopTransport required String sessionId, required String threadId, }) async { - final endpoint = await _resolveEndpoint(target); - if (endpoint == null) { - return; - } - await _acpClient.closeSession( - sessionId: sessionId, - threadId: threadId, - endpointOverride: endpoint, + await _bridge.request( + method: 'session.close', + params: {'sessionId': sessionId, 'threadId': threadId}, ); } @override - Future dispose() async { - final process = _localProcess; - _localProcess = null; - _localEndpoint = null; - _localEndpointFuture = null; - if (process != null) { - try { - process.kill(); - } catch (_) { - // Best effort only. - } - } - } + Future dispose() => _bridge.dispose(); - Future _resolveEndpoint(AssistantExecutionTarget target) async { - if (target == AssistantExecutionTarget.singleAgent) { - return _ensureLocalEndpoint(); - } - return _endpointResolver(target); - } - - Future _ensureLocalEndpoint() async { - if (_localEndpoint != null) { - return _localEndpoint; - } - final inFlight = _localEndpointFuture; - if (inFlight != null) { - return inFlight; - } - final next = _startLocalProcess(); - _localEndpointFuture = next; - try { - _localEndpoint = await next; - return _localEndpoint; - } finally { - _localEndpointFuture = null; - } - } - - Future _startLocalProcess() async { - final launch = await _goCoreLocator.locate(); - if (launch == null) { - return null; - } - if (shouldBlockGoCoreLaunch( - launch, - isAppleHost: Platform.isIOS || Platform.isMacOS, - )) { - return null; - } - final reservedSocket = await ServerSocket.bind( - InternetAddress.loopbackIPv4, - 0, - ); - final port = reservedSocket.port; - await reservedSocket.close(); - final listenAddress = '127.0.0.1:$port'; - final process = await _processStarter( - launch.executable, - [...launch.arguments, 'serve', '--listen', listenAddress], - environment: Platform.environment, - workingDirectory: launch.workingDirectory, - ); - _localProcess = process; - unawaited(process.stdout.drain()); - unawaited(process.stderr.drain()); - final endpoint = Uri(scheme: 'http', host: '127.0.0.1', port: port); - final deadline = DateTime.now().add(const Duration(seconds: 8)); - while (DateTime.now().isBefore(deadline)) { - if (_localProcess != process) { - break; - } - final exitCode = await process.exitCode.timeout( - const Duration(milliseconds: 20), - onTimeout: () => -1, - ); - if (exitCode != -1) { - break; - } - try { - await _acpClient.request( - method: 'acp.capabilities', - params: const {}, - endpointOverride: endpoint, - ); - return endpoint; - } catch (_) { - await Future.delayed(const Duration(milliseconds: 120)); - } - } - await dispose(); - return null; - } - - Future _syncProvidersToEndpoint( - Uri endpoint, - List providers, - ) async { - if (providers.isEmpty) { - return; - } - await _acpClient.request( + Future _syncProviders() async { + await _bridge.request( method: 'xworkmate.providers.sync', params: { - 'providers': providers - .map((item) => item.toJson()) + 'providers': _syncedProviders + .map( + (item) => { + 'providerId': item.providerId, + 'endpoint': item.endpoint, + 'label': item.label, + 'authorizationHeader': item.authorizationHeader, + 'enabled': item.enabled, + }, + ) .toList(growable: false), }, - endpointOverride: endpoint, ); } + + Map _castMap(Object? value) { + if (value is Map) { + return value; + } + if (value is Map) { + return value.cast(); + } + return const {}; + } + + List _asList(Object? raw) { + if (raw is List) { + return raw; + } + if (raw is List) { + return raw.cast(); + } + return const []; + } + + bool? _boolValue(Object? raw) { + if (raw is bool) { + return raw; + } + if (raw is num) { + return raw != 0; + } + final text = raw?.toString().trim().toLowerCase(); + if (text == null || text.isEmpty) { + return null; + } + if (text == 'true' || text == '1' || text == 'yes') { + return true; + } + if (text == 'false' || text == '0' || text == 'no') { + return false; + } + return null; + } } diff --git a/lib/runtime/go_acp_stdio_bridge.dart b/lib/runtime/go_acp_stdio_bridge.dart new file mode 100644 index 00000000..c27b129e --- /dev/null +++ b/lib/runtime/go_acp_stdio_bridge.dart @@ -0,0 +1,238 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'embedded_agent_launch_policy.dart'; +import 'go_core.dart'; + +typedef GoAcpStdioProcessStarter = + Future Function( + String executable, + List arguments, { + Map? environment, + String? workingDirectory, + }); + +class GoAcpStdioBridge { + GoAcpStdioBridge({ + GoCoreLocator? goCoreLocator, + GoAcpStdioProcessStarter? processStarter, + }) : _goCoreLocator = goCoreLocator ?? GoCoreLocator(), + _processStarter = + processStarter ?? + ((executable, arguments, {environment, workingDirectory}) { + return Process.start( + executable, + arguments, + environment: environment, + workingDirectory: workingDirectory, + ); + }); + + final GoCoreLocator _goCoreLocator; + final GoAcpStdioProcessStarter _processStarter; + + final StreamController> _notificationsController = + StreamController>.broadcast(); + final Map>> _pending = + >>{}; + + Process? _process; + StreamSubscription? _stdoutSubscription; + StreamSubscription? _stderrSubscription; + Future? _startupFuture; + int _requestCounter = 0; + + Stream> get notifications => + _notificationsController.stream; + + Future> request({ + required String method, + required Map params, + Duration timeout = const Duration(seconds: 120), + }) async { + await _ensureStarted(); + final process = _process; + if (process == null) { + throw StateError('Missing Go ACP stdio process.'); + } + final id = + '${DateTime.now().microsecondsSinceEpoch}-$method-${_requestCounter++}'; + final completer = Completer>(); + _pending[id] = completer; + process.stdin.writeln( + jsonEncode({ + 'jsonrpc': '2.0', + 'id': id, + 'method': method, + 'params': params, + }), + ); + try { + return await completer.future.timeout( + timeout, + onTimeout: () => throw TimeoutException( + 'Go ACP stdio request timed out: $method', + timeout, + ), + ); + } finally { + _pending.remove(id); + } + } + + Future dispose() async { + final process = _process; + _process = null; + _startupFuture = null; + for (final completer in _pending.values) { + if (!completer.isCompleted) { + completer.completeError( + StateError('Go ACP stdio bridge disposed before response.'), + ); + } + } + _pending.clear(); + await _stdoutSubscription?.cancel(); + await _stderrSubscription?.cancel(); + _stdoutSubscription = null; + _stderrSubscription = null; + if (process != null) { + try { + await process.stdin.close(); + } catch (_) { + // Ignore broken pipes during disposal. + } + try { + process.kill(); + } catch (_) { + // Best effort only. + } + } + await _notificationsController.close(); + } + + Future _ensureStarted() async { + if (_process != null) { + return; + } + final inFlight = _startupFuture; + if (inFlight != null) { + return inFlight; + } + final next = _start(); + _startupFuture = next; + try { + await next; + } finally { + _startupFuture = null; + } + } + + Future _start() async { + final launch = await _goCoreLocator.locate(); + if (launch == null) { + throw StateError('Go core is unavailable.'); + } + if (shouldBlockGoCoreLaunch( + launch, + isAppleHost: Platform.isIOS || Platform.isMacOS, + )) { + throw UnsupportedError( + 'App Store builds only allow the bundled Go core helper inside the app bundle.', + ); + } + final process = await _processStarter( + launch.executable, + [...launch.arguments, 'acp-stdio'], + environment: Platform.environment, + workingDirectory: launch.workingDirectory, + ); + _process = process; + _stdoutSubscription = process.stdout + .transform(utf8.decoder) + .transform(const LineSplitter()) + .listen(_handleStdoutLine, onError: _handleProcessError); + _stderrSubscription = process.stderr + .transform(utf8.decoder) + .transform(const LineSplitter()) + .listen((_) {}, onError: _handleProcessError); + unawaited( + process.exitCode.then((exitCode) { + if (_process != process) { + return; + } + _process = null; + _failPending( + StateError('Go ACP stdio process exited with code $exitCode'), + ); + }), + ); + await request(method: 'acp.capabilities', params: const {}); + } + + void _handleStdoutLine(String line) { + final trimmed = line.trim(); + if (trimmed.isEmpty || !trimmed.startsWith('{')) { + return; + } + final json = _decodeMap(trimmed); + final id = json['id']?.toString().trim(); + if (id != null && id.isNotEmpty) { + final completer = _pending[id]; + if (completer == null || completer.isCompleted) { + return; + } + final error = _castMap(json['error']); + if (error.isNotEmpty) { + completer.completeError( + StateError( + error['message']?.toString() ?? 'Go ACP stdio request failed', + ), + ); + return; + } + completer.complete(json); + return; + } + if ((json['method']?.toString().trim() ?? '').isNotEmpty && + !_notificationsController.isClosed) { + _notificationsController.add(json); + } + } + + void _handleProcessError(Object error) { + _failPending(error); + } + + void _failPending(Object error) { + final pending = Map>>.from(_pending); + _pending.clear(); + for (final completer in pending.values) { + if (!completer.isCompleted) { + completer.completeError(error); + } + } + } + + Map _decodeMap(String raw) { + final decoded = jsonDecode(raw); + if (decoded is Map) { + return decoded; + } + if (decoded is Map) { + return decoded.cast(); + } + return const {}; + } + + Map _castMap(Object? value) { + if (value is Map) { + return value; + } + if (value is Map) { + return value.cast(); + } + return const {}; + } +} diff --git a/lib/runtime/go_gateway_runtime_desktop_client.dart b/lib/runtime/go_gateway_runtime_desktop_client.dart index e370cde3..e8c29db2 100644 --- a/lib/runtime/go_gateway_runtime_desktop_client.dart +++ b/lib/runtime/go_gateway_runtime_desktop_client.dart @@ -1,52 +1,24 @@ import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; -import 'embedded_agent_launch_policy.dart'; import 'gateway_runtime_errors.dart'; -import 'gateway_runtime_helpers.dart'; import 'gateway_runtime_session_client.dart'; -import 'go_core.dart'; - -typedef GoGatewayRuntimeProcessStarter = - Future Function( - String executable, - List arguments, { - Map? environment, - String? workingDirectory, - }); +import 'go_acp_stdio_bridge.dart'; class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { - GoGatewayRuntimeDesktopClient({ - GoCoreLocator? goCoreLocator, - GoGatewayRuntimeProcessStarter? processStarter, - }) : _goCoreLocator = goCoreLocator ?? GoCoreLocator(), - _processStarter = - processStarter ?? - ((executable, arguments, {environment, workingDirectory}) { - return Process.start( - executable, - arguments, - environment: environment, - workingDirectory: workingDirectory, - ); - }); - - final GoCoreLocator _goCoreLocator; - final GoGatewayRuntimeProcessStarter _processStarter; + GoGatewayRuntimeDesktopClient({GoAcpStdioBridge? bridge}) + : _bridge = bridge ?? GoAcpStdioBridge() { + _notificationsSubscription = _bridge.notifications.listen( + _handleNotification, + onError: (Object error, StackTrace stackTrace) { + _updatesController.addError(error, stackTrace); + }, + ); + } + final GoAcpStdioBridge _bridge; + late final StreamSubscription> _notificationsSubscription; final StreamController _updatesController = StreamController.broadcast(); - final Map>> _pending = - >>{}; - - Process? _localProcess; - Uri? _localEndpoint; - Future? _localEndpointFuture; - WebSocket? _socket; - StreamSubscription? _socketSubscription; - Future? _socketReadyFuture; - int _requestCounter = 0; @override Stream get updates => _updatesController.stream; @@ -59,7 +31,7 @@ class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { method: 'xworkmate.gateway.connect', params: request.toJson(), ); - if (boolValue(result['ok']) != true) { + if (_boolValue(result['ok']) != true) { throw _gatewayErrorFromResult( result, fallbackMessage: 'Gateway connect failed', @@ -84,7 +56,7 @@ class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { 'timeoutMs': timeout.inMilliseconds, }, ); - if (boolValue(result['ok']) != true) { + if (_boolValue(result['ok']) != true) { throw _gatewayErrorFromResult( result, fallbackMessage: '$method request failed', @@ -103,37 +75,8 @@ class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { @override Future dispose() async { - for (final completer in _pending.values) { - if (!completer.isCompleted) { - completer.completeError( - GatewayRuntimeException( - 'Go gateway runtime transport disposed', - code: 'GO_GATEWAY_RUNTIME_TRANSPORT_DISPOSED', - ), - ); - } - } - _pending.clear(); - await _socketSubscription?.cancel(); - _socketSubscription = null; - try { - await _socket?.close(); - } catch (_) { - // Best effort only. - } - _socket = null; - _socketReadyFuture = null; - final process = _localProcess; - _localProcess = null; - _localEndpoint = null; - _localEndpointFuture = null; - if (process != null) { - try { - process.kill(); - } catch (_) { - // Best effort only. - } - } + await _notificationsSubscription.cancel(); + await _bridge.dispose(); await _updatesController.close(); } @@ -141,208 +84,24 @@ class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { required String method, required Map params, }) async { - await _ensureSocketReady(); - final socket = _socket; - if (socket == null) { - throw GatewayRuntimeException( - 'Missing Go gateway runtime transport', - code: 'GO_GATEWAY_RUNTIME_TRANSPORT_UNAVAILABLE', - ); - } - final requestId = - '${DateTime.now().microsecondsSinceEpoch}-$method-${_requestCounter++}'; - final completer = Completer>(); - _pending[requestId] = completer; - socket.add( - jsonEncode({ - 'jsonrpc': '2.0', - 'id': requestId, - 'method': method, - 'params': params, - }), - ); - try { - return await completer.future.timeout(const Duration(seconds: 120)); - } finally { - _pending.remove(requestId); - } + final response = await _bridge.request(method: method, params: params); + return _castMap(response['result']); } - Future _ensureSocketReady() async { - final inFlight = _socketReadyFuture; - if (inFlight != null) { - return inFlight; - } - final next = _openSocket(); - _socketReadyFuture = next; - try { - await next; - } finally { - _socketReadyFuture = null; - } - } - - Future _openSocket() async { - if (_socket != null) { - return; - } - final endpoint = await _ensureLocalEndpoint(); - if (endpoint == null) { - throw GatewayRuntimeException( - 'Missing Go gateway runtime endpoint', - code: 'GO_GATEWAY_RUNTIME_ENDPOINT_MISSING', - ); - } - final wsEndpoint = endpoint.replace( - scheme: endpoint.scheme == 'https' ? 'wss' : 'ws', - path: '/acp', - ); - final socket = await WebSocket.connect(wsEndpoint.toString()).timeout( - const Duration(seconds: 6), - onTimeout: () => throw GatewayRuntimeException( - 'Go gateway runtime websocket connect timeout', - code: 'GO_GATEWAY_RUNTIME_WS_CONNECT_TIMEOUT', - ), - ); - _socket = socket; - _socketSubscription = socket.listen( - _handleSocketMessage, - onError: (Object error, StackTrace stackTrace) { - _failPending( - GatewayRuntimeException( - error.toString(), - code: 'GO_GATEWAY_RUNTIME_WS_ERROR', - ), - ); - }, - onDone: () { - _socket = null; - _socketSubscription = null; - _failPending( - GatewayRuntimeException( - 'Go gateway runtime websocket closed', - code: 'GO_GATEWAY_RUNTIME_WS_CLOSED', - ), - ); - }, - cancelOnError: true, - ); - } - - void _handleSocketMessage(dynamic raw) { - final json = _decodeMap(raw); - final id = json['id']?.toString().trim(); - if (id != null && id.isNotEmpty) { - final completer = _pending[id]; - if (completer != null && !completer.isCompleted) { - final error = _castMap(json['error']); - if (error.isNotEmpty) { - completer.completeError( - GatewayRuntimeException( - error['message']?.toString() ?? - 'Go gateway runtime request failed', - code: error['code']?.toString(), - ), - ); - } else { - completer.complete(_castMap(json['result'])); - } - } - return; - } - final method = json['method']?.toString().trim() ?? ''; + void _handleNotification(Map notification) { + final method = notification['method']?.toString().trim() ?? ''; if (method.isEmpty) { return; } try { _updatesController.add( - GatewayRuntimeSessionUpdate.fromNotification(json), + GatewayRuntimeSessionUpdate.fromNotification(notification), ); } catch (_) { // Ignore unrelated ACP notifications. } } - void _failPending(GatewayRuntimeException error) { - for (final completer in _pending.values) { - if (!completer.isCompleted) { - completer.completeError(error); - } - } - _pending.clear(); - } - - Future _ensureLocalEndpoint() async { - if (_localEndpoint != null) { - return _localEndpoint; - } - final inFlight = _localEndpointFuture; - if (inFlight != null) { - return inFlight; - } - final next = _startLocalProcess(); - _localEndpointFuture = next; - try { - _localEndpoint = await next; - return _localEndpoint; - } finally { - _localEndpointFuture = null; - } - } - - Future _startLocalProcess() async { - final launch = await _goCoreLocator.locate(); - if (launch == null) { - return null; - } - if (shouldBlockGoCoreLaunch( - launch, - isAppleHost: Platform.isIOS || Platform.isMacOS, - )) { - return null; - } - final reservedSocket = await ServerSocket.bind( - InternetAddress.loopbackIPv4, - 0, - ); - final port = reservedSocket.port; - await reservedSocket.close(); - final listenAddress = '127.0.0.1:$port'; - final process = await _processStarter( - launch.executable, - [...launch.arguments, 'serve', '--listen', listenAddress], - environment: Platform.environment, - workingDirectory: launch.workingDirectory, - ); - _localProcess = process; - unawaited(process.stdout.drain()); - unawaited(process.stderr.drain()); - final endpoint = Uri(scheme: 'http', host: '127.0.0.1', port: port); - final deadline = DateTime.now().add(const Duration(seconds: 8)); - while (DateTime.now().isBefore(deadline)) { - if (_localProcess != process) { - break; - } - final exitCode = await process.exitCode.timeout( - const Duration(milliseconds: 20), - onTimeout: () => -1, - ); - if (exitCode != -1) { - break; - } - try { - final probe = await WebSocket.connect( - endpoint.replace(scheme: 'ws', path: '/acp').toString(), - ).timeout(const Duration(milliseconds: 300)); - await probe.close(); - return endpoint; - } catch (_) { - await Future.delayed(const Duration(milliseconds: 120)); - } - } - return null; - } - GatewayRuntimeException _gatewayErrorFromResult( Map result, { required String fallbackMessage, @@ -355,22 +114,6 @@ class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { ); } - Map _decodeMap(dynamic raw) { - if (raw is Map) { - return raw; - } - if (raw is Map) { - return raw.cast(); - } - if (raw is String) { - return _castMap(jsonDecode(raw)); - } - if (raw is List) { - return _castMap(jsonDecode(utf8.decode(raw))); - } - return const {}; - } - Map _castMap(Object? value) { if (value is Map) { return value; @@ -380,4 +123,24 @@ class GoGatewayRuntimeDesktopClient implements GatewayRuntimeSessionClient { } return const {}; } + + bool? _boolValue(Object? raw) { + if (raw is bool) { + return raw; + } + if (raw is num) { + return raw != 0; + } + final text = raw?.toString().trim().toLowerCase(); + if (text == null || text.isEmpty) { + return null; + } + if (text == 'true' || text == '1' || text == 'yes') { + return true; + } + if (text == 'false' || text == '0' || text == 'no') { + return false; + } + return null; + } } diff --git a/lib/runtime/go_multi_agent_mount_desktop_client.dart b/lib/runtime/go_multi_agent_mount_desktop_client.dart index 9972871a..393267d0 100644 --- a/lib/runtime/go_multi_agent_mount_desktop_client.dart +++ b/lib/runtime/go_multi_agent_mount_desktop_client.dart @@ -1,45 +1,12 @@ -import 'dart:async'; -import 'dart:io'; - -import 'embedded_agent_launch_policy.dart'; -import 'gateway_acp_client.dart'; -import 'go_core.dart'; +import 'go_acp_stdio_bridge.dart'; import 'multi_agent_mount_resolver.dart'; import 'runtime_models.dart'; -typedef GoMultiAgentMountProcessStarter = - Future Function( - String executable, - List arguments, { - Map? environment, - String? workingDirectory, - }); - class GoMultiAgentMountDesktopClient implements MultiAgentMountResolver { - GoMultiAgentMountDesktopClient({ - GatewayAcpClient? acpClient, - GoCoreLocator? goCoreLocator, - GoMultiAgentMountProcessStarter? processStarter, - }) : _acpClient = acpClient ?? GatewayAcpClient(endpointResolver: () => null), - _goCoreLocator = goCoreLocator ?? GoCoreLocator(), - _processStarter = - processStarter ?? - ((executable, arguments, {environment, workingDirectory}) { - return Process.start( - executable, - arguments, - environment: environment, - workingDirectory: workingDirectory, - ); - }); + GoMultiAgentMountDesktopClient({GoAcpStdioBridge? bridge}) + : _bridge = bridge ?? GoAcpStdioBridge(); - final GatewayAcpClient _acpClient; - final GoCoreLocator _goCoreLocator; - final GoMultiAgentMountProcessStarter _processStarter; - - Process? _localProcess; - Uri? _localEndpoint; - Future? _localEndpointFuture; + final GoAcpStdioBridge _bridge; @override Future reconcile({ @@ -50,11 +17,7 @@ class GoMultiAgentMountDesktopClient implements MultiAgentMountResolver { required String opencodeHome, required ArisMountProbe arisProbe, }) async { - final endpoint = await _ensureLocalEndpoint(); - if (endpoint == null) { - return null; - } - final response = await _acpClient.request( + final response = await _bridge.request( method: 'xworkmate.mounts.reconcile', params: { 'config': { @@ -70,7 +33,6 @@ class GoMultiAgentMountDesktopClient implements MultiAgentMountResolver { 'opencodeHome': opencodeHome.trim(), 'aris': arisProbe.toJson(), }, - endpointOverride: endpoint, ); final result = _castMap(response['result']); final rawTargets = result['mountTargets']; @@ -98,93 +60,7 @@ class GoMultiAgentMountDesktopClient implements MultiAgentMountResolver { } @override - Future dispose() async { - final process = _localProcess; - _localProcess = null; - _localEndpoint = null; - _localEndpointFuture = null; - if (process != null) { - try { - process.kill(); - } catch (_) { - // Best effort only. - } - } - await _acpClient.dispose(); - } - - Future _ensureLocalEndpoint() async { - if (_localEndpoint != null) { - return _localEndpoint; - } - final inFlight = _localEndpointFuture; - if (inFlight != null) { - return inFlight; - } - final next = _startLocalProcess(); - _localEndpointFuture = next; - try { - _localEndpoint = await next; - return _localEndpoint; - } finally { - _localEndpointFuture = null; - } - } - - Future _startLocalProcess() async { - final launch = await _goCoreLocator.locate(); - if (launch == null) { - return null; - } - if (shouldBlockGoCoreLaunch( - launch, - isAppleHost: Platform.isIOS || Platform.isMacOS, - )) { - return null; - } - final reservedSocket = await ServerSocket.bind( - InternetAddress.loopbackIPv4, - 0, - ); - final port = reservedSocket.port; - await reservedSocket.close(); - final listenAddress = '127.0.0.1:$port'; - final process = await _processStarter( - launch.executable, - [...launch.arguments, 'serve', '--listen', listenAddress], - environment: Platform.environment, - workingDirectory: launch.workingDirectory, - ); - _localProcess = process; - unawaited(process.stdout.drain()); - unawaited(process.stderr.drain()); - final endpoint = Uri(scheme: 'http', host: '127.0.0.1', port: port); - final deadline = DateTime.now().add(const Duration(seconds: 8)); - while (DateTime.now().isBefore(deadline)) { - if (_localProcess != process) { - break; - } - final exitCode = await process.exitCode.timeout( - const Duration(milliseconds: 20), - onTimeout: () => -1, - ); - if (exitCode != -1) { - break; - } - try { - await _acpClient.request( - method: 'acp.capabilities', - params: const {}, - endpointOverride: endpoint, - ); - return endpoint; - } catch (_) { - await Future.delayed(const Duration(milliseconds: 120)); - } - } - await dispose(); - return null; - } + Future dispose() => _bridge.dispose(); Map _castMap(Object? value) { if (value is Map) { diff --git a/lib/runtime/go_runtime_dispatch_desktop_client.dart b/lib/runtime/go_runtime_dispatch_desktop_client.dart index 08e3ac14..8d055947 100644 --- a/lib/runtime/go_runtime_dispatch_desktop_client.dart +++ b/lib/runtime/go_runtime_dispatch_desktop_client.dart @@ -1,45 +1,12 @@ -import 'dart:async'; -import 'dart:io'; - -import 'embedded_agent_launch_policy.dart'; -import 'gateway_acp_client.dart'; -import 'go_core.dart'; +import 'go_acp_stdio_bridge.dart'; import 'runtime_dispatch_resolver.dart'; import 'runtime_external_code_agents.dart'; -typedef GoRuntimeDispatchProcessStarter = - Future Function( - String executable, - List arguments, { - Map? environment, - String? workingDirectory, - }); - class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { - GoRuntimeDispatchDesktopClient({ - GatewayAcpClient? acpClient, - GoCoreLocator? goCoreLocator, - GoRuntimeDispatchProcessStarter? processStarter, - }) : _acpClient = acpClient ?? GatewayAcpClient(endpointResolver: () => null), - _goCoreLocator = goCoreLocator ?? GoCoreLocator(), - _processStarter = - processStarter ?? - ((executable, arguments, {environment, workingDirectory}) { - return Process.start( - executable, - arguments, - environment: environment, - workingDirectory: workingDirectory, - ); - }); + GoRuntimeDispatchDesktopClient({GoAcpStdioBridge? bridge}) + : _bridge = bridge ?? GoAcpStdioBridge(); - final GatewayAcpClient _acpClient; - final GoCoreLocator _goCoreLocator; - final GoRuntimeDispatchProcessStarter _processStarter; - - Process? _localProcess; - Uri? _localEndpoint; - Future? _localEndpointFuture; + final GoAcpStdioBridge _bridge; @override Future selectProviderId({ @@ -47,11 +14,7 @@ class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { String preferredProviderId = '', Iterable requiredCapabilities = const [], }) async { - final endpoint = await _ensureLocalEndpoint(); - if (endpoint == null) { - return null; - } - final response = await _acpClient.request( + final response = await _bridge.request( method: 'xworkmate.dispatch.resolve', params: { 'preferredProviderId': preferredProviderId.trim(), @@ -61,7 +24,6 @@ class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { .toList(growable: false), 'providers': providers.map(_providerToJson).toList(growable: false), }, - endpointOverride: endpoint, ); final result = _castMap(response['result']); return result['providerId']?.toString().trim().isNotEmpty == true @@ -77,11 +39,7 @@ class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { required Map nodeState, required Map nodeInfo, }) async { - final endpoint = await _ensureLocalEndpoint(); - if (endpoint == null) { - return const RuntimeDispatchResolution(metadata: {}); - } - final response = await _acpClient.request( + final response = await _bridge.request( method: 'xworkmate.dispatch.resolve', params: { 'preferredProviderId': preferredProviderId.trim(), @@ -93,7 +51,6 @@ class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { 'nodeState': nodeState, 'nodeInfo': nodeInfo, }, - endpointOverride: endpoint, ); final result = _castMap(response['result']); return RuntimeDispatchResolution( @@ -109,20 +66,7 @@ class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { } @override - Future dispose() async { - final process = _localProcess; - _localProcess = null; - _localEndpoint = null; - _localEndpointFuture = null; - if (process != null) { - try { - process.kill(); - } catch (_) { - // Best effort only. - } - } - await _acpClient.dispose(); - } + Future dispose() => _bridge.dispose(); Map _providerToJson(ExternalCodeAgentProvider provider) { return { @@ -133,79 +77,6 @@ class GoRuntimeDispatchDesktopClient implements RuntimeDispatchResolver { }; } - Future _ensureLocalEndpoint() async { - if (_localEndpoint != null) { - return _localEndpoint; - } - final inFlight = _localEndpointFuture; - if (inFlight != null) { - return inFlight; - } - final next = _startLocalProcess(); - _localEndpointFuture = next; - try { - _localEndpoint = await next; - return _localEndpoint; - } finally { - _localEndpointFuture = null; - } - } - - Future _startLocalProcess() async { - final launch = await _goCoreLocator.locate(); - if (launch == null) { - return null; - } - if (shouldBlockGoCoreLaunch( - launch, - isAppleHost: Platform.isIOS || Platform.isMacOS, - )) { - return null; - } - final reservedSocket = await ServerSocket.bind( - InternetAddress.loopbackIPv4, - 0, - ); - final port = reservedSocket.port; - await reservedSocket.close(); - final listenAddress = '127.0.0.1:$port'; - final process = await _processStarter( - launch.executable, - [...launch.arguments, 'serve', '--listen', listenAddress], - environment: Platform.environment, - workingDirectory: launch.workingDirectory, - ); - _localProcess = process; - unawaited(process.stdout.drain()); - unawaited(process.stderr.drain()); - final endpoint = Uri(scheme: 'http', host: '127.0.0.1', port: port); - final deadline = DateTime.now().add(const Duration(seconds: 8)); - while (DateTime.now().isBefore(deadline)) { - if (_localProcess != process) { - break; - } - final exitCode = await process.exitCode.timeout( - const Duration(milliseconds: 20), - onTimeout: () => -1, - ); - if (exitCode != -1) { - break; - } - try { - await _acpClient.request( - method: 'acp.capabilities', - params: const {}, - endpointOverride: endpoint, - ); - return endpoint; - } catch (_) { - await Future.delayed(const Duration(milliseconds: 120)); - } - } - await dispose(); - return null; - } - Map _castMap(Object? value) { if (value is Map) { return value; diff --git a/lib/runtime/go_task_service_client.dart b/lib/runtime/go_task_service_client.dart index 4e12a9ec..57d068ea 100644 --- a/lib/runtime/go_task_service_client.dart +++ b/lib/runtime/go_task_service_client.dart @@ -579,16 +579,21 @@ GoTaskServiceResult goTaskServiceResultFromAcpResponse( String? completedMessage, }) { final result = _castMap(response['result']); + final responseText = + (result['output']?.toString().trim().isNotEmpty == true + ? result['output'].toString().trim() + : result['summary']?.toString().trim().isNotEmpty == true + ? result['summary'].toString().trim() + : result['message']?.toString().trim() ?? '') + .trim(); final primaryText = (completedMessage?.trim().isNotEmpty == true ? completedMessage!.trim() + : responseText.isNotEmpty + ? responseText : streamedText.trim().isNotEmpty ? streamedText.trim() - : (result['output']?.toString().trim().isNotEmpty == true - ? result['output'].toString().trim() - : result['summary']?.toString().trim().isNotEmpty == true - ? result['summary'].toString().trim() - : result['message']?.toString().trim() ?? '')) + : '') .trim(); return GoTaskServiceResult( success: _boolValue(result['success']) ?? true, diff --git a/test/runtime/external_code_agent_acp_desktop_transport_test.dart b/test/runtime/external_code_agent_acp_desktop_transport_test.dart index 47c44903..f887d48b 100644 --- a/test/runtime/external_code_agent_acp_desktop_transport_test.dart +++ b/test/runtime/external_code_agent_acp_desktop_transport_test.dart @@ -2,29 +2,70 @@ library; import 'dart:async'; -import 'dart:convert'; -import 'dart:io'; import 'package:flutter_test/flutter_test.dart'; import 'package:xworkmate/runtime/external_code_agent_acp_desktop_transport.dart'; -import 'package:xworkmate/runtime/gateway_acp_client.dart'; +import 'package:xworkmate/runtime/go_acp_stdio_bridge.dart'; import 'package:xworkmate/runtime/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; void main() { group('ExternalCodeAgentAcpDesktopTransport', () { - test('uses resolved gateway endpoint for local gateway sessions', () async { - final server = await _AcpFakeServer.start(); - addTearDown(server.close); - - final transport = ExternalCodeAgentAcpDesktopTransport( - acpClient: GatewayAcpClient(endpointResolver: () => null), - endpointResolver: (target) => switch (target) { - AssistantExecutionTarget.local => server.baseHttpUri, - _ => null, + test('uses direct Go ACP stdio bridge for desktop task execution', () async { + late final _FakeGoAcpStdioBridge bridge; + bridge = _FakeGoAcpStdioBridge( + handler: (method, params) async { + switch (method) { + case 'acp.capabilities': + return { + 'jsonrpc': '2.0', + 'id': 'capabilities', + 'result': { + 'singleAgent': true, + 'multiAgent': true, + 'providers': ['codex'], + 'capabilities': { + 'single_agent': true, + 'multi_agent': true, + 'providers': ['codex'], + }, + }, + }; + case 'xworkmate.providers.sync': + return { + 'jsonrpc': '2.0', + 'id': 'sync', + 'result': {'ok': true}, + }; + case 'session.start': + bridge.emit({ + 'jsonrpc': '2.0', + 'method': 'session.update', + 'params': { + 'sessionId': 'session-local', + 'threadId': 'thread-local', + 'turnId': 'turn-1', + 'type': 'delta', + 'delta': 'gateway-', + }, + }); + return { + 'jsonrpc': '2.0', + 'id': 'start', + 'result': { + 'success': true, + 'message': 'gateway-ok', + 'summary': 'gateway-ok', + 'turnId': 'turn-1', + }, + }; + } + throw StateError('Unexpected method: $method'); }, ); + final transport = ExternalCodeAgentAcpDesktopTransport(bridge: bridge); + final updates = []; final result = await transport.executeTask( const GoTaskServiceRequest( sessionId: 'session-local', @@ -42,112 +83,52 @@ void main() { agentId: '', metadata: {}, ), - onUpdate: (_) {}, + onUpdate: updates.add, ); 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 = ExternalCodeAgentAcpDesktopTransport( - acpClient: GatewayAcpClient(endpointResolver: () => null), - endpointResolver: (_) => null, - ); - - await expectLater( - () => transport.executeTask( - const GoTaskServiceRequest( - 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', - 'EXTERNAL_ACP_ENDPOINT_MISSING', - ), - ), + expect( + bridge.recordedMethods, + containsAll(['xworkmate.providers.sync', 'session.start']), ); + expect(updates.single.text, 'gateway-'); }); }); } -class _AcpFakeServer { - _AcpFakeServer._(this._server); +class _FakeGoAcpStdioBridge extends GoAcpStdioBridge { + _FakeGoAcpStdioBridge({required this.handler}); - final HttpServer _server; - final List rpcMethods = []; - String? lastHttpRequestPath; - String? lastSessionMode; + final Future> Function( + String method, + Map params, + ) + handler; - Uri get baseHttpUri => Uri.parse('http://127.0.0.1:${_server.port}'); + final StreamController> _notificationsController = + StreamController>.broadcast(); + final List recordedMethods = []; - static Future<_AcpFakeServer> start() async { - final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); - final fake = _AcpFakeServer._(server); - unawaited(fake._listen()); - return fake; + @override + Stream> get notifications => _notificationsController.stream; + + void emit(Map notification) { + _notificationsController.add(notification); } - Future close() async { - await _server.close(force: true); + @override + Future> request({ + required String method, + required Map params, + Duration timeout = const Duration(seconds: 120), + }) async { + recordedMethods.add(method); + return handler(method, params); } - 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(); + @override + Future dispose() async { + await _notificationsController.close(); } }