refactor: route task threads through go task service

This commit is contained in:
Haitao Pan 2026-04-06 11:58:54 +08:00
parent dea09ea464
commit a486eb58d7
21 changed files with 1633 additions and 206 deletions

View File

@ -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<String>? singleAgentSharedSkillScanRootOverrides,
List<SingleAgentProvider>? 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<SingleAgentProvider, DirectSingleAgentCapabilities>
@ -319,8 +323,8 @@ class AppController extends ChangeNotifier {
final Map<String, Map<String, dynamic>>
latestRoutingResolutionBySessionInternal =
<String, Map<String, dynamic>>{};
final Map<String, GoAgentCoreSyncedProvider> syncedGoAgentProvidersInternal =
<String, GoAgentCoreSyncedProvider>{};
final Map<String, ExternalCodeAgentAcpSyncedProvider>
syncedGoAgentProvidersInternal = <String, ExternalCodeAgentAcpSyncedProvider>{};
final DesktopThreadArtifactService threadArtifactServiceInternal =
DesktopThreadArtifactService();
List<AssistantThreadSkillEntry> singleAgentSharedImportedSkillsInternal =

View File

@ -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<List<GoAgentCoreSyncedProvider>>
buildGoAgentCoreSyncedProvidersInternal() async {
final providers = <GoAgentCoreSyncedProvider>[];
Future<List<ExternalCodeAgentAcpSyncedProvider>>
buildExternalAcpSyncedProvidersInternal() async {
final providers = <ExternalCodeAgentAcpSyncedProvider>[];
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<void> syncGoAgentCoreProvidersInternal() async {
final providers = await buildGoAgentCoreSyncedProvidersInternal();
Future<void> 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,

View File

@ -84,9 +84,9 @@ Future<void> 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,
);

View File

@ -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,

View File

@ -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,

View File

@ -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<String> pendingSessionKeysInternal = <String>{};
final Set<String> goTaskServiceManagedRelaySessionsInternal = <String>{};
final Map<String, String> streamingTextBySessionInternal = <String, String>{};
final Map<String, Future<void>> threadTurnQueuesInternal =
<String, Future<void>>{};
@ -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();
}

View File

@ -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 <String, dynamic>{},
routing: buildWebGoAgentCoreRoutingForSessionInternal(
routing: buildWebExternalAcpRoutingForSessionInternal(
sessionKey,
explicitExecutionTarget: 'multiAgent',
),
@ -230,7 +230,7 @@ extension AppControllerWebGatewayChat on AppController {
notifyChangedInternal();
}
Future<void> executeGoAgentCoreRunInternal({
Future<void> 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: <String, dynamic>{
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,

View File

@ -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);

View File

@ -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<Process> Function(
String executable,
List<String> 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<Uri?>? _localEndpointFuture;
@override
Future<void> syncProviders(List<GoAgentCoreSyncedProvider> providers) async {
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) async {
final endpoint = await _ensureLocalEndpoint();
if (endpoint == null) {
return;
@ -60,19 +62,19 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient {
}
@override
Future<GoAgentCoreCapabilities> loadCapabilities({
Future<ExternalCodeAgentAcpCapabilities> 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<GoAgentCoreRunResult> executeSession(
GoAgentCoreSessionRequest request, {
required void Function(GoAgentCoreSessionUpdate update) onUpdate,
Future<GoTaskServiceResult> 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<void> cancelSession({
Future<void> cancelTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
@ -137,7 +140,7 @@ class GoAgentCoreDesktopTransport implements GoAgentCoreClient {
}
@override
Future<void> closeSession({
Future<void> closeTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,

View File

@ -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 <SingleAgentProvider>{},
raw = const <String, dynamic>{};
final bool singleAgent;
final bool multiAgent;
final Set<SingleAgentProvider> providers;
final Map<String, dynamic> 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<String, dynamic> toJson() {
return <String, dynamic>{
'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<String, dynamic> toJson() {
return <String, dynamic>{
'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 <ExternalCodeAgentAcpAvailableSkill>[],
}) : mode = ExternalCodeAgentAcpRoutingMode.auto,
explicitExecutionTarget = '',
explicitProviderId = '',
explicitModel = '',
explicitSkills = const <String>[],
allowSkillInstall = false,
installApproval = null;
final ExternalCodeAgentAcpRoutingMode mode;
final String preferredGatewayTarget;
final String explicitExecutionTarget;
final String explicitProviderId;
final String explicitModel;
final List<String> explicitSkills;
final bool allowSkillInstall;
final List<ExternalCodeAgentAcpAvailableSkill> availableSkills;
final ExternalCodeAgentAcpSkillInstallApproval? installApproval;
bool get isAuto => mode == ExternalCodeAgentAcpRoutingMode.auto;
Map<String, dynamic> toJson() {
return <String, dynamic>{
'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<String> approvedSkillKeys;
Map<String, dynamic> toJson() {
return <String, dynamic>{
'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<String> selectedSkills;
final List<GatewayChatAttachmentPayload> inlineAttachments;
final List<CollaborationAttachment> localAttachments;
final String aiGatewayBaseUrl;
final String aiGatewayApiKey;
final String agentId;
final Map<String, dynamic> 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<String, dynamic> toExternalAcpParams() {
final resolvedRouting = effectiveRouting;
final params = <String, dynamic>{
'sessionId': sessionId,
'threadId': threadId,
'mode': acpMode,
'taskPrompt': prompt,
'workingDirectory': workingDirectory.trim(),
'selectedSkills': selectedSkills,
'attachments': <Map<String, dynamic>>[
...localAttachments.map(
(item) => <String, dynamic>{
'name': item.name,
'description': item.description,
'path': item.path,
},
),
...inlineAttachments.map(
(item) => <String, dynamic>{
'name': item.fileName,
'description': item.mimeType,
'path': '',
},
),
],
if (inlineAttachments.isNotEmpty)
'inlineAttachments': inlineAttachments
.map(
(item) => <String, dynamic>{
'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)) ...<String, dynamic>{
'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 <ExternalCodeAgentAcpAvailableSkill>[],
);
}
}
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<String, dynamic> 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<String, dynamic> 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<String> get resolvedSkills {
final rawList = raw['resolvedSkills'];
if (rawList is! List) {
return const <String>[];
}
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<Map<String, dynamic>> get skillCandidates =>
_castMapList(raw['skillCandidates']);
List<Map<String, dynamic>> 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<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
);
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
});
Future<GoTaskServiceResult> executeTask(
GoTaskServiceRequest request, {
required void Function(GoTaskServiceUpdate update) onUpdate,
});
Future<void> cancelTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
});
Future<void> closeTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
});
Future<void> dispose();
}
abstract class GoTaskServiceClient {
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
);
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
});
Future<GoTaskServiceResult> executeTask(
GoTaskServiceRequest request, {
required void Function(GoTaskServiceUpdate update) onUpdate,
});
Future<void> cancelTask({
required GoTaskServiceRoute route,
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
});
Future<void> closeTask({
required GoTaskServiceRoute route,
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
});
Future<void> dispose();
}
GoTaskServiceUpdate? goTaskServiceUpdateFromAcpNotification(
Map<String, dynamic> 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<String, dynamic> 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<String, dynamic> mergeGoTaskServiceResponseResult(
Map<String, dynamic> response,
Map<String, dynamic> overlay,
) {
if (overlay.isEmpty) {
return response;
}
final next = Map<String, dynamic>.from(response);
final result = Map<String, dynamic>.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<String, dynamic> _castMap(Object? value) {
if (value is Map<String, dynamic>) {
return value;
}
if (value is Map) {
return value.cast<String, dynamic>();
}
return const <String, dynamic>{};
}
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<Map<String, dynamic>> _castMapList(Object? raw) {
if (raw is! List) {
return const <Map<String, dynamic>>[];
}
return raw.map(_castMap).toList(growable: false);
}

View File

@ -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<GatewayPushEvent> _gatewayEventsSubscription;
final Map<String, _PendingOpenClawTask> _pendingOpenClawTasksByRunId =
<String, _PendingOpenClawTask>{};
final Map<String, String> _openClawRunIdsBySession = <String, String>{};
@override
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) => _acpTransport.syncExternalProviders(providers);
@override
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
}) => _acpTransport.loadExternalAcpCapabilities(
target: target,
forceRefresh: forceRefresh,
);
@override
Future<GoTaskServiceResult> 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<void> 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<void> 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<void> 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<GoTaskServiceResult> _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<GoTaskServiceResult>(),
);
_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(<String, dynamic>{
'runId': payload['runId'],
'sessionKey': payload['sessionKey'],
'state': payload['state'],
'assistantText': extractMessageText(message),
'errorMessage': payload['errorMessage'],
});
}
}
void _handleOpenClawRunPayload(Map<String, dynamic> 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<GoTaskServiceResult> completer;
String streamedText = '';
}

View File

@ -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<void> syncProviders(List<GoAgentCoreSyncedProvider> providers) async {
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) async {
final endpoint = _goCoreEndpoint;
if (endpoint == null) {
return;
@ -30,16 +32,16 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient {
}
@override
Future<GoAgentCoreCapabilities> loadCapabilities({
Future<ExternalCodeAgentAcpCapabilities> 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<GoAgentCoreRunResult> executeSession(
GoAgentCoreSessionRequest request, {
required void Function(GoAgentCoreSessionUpdate update) onUpdate,
Future<GoTaskServiceResult> 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<void> cancelSession({
Future<void> cancelTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
@ -104,7 +107,7 @@ class GoAgentCoreWebTransport implements GoAgentCoreClient {
}
@override
Future<void> closeSession({
Future<void> closeTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,

View File

@ -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<GatewayPushEvent> _relayEventsSubscription;
final Map<String, _PendingOpenClawTask> _pendingOpenClawTasksByRunId =
<String, _PendingOpenClawTask>{};
final Map<String, String> _openClawRunIdsBySession = <String, String>{};
@override
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) => _acpTransport.syncExternalProviders(providers);
@override
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
}) => _acpTransport.loadExternalAcpCapabilities(
target: target,
forceRefresh: forceRefresh,
);
@override
Future<GoTaskServiceResult> 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<void> 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: <String, dynamic>{'sessionKey': sessionId, 'runId': runId},
);
return;
}
await _acpTransport.cancelTask(
target: target,
sessionId: sessionId,
threadId: threadId,
);
}
@override
Future<void> 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<void> 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<GoTaskServiceResult> _executeOpenClawTask(
GoTaskServiceRequest request, {
required void Function(GoTaskServiceUpdate update) onUpdate,
}) async {
if (!_relayClient.isConnected) {
throw GatewayRuntimeException('gateway not connected');
}
final metadata = <String, dynamic>{
...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<GoTaskServiceResult>(),
);
_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(<String, dynamic>{
'runId': payload['runId'],
'sessionKey': payload['sessionKey'],
'state': payload['state'],
'assistantText': extractMessageText(message),
'errorMessage': payload['errorMessage'],
});
}
}
void _handleOpenClawRunPayload(Map<String, dynamic> 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<GoTaskServiceResult> completer;
String streamedText = '';
}

View File

@ -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');

View File

@ -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<void> 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: <String, dynamic>{},
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<GoAgentCoreSessionRequest> requests =
<GoAgentCoreSessionRequest>[];
GoTaskServiceRequest? lastRequest;
final List<GoTaskServiceRequest> requests = <GoTaskServiceRequest>[];
@override
Future<void> syncProviders(List<GoAgentCoreSyncedProvider> providers) async {}
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) async {}
@override
Future<GoAgentCoreCapabilities> loadCapabilities({
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
}) async {
@ -148,16 +150,16 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient {
}
@override
Future<GoAgentCoreRunResult> executeSession(
GoAgentCoreSessionRequest request, {
required void Function(GoAgentCoreSessionUpdate update) onUpdate,
Future<GoTaskServiceResult> 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 <String, dynamic>{'type': 'delta'},
),
);
@ -174,7 +177,8 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient {
}
@override
Future<void> cancelSession({
Future<void> cancelTask({
required GoTaskServiceRoute route,
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
@ -183,7 +187,8 @@ class FakeGoAgentCoreClientInternal implements GoAgentCoreClient {
}
@override
Future<void> closeSession({
Future<void> 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 {

View File

@ -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<AppController> createAppControllerInternal({
List<SingleAgentProvider> availableSingleAgentProvidersOverride =
const <SingleAgentProvider>[],
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);

View File

@ -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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
result: const GoAgentCoreRunResult(
result: const GoTaskServiceResult(
success: true,
message: 'CODEX_REPLY',
turnId: 'turn-1',
raw: <String, dynamic>{},
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>{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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
result: const GoAgentCoreRunResult(
result: const GoTaskServiceResult(
success: true,
message: 'CODEX_REPLY',
turnId: 'turn-1',
raw: <String, dynamic>{},
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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
result: const GoAgentCoreRunResult(
result: const GoTaskServiceResult(
success: true,
message: 'WORKSPACE_OK',
turnId: 'turn-1',
raw: <String, dynamic>{},
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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
result: const GoAgentCoreRunResult(
result: const GoTaskServiceResult(
success: true,
message: 'WORKSPACE_PLACEHOLDER_OK',
turnId: 'turn-1',
raw: <String, dynamic>{},
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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
result: const GoAgentCoreRunResult(
result: const GoTaskServiceResult(
success: true,
message: 'THREAD_OK',
turnId: 'turn-1',
raw: <String, dynamic>{},
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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
result: const GoAgentCoreRunResult(
result: const GoTaskServiceResult(
success: true,
message: 'THREAD_OK',
turnId: 'turn-1',
raw: <String, dynamic>{},
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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
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>{SingleAgentProvider.opencode},
raw: <String, dynamic>{},
),
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(

View File

@ -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<void> syncProviders(List<GoAgentCoreSyncedProvider> providers) async {}
Future<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) async {}
@override
Future<GoAgentCoreCapabilities> loadCapabilities({
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
}) async {
return GoAgentCoreCapabilities(
return ExternalCodeAgentAcpCapabilities(
singleAgent: true,
multiAgent: true,
providers: <SingleAgentProvider>{
@ -651,14 +653,14 @@ class _FakeGoAgentCoreClient implements GoAgentCoreClient {
}
@override
Future<GoAgentCoreRunResult> executeSession(
GoAgentCoreSessionRequest request, {
required void Function(GoAgentCoreSessionUpdate update) onUpdate,
Future<GoTaskServiceResult> 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 <String, dynamic>{'type': 'delta'},
),
);
return const GoAgentCoreRunResult(
return GoTaskServiceResult(
success: true,
message: 'XWORKMATE_OK',
turnId: 'turn-1',
raw: <String, dynamic>{},
errorMessage: '',
resolvedModel: '',
route: request.route,
);
}
@override
Future<void> cancelSession({
Future<void> cancelTask({
required GoTaskServiceRoute route,
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
}) async {}
@override
Future<void> closeSession({
Future<void> closeTask({
required GoTaskServiceRoute route,
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,

View File

@ -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,

View File

@ -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 <String>[],
inlineAttachments: const <GatewayChatAttachmentPayload>[],
localAttachments: const <CollaborationAttachment>[],
aiGatewayBaseUrl: '',
aiGatewayApiKey: '',
agentId: '',
metadata: const <String, dynamic>{},
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,
);
});
});
}

View File

@ -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<GatewayPushEvent> controller =
StreamController<GatewayPushEvent>.broadcast();
final List<Map<String, Object?>> sendChatCalls = <Map<String, Object?>>[];
final List<Map<String, String>> abortChatCalls = <Map<String, String>>[];
@override
Stream<GatewayPushEvent> get events => controller.stream;
@override
bool get isConnected => true;
@override
Future<String> sendChat({
required String sessionKey,
required String message,
required String thinking,
List<GatewayChatAttachmentPayload> attachments =
const <GatewayChatAttachmentPayload>[],
String? agentId,
Map<String, dynamic>? metadata,
}) async {
sendChatCalls.add(<String, Object?>{
'sessionKey': sessionKey,
'message': message,
'thinking': thinking,
'agentId': agentId,
'metadata': metadata,
});
return 'run-1';
}
@override
Future<void> abortChat({required String sessionKey, required String runId}) async {
abortChatCalls.add(<String, String>{
'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<void> syncExternalProviders(
List<ExternalCodeAgentAcpSyncedProvider> providers,
) async {}
@override
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
required AssistantExecutionTarget target,
bool forceRefresh = false,
}) async {
return const ExternalCodeAgentAcpCapabilities.empty();
}
@override
Future<GoTaskServiceResult> 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 <String, dynamic>{'type': 'delta'},
),
);
return GoTaskServiceResult(
success: true,
message: 'ACP_OK',
turnId: 'turn-1',
raw: <String, dynamic>{},
errorMessage: '',
resolvedModel: '',
route: request.route,
);
}
@override
Future<void> cancelTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
}) async {
cancelCalls += 1;
}
@override
Future<void> closeTask({
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
}) async {}
@override
Future<void> 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 <String>[],
inlineAttachments: const <GatewayChatAttachmentPayload>[],
localAttachments: const <CollaborationAttachment>[],
aiGatewayBaseUrl: '',
aiGatewayApiKey: '',
agentId: 'agent-1',
metadata: const <String, dynamic>{'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 = <GoTaskServiceUpdate>[];
final future = service.executeTask(
_request(target: AssistantExecutionTarget.local),
onUpdate: updates.add,
);
await Future<void>.delayed(Duration.zero);
gateway.controller.add(
GatewayPushEvent(
event: 'chat',
payload: <String, dynamic>{
'runId': 'run-1',
'sessionKey': 'thread-1',
'state': 'final',
'message': <String, dynamic>{
'role': 'assistant',
'content': <Map<String, dynamic>>[
<String, dynamic>{'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);
});
});
}