From 2d1d8ecb42dda032cf934f2a4eff5ff7be9410d0 Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Fri, 27 Mar 2026 11:21:05 +0800 Subject: [PATCH] Fix OpenCode single-agent ACP transport --- ...direct_single_agent_app_server_client.dart | 487 ++++++++++++++++++ .../direct_single_agent_app_server_suite.dart | 255 ++++++++- 2 files changed, 741 insertions(+), 1 deletion(-) diff --git a/lib/runtime/direct_single_agent_app_server_client.dart b/lib/runtime/direct_single_agent_app_server_client.dart index 3d174393..2d05329e 100644 --- a/lib/runtime/direct_single_agent_app_server_client.dart +++ b/lib/runtime/direct_single_agent_app_server_client.dart @@ -75,6 +75,7 @@ class DirectSingleAgentAppServerClient { final Map _activeConnections = {}; final Map _threadIds = {}; + final Map _restSessionIds = {}; final Set _abortedSessions = {}; final Map @@ -97,6 +98,37 @@ class DirectSingleAgentAppServerClient { } final endpoint = _resolveWebSocketEndpoint(provider); + if (_usesRestSessionApi(provider)) { + final base = endpointResolver(provider); + if (base == null) { + final unavailable = const DirectSingleAgentCapabilities.unavailable( + endpoint: '', + errorMessage: 'Single-agent app-server endpoint is not configured.', + ); + _cachedCapabilities[provider] = unavailable; + _capabilitiesRefreshedAt[provider] = DateTime.now(); + return unavailable; + } + try { + await _fetchJson( + _buildRestUri(base, '/global/health'), + gatewayToken: gatewayToken, + ); + _cachedCapabilities[provider] = DirectSingleAgentCapabilities( + available: true, + supportedProviders: [provider], + endpoint: base.toString(), + ); + } catch (error) { + _cachedCapabilities[provider] = DirectSingleAgentCapabilities.unavailable( + endpoint: base.toString(), + errorMessage: error.toString(), + ); + } finally { + _capabilitiesRefreshedAt[provider] = DateTime.now(); + } + return _cachedCapabilities[provider]!; + } if (endpoint == null) { final unavailable = const DirectSingleAgentCapabilities.unavailable( endpoint: '', @@ -135,6 +167,9 @@ class DirectSingleAgentAppServerClient { Future run( DirectSingleAgentRunRequest request, ) async { + if (_usesRestSessionApi(request.provider)) { + return _runViaRestApi(request); + } final endpoint = _resolveWebSocketEndpoint(request.provider); if (endpoint == null) { return const DirectSingleAgentRunResult( @@ -302,6 +337,22 @@ class DirectSingleAgentAppServerClient { return; } _abortedSessions.add(normalizedSessionId); + final restSessionId = _restSessionIds[normalizedSessionId]?.trim() ?? ''; + if (restSessionId.isNotEmpty) { + final provider = SingleAgentProvider.opencode; + final base = endpointResolver(provider); + if (base != null) { + try { + await _postJson( + _buildRestUri(base, '/session/$restSessionId/abort'), + body: null, + gatewayToken: '', + ); + } catch (_) { + // Best effort only. + } + } + } final connection = _activeConnections[normalizedSessionId]; final threadId = _threadIds[normalizedSessionId]; if (connection == null || threadId == null || threadId.isEmpty) { @@ -365,6 +416,442 @@ class DirectSingleAgentAppServerClient { return threadId; } + Future _runViaRestApi( + DirectSingleAgentRunRequest request, + ) async { + final base = endpointResolver(request.provider); + if (base == null) { + return const DirectSingleAgentRunResult( + success: false, + output: '', + errorMessage: 'Single-agent REST endpoint is missing.', + ); + } + final normalizedSessionId = request.sessionId.trim(); + if (normalizedSessionId.isEmpty) { + return const DirectSingleAgentRunResult( + success: false, + output: '', + errorMessage: 'Single-agent session id is missing.', + ); + } + + _abortedSessions.remove(normalizedSessionId); + final remoteSessionId = await _ensureRestSession( + base, + sessionId: normalizedSessionId, + workingDirectory: request.workingDirectory, + gatewayToken: request.gatewayToken, + ); + + final output = StringBuffer(); + final completion = Completer(); + String? activeAssistantMessageId; + String? lastAssistantText; + var busySeen = false; + + 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 ?? ''); + completion.complete( + DirectSingleAgentRunResult( + success: true, + output: resolvedOutput, + errorMessage: '', + resolvedModel: request.model, + ), + ); + } + + try { + final eventUri = _buildRestUri(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 = _decodeMap(line.substring(6)); + final payload = _asMap(event['payload']); + final type = payload['type']?.toString().trim() ?? ''; + final properties = _asMap(payload['properties']); + if (properties['sessionID']?.toString().trim() != + remoteSessionId) { + return; + } + if (type == 'session.status') { + final status = _asMap(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 = _asMap(properties['error']); + completion.complete( + DirectSingleAgentRunResult( + success: false, + output: output.toString(), + errorMessage: + error['message']?.toString() ?? + error['name']?.toString() ?? + 'OpenCode session failed.', + aborted: _abortedSessions.contains(normalizedSessionId), + resolvedModel: request.model, + ), + ); + return; + } + if (type == 'message.updated') { + final info = _asMap(properties['info']); + if (info['role']?.toString().trim() == 'assistant') { + activeAssistantMessageId = info['id']?.toString().trim(); + } + return; + } + if (type == 'message.part.delta') { + final part = _asMap(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 = _asMap(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) { + // OpenCode event streams can disconnect independently from the + // backing session lifecycle. Keep polling session state instead. + }, + onDone: () {}, + cancelOnError: true, + ); + + await _postJson( + _buildRestUri( + base, + '/session/$remoteSessionId/message', + queryParameters: { + 'directory': request.workingDirectory, + }, + ), + body: { + 'agent': 'build', + 'parts': >[ + {'type': 'text', 'text': request.prompt}, + ], + }, + gatewayToken: request.gatewayToken, + ); + unawaited( + _pollRestAssistantMessage( + 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: (message) { + if (!completion.isCompleted) { + completion.complete( + DirectSingleAgentRunResult( + success: false, + output: output.toString(), + errorMessage: message, + aborted: _abortedSessions.contains(normalizedSessionId), + resolvedModel: request.model, + ), + ); + } + }, + ), + ); + + return await completion.future.timeout( + const Duration(minutes: 10), + onTimeout: () => DirectSingleAgentRunResult( + success: false, + output: output.toString(), + errorMessage: 'OpenCode REST request timed out.', + aborted: _abortedSessions.contains(normalizedSessionId), + resolvedModel: request.model, + ), + ); + } catch (error) { + return DirectSingleAgentRunResult( + success: false, + output: output.toString(), + errorMessage: error.toString(), + aborted: _abortedSessions.contains(normalizedSessionId), + resolvedModel: request.model, + ); + } finally { + unawaited(lineSubscription?.cancel()); + eventClient.close(force: true); + _abortedSessions.remove(normalizedSessionId); + } + } + + Future _ensureRestSession( + Uri base, { + required String sessionId, + required String workingDirectory, + required String gatewayToken, + }) async { + final existing = _restSessionIds[sessionId]?.trim() ?? ''; + if (existing.isNotEmpty) { + return existing; + } + final created = await _postJson( + _buildRestUri( + base, + '/session', + queryParameters: {'directory': workingDirectory}, + ), + body: {'title': sessionId}, + gatewayToken: gatewayToken, + ); + final createdId = created['id']?.toString().trim() ?? ''; + if (createdId.isEmpty) { + throw StateError('OpenCode REST endpoint returned an empty session id.'); + } + _restSessionIds[sessionId] = createdId; + return createdId; + } + + Future _pollRestAssistantMessage( + 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 _fetchJsonList( + _buildRestUri( + base, + '/session/$remoteSessionId/message', + queryParameters: { + 'directory': workingDirectory, + 'limit': '20', + }, + ), + gatewayToken: gatewayToken, + ); + final text = _latestAssistantTextFromRestMessages(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)); + } + } + + String _latestAssistantTextFromRestMessages(List items) { + for (final raw in items.reversed) { + final item = _asMap(raw); + final info = _asMap(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 = _asMap(rawPart); + if (part['type']?.toString().trim() == 'text') { + final text = part['text']?.toString() ?? ''; + if (text.trim().isNotEmpty) { + return text; + } + } + } + } + return ''; + } + + bool _usesRestSessionApi(SingleAgentProvider provider) { + if (provider.providerId != SingleAgentProvider.opencode.providerId) { + return false; + } + final base = endpointResolver(provider); + final scheme = base?.scheme.toLowerCase() ?? ''; + return scheme == 'http' || scheme == 'https'; + } + + Uri _buildRestUri( + Uri base, + String path, { + Map? queryParameters, + }) { + final normalizedPath = path.startsWith('/') ? path : '/$path'; + return base.replace( + path: normalizedPath, + queryParameters: queryParameters, + fragment: null, + ); + } + + Future> _fetchJson( + 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 _decodeMap(body); + } finally { + client.close(force: true); + } + } + + Future> _postJson( + 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 _decodeMap(text); + } finally { + client.close(force: true); + } + } + + Future> _fetchJsonList( + 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? _extractThreadId(Map payload) { final topLevelId = payload['id']?.toString().trim() ?? ''; if (topLevelId.isNotEmpty) { diff --git a/test/runtime/direct_single_agent_app_server_suite.dart b/test/runtime/direct_single_agent_app_server_suite.dart index c688f412..42b73076 100644 --- a/test/runtime/direct_single_agent_app_server_suite.dart +++ b/test/runtime/direct_single_agent_app_server_suite.dart @@ -50,7 +50,7 @@ void main() { ).copyWith(onOutput: deltas.add), ); - expect(result.success, isTrue); + expect(result.success, isTrue, reason: result.errorMessage); expect(result.output, 'hello world from app server'); expect(result.resolvedModel, 'codex-sonnet'); expect(server.lastTurnInput, [ @@ -175,6 +175,54 @@ void main() { expect(result.resolvedModel, 'codex-sonnet'); }, ); + + test('probes OpenCode REST endpoint and reports provider support', () async { + final server = await _FakeOpenCodeRestServer.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 _FakeOpenCodeRestServer.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'); + }); }); } @@ -424,6 +472,211 @@ class _FakeAppServer { } } +class _FakeOpenCodeRestServer { + _FakeOpenCodeRestServer._(this._server); + + final HttpServer _server; + final List _eventResponses = []; + var _sessionCounter = 0; + var _messageCounter = 0; + bool healthRequested = false; + int createdSessionCount = 0; + String lastPromptText = ''; + final Map _assistantTextBySession = {}; + + Uri get baseHttpUri => Uri.parse('http://127.0.0.1:${_server.port}'); + + static Future<_FakeOpenCodeRestServer> start() async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + final fake = _FakeOpenCodeRestServer._(server); + unawaited(fake._listen()); + return fake; + } + + Future close() async { + for (final response in _eventResponses.toList(growable: false)) { + try { + await response.close(); + } catch (_) { + // Best effort. + } + } + await _server.close(force: true); + } + + Future _listen() async { + await for (final request in _server) { + 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(); + _eventResponses.add(request.response); + continue; + } + if (request.uri.path == '/session' && request.method == 'POST') { + createdSessionCount += 1; + final sessionId = 'ses-${_sessionCounter++}'; + 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 = _assistantTextBySession[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-${_messageCounter++}'; + await _broadcastEvent( + { + 'payload': { + 'type': 'session.status', + 'properties': { + 'sessionID': sessionId, + 'status': {'type': 'busy'}, + }, + }, + }, + ); + await _broadcastEvent( + { + 'payload': { + 'type': 'message.updated', + 'properties': { + 'sessionID': sessionId, + 'info': { + 'id': assistantMessageId, + 'role': 'assistant', + }, + }, + }, + }, + ); + for (final delta in ['hello ', 'world ', 'from ', 'opencode']) { + await _broadcastEvent( + { + 'payload': { + 'type': 'message.part.delta', + 'properties': { + 'sessionID': sessionId, + 'part': {'messageID': assistantMessageId}, + 'text': delta, + }, + }, + }, + ); + } + await _broadcastEvent( + { + 'payload': { + 'type': 'message.part.updated', + 'properties': { + 'sessionID': sessionId, + 'part': { + 'messageID': assistantMessageId, + 'type': 'text', + 'text': 'hello world from opencode', + }, + }, + }, + }, + ); + _assistantTextBySession[sessionId] = 'hello world from opencode'; + await _broadcastEvent( + { + '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 _broadcastEvent(Map event) async { + final payload = 'data: ${jsonEncode(event)}\n\n'; + for (final response in _eventResponses.toList(growable: false)) { + response.write(payload); + await response.flush(); + } + } +} + Map _decodeMap(Object raw) { if (raw is Map) { return raw;