diff --git a/lib/app/app_controller_desktop_core.dart b/lib/app/app_controller_desktop_core.dart index 10a2ba26..e8407da7 100644 --- a/lib/app/app_controller_desktop_core.dart +++ b/lib/app/app_controller_desktop_core.dart @@ -29,8 +29,9 @@ import '../runtime/codex_config_bridge.dart'; import '../runtime/code_agent_node_orchestrator.dart'; import '../runtime/assistant_artifacts.dart'; import '../runtime/desktop_thread_artifact_service.dart'; -import '../runtime/go_agent_core_client.dart'; import '../runtime/go_agent_core_desktop_transport.dart'; +import '../runtime/go_task_service_client.dart'; +import '../runtime/go_task_service_desktop_service.dart'; import '../runtime/go_gateway_runtime_desktop_client.dart'; import '../runtime/go_multi_agent_mount_desktop_client.dart'; import '../runtime/go_runtime_dispatch_desktop_client.dart'; @@ -127,7 +128,7 @@ class AppController extends ChangeNotifier { List? singleAgentSharedSkillScanRootOverrides, List? availableSingleAgentProvidersOverride, ArisBundleRepository? arisBundleRepository, - GoAgentCoreClient? goAgentCoreClient, + GoTaskServiceClient? goTaskServiceClient, MultiAgentMountManager? multiAgentMountManager, }) { storeInternal = store ?? SecureConfigStore(); @@ -214,12 +215,15 @@ class AppController extends ChangeNotifier { goCoreLocator: goCoreLocatorInternal, ), ); - goAgentCoreClientInternal = - goAgentCoreClient ?? - GoAgentCoreDesktopTransport( - acpClient: gatewayAcpClientInternal, - endpointResolver: resolveGoAgentCoreEndpointForTargetInternal, - goCoreLocator: goCoreLocatorInternal, + goTaskServiceClientInternal = + goTaskServiceClient ?? + DesktopGoTaskService( + gateway: runtimeCoordinatorInternal.gateway, + acpTransport: ExternalCodeAgentAcpDesktopTransport( + acpClient: gatewayAcpClientInternal, + endpointResolver: resolveGoAgentCoreEndpointForTargetInternal, + goCoreLocator: goCoreLocatorInternal, + ), ); multiAgentOrchestratorInternal = MultiAgentOrchestrator( config: resolveMultiAgentConfigInternal( @@ -267,7 +271,7 @@ class AppController extends ChangeNotifier { storeInternal.dispose(); desktopPlatformServiceInternal.dispose(); unawaited(multiAgentMountManagerInternal.dispose()); - unawaited(goAgentCoreClientInternal.dispose()); + unawaited(goTaskServiceClientInternal.dispose()); unawaited(gatewayAcpClientInternal.dispose()); super.dispose(); } @@ -298,7 +302,7 @@ class AppController extends ChangeNotifier { availableSingleAgentProvidersOverrideInternal; late final ArisBundleRepository arisBundleRepositoryInternal; late final GoCoreLocator goCoreLocatorInternal; - late final GoAgentCoreClient goAgentCoreClientInternal; + late final GoTaskServiceClient goTaskServiceClientInternal; late final MultiAgentOrchestrator multiAgentOrchestratorInternal; late final MultiAgentMountManager multiAgentMountManagerInternal; Map @@ -319,8 +323,8 @@ class AppController extends ChangeNotifier { final Map> latestRoutingResolutionBySessionInternal = >{}; - final Map syncedGoAgentProvidersInternal = - {}; + final Map + syncedGoAgentProvidersInternal = {}; final DesktopThreadArtifactService threadArtifactServiceInternal = DesktopThreadArtifactService(); List singleAgentSharedImportedSkillsInternal = diff --git a/lib/app/app_controller_desktop_go_agent_core_routing.dart b/lib/app/app_controller_desktop_go_agent_core_routing.dart index e7db8be5..068a278e 100644 --- a/lib/app/app_controller_desktop_go_agent_core_routing.dart +++ b/lib/app/app_controller_desktop_go_agent_core_routing.dart @@ -28,7 +28,7 @@ import '../runtime/codex_config_bridge.dart'; import '../runtime/code_agent_node_orchestrator.dart'; import '../runtime/assistant_artifacts.dart'; import '../runtime/desktop_thread_artifact_service.dart'; -import '../runtime/go_agent_core_client.dart'; +import '../runtime/go_task_service_client.dart'; import '../runtime/mode_switcher.dart'; import '../runtime/agent_registry.dart'; import '../runtime/multi_agent_orchestrator.dart'; @@ -39,9 +39,9 @@ import 'app_controller_desktop_core.dart'; import 'app_controller_desktop_thread_sessions.dart'; extension AppControllerDesktopGoAgentCoreRouting on AppController { - Future> - buildGoAgentCoreSyncedProvidersInternal() async { - final providers = []; + Future> + buildExternalAcpSyncedProvidersInternal() async { + final providers = []; for (final profile in settings.externalAcpEndpoints) { final providerId = profile.providerKey.trim(); final endpoint = profile.endpoint.trim(); @@ -54,7 +54,7 @@ extension AppControllerDesktopGoAgentCoreRouting on AppController { refName: profile.authRef.trim(), ); providers.add( - GoAgentCoreSyncedProvider( + ExternalCodeAgentAcpSyncedProvider( providerId: providerId, label: profile.label, endpoint: endpoint, @@ -66,19 +66,19 @@ extension AppControllerDesktopGoAgentCoreRouting on AppController { return providers; } - Future syncGoAgentCoreProvidersInternal() async { - final providers = await buildGoAgentCoreSyncedProvidersInternal(); + Future syncExternalAcpProvidersInternal() async { + final providers = await buildExternalAcpSyncedProvidersInternal(); syncedGoAgentProvidersInternal ..clear() ..addEntries( providers.map((item) => MapEntry(item.providerId.trim(), item)), ); - await goAgentCoreClientInternal.syncProviders(providers); + await goTaskServiceClientInternal.syncExternalProviders(providers); } void updateLatestRoutingResolutionInternal( String sessionKey, - GoAgentCoreRunResult result, + GoTaskServiceResult result, ) { final normalizedSessionKey = normalizedAssistantSessionKeyInternal( sessionKey, diff --git a/lib/app/app_controller_desktop_runtime_coordination_impl.dart b/lib/app/app_controller_desktop_runtime_coordination_impl.dart index ccbbfb55..f19e4c95 100644 --- a/lib/app/app_controller_desktop_runtime_coordination_impl.dart +++ b/lib/app/app_controller_desktop_runtime_coordination_impl.dart @@ -84,9 +84,9 @@ Future refreshSingleAgentCapabilitiesRuntimeInternal( AppController controller, { bool forceRefresh = false, }) async { - await controller.syncGoAgentCoreProvidersInternal(); - final capabilities = await controller.goAgentCoreClientInternal - .loadCapabilities( + await controller.syncExternalAcpProvidersInternal(); + final capabilities = await controller.goTaskServiceClientInternal + .loadExternalAcpCapabilities( target: AssistantExecutionTarget.singleAgent, forceRefresh: forceRefresh, ); diff --git a/lib/app/app_controller_desktop_single_agent.dart b/lib/app/app_controller_desktop_single_agent.dart index e3376d39..1755176e 100644 --- a/lib/app/app_controller_desktop_single_agent.dart +++ b/lib/app/app_controller_desktop_single_agent.dart @@ -28,7 +28,7 @@ import '../runtime/codex_config_bridge.dart'; import '../runtime/code_agent_node_orchestrator.dart'; import '../runtime/assistant_artifacts.dart'; import '../runtime/desktop_thread_artifact_service.dart'; -import '../runtime/go_agent_core_client.dart'; +import '../runtime/go_task_service_client.dart'; import '../runtime/mode_switcher.dart'; import '../runtime/agent_registry.dart'; import '../runtime/multi_agent_orchestrator.dart'; @@ -84,9 +84,10 @@ extension AppControllerDesktopSingleAgent on AppController { notifyIfActiveInternal(); try { - final routing = buildGoAgentCoreRoutingForSessionInternal(sessionKey); + final routing = buildExternalAcpRoutingForSessionInternal(sessionKey); final selection = singleAgentProviderForSession(sessionKey); - final capabilities = await goAgentCoreClientInternal.loadCapabilities( + final capabilities = await goTaskServiceClientInternal + .loadExternalAcpCapabilities( target: AssistantExecutionTarget.singleAgent, forceRefresh: true, ); @@ -175,8 +176,8 @@ extension AppControllerDesktopSingleAgent on AppController { .map((item) => item.label.trim().isNotEmpty ? item.label : item.key) .where((item) => item.trim().isNotEmpty) .toList(growable: false); - final result = await goAgentCoreClientInternal.executeSession( - GoAgentCoreSessionRequest( + final result = await goTaskServiceClientInternal.executeTask( + GoTaskServiceRequest( sessionId: sessionKey, threadId: sessionKey, target: AssistantExecutionTarget.singleAgent, @@ -838,7 +839,7 @@ extension AppControllerDesktopSingleAgent on AppController { ); } - GoAgentCoreRoutingConfig buildGoAgentCoreRoutingForSessionInternal( + ExternalCodeAgentAcpRoutingConfig buildExternalAcpRoutingForSessionInternal( String sessionKey, { String? explicitExecutionTarget, }) { @@ -861,7 +862,7 @@ extension AppControllerDesktopSingleAgent on AppController { final availableSkills = assistantImportedSkillsForSession(normalizedSessionKey) .map( - (item) => GoAgentCoreAvailableSkill( + (item) => ExternalCodeAgentAcpAvailableSkill( id: item.key, label: item.label, description: item.description, @@ -905,14 +906,14 @@ extension AppControllerDesktopSingleAgent on AppController { resolvedExplicitSkills.isNotEmpty; if (!hasExplicitSelection) { - return GoAgentCoreRoutingConfig.auto( + return ExternalCodeAgentAcpRoutingConfig.auto( preferredGatewayTarget: preferredGatewayTarget, availableSkills: availableSkills, ); } - return GoAgentCoreRoutingConfig( - mode: GoAgentCoreRoutingMode.explicit, + return ExternalCodeAgentAcpRoutingConfig( + mode: ExternalCodeAgentAcpRoutingMode.explicit, preferredGatewayTarget: preferredGatewayTarget, explicitExecutionTarget: resolvedExplicitExecutionTarget, explicitProviderId: resolvedExplicitProviderId, diff --git a/lib/app/app_controller_desktop_thread_actions.dart b/lib/app/app_controller_desktop_thread_actions.dart index a0d21a6e..7d891b9a 100644 --- a/lib/app/app_controller_desktop_thread_actions.dart +++ b/lib/app/app_controller_desktop_thread_actions.dart @@ -28,7 +28,7 @@ import '../runtime/codex_config_bridge.dart'; import '../runtime/code_agent_node_orchestrator.dart'; import '../runtime/assistant_artifacts.dart'; import '../runtime/desktop_thread_artifact_service.dart'; -import '../runtime/go_agent_core_client.dart'; +import '../runtime/go_task_service_client.dart'; import '../runtime/mode_switcher.dart'; import '../runtime/agent_registry.dart'; import '../runtime/multi_agent_orchestrator.dart'; @@ -313,8 +313,8 @@ extension AppControllerDesktopThreadActions on AppController { try { final dispatch = await codeAgentNodeOrchestratorInternal .buildGatewayDispatch(buildCodeAgentNodeStateInternal()); - final result = await goAgentCoreClientInternal.executeSession( - GoAgentCoreSessionRequest( + final result = await goTaskServiceClientInternal.executeTask( + GoTaskServiceRequest( sessionId: sessionKey, threadId: sessionKey, target: currentTarget, @@ -331,7 +331,7 @@ extension AppControllerDesktopThreadActions on AppController { aiGatewayApiKey: await loadAiGatewayApiKey(), agentId: dispatch.agentId ?? '', metadata: dispatch.metadata, - routing: buildGoAgentCoreRoutingForSessionInternal(sessionKey), + routing: buildExternalAcpRoutingForSessionInternal(sessionKey), ), onUpdate: (update) { if (update.isDelta) { @@ -426,7 +426,8 @@ extension AppControllerDesktopThreadActions on AppController { sessionsControllerInternal.currentSessionKey, ); if (aiGatewayPendingSessionKeysInternal.contains(sessionKey)) { - await goAgentCoreClientInternal.cancelSession( + await goTaskServiceClientInternal.cancelTask( + route: GoTaskServiceRoute.externalAcpSingle, target: AssistantExecutionTarget.singleAgent, sessionId: sessionKey, threadId: sessionKey, @@ -444,7 +445,13 @@ extension AppControllerDesktopThreadActions on AppController { sessionsControllerInternal.currentSessionKey, ); if (aiGatewayPendingSessionKeysInternal.contains(sessionKey)) { - await goAgentCoreClientInternal.cancelSession( + await goTaskServiceClientInternal.cancelTask( + route: assistantExecutionTargetForSession(sessionKey) == + AssistantExecutionTarget.singleAgent || + assistantExecutionTargetForSession(sessionKey) == + AssistantExecutionTarget.auto + ? GoTaskServiceRoute.externalAcpSingle + : GoTaskServiceRoute.openClawTask, target: assistantExecutionTargetForSession(sessionKey), sessionId: sessionKey, threadId: sessionKey, diff --git a/lib/app/app_controller_web_core.dart b/lib/app/app_controller_web_core.dart index 9f85ed08..c799c100 100644 --- a/lib/app/app_controller_web_core.dart +++ b/lib/app/app_controller_web_core.dart @@ -5,10 +5,11 @@ import 'package:flutter/material.dart'; import '../i18n/app_language.dart'; import '../models/app_models.dart'; import '../runtime/assistant_artifacts.dart'; -import '../runtime/go_agent_core_client.dart'; +import '../runtime/go_task_service_client.dart'; import '../runtime/runtime_models.dart'; import '../web/web_acp_client.dart'; import '../web/go_agent_core_web_transport.dart'; +import '../web/go_task_service_web_service.dart'; import '../web/web_ai_gateway_client.dart'; import '../web/web_artifact_proxy_client.dart'; import '../web/web_relay_gateway_client.dart'; @@ -38,7 +39,7 @@ class AppController extends ChangeNotifier { WebStore? store, WebAiGatewayClient? aiGatewayClient, WebAcpClient? acpClient, - GoAgentCoreClient? goAgentCoreClient, + GoTaskServiceClient? goTaskServiceClient, WebRelayGatewayClient? relayClient, RemoteWebSessionRepositoryBuilder? remoteSessionRepositoryBuilder, UiFeatureManifest? uiFeatureManifest, @@ -51,11 +52,14 @@ class AppController extends ChangeNotifier { remoteSessionRepositoryBuilder ?? defaultRemoteSessionRepositoryInternal { relayClientInternal = relayClient ?? WebRelayGatewayClient(storeInternal); - goAgentCoreClientInternal = - goAgentCoreClient ?? - GoAgentCoreWebTransport( - acpClient: acpClientInternal, - endpointResolver: acpEndpointForTargetInternal, + goTaskServiceClientInternal = + goTaskServiceClient ?? + WebGoTaskService( + relayClient: relayClientInternal, + acpTransport: ExternalCodeAgentAcpWebTransport( + acpClient: acpClientInternal, + endpointResolver: acpEndpointForTargetInternal, + ), ); artifactProxyClientInternal = WebArtifactProxyClient(relayClientInternal); relayEventsSubscriptionInternal = relayClientInternal.events.listen( @@ -68,7 +72,7 @@ class AppController extends ChangeNotifier { final UiFeatureManifest uiFeatureManifestInternal; final WebAiGatewayClient aiGatewayClientInternal; final WebAcpClient acpClientInternal; - late final GoAgentCoreClient goAgentCoreClientInternal; + late final GoTaskServiceClient goTaskServiceClientInternal; final RemoteWebSessionRepositoryBuilder remoteSessionRepositoryBuilderInternal; late final WebRelayGatewayClient relayClientInternal; @@ -97,6 +101,7 @@ class AppController extends ChangeNotifier { final WebTaskThreadRepository threadRepositoryInternal = WebTaskThreadRepository(); final Set pendingSessionKeysInternal = {}; + final Set goTaskServiceManagedRelaySessionsInternal = {}; final Map streamingTextBySessionInternal = {}; final Map> threadTurnQueuesInternal = >{}; @@ -330,7 +335,7 @@ class AppController extends ChangeNotifier { @override void dispose() { unawaited(relayEventsSubscriptionInternal.cancel()); - unawaited(goAgentCoreClientInternal.dispose()); + unawaited(goTaskServiceClientInternal.dispose()); unawaited(relayClientInternal.dispose()); super.dispose(); } diff --git a/lib/app/app_controller_web_gateway_chat.dart b/lib/app/app_controller_web_gateway_chat.dart index b045b90d..4e43fa83 100644 --- a/lib/app/app_controller_web_gateway_chat.dart +++ b/lib/app/app_controller_web_gateway_chat.dart @@ -5,7 +5,7 @@ import 'package:flutter/material.dart'; import '../i18n/app_language.dart'; import '../models/app_models.dart'; import '../runtime/assistant_artifacts.dart'; -import '../runtime/go_agent_core_client.dart'; +import '../runtime/go_task_service_client.dart'; import '../runtime/runtime_models.dart'; import '../web/web_acp_client.dart'; import '../web/web_ai_gateway_client.dart'; @@ -108,7 +108,7 @@ extension AppControllerWebGatewayChat on AppController { } if (target == AssistantExecutionTarget.singleAgent || target == AssistantExecutionTarget.auto) { - await executeGoAgentCoreRunInternal( + await executeGoTaskServiceRunInternal( sessionKey: sessionKey, prompt: trimmed, target: target == AssistantExecutionTarget.auto @@ -125,7 +125,7 @@ extension AppControllerWebGatewayChat on AppController { selectedSkillLabels: selectedSkillLabels, ); } else { - await executeGoAgentCoreRunInternal( + await executeGoTaskServiceRunInternal( sessionKey: sessionKey, prompt: trimmed, target: target, @@ -164,8 +164,8 @@ extension AppControllerWebGatewayChat on AppController { pendingSessionKeysInternal.add(sessionKey); notifyChangedInternal(); try { - final result = await goAgentCoreClientInternal.executeSession( - GoAgentCoreSessionRequest( + final result = await goTaskServiceClientInternal.executeTask( + GoTaskServiceRequest( sessionId: sessionKey, threadId: sessionKey, target: assistantExecutionTargetForSession(sessionKey), @@ -180,7 +180,7 @@ extension AppControllerWebGatewayChat on AppController { aiGatewayApiKey: aiGatewayApiKeyCacheInternal.trim(), agentId: selectedAgentId, metadata: const {}, - routing: buildWebGoAgentCoreRoutingForSessionInternal( + routing: buildWebExternalAcpRoutingForSessionInternal( sessionKey, explicitExecutionTarget: 'multiAgent', ), @@ -230,7 +230,7 @@ extension AppControllerWebGatewayChat on AppController { notifyChangedInternal(); } - Future executeGoAgentCoreRunInternal({ + Future executeGoTaskServiceRunInternal({ required String sessionKey, required String prompt, required AssistantExecutionTarget target, @@ -244,8 +244,7 @@ extension AppControllerWebGatewayChat on AppController { .map((item) => item.trim()) .where((item) => item.isNotEmpty) .toList(growable: false); - final result = await goAgentCoreClientInternal.executeSession( - GoAgentCoreSessionRequest( + final request = GoTaskServiceRequest( sessionId: sessionKey, threadId: sessionKey, target: target, @@ -262,37 +261,53 @@ extension AppControllerWebGatewayChat on AppController { metadata: { if (selectedSkills.isNotEmpty) 'selectedSkills': selectedSkills, }, - routing: buildWebGoAgentCoreRoutingForSessionInternal(sessionKey), + routing: buildWebExternalAcpRoutingForSessionInternal(sessionKey), provider: provider, - ), - onUpdate: (update) { - if (update.isDelta) { - appendStreamingTextInternal(sessionKey, update.text); - notifyChangedInternal(); - } - }, - ); - final message = result.message.trim(); - if (!result.success && result.errorMessage.trim().isNotEmpty) { - throw Exception(result.errorMessage.trim()); - } - if (message.isEmpty) { - throw Exception( - appText( - 'Go Agent-core 没有返回可显示的输出。', - 'Go Agent-core returned no displayable output.', - ), ); + final route = request.route; + if (route == GoTaskServiceRoute.openClawTask) { + goTaskServiceManagedRelaySessionsInternal.add(sessionKey); + } + try { + final result = await goTaskServiceClientInternal.executeTask( + request, + onUpdate: (update) { + if (update.isDelta) { + appendStreamingTextInternal(sessionKey, update.text); + notifyChangedInternal(); + } + }, + ); + final message = result.message.trim(); + if (!result.success && result.errorMessage.trim().isNotEmpty) { + throw Exception(result.errorMessage.trim()); + } + if (message.isEmpty) { + throw Exception( + appText( + route == GoTaskServiceRoute.openClawTask + ? 'OpenClaw task 没有返回可显示的输出。' + : 'Go Task Service 没有返回可显示的输出。', + route == GoTaskServiceRoute.openClawTask + ? 'OpenClaw task returned no displayable output.' + : 'Go Task Service returned no displayable output.', + ), + ); + } + appendAssistantMessageInternal( + sessionKey: sessionKey, + text: message, + error: false, + ); + } finally { + clearStreamingTextInternal(sessionKey); + if (route == GoTaskServiceRoute.openClawTask) { + goTaskServiceManagedRelaySessionsInternal.remove(sessionKey); + } } - appendAssistantMessageInternal( - sessionKey: sessionKey, - text: message, - error: false, - ); - clearStreamingTextInternal(sessionKey); } - GoAgentCoreRoutingConfig buildWebGoAgentCoreRoutingForSessionInternal( + ExternalCodeAgentAcpRoutingConfig buildWebExternalAcpRoutingForSessionInternal( String sessionKey, { String? explicitExecutionTarget, }) { @@ -310,7 +325,7 @@ extension AppControllerWebGatewayChat on AppController { final availableSkills = assistantImportedSkillsForSession(normalizedSessionKey) .map( - (item) => GoAgentCoreAvailableSkill( + (item) => ExternalCodeAgentAcpAvailableSkill( id: item.key, label: item.label, description: item.description, @@ -359,14 +374,14 @@ extension AppControllerWebGatewayChat on AppController { resolvedExplicitSkills.isNotEmpty; if (!hasExplicitSelection) { - return GoAgentCoreRoutingConfig.auto( + return ExternalCodeAgentAcpRoutingConfig.auto( preferredGatewayTarget: preferredGatewayTarget, availableSkills: availableSkills, ); } - return GoAgentCoreRoutingConfig( - mode: GoAgentCoreRoutingMode.explicit, + return ExternalCodeAgentAcpRoutingConfig( + mode: ExternalCodeAgentAcpRoutingMode.explicit, preferredGatewayTarget: preferredGatewayTarget, explicitExecutionTarget: resolvedExplicitExecutionTarget, explicitProviderId: resolvedExplicitProviderId, diff --git a/lib/app/app_controller_web_helpers.dart b/lib/app/app_controller_web_helpers.dart index b40bb2d6..06195746 100644 --- a/lib/app/app_controller_web_helpers.dart +++ b/lib/app/app_controller_web_helpers.dart @@ -369,6 +369,9 @@ extension AppControllerWebHelpers on AppController { if (sessionKey.isEmpty) { return; } + if (goTaskServiceManagedRelaySessionsInternal.contains(sessionKey)) { + return; + } final state = payload['state']?.toString().trim() ?? ''; final message = castMapInternal(payload['message']); final text = extractMessageTextInternal(message); diff --git a/lib/runtime/go_agent_core_desktop_transport.dart b/lib/runtime/go_agent_core_desktop_transport.dart index 1035a0f3..68c9bca0 100644 --- a/lib/runtime/go_agent_core_desktop_transport.dart +++ b/lib/runtime/go_agent_core_desktop_transport.dart @@ -3,11 +3,11 @@ import 'dart:io'; import 'embedded_agent_launch_policy.dart'; import 'gateway_acp_client.dart'; -import 'go_agent_core_client.dart'; import 'go_core.dart'; +import 'go_task_service_client.dart'; import 'runtime_models.dart'; -typedef GoAgentCoreProcessStarter = +typedef ExternalCodeAgentAcpProcessStarter = Future Function( String executable, List arguments, { @@ -15,12 +15,12 @@ typedef GoAgentCoreProcessStarter = String? workingDirectory, }); -class GoAgentCoreDesktopTransport implements GoAgentCoreClient { - GoAgentCoreDesktopTransport({ +class ExternalCodeAgentAcpDesktopTransport implements ExternalCodeAgentAcpTransport { + ExternalCodeAgentAcpDesktopTransport({ required GatewayAcpClient acpClient, required Uri? Function(AssistantExecutionTarget target) endpointResolver, GoCoreLocator? goCoreLocator, - GoAgentCoreProcessStarter? processStarter, + ExternalCodeAgentAcpProcessStarter? processStarter, }) : _acpClient = acpClient, _endpointResolver = endpointResolver, _goCoreLocator = goCoreLocator ?? GoCoreLocator(), @@ -38,14 +38,16 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { final GatewayAcpClient _acpClient; final Uri? Function(AssistantExecutionTarget target) _endpointResolver; final GoCoreLocator _goCoreLocator; - final GoAgentCoreProcessStarter _processStarter; + final ExternalCodeAgentAcpProcessStarter _processStarter; Process? _localProcess; Uri? _localEndpoint; Future? _localEndpointFuture; @override - Future syncProviders(List providers) async { + Future syncExternalProviders( + List providers, + ) async { final endpoint = await _ensureLocalEndpoint(); if (endpoint == null) { return; @@ -60,19 +62,19 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { } @override - Future loadCapabilities({ + Future loadExternalAcpCapabilities({ required AssistantExecutionTarget target, bool forceRefresh = false, }) async { final endpoint = await _resolveEndpoint(target); if (endpoint == null) { - return const GoAgentCoreCapabilities.empty(); + return const ExternalCodeAgentAcpCapabilities.empty(); } final capabilities = await _acpClient.loadCapabilities( forceRefresh: forceRefresh, endpointOverride: endpoint, ); - return GoAgentCoreCapabilities( + return ExternalCodeAgentAcpCapabilities( singleAgent: capabilities.singleAgent, multiAgent: capabilities.multiAgent, providers: capabilities.providers, @@ -81,9 +83,9 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { } @override - Future executeSession( - GoAgentCoreSessionRequest request, { - required void Function(GoAgentCoreSessionUpdate update) onUpdate, + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, }) async { final endpoint = await _resolveEndpoint(request.target); if (endpoint == null) { @@ -96,10 +98,10 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { String? completedMessage; final response = await _acpClient.request( method: request.resumeSession ? 'session.message' : 'session.start', - params: request.toAcpParams(), + params: request.toExternalAcpParams(), endpointOverride: endpoint, onNotification: (notification) { - final update = goAgentCoreUpdateFromNotification(notification); + final update = goTaskServiceUpdateFromAcpNotification(notification); if (update == null) { return; } @@ -112,15 +114,16 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { onUpdate(update); }, ); - return goAgentCoreRunResultFromResponse( + return goTaskServiceResultFromAcpResponse( response, + route: request.route, streamedText: streamedText, completedMessage: completedMessage, ); } @override - Future cancelSession({ + Future cancelTask({ required AssistantExecutionTarget target, required String sessionId, required String threadId, @@ -137,7 +140,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient { } @override - Future closeSession({ + Future closeTask({ required AssistantExecutionTarget target, required String sessionId, required String threadId, diff --git a/lib/runtime/go_task_service_client.dart b/lib/runtime/go_task_service_client.dart new file mode 100644 index 00000000..e8186f01 --- /dev/null +++ b/lib/runtime/go_task_service_client.dart @@ -0,0 +1,635 @@ +import 'runtime_models.dart'; + +enum GoTaskServiceRoute { openClawTask, externalAcpSingle, externalAcpMulti } + +class ExternalCodeAgentAcpCapabilities { + const ExternalCodeAgentAcpCapabilities({ + required this.singleAgent, + required this.multiAgent, + required this.providers, + required this.raw, + }); + + const ExternalCodeAgentAcpCapabilities.empty() + : singleAgent = false, + multiAgent = false, + providers = const {}, + raw = const {}; + + final bool singleAgent; + final bool multiAgent; + final Set providers; + final Map raw; +} + +class ExternalCodeAgentAcpSyncedProvider { + const ExternalCodeAgentAcpSyncedProvider({ + required this.providerId, + required this.label, + required this.endpoint, + required this.authorizationHeader, + required this.enabled, + }); + + final String providerId; + final String label; + final String endpoint; + final String authorizationHeader; + final bool enabled; + + Map toJson() { + return { + 'providerId': providerId.trim(), + 'label': label.trim(), + 'endpoint': endpoint.trim(), + 'authorizationHeader': authorizationHeader.trim(), + 'enabled': enabled, + }; + } +} + +enum ExternalCodeAgentAcpRoutingMode { auto, explicit } + +class ExternalCodeAgentAcpAvailableSkill { + const ExternalCodeAgentAcpAvailableSkill({ + required this.id, + required this.label, + required this.description, + this.installed = true, + }); + + final String id; + final String label; + final String description; + final bool installed; + + Map toJson() { + return { + 'id': id.trim(), + 'label': label.trim(), + 'description': description.trim(), + 'installed': installed, + }; + } +} + +class ExternalCodeAgentAcpRoutingConfig { + const ExternalCodeAgentAcpRoutingConfig({ + required this.mode, + required this.preferredGatewayTarget, + required this.explicitExecutionTarget, + required this.explicitProviderId, + required this.explicitModel, + required this.explicitSkills, + required this.allowSkillInstall, + required this.availableSkills, + this.installApproval, + }); + + const ExternalCodeAgentAcpRoutingConfig.auto({ + this.preferredGatewayTarget = '', + this.availableSkills = const [], + }) : mode = ExternalCodeAgentAcpRoutingMode.auto, + explicitExecutionTarget = '', + explicitProviderId = '', + explicitModel = '', + explicitSkills = const [], + allowSkillInstall = false, + installApproval = null; + + final ExternalCodeAgentAcpRoutingMode mode; + final String preferredGatewayTarget; + final String explicitExecutionTarget; + final String explicitProviderId; + final String explicitModel; + final List explicitSkills; + final bool allowSkillInstall; + final List availableSkills; + final ExternalCodeAgentAcpSkillInstallApproval? installApproval; + + bool get isAuto => mode == ExternalCodeAgentAcpRoutingMode.auto; + + Map toJson() { + return { + 'routingMode': mode.name, + if (preferredGatewayTarget.trim().isNotEmpty) + 'preferredGatewayTarget': preferredGatewayTarget.trim(), + if (explicitExecutionTarget.trim().isNotEmpty) + 'explicitExecutionTarget': explicitExecutionTarget.trim(), + if (explicitProviderId.trim().isNotEmpty) + 'explicitProviderId': explicitProviderId.trim(), + if (explicitModel.trim().isNotEmpty) + 'explicitModel': explicitModel.trim(), + 'explicitSkills': explicitSkills + .map((item) => item.trim()) + .where((item) => item.isNotEmpty) + .toList(growable: false), + 'allowSkillInstall': allowSkillInstall, + 'availableSkills': availableSkills + .map((item) => item.toJson()) + .toList(growable: false), + if (installApproval != null) 'installApproval': installApproval!.toJson(), + }; + } +} + +class ExternalCodeAgentAcpSkillInstallApproval { + const ExternalCodeAgentAcpSkillInstallApproval({ + required this.requestId, + required this.approvedSkillKeys, + }); + + final String requestId; + final List approvedSkillKeys; + + Map toJson() { + return { + 'requestId': requestId.trim(), + 'approvedSkillKeys': approvedSkillKeys + .map((item) => item.trim()) + .where((item) => item.isNotEmpty) + .toList(growable: false), + }; + } +} + +class GoTaskServiceRequest { + const GoTaskServiceRequest({ + required this.sessionId, + required this.threadId, + required this.target, + required this.prompt, + required this.workingDirectory, + required this.model, + required this.thinking, + required this.selectedSkills, + required this.inlineAttachments, + required this.localAttachments, + required this.aiGatewayBaseUrl, + required this.aiGatewayApiKey, + required this.agentId, + required this.metadata, + this.routing, + this.provider = SingleAgentProvider.auto, + this.resumeSession = false, + this.multiAgent = false, + }); + + final String sessionId; + final String threadId; + final AssistantExecutionTarget target; + final String prompt; + final String workingDirectory; + final String model; + final String thinking; + final List selectedSkills; + final List inlineAttachments; + final List localAttachments; + final String aiGatewayBaseUrl; + final String aiGatewayApiKey; + final String agentId; + final Map metadata; + final ExternalCodeAgentAcpRoutingConfig? routing; + final SingleAgentProvider provider; + final bool resumeSession; + final bool multiAgent; + + GoTaskServiceRoute get route { + if (multiAgent) { + return GoTaskServiceRoute.externalAcpMulti; + } + return switch (target) { + AssistantExecutionTarget.local => GoTaskServiceRoute.openClawTask, + AssistantExecutionTarget.remote => GoTaskServiceRoute.openClawTask, + AssistantExecutionTarget.singleAgent => GoTaskServiceRoute.externalAcpSingle, + AssistantExecutionTarget.auto => GoTaskServiceRoute.externalAcpSingle, + }; + } + + String get acpMode { + if (route == GoTaskServiceRoute.externalAcpMulti) { + return 'multi-agent'; + } + return switch (target) { + AssistantExecutionTarget.auto => 'single-agent', + AssistantExecutionTarget.singleAgent => 'single-agent', + AssistantExecutionTarget.local => _gatewaySessionMode, + AssistantExecutionTarget.remote => _gatewaySessionMode, + }; + } + + String get routingExecutionTarget { + if (route == GoTaskServiceRoute.externalAcpMulti) { + return 'multi-agent'; + } + return switch (target) { + AssistantExecutionTarget.auto => 'single-agent', + AssistantExecutionTarget.singleAgent => 'single-agent', + AssistantExecutionTarget.local => 'gateway', + AssistantExecutionTarget.remote => 'gateway', + }; + } + + bool get hasInlineAttachments => inlineAttachments.isNotEmpty; + + ExternalCodeAgentAcpRoutingConfig get effectiveRouting => + routing ?? _synthesizedRouting(); + + Map toExternalAcpParams() { + final resolvedRouting = effectiveRouting; + final params = { + 'sessionId': sessionId, + 'threadId': threadId, + 'mode': acpMode, + 'taskPrompt': prompt, + 'workingDirectory': workingDirectory.trim(), + 'selectedSkills': selectedSkills, + 'attachments': >[ + ...localAttachments.map( + (item) => { + 'name': item.name, + 'description': item.description, + 'path': item.path, + }, + ), + ...inlineAttachments.map( + (item) => { + 'name': item.fileName, + 'description': item.mimeType, + 'path': '', + }, + ), + ], + if (inlineAttachments.isNotEmpty) + 'inlineAttachments': inlineAttachments + .map( + (item) => { + 'name': item.fileName, + 'mimeType': item.mimeType, + 'content': item.content, + 'sizeBytes': goTaskServiceBase64Size(item.content), + }, + ) + .toList(growable: false), + if (provider != SingleAgentProvider.auto) 'provider': provider.providerId, + if (model.trim().isNotEmpty) 'model': model.trim(), + if (thinking.trim().isNotEmpty) 'thinking': thinking.trim(), + if (aiGatewayBaseUrl.trim().isNotEmpty) + 'aiGatewayBaseUrl': aiGatewayBaseUrl.trim(), + if (aiGatewayApiKey.trim().isNotEmpty) + 'aiGatewayApiKey': aiGatewayApiKey.trim(), + 'routing': resolvedRouting.toJson(), + if (_usesGatewaySessionMode(acpMode)) ...{ + 'executionTarget': target.promptValue, + if (agentId.trim().isNotEmpty) 'agentId': agentId.trim(), + if (metadata.isNotEmpty) 'metadata': metadata, + }, + }; + return params; + } + + ExternalCodeAgentAcpRoutingConfig _synthesizedRouting() { + final preferredGatewayTarget = switch (target) { + AssistantExecutionTarget.remote => 'remote', + _ => 'local', + }; + final explicitExecutionTarget = switch (target) { + AssistantExecutionTarget.local => 'local', + AssistantExecutionTarget.remote => 'remote', + AssistantExecutionTarget.singleAgent => 'singleAgent', + AssistantExecutionTarget.auto => '', + }; + final explicitProviderId = provider == SingleAgentProvider.auto + ? '' + : provider.providerId; + final explicitModelValue = model.trim(); + final explicitSkillsValue = selectedSkills + .map((item) => item.trim()) + .where((item) => item.isNotEmpty) + .toList(growable: false); + final hasExplicitSelection = + explicitExecutionTarget.isNotEmpty || + explicitProviderId.isNotEmpty || + explicitModelValue.isNotEmpty || + explicitSkillsValue.isNotEmpty; + if (!hasExplicitSelection) { + return ExternalCodeAgentAcpRoutingConfig.auto( + preferredGatewayTarget: preferredGatewayTarget, + ); + } + return ExternalCodeAgentAcpRoutingConfig( + mode: ExternalCodeAgentAcpRoutingMode.explicit, + preferredGatewayTarget: preferredGatewayTarget, + explicitExecutionTarget: explicitExecutionTarget, + explicitProviderId: explicitProviderId, + explicitModel: explicitModelValue, + explicitSkills: explicitSkillsValue, + allowSkillInstall: false, + availableSkills: const [], + ); + } +} + +const String _gatewaySessionMode = 'gateway-chat'; + +bool _usesGatewaySessionMode(String mode) { + final normalized = mode.trim(); + return normalized == 'gateway' || normalized == _gatewaySessionMode; +} + +class GoTaskServiceUpdate { + const GoTaskServiceUpdate({ + required this.sessionId, + required this.threadId, + required this.turnId, + required this.type, + required this.text, + required this.message, + required this.pending, + required this.error, + required this.route, + required this.payload, + }); + + final String sessionId; + final String threadId; + final String turnId; + final String type; + final String text; + final String message; + final bool pending; + final bool error; + final GoTaskServiceRoute route; + final Map payload; + + bool get isDelta => type == 'delta' && text.isNotEmpty; + bool get isDone => type == 'done' || payload['event'] == 'completed'; +} + +class GoTaskServiceResult { + const GoTaskServiceResult({ + required this.success, + required this.message, + required this.turnId, + required this.raw, + required this.errorMessage, + required this.resolvedModel, + required this.route, + }); + + final bool success; + final String message; + final String turnId; + final Map raw; + final String errorMessage; + final String resolvedModel; + final GoTaskServiceRoute route; + + String get resolvedWorkingDirectory => + raw['resolvedWorkingDirectory']?.toString().trim() ?? + raw['workingDirectory']?.toString().trim() ?? + ''; + + String get resolvedExecutionTarget => + raw['resolvedExecutionTarget']?.toString().trim() ?? ''; + + String get resolvedEndpointTarget => + raw['resolvedEndpointTarget']?.toString().trim() ?? ''; + + String get resolvedProviderId => + raw['resolvedProviderId']?.toString().trim() ?? ''; + + List get resolvedSkills { + final rawList = raw['resolvedSkills']; + if (rawList is! List) { + return const []; + } + return rawList + .map((item) => item?.toString().trim() ?? '') + .where((item) => item.isNotEmpty) + .toList(growable: false); + } + + String get skillResolutionSource => + raw['skillResolutionSource']?.toString().trim() ?? ''; + + bool get needsSkillInstall => _boolValue(raw['needsSkillInstall']) ?? false; + + String get skillInstallRequestId => + raw['skillInstallRequestId']?.toString().trim() ?? ''; + + List> get skillCandidates => + _castMapList(raw['skillCandidates']); + + List> get memorySources => + _castMapList(raw['memorySources']); + + WorkspaceRefKind? get resolvedWorkspaceRefKind { + final rawValue = raw['resolvedWorkspaceRefKind']?.toString().trim() ?? ''; + if (rawValue.isEmpty) { + return null; + } + return WorkspaceRefKindCopy.fromJsonValue(rawValue); + } +} + +abstract class ExternalCodeAgentAcpTransport { + Future syncExternalProviders( + List providers, + ); + + Future loadExternalAcpCapabilities({ + required AssistantExecutionTarget target, + bool forceRefresh = false, + }); + + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }); + + Future cancelTask({ + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }); + + Future closeTask({ + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }); + + Future dispose(); +} + +abstract class GoTaskServiceClient { + Future syncExternalProviders( + List providers, + ); + + Future loadExternalAcpCapabilities({ + required AssistantExecutionTarget target, + bool forceRefresh = false, + }); + + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }); + + Future cancelTask({ + required GoTaskServiceRoute route, + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }); + + Future closeTask({ + required GoTaskServiceRoute route, + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }); + + Future dispose(); +} + +GoTaskServiceUpdate? goTaskServiceUpdateFromAcpNotification( + Map notification, +) { + final method = notification['method']?.toString().trim().toLowerCase() ?? ''; + if (method != 'session.update' && method != 'acp.session.update') { + return null; + } + final params = _castMap(notification['params']); + final payload = params.isNotEmpty + ? params + : _castMap(notification['payload']); + final type = + payload['type']?.toString().trim().toLowerCase() ?? + payload['state']?.toString().trim().toLowerCase() ?? + payload['event']?.toString().trim().toLowerCase() ?? + 'status'; + return GoTaskServiceUpdate( + sessionId: payload['sessionId']?.toString().trim().isNotEmpty == true + ? payload['sessionId'].toString().trim() + : payload['threadId']?.toString().trim() ?? '', + threadId: payload['threadId']?.toString().trim() ?? '', + turnId: payload['turnId']?.toString().trim() ?? '', + type: type, + text: + payload['delta']?.toString() ?? + payload['text']?.toString() ?? + _castMap(payload['message'])['content']?.toString() ?? + '', + message: payload['message']?.toString() ?? '', + pending: _boolValue(payload['pending']) ?? false, + error: _boolValue(payload['error']) ?? false, + route: GoTaskServiceRoute.externalAcpSingle, + payload: payload, + ); +} + +GoTaskServiceResult goTaskServiceResultFromAcpResponse( + Map response, { + required GoTaskServiceRoute route, + String streamedText = '', + String? completedMessage, +}) { + final result = _castMap(response['result']); + final primaryText = + (completedMessage?.trim().isNotEmpty == true + ? completedMessage!.trim() + : 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, + message: primaryText, + turnId: result['turnId']?.toString().trim() ?? '', + raw: result, + errorMessage: result['error']?.toString() ?? '', + resolvedModel: + result['model']?.toString().trim() ?? + result['resolvedModel']?.toString().trim() ?? + '', + route: route, + ); +} + +Map mergeGoTaskServiceResponseResult( + Map response, + Map overlay, +) { + if (overlay.isEmpty) { + return response; + } + final next = Map.from(response); + final result = Map.from(_castMap(next['result'])); + overlay.forEach((key, value) { + if (value == null) { + return; + } + if (value is String && value.trim().isEmpty) { + if (result.containsKey(key)) { + return; + } + } + result[key] = value; + }); + next['result'] = result; + return next; +} + +int goTaskServiceBase64Size(String base64) { + final normalized = base64.trim().split(',').last.trim(); + if (normalized.isEmpty) { + return 0; + } + final padding = normalized.endsWith('==') + ? 2 + : (normalized.endsWith('=') ? 1 : 0); + return (normalized.length * 3 ~/ 4) - padding; +} + +Map _castMap(Object? value) { + if (value is Map) { + return value; + } + if (value is Map) { + return value.cast(); + } + return const {}; +} + +bool? _boolValue(Object? raw) { + if (raw is bool) { + return raw; + } + if (raw is num) { + return raw != 0; + } + if (raw is String) { + final normalized = raw.trim().toLowerCase(); + if (normalized == 'true') { + return true; + } + if (normalized == 'false') { + return false; + } + } + return null; +} + +List> _castMapList(Object? raw) { + if (raw is! List) { + return const >[]; + } + return raw.map(_castMap).toList(growable: false); +} diff --git a/lib/runtime/go_task_service_desktop_service.dart b/lib/runtime/go_task_service_desktop_service.dart new file mode 100644 index 00000000..c905242d --- /dev/null +++ b/lib/runtime/go_task_service_desktop_service.dart @@ -0,0 +1,221 @@ +import 'dart:async'; + +import 'gateway_runtime.dart'; +import 'go_task_service_client.dart'; +import 'runtime_models.dart'; + +class DesktopGoTaskService implements GoTaskServiceClient { + DesktopGoTaskService({ + required GatewayRuntime gateway, + required ExternalCodeAgentAcpTransport acpTransport, + }) : _gateway = gateway, + _acpTransport = acpTransport { + _gatewayEventsSubscription = _gateway.events.listen(_handleGatewayEvent); + } + + final GatewayRuntime _gateway; + final ExternalCodeAgentAcpTransport _acpTransport; + + late final StreamSubscription _gatewayEventsSubscription; + final Map _pendingOpenClawTasksByRunId = + {}; + final Map _openClawRunIdsBySession = {}; + + @override + Future syncExternalProviders( + List providers, + ) => _acpTransport.syncExternalProviders(providers); + + @override + Future loadExternalAcpCapabilities({ + required AssistantExecutionTarget target, + bool forceRefresh = false, + }) => _acpTransport.loadExternalAcpCapabilities( + target: target, + forceRefresh: forceRefresh, + ); + + @override + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }) async { + switch (request.route) { + case GoTaskServiceRoute.openClawTask: + return _executeOpenClawTask(request, onUpdate: onUpdate); + case GoTaskServiceRoute.externalAcpSingle: + case GoTaskServiceRoute.externalAcpMulti: + return _acpTransport.executeTask(request, onUpdate: onUpdate); + } + } + + @override + Future cancelTask({ + required GoTaskServiceRoute route, + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }) async { + if (route == GoTaskServiceRoute.openClawTask) { + final runId = _openClawRunIdsBySession[sessionId]; + if (runId == null || runId.trim().isEmpty) { + return; + } + await _gateway.abortChat(sessionKey: sessionId, runId: runId); + return; + } + await _acpTransport.cancelTask( + target: target, + sessionId: sessionId, + threadId: threadId, + ); + } + + @override + Future closeTask({ + required GoTaskServiceRoute route, + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }) async { + if (route == GoTaskServiceRoute.openClawTask) { + _openClawRunIdsBySession.remove(sessionId); + return; + } + await _acpTransport.closeTask( + target: target, + sessionId: sessionId, + threadId: threadId, + ); + } + + @override + Future dispose() async { + for (final pending in _pendingOpenClawTasksByRunId.values) { + if (!pending.completer.isCompleted) { + pending.completer.completeError( + GatewayRuntimeException('task service disposed'), + ); + } + } + _pendingOpenClawTasksByRunId.clear(); + _openClawRunIdsBySession.clear(); + await _gatewayEventsSubscription.cancel(); + await _acpTransport.dispose(); + } + + Future _executeOpenClawTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }) async { + if (!_gateway.isConnected) { + throw GatewayRuntimeException('gateway not connected'); + } + final runId = await _gateway.sendChat( + sessionKey: request.sessionId, + message: request.prompt, + thinking: request.thinking, + attachments: request.inlineAttachments, + agentId: request.agentId.trim().isEmpty ? null : request.agentId.trim(), + metadata: request.metadata.isEmpty ? null : request.metadata, + ); + final pending = _PendingOpenClawTask( + request: request, + runId: runId, + onUpdate: onUpdate, + completer: Completer(), + ); + _pendingOpenClawTasksByRunId[runId] = pending; + _openClawRunIdsBySession[request.sessionId] = runId; + return pending.completer.future; + } + + void _handleGatewayEvent(GatewayPushEvent event) { + if (event.event == 'chat.run') { + _handleOpenClawRunPayload(asMap(event.payload)); + return; + } + if (event.event == 'chat') { + final payload = asMap(event.payload); + final message = asMap(payload['message']); + _handleOpenClawRunPayload({ + 'runId': payload['runId'], + 'sessionKey': payload['sessionKey'], + 'state': payload['state'], + 'assistantText': extractMessageText(message), + 'errorMessage': payload['errorMessage'], + }); + } + } + + void _handleOpenClawRunPayload(Map payload) { + final runId = stringValue(payload['runId']) ?? ''; + if (runId.isEmpty) { + return; + } + final pending = _pendingOpenClawTasksByRunId[runId]; + if (pending == null) { + return; + } + final state = stringValue(payload['state']) ?? ''; + final assistantText = stringValue(payload['assistantText']) ?? ''; + final errorMessage = stringValue(payload['errorMessage']) ?? ''; + if (assistantText.isNotEmpty && (state == 'delta' || state == 'final')) { + pending.streamedText = assistantText; + pending.onUpdate( + GoTaskServiceUpdate( + sessionId: pending.request.sessionId, + threadId: pending.request.threadId, + turnId: runId, + type: 'delta', + text: assistantText, + message: '', + pending: state != 'final', + error: false, + route: GoTaskServiceRoute.openClawTask, + payload: payload, + ), + ); + } + final terminal = + boolValue(payload['terminal']) ?? false || + state == 'final' || + state == 'aborted' || + state == 'error'; + if (!terminal || pending.completer.isCompleted) { + return; + } + _pendingOpenClawTasksByRunId.remove(runId); + _openClawRunIdsBySession.remove(pending.request.sessionId); + final success = state != 'error' && state != 'aborted'; + pending.completer.complete( + GoTaskServiceResult( + success: success, + message: pending.streamedText.trim(), + turnId: runId, + raw: payload, + errorMessage: errorMessage, + resolvedModel: + stringValue(payload['model']) ?? + stringValue(payload['resolvedModel']) ?? + '', + route: GoTaskServiceRoute.openClawTask, + ), + ); + } +} + +class _PendingOpenClawTask { + _PendingOpenClawTask({ + required this.request, + required this.runId, + required this.onUpdate, + required this.completer, + }); + + final GoTaskServiceRequest request; + final String runId; + final void Function(GoTaskServiceUpdate update) onUpdate; + final Completer completer; + String streamedText = ''; +} diff --git a/lib/web/go_agent_core_web_transport.dart b/lib/web/go_agent_core_web_transport.dart index aad6c2df..aaca1a84 100644 --- a/lib/web/go_agent_core_web_transport.dart +++ b/lib/web/go_agent_core_web_transport.dart @@ -1,9 +1,9 @@ -import '../runtime/go_agent_core_client.dart'; +import '../runtime/go_task_service_client.dart'; import '../runtime/runtime_models.dart'; import 'web_acp_client.dart'; -class GoAgentCoreWebTransport implements GoAgentCoreClient { - const GoAgentCoreWebTransport({ +class ExternalCodeAgentAcpWebTransport implements ExternalCodeAgentAcpTransport { + const ExternalCodeAgentAcpWebTransport({ required WebAcpClient acpClient, required Uri? Function(AssistantExecutionTarget target) endpointResolver, }) : _acpClient = acpClient, @@ -15,7 +15,9 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient { Uri? get _goCoreEndpoint => _endpointResolver(AssistantExecutionTarget.singleAgent); @override - Future syncProviders(List providers) async { + Future syncExternalProviders( + List providers, + ) async { final endpoint = _goCoreEndpoint; if (endpoint == null) { return; @@ -30,16 +32,16 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient { } @override - Future loadCapabilities({ + Future loadExternalAcpCapabilities({ required AssistantExecutionTarget target, bool forceRefresh = false, }) async { final endpoint = _goCoreEndpoint; if (endpoint == null) { - return const GoAgentCoreCapabilities.empty(); + return const ExternalCodeAgentAcpCapabilities.empty(); } final capabilities = await _acpClient.loadCapabilities(endpoint: endpoint); - return GoAgentCoreCapabilities( + return ExternalCodeAgentAcpCapabilities( singleAgent: capabilities.singleAgent, multiAgent: capabilities.multiAgent, providers: capabilities.providers, @@ -48,9 +50,9 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient { } @override - Future executeSession( - GoAgentCoreSessionRequest request, { - required void Function(GoAgentCoreSessionUpdate update) onUpdate, + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, }) async { final endpoint = _goCoreEndpoint; if (endpoint == null) { @@ -64,9 +66,9 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient { final response = await _acpClient.request( endpoint: endpoint, method: request.resumeSession ? 'session.message' : 'session.start', - params: request.toAcpParams(), + params: request.toExternalAcpParams(), onNotification: (notification) { - final update = goAgentCoreUpdateFromNotification(notification); + final update = goTaskServiceUpdateFromAcpNotification(notification); if (update == null) { return; } @@ -79,15 +81,16 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient { onUpdate(update); }, ); - return goAgentCoreRunResultFromResponse( + return goTaskServiceResultFromAcpResponse( response, + route: request.route, streamedText: streamedText, completedMessage: completedMessage, ); } @override - Future cancelSession({ + Future cancelTask({ required AssistantExecutionTarget target, required String sessionId, required String threadId, @@ -104,7 +107,7 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient { } @override - Future closeSession({ + Future closeTask({ required AssistantExecutionTarget target, required String sessionId, required String threadId, diff --git a/lib/web/go_task_service_web_service.dart b/lib/web/go_task_service_web_service.dart new file mode 100644 index 00000000..c1122ee3 --- /dev/null +++ b/lib/web/go_task_service_web_service.dart @@ -0,0 +1,231 @@ +import 'dart:async'; + +import '../runtime/gateway_runtime_errors.dart'; +import '../runtime/gateway_runtime_helpers.dart'; +import '../runtime/go_task_service_client.dart'; +import '../runtime/runtime_models.dart'; +import 'web_relay_gateway_client.dart'; + +class WebGoTaskService implements GoTaskServiceClient { + WebGoTaskService({ + required WebRelayGatewayClient relayClient, + required ExternalCodeAgentAcpTransport acpTransport, + }) : _relayClient = relayClient, + _acpTransport = acpTransport { + _relayEventsSubscription = _relayClient.events.listen(_handleRelayEvent); + } + + final WebRelayGatewayClient _relayClient; + final ExternalCodeAgentAcpTransport _acpTransport; + + late final StreamSubscription _relayEventsSubscription; + final Map _pendingOpenClawTasksByRunId = + {}; + final Map _openClawRunIdsBySession = {}; + + @override + Future syncExternalProviders( + List providers, + ) => _acpTransport.syncExternalProviders(providers); + + @override + Future loadExternalAcpCapabilities({ + required AssistantExecutionTarget target, + bool forceRefresh = false, + }) => _acpTransport.loadExternalAcpCapabilities( + target: target, + forceRefresh: forceRefresh, + ); + + @override + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }) async { + switch (request.route) { + case GoTaskServiceRoute.openClawTask: + return _executeOpenClawTask(request, onUpdate: onUpdate); + case GoTaskServiceRoute.externalAcpSingle: + case GoTaskServiceRoute.externalAcpMulti: + return _acpTransport.executeTask(request, onUpdate: onUpdate); + } + } + + @override + Future cancelTask({ + required GoTaskServiceRoute route, + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }) async { + if (route == GoTaskServiceRoute.openClawTask) { + final runId = _openClawRunIdsBySession[sessionId]; + if (runId == null || runId.trim().isEmpty) { + return; + } + await _relayClient.request( + 'chat.abort', + params: {'sessionKey': sessionId, 'runId': runId}, + ); + return; + } + await _acpTransport.cancelTask( + target: target, + sessionId: sessionId, + threadId: threadId, + ); + } + + @override + Future closeTask({ + required GoTaskServiceRoute route, + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }) async { + if (route == GoTaskServiceRoute.openClawTask) { + _openClawRunIdsBySession.remove(sessionId); + return; + } + await _acpTransport.closeTask( + target: target, + sessionId: sessionId, + threadId: threadId, + ); + } + + @override + Future dispose() async { + for (final pending in _pendingOpenClawTasksByRunId.values) { + if (!pending.completer.isCompleted) { + pending.completer.completeError( + GatewayRuntimeException('task service disposed'), + ); + } + } + _pendingOpenClawTasksByRunId.clear(); + _openClawRunIdsBySession.clear(); + await _relayEventsSubscription.cancel(); + await _acpTransport.dispose(); + } + + Future _executeOpenClawTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }) async { + if (!_relayClient.isConnected) { + throw GatewayRuntimeException('gateway not connected'); + } + final metadata = { + ...request.metadata, + if (request.agentId.trim().isNotEmpty && + !request.metadata.containsKey('agentId')) + 'agentId': request.agentId.trim(), + }; + final runId = await _relayClient.sendChat( + sessionKey: request.sessionId, + message: request.prompt, + thinking: request.thinking, + attachments: request.inlineAttachments, + metadata: metadata, + ); + final pending = _PendingOpenClawTask( + request: request, + runId: runId, + onUpdate: onUpdate, + completer: Completer(), + ); + _pendingOpenClawTasksByRunId[runId] = pending; + _openClawRunIdsBySession[request.sessionId] = runId; + return pending.completer.future; + } + + void _handleRelayEvent(GatewayPushEvent event) { + if (event.event == 'chat.run') { + _handleOpenClawRunPayload(asMap(event.payload)); + return; + } + if (event.event == 'chat') { + final payload = asMap(event.payload); + final message = asMap(payload['message']); + _handleOpenClawRunPayload({ + 'runId': payload['runId'], + 'sessionKey': payload['sessionKey'], + 'state': payload['state'], + 'assistantText': extractMessageText(message), + 'errorMessage': payload['errorMessage'], + }); + } + } + + void _handleOpenClawRunPayload(Map payload) { + final runId = stringValue(payload['runId']) ?? ''; + if (runId.isEmpty) { + return; + } + final pending = _pendingOpenClawTasksByRunId[runId]; + if (pending == null) { + return; + } + final state = stringValue(payload['state']) ?? ''; + final assistantText = stringValue(payload['assistantText']) ?? ''; + final errorMessage = stringValue(payload['errorMessage']) ?? ''; + if (assistantText.isNotEmpty && (state == 'delta' || state == 'final')) { + pending.streamedText = assistantText; + pending.onUpdate( + GoTaskServiceUpdate( + sessionId: pending.request.sessionId, + threadId: pending.request.threadId, + turnId: runId, + type: 'delta', + text: assistantText, + message: '', + pending: state != 'final', + error: false, + route: GoTaskServiceRoute.openClawTask, + payload: payload, + ), + ); + } + final terminal = + boolValue(payload['terminal']) ?? false || + state == 'final' || + state == 'aborted' || + state == 'error'; + if (!terminal || pending.completer.isCompleted) { + return; + } + _pendingOpenClawTasksByRunId.remove(runId); + _openClawRunIdsBySession.remove(pending.request.sessionId); + final success = state != 'error' && state != 'aborted'; + pending.completer.complete( + GoTaskServiceResult( + success: success, + message: pending.streamedText.trim(), + turnId: runId, + raw: payload, + errorMessage: errorMessage, + resolvedModel: + stringValue(payload['model']) ?? + stringValue(payload['resolvedModel']) ?? + '', + route: GoTaskServiceRoute.openClawTask, + ), + ); + } +} + +class _PendingOpenClawTask { + _PendingOpenClawTask({ + required this.request, + required this.runId, + required this.onUpdate, + required this.completer, + }); + + final GoTaskServiceRequest request; + final String runId; + final void Function(GoTaskServiceUpdate update) onUpdate; + final Completer completer; + String streamedText = ''; +} diff --git a/test/runtime/app_controller_ai_gateway_chat_suite_chat.dart b/test/runtime/app_controller_ai_gateway_chat_suite_chat.dart index 844486a0..28fe192a 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite_chat.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite_chat.dart @@ -42,7 +42,7 @@ void registerAppControllerAiGatewayChatSuiteChatTestsInternal() { gateway: gateway, codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: FallbackOnlyGoAgentCoreClientInternal(), + goTaskServiceClient: FallbackOnlyGoAgentCoreClientInternal(), ); await controller.settingsController.saveAiGatewayApiKey('live-key'); @@ -104,7 +104,7 @@ void registerAppControllerAiGatewayChatSuiteChatTestsInternal() { gateway: secondGateway, codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: FallbackOnlyGoAgentCoreClientInternal(), + goTaskServiceClient: FallbackOnlyGoAgentCoreClientInternal(), ); await secondController.settingsController.saveAiGatewayApiKey( @@ -182,7 +182,7 @@ void registerAppControllerAiGatewayChatSuiteChatTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: FallbackOnlyGoAgentCoreClientInternal(), + goTaskServiceClient: FallbackOnlyGoAgentCoreClientInternal(), ); await controller.settingsController.saveAiGatewayApiKey('live-key'); @@ -242,7 +242,7 @@ void registerAppControllerAiGatewayChatSuiteChatTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: FallbackOnlyGoAgentCoreClientInternal(), + goTaskServiceClient: FallbackOnlyGoAgentCoreClientInternal(), ); await controller.settingsController.saveAiGatewayApiKey('live-key'); 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 af0ebd6f..628b9805 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite_fakes.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite_fakes.dart @@ -9,7 +9,7 @@ import 'package:xworkmate/app/app_controller.dart'; import 'package:xworkmate/runtime/codex_runtime.dart'; import 'package:xworkmate/runtime/device_identity_store.dart'; import 'package:xworkmate/runtime/gateway_runtime.dart'; -import 'package:xworkmate/runtime/go_agent_core_client.dart'; +import 'package:xworkmate/runtime/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_coordinator.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; @@ -112,34 +112,36 @@ class FakeCodexRuntimeInternal extends CodexRuntime { Future stop() async {} } -class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { +class FakeGoAgentCoreClientInternal implements GoTaskServiceClient { FakeGoAgentCoreClientInternal({ - this.capabilities = const GoAgentCoreCapabilities.empty(), - this.result = const GoAgentCoreRunResult( + this.capabilities = const ExternalCodeAgentAcpCapabilities.empty(), + this.result = const GoTaskServiceResult( success: false, message: '', turnId: '', raw: {}, errorMessage: 'no result configured', resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, ), }); - final GoAgentCoreCapabilities capabilities; - final GoAgentCoreRunResult result; + final ExternalCodeAgentAcpCapabilities capabilities; + final GoTaskServiceResult result; int capabilitiesCalls = 0; int executeCalls = 0; int cancelCalls = 0; - GoAgentCoreSessionRequest? lastRequest; - final List requests = - []; + GoTaskServiceRequest? lastRequest; + final List requests = []; @override - Future syncProviders(List providers) async {} + Future syncExternalProviders( + List providers, + ) async {} @override - Future loadCapabilities({ + Future loadExternalAcpCapabilities({ required AssistantExecutionTarget target, bool forceRefresh = false, }) async { @@ -148,16 +150,16 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { } @override - Future executeSession( - GoAgentCoreSessionRequest request, { - required void Function(GoAgentCoreSessionUpdate update) onUpdate, + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, }) async { executeCalls += 1; lastRequest = request; requests.add(request); if (result.message.trim().isNotEmpty) { onUpdate( - GoAgentCoreSessionUpdate( + GoTaskServiceUpdate( sessionId: request.sessionId, threadId: request.threadId, turnId: result.turnId, @@ -166,6 +168,7 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { message: '', pending: false, error: false, + route: result.route, payload: const {'type': 'delta'}, ), ); @@ -174,7 +177,8 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { } @override - Future cancelSession({ + Future cancelTask({ + required GoTaskServiceRoute route, required AssistantExecutionTarget target, required String sessionId, required String threadId, @@ -183,7 +187,8 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { } @override - Future closeSession({ + Future closeTask({ + required GoTaskServiceRoute route, required AssistantExecutionTarget target, required String sessionId, required String threadId, @@ -196,7 +201,7 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient { class FallbackOnlyGoAgentCoreClientInternal extends FakeGoAgentCoreClientInternal { FallbackOnlyGoAgentCoreClientInternal() - : super(capabilities: const GoAgentCoreCapabilities.empty()); + : super(capabilities: const ExternalCodeAgentAcpCapabilities.empty()); } class FakeAiGatewayServerInternal { diff --git a/test/runtime/app_controller_ai_gateway_chat_suite_fixtures.dart b/test/runtime/app_controller_ai_gateway_chat_suite_fixtures.dart index f1a9da6b..d3734b50 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite_fixtures.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite_fixtures.dart @@ -9,7 +9,7 @@ import 'package:xworkmate/app/app_controller.dart'; import 'package:xworkmate/runtime/codex_runtime.dart'; import 'package:xworkmate/runtime/device_identity_store.dart'; import 'package:xworkmate/runtime/gateway_runtime.dart'; -import 'package:xworkmate/runtime/go_agent_core_client.dart'; +import 'package:xworkmate/runtime/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_coordinator.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; @@ -23,14 +23,14 @@ Future createAppControllerInternal({ List availableSingleAgentProvidersOverride = const [], RuntimeCoordinator? runtimeCoordinator, - GoAgentCoreClient? goAgentCoreClient, + GoTaskServiceClient? goTaskServiceClient, }) async { final controller = AppController( store: store, availableSingleAgentProvidersOverride: availableSingleAgentProvidersOverride, runtimeCoordinator: runtimeCoordinator, - goAgentCoreClient: goAgentCoreClient, + goTaskServiceClient: goTaskServiceClient, ); addTearDown(controller.dispose); await waitForInternal(() => !controller.initializing); diff --git a/test/runtime/app_controller_ai_gateway_chat_suite_single_agent.dart b/test/runtime/app_controller_ai_gateway_chat_suite_single_agent.dart index a141acbe..e073eb1c 100644 --- a/test/runtime/app_controller_ai_gateway_chat_suite_single_agent.dart +++ b/test/runtime/app_controller_ai_gateway_chat_suite_single_agent.dart @@ -9,7 +9,7 @@ import 'package:xworkmate/app/app_controller.dart'; import 'package:xworkmate/runtime/codex_runtime.dart'; import 'package:xworkmate/runtime/device_identity_store.dart'; import 'package:xworkmate/runtime/gateway_runtime.dart'; -import 'package:xworkmate/runtime/go_agent_core_client.dart'; +import 'package:xworkmate/runtime/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_coordinator.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; @@ -28,19 +28,20 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ); final store = createStoreFromTempDirectoryInternal(tempDirectory); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'CODEX_REPLY', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: 'codex-sonnet', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -52,7 +53,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.saveSettings( controller.settings.copyWith( @@ -107,7 +108,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ); final store = createStoreFromTempDirectoryInternal(tempDirectory); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, @@ -123,7 +124,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.saveSettings( controller.settings.copyWith( @@ -159,19 +160,20 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ); final store = createStoreFromTempDirectoryInternal(tempDirectory); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'CODEX_REPLY', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: 'codex-sonnet', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -183,7 +185,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.saveSettings( @@ -221,19 +223,20 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { SettingsSnapshot.defaults().copyWith(workspacePath: tempDirectory.path), ); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'WORKSPACE_OK', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: 'codex-sonnet', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -245,7 +248,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.saveSettings( controller.settings.copyWith( @@ -295,19 +298,20 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ), ); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'WORKSPACE_PLACEHOLDER_OK', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: 'codex-sonnet', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -319,7 +323,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.setAssistantExecutionTarget( @@ -373,7 +377,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.settingsController.saveAiGatewayApiKey('live-key'); @@ -437,7 +441,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.settingsController.saveAiGatewayApiKey('live-key'); @@ -506,7 +510,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.settingsController.saveAiGatewayApiKey('live-key'); @@ -590,19 +594,20 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ]); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'THREAD_OK', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -614,7 +619,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.sendChatMessage('检查当前线程目录', thinking: 'low'); @@ -649,19 +654,20 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'THREAD_OK', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -673,7 +679,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); controller.initializeAssistantThreadContext( @@ -725,13 +731,13 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'THREAD_OK', turnId: 'turn-1', @@ -742,6 +748,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { }, errorMessage: '', resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -753,7 +760,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); controller.initializeAssistantThreadContext( @@ -816,13 +823,13 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { ); final client = FakeGoAgentCoreClientInternal( - capabilities: GoAgentCoreCapabilities( + capabilities: ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: false, providers: {SingleAgentProvider.opencode}, raw: {}, ), - result: const GoAgentCoreRunResult( + result: const GoTaskServiceResult( success: true, message: 'THREAD_OK', turnId: 'turn-1', @@ -832,6 +839,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { }, errorMessage: '', resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, ), ); final controller = await createAppControllerInternal( @@ -843,7 +851,7 @@ void registerAppControllerAiGatewayChatSuiteSingleAgentTestsInternal() { gateway: FakeGatewayRuntimeInternal(store: store), codex: FakeCodexRuntimeInternal(), ), - goAgentCoreClient: client, + goTaskServiceClient: client, ); await controller.setAssistantExecutionTarget( diff --git a/test/runtime/app_controller_assistant_flow_suite.dart b/test/runtime/app_controller_assistant_flow_suite.dart index 1055bbb1..c32d3408 100644 --- a/test/runtime/app_controller_assistant_flow_suite.dart +++ b/test/runtime/app_controller_assistant_flow_suite.dart @@ -8,7 +8,7 @@ import 'dart:io'; import 'package:flutter_test/flutter_test.dart'; import 'package:shared_preferences/shared_preferences.dart'; import 'package:xworkmate/app/app_controller.dart'; -import 'package:xworkmate/runtime/go_agent_core_client.dart'; +import 'package:xworkmate/runtime/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; import 'package:xworkmate/runtime/secure_config_store.dart'; @@ -34,7 +34,7 @@ void main() { ); final controller = AppController( store: store, - goAgentCoreClient: goCoreClient, + goTaskServiceClient: goCoreClient, ); addTearDown(() async { controller.dispose(); @@ -89,7 +89,7 @@ void main() { ); expect( goCoreClient.lastRequest?.routing?.mode, - GoAgentCoreRoutingMode.auto, + ExternalCodeAgentAcpRoutingMode.auto, ); expect( goCoreClient.lastRequest?.routing?.preferredGatewayTarget, @@ -116,7 +116,7 @@ void main() { final goCoreClient = _FakeGoAgentCoreClient(); final controller = AppController( store: store, - goAgentCoreClient: goCoreClient, + goTaskServiceClient: goCoreClient, ); addTearDown(controller.dispose); @@ -139,7 +139,7 @@ void main() { expect( goCoreClient.lastRequest?.routing?.mode, - GoAgentCoreRoutingMode.explicit, + ExternalCodeAgentAcpRoutingMode.explicit, ); expect( goCoreClient.lastRequest?.routing?.explicitExecutionTarget, @@ -167,7 +167,7 @@ void main() { ); final controller = AppController( store: store, - goAgentCoreClient: _FakeGoAgentCoreClient( + goTaskServiceClient: _FakeGoAgentCoreClient( onExecute: gateway.recordGoCoreTurn, ), ); @@ -217,7 +217,7 @@ void main() { ); final controller = AppController( store: store, - goAgentCoreClient: _FakeGoAgentCoreClient( + goTaskServiceClient: _FakeGoAgentCoreClient( onExecute: gateway.recordGoCoreTurn, ), ); @@ -615,8 +615,8 @@ class _FakeGatewayServer { socket.add(jsonEncode(frame)); } - void recordGoCoreTurn(GoAgentCoreSessionRequest request) { - lastChatSendParams = request.toAcpParams(); + void recordGoCoreTurn(GoTaskServiceRequest request) { + lastChatSendParams = request.toExternalAcpParams(); final prompt = request.prompt.trim(); if (prompt.isNotEmpty) { _appendMessage(role: 'user', text: prompt); @@ -625,21 +625,23 @@ class _FakeGatewayServer { } } -class _FakeGoAgentCoreClient implements GoAgentCoreClient { +class _FakeGoAgentCoreClient implements GoTaskServiceClient { _FakeGoAgentCoreClient({this.onExecute}); - GoAgentCoreSessionRequest? lastRequest; - final void Function(GoAgentCoreSessionRequest request)? onExecute; + GoTaskServiceRequest? lastRequest; + final void Function(GoTaskServiceRequest request)? onExecute; @override - Future syncProviders(List providers) async {} + Future syncExternalProviders( + List providers, + ) async {} @override - Future loadCapabilities({ + Future loadExternalAcpCapabilities({ required AssistantExecutionTarget target, bool forceRefresh = false, }) async { - return GoAgentCoreCapabilities( + return ExternalCodeAgentAcpCapabilities( singleAgent: true, multiAgent: true, providers: { @@ -651,14 +653,14 @@ class _FakeGoAgentCoreClient implements GoAgentCoreClient { } @override - Future executeSession( - GoAgentCoreSessionRequest request, { - required void Function(GoAgentCoreSessionUpdate update) onUpdate, + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, }) async { lastRequest = request; onExecute?.call(request); onUpdate( - GoAgentCoreSessionUpdate( + GoTaskServiceUpdate( sessionId: request.sessionId, threadId: request.threadId, turnId: 'turn-1', @@ -667,28 +669,32 @@ class _FakeGoAgentCoreClient implements GoAgentCoreClient { message: '', pending: false, error: false, + route: request.route, payload: const {'type': 'delta'}, ), ); - return const GoAgentCoreRunResult( + return GoTaskServiceResult( success: true, message: 'XWORKMATE_OK', turnId: 'turn-1', raw: {}, errorMessage: '', resolvedModel: '', + route: request.route, ); } @override - Future cancelSession({ + Future cancelTask({ + required GoTaskServiceRoute route, required AssistantExecutionTarget target, required String sessionId, required String threadId, }) async {} @override - Future closeSession({ + Future closeTask({ + required GoTaskServiceRoute route, required AssistantExecutionTarget target, required String sessionId, required String threadId, diff --git a/test/runtime/go_agent_core_desktop_transport_suite.dart b/test/runtime/go_agent_core_desktop_transport_suite.dart index 43322beb..9a5d565d 100644 --- a/test/runtime/go_agent_core_desktop_transport_suite.dart +++ b/test/runtime/go_agent_core_desktop_transport_suite.dart @@ -7,17 +7,17 @@ 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/go_task_service_client.dart'; import 'package:xworkmate/runtime/runtime_models.dart'; void main() { - group('GoAgentCoreDesktopTransport', () { + group('ExternalCodeAgentAcpDesktopTransport', () { test('uses resolved gateway endpoint for local gateway sessions', () async { final server = await _AcpFakeServer.start(); addTearDown(server.close); - final transport = GoAgentCoreDesktopTransport( + final transport = ExternalCodeAgentAcpDesktopTransport( acpClient: GatewayAcpClient(endpointResolver: () => null), endpointResolver: (target) => switch (target) { AssistantExecutionTarget.local => server.baseHttpUri, @@ -25,8 +25,8 @@ void main() { }, ); - final result = await transport.executeSession( - const GoAgentCoreSessionRequest( + final result = await transport.executeTask( + const GoTaskServiceRequest( sessionId: 'session-local', threadId: 'thread-local', target: AssistantExecutionTarget.local, @@ -53,14 +53,14 @@ void main() { }); test('reports missing endpoint when gateway target cannot resolve', () async { - final transport = GoAgentCoreDesktopTransport( + final transport = ExternalCodeAgentAcpDesktopTransport( acpClient: GatewayAcpClient(endpointResolver: () => null), endpointResolver: (_) => null, ); await expectLater( - () => transport.executeSession( - const GoAgentCoreSessionRequest( + () => transport.executeTask( + const GoTaskServiceRequest( sessionId: 'session-local', threadId: 'thread-local', target: AssistantExecutionTarget.local, diff --git a/test/runtime/go_task_service_client_test.dart b/test/runtime/go_task_service_client_test.dart new file mode 100644 index 00000000..51a05617 --- /dev/null +++ b/test/runtime/go_task_service_client_test.dart @@ -0,0 +1,62 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:xworkmate/runtime/go_task_service_client.dart'; +import 'package:xworkmate/runtime/runtime_models.dart'; + +void main() { + group('GoTaskServiceRequest routing', () { + GoTaskServiceRequest buildRequest({ + required AssistantExecutionTarget target, + bool multiAgent = false, + }) { + return GoTaskServiceRequest( + sessionId: 'thread-1', + threadId: 'thread-1', + target: target, + prompt: 'hello', + workingDirectory: '/tmp/workspace', + model: '', + thinking: 'medium', + selectedSkills: const [], + inlineAttachments: const [], + localAttachments: const [], + aiGatewayBaseUrl: '', + aiGatewayApiKey: '', + agentId: '', + metadata: const {}, + multiAgent: multiAgent, + ); + } + + test('routes local and remote targets to the OpenClaw lane', () { + expect( + buildRequest(target: AssistantExecutionTarget.local).route, + GoTaskServiceRoute.openClawTask, + ); + expect( + buildRequest(target: AssistantExecutionTarget.remote).route, + GoTaskServiceRoute.openClawTask, + ); + }); + + test('routes single-agent and auto targets to the ACP single lane', () { + expect( + buildRequest(target: AssistantExecutionTarget.singleAgent).route, + GoTaskServiceRoute.externalAcpSingle, + ); + expect( + buildRequest(target: AssistantExecutionTarget.auto).route, + GoTaskServiceRoute.externalAcpSingle, + ); + }); + + test('routes multi-agent requests to the ACP multi lane', () { + expect( + buildRequest( + target: AssistantExecutionTarget.remote, + multiAgent: true, + ).route, + GoTaskServiceRoute.externalAcpMulti, + ); + }); + }); +} diff --git a/test/runtime/go_task_service_desktop_service_test.dart b/test/runtime/go_task_service_desktop_service_test.dart new file mode 100644 index 00000000..a478e261 --- /dev/null +++ b/test/runtime/go_task_service_desktop_service_test.dart @@ -0,0 +1,218 @@ +import 'dart:async'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:xworkmate/runtime/device_identity_store.dart'; +import 'package:xworkmate/runtime/gateway_runtime.dart'; +import 'package:xworkmate/runtime/go_task_service_client.dart'; +import 'package:xworkmate/runtime/go_task_service_desktop_service.dart'; +import 'package:xworkmate/runtime/runtime_models.dart'; +import 'package:xworkmate/runtime/secure_config_store.dart'; + +class _FakeGatewayRuntime extends GatewayRuntime { + _FakeGatewayRuntime() + : super( + store: SecureConfigStore(), + identityStore: DeviceIdentityStoreForTest(), + ); + + final StreamController controller = + StreamController.broadcast(); + final List> sendChatCalls = >[]; + final List> abortChatCalls = >[]; + + @override + Stream get events => controller.stream; + + @override + bool get isConnected => true; + + @override + Future sendChat({ + required String sessionKey, + required String message, + required String thinking, + List attachments = + const [], + String? agentId, + Map? metadata, + }) async { + sendChatCalls.add({ + 'sessionKey': sessionKey, + 'message': message, + 'thinking': thinking, + 'agentId': agentId, + 'metadata': metadata, + }); + return 'run-1'; + } + + @override + Future abortChat({required String sessionKey, required String runId}) async { + abortChatCalls.add({ + 'sessionKey': sessionKey, + 'runId': runId, + }); + } +} + +class DeviceIdentityStoreForTest extends DeviceIdentityStore { + DeviceIdentityStoreForTest() : super(SecureConfigStore()); +} + +class _FakeExternalAcpTransport implements ExternalCodeAgentAcpTransport { + int executeCalls = 0; + int cancelCalls = 0; + GoTaskServiceRequest? lastRequest; + + @override + Future syncExternalProviders( + List providers, + ) async {} + + @override + Future loadExternalAcpCapabilities({ + required AssistantExecutionTarget target, + bool forceRefresh = false, + }) async { + return const ExternalCodeAgentAcpCapabilities.empty(); + } + + @override + Future executeTask( + GoTaskServiceRequest request, { + required void Function(GoTaskServiceUpdate update) onUpdate, + }) async { + executeCalls += 1; + lastRequest = request; + onUpdate( + GoTaskServiceUpdate( + sessionId: request.sessionId, + threadId: request.threadId, + turnId: 'turn-1', + type: 'delta', + text: 'ACP_OK', + message: '', + pending: false, + error: false, + route: GoTaskServiceRoute.externalAcpSingle, + payload: const {'type': 'delta'}, + ), + ); + return GoTaskServiceResult( + success: true, + message: 'ACP_OK', + turnId: 'turn-1', + raw: {}, + errorMessage: '', + resolvedModel: '', + route: request.route, + ); + } + + @override + Future cancelTask({ + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }) async { + cancelCalls += 1; + } + + @override + Future closeTask({ + required AssistantExecutionTarget target, + required String sessionId, + required String threadId, + }) async {} + + @override + Future dispose() async {} +} + +GoTaskServiceRequest _request({ + required AssistantExecutionTarget target, + bool multiAgent = false, +}) { + return GoTaskServiceRequest( + sessionId: 'thread-1', + threadId: 'thread-1', + target: target, + prompt: 'hello', + workingDirectory: '/tmp/workspace', + model: '', + thinking: 'medium', + selectedSkills: const [], + inlineAttachments: const [], + localAttachments: const [], + aiGatewayBaseUrl: '', + aiGatewayApiKey: '', + agentId: 'agent-1', + metadata: const {'threadMode': 'test'}, + multiAgent: multiAgent, + ); +} + +void main() { + group('DesktopGoTaskService', () { + test('routes OpenClaw tasks through GatewayRuntime', () async { + final gateway = _FakeGatewayRuntime(); + final acp = _FakeExternalAcpTransport(); + final service = DesktopGoTaskService(gateway: gateway, acpTransport: acp); + + final updates = []; + final future = service.executeTask( + _request(target: AssistantExecutionTarget.local), + onUpdate: updates.add, + ); + + await Future.delayed(Duration.zero); + gateway.controller.add( + GatewayPushEvent( + event: 'chat', + payload: { + 'runId': 'run-1', + 'sessionKey': 'thread-1', + 'state': 'final', + 'message': { + 'role': 'assistant', + 'content': >[ + {'type': 'text', 'text': 'OPENCLAW_OK'}, + ], + }, + }, + ), + ); + + final result = await future; + + expect(acp.executeCalls, 0); + expect(gateway.sendChatCalls, hasLength(1)); + expect(result.route, GoTaskServiceRoute.openClawTask); + expect(result.message, 'OPENCLAW_OK'); + expect(updates.last.route, GoTaskServiceRoute.openClawTask); + }); + + test('routes single-agent and multi-agent tasks through ACP transport', () async { + final gateway = _FakeGatewayRuntime(); + final acp = _FakeExternalAcpTransport(); + final service = DesktopGoTaskService(gateway: gateway, acpTransport: acp); + + final singleResult = await service.executeTask( + _request(target: AssistantExecutionTarget.singleAgent), + onUpdate: (_) {}, + ); + final multiResult = await service.executeTask( + _request( + target: AssistantExecutionTarget.remote, + multiAgent: true, + ), + onUpdate: (_) {}, + ); + + expect(gateway.sendChatCalls, isEmpty); + expect(acp.executeCalls, 2); + expect(singleResult.route, GoTaskServiceRoute.externalAcpSingle); + expect(multiResult.route, GoTaskServiceRoute.externalAcpMulti); + }); + }); +}