diff --git a/lib/runtime/direct_single_agent_app_server_client.dart b/lib/runtime/direct_single_agent_app_server_client.dart index e6fc0067..2d789953 100644 --- a/lib/runtime/direct_single_agent_app_server_client.dart +++ b/lib/runtime/direct_single_agent_app_server_client.dart @@ -106,7 +106,11 @@ class DirectSingleAgentEndpointDescriptor { ); } final scheme = endpoint.scheme.toLowerCase(); - final normalizedBase = endpoint.replace(path: '', query: null, fragment: null); + final normalizedBase = endpoint.replace( + path: '', + query: null, + fragment: null, + ); final isLocal = _isLocalHost(endpoint.host); if (scheme == 'ws' && isLocal) { return DirectSingleAgentEndpointDescriptor( @@ -245,11 +249,11 @@ class DirectSingleAgentAppServerClient { sessionId, candidateBases: [ for (final entry in _transportKinds.entries) - if (entry.value == _DirectSingleAgentTransportKind.restSessionApi) - ...[ - if (_describeEndpoint(entry.key).baseUri != null) - _describeEndpoint(entry.key).baseUri!, - ], + if (entry.value == + _DirectSingleAgentTransportKind.restSessionApi) ...[ + if (_describeEndpoint(entry.key).baseUri != null) + _describeEndpoint(entry.key).baseUri!, + ], ], ); await _webSocketTransport.abort(sessionId); @@ -262,7 +266,9 @@ class DirectSingleAgentAppServerClient { DirectSingleAgentEndpointDescriptor _describeEndpoint( SingleAgentProvider provider, ) { - return DirectSingleAgentEndpointDescriptor.describe(endpointResolver(provider)); + return DirectSingleAgentEndpointDescriptor.describe( + endpointResolver(provider), + ); } Future<_ResolvedSingleAgentTransport> _resolveTransport( @@ -272,16 +278,16 @@ class DirectSingleAgentAppServerClient { }) async { final cachedKind = _transportKinds[provider]; if (cachedKind != null) { - final cachedEndpoint = cachedKind == - _DirectSingleAgentTransportKind.websocketAppServer + final cachedEndpoint = + cachedKind == _DirectSingleAgentTransportKind.websocketAppServer ? descriptor.websocketUri : descriptor.baseUri; if (cachedEndpoint != null) { return _ResolvedSingleAgentTransport( kind: cachedKind, endpoint: cachedEndpoint, - websocket: cachedKind == - _DirectSingleAgentTransportKind.websocketAppServer + websocket: + cachedKind == _DirectSingleAgentTransportKind.websocketAppServer ? _webSocketTransport : null, rest: cachedKind == _DirectSingleAgentTransportKind.restSessionApi @@ -356,10 +362,7 @@ class _DirectSingleAgentWebSocketTransport { final Map _threadIds = {}; final Set _abortedSessions = {}; - Future probe( - Uri endpoint, { - required String gatewayToken, - }) async { + Future probe(Uri endpoint, {required String gatewayToken}) async { _DirectAppServerConnection? connection; try { connection = await _DirectAppServerConnection.connect( @@ -602,10 +605,7 @@ class _DirectSingleAgentRestTransport { final Map _restSessionIds = {}; final Set _abortedSessions = {}; - Future probe( - Uri base, { - required String gatewayToken, - }) async { + Future probe(Uri base, {required String gatewayToken}) async { await _fetchJson( _buildRestUri(base, '/global/health'), gatewayToken: gatewayToken, @@ -639,6 +639,26 @@ class _DirectSingleAgentRestTransport { String? lastAssistantText; var busySeen = false; + bool hasResolvedAssistantContent() { + return output.toString().trim().isNotEmpty || + (lastAssistantText?.trim().isNotEmpty ?? false); + } + + void completeFailure(String message) { + if (completion.isCompleted) { + return; + } + completion.complete( + DirectSingleAgentRunResult( + success: false, + output: output.toString(), + errorMessage: message, + aborted: _abortedSessions.contains(normalizedSessionId), + resolvedModel: request.model, + ), + ); + } + final eventClient = HttpClient() ..connectionTimeout = const Duration(seconds: 8); late final HttpClientRequest eventRequest; @@ -652,6 +672,12 @@ class _DirectSingleAgentRestTransport { 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, @@ -707,17 +733,10 @@ class _DirectSingleAgentRestTransport { } 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, - ), + completeFailure( + error['message']?.toString() ?? + error['name']?.toString() ?? + 'OpenCode session failed.', ); return; } @@ -733,7 +752,8 @@ class _DirectSingleAgentRestTransport { if (activeAssistantMessageId != null && part['messageID']?.toString().trim() == activeAssistantMessageId) { - final delta = properties['text']?.toString() ?? + final delta = + properties['text']?.toString() ?? properties['delta']?.toString() ?? ''; if (delta.isNotEmpty) { @@ -793,19 +813,7 @@ class _DirectSingleAgentRestTransport { completeSuccess(); } }, - onError: (message) { - if (!completion.isCompleted) { - completion.complete( - DirectSingleAgentRunResult( - success: false, - output: output.toString(), - errorMessage: message, - aborted: _abortedSessions.contains(normalizedSessionId), - resolvedModel: request.model, - ), - ); - } - }, + onError: completeFailure, ), ); @@ -930,6 +938,7 @@ class _DirectSingleAgentRestTransport { } await Future.delayed(const Duration(milliseconds: 200)); } + onError('OpenCode REST session completed without assistant content.'); } String _latestAssistantTextFromRestMessages(List items) { diff --git a/test/runtime/direct_single_agent_app_server_suite.dart b/test/runtime/direct_single_agent_app_server_suite.dart index 9e760ef1..c2187f8f 100644 --- a/test/runtime/direct_single_agent_app_server_suite.dart +++ b/test/runtime/direct_single_agent_app_server_suite.dart @@ -256,6 +256,36 @@ void main() { expect(server.createdSessionCount, 1); expect(server.lastPromptText, 'hello opencode'); }); + + test( + 'fails OpenCode REST turns that complete without assistant content', + () async { + final server = await _FakeOpenCodeRestServer.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')); + }, + ); }); } @@ -506,9 +536,10 @@ class _FakeAppServer { } class _FakeOpenCodeRestServer { - _FakeOpenCodeRestServer._(this._server); + _FakeOpenCodeRestServer._(this._server, {required this.emitAssistantContent}); final HttpServer _server; + final bool emitAssistantContent; final List _eventResponses = []; var _sessionCounter = 0; var _messageCounter = 0; @@ -519,9 +550,14 @@ class _FakeOpenCodeRestServer { Uri get baseHttpUri => Uri.parse('http://127.0.0.1:${_server.port}'); - static Future<_FakeOpenCodeRestServer> start() async { + static Future<_FakeOpenCodeRestServer> start({ + bool emitAssistantContent = true, + }) async { final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); - final fake = _FakeOpenCodeRestServer._(server); + final fake = _FakeOpenCodeRestServer._( + server, + emitAssistantContent: emitAssistantContent, + ); unawaited(fake._listen()); return fake; } @@ -644,32 +680,39 @@ class _FakeOpenCodeRestServer { }, }, }); - for (final delta in ['hello ', 'world ', 'from ', 'opencode']) { + if (emitAssistantContent) { + 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.delta', + 'type': 'message.part.updated', 'properties': { 'sessionID': sessionId, - 'part': {'messageID': assistantMessageId}, - 'text': delta, + 'part': { + 'messageID': assistantMessageId, + 'type': 'text', + 'text': 'hello world from opencode', + }, }, }, }); + _assistantTextBySession[sessionId] = 'hello world from opencode'; } - 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',