feat: add assistant task runtime controls

This commit is contained in:
Haitao Pan 2026-05-08 17:16:50 +08:00
parent f9fb3526b4
commit 1239c59d42
9 changed files with 586 additions and 145 deletions

View File

@ -3,6 +3,7 @@
import 'dart:async';
import 'dart:convert';
import 'dart:io';
import 'dart:math' as math;
import 'package:crypto/crypto.dart' as crypto;
import 'package:flutter/material.dart';
import 'app_metadata.dart';
@ -786,7 +787,8 @@ extension AppControllerDesktopRuntimeHelpers on AppController {
String authorization,
) async {
var bytes = <int>[];
for (var attempt = 1; attempt <= 3; attempt++) {
const maxAttempts = 5;
for (var attempt = 1; attempt <= maxAttempts; attempt++) {
final result = await _downloadBridgeArtifactBytesOnceInternal(
uri,
authorization,
@ -801,8 +803,9 @@ extension AppControllerDesktopRuntimeHelpers on AppController {
if (result.completed) {
return bytes;
}
if (attempt < 3) {
await Future<void>.delayed(Duration(milliseconds: attempt * 250));
if (attempt < maxAttempts) {
final delayMs = math.min(2000, 250 * (1 << (attempt - 1)));
await Future<void>.delayed(Duration(milliseconds: delayMs));
}
}
return null;
@ -948,10 +951,7 @@ extension AppControllerDesktopRuntimeHelpers on AppController {
if (bridgeEndpoint == null) {
return null;
}
if (_usesOpenClawTaskSubmitEndpointInternal(request)) {
return bridgeEndpoint.replace(path: '/gateway/openclaw');
}
return bridgeEndpoint;
return resolveAcpHttpRpcEndpoint(bridgeEndpoint);
}
Uri? gatewayProfileBaseUriInternal(GatewayConnectionProfile profile) {
@ -1051,22 +1051,6 @@ extension AppControllerDesktopRuntimeHelpers on AppController {
) => kGatewayRemoteProfileIndex;
}
bool _usesOpenClawTaskSubmitEndpointInternal(GoTaskServiceRequest request) {
if (request.isMultiAgentRequest || !request.target.isGateway) {
return false;
}
final providerId = normalizeSingleAgentProviderId(
request.provider.providerId,
);
if (providerId == kCanonicalGatewayProviderId) {
return true;
}
return normalizeSingleAgentProviderId(
request.effectiveRouting.preferredGatewayTarget,
) ==
kCanonicalGatewayProviderId;
}
String _normalizeAuthorizationHeaderInternal(String raw) {
final trimmed = raw.trim();
if (trimmed.isEmpty) {

View File

@ -402,6 +402,10 @@ extension AppControllerDesktopThreadActions on AppController {
}
},
);
if (!aiGatewayPendingSessionKeysInternal.contains(sessionKey)) {
clearAiGatewayStreamingTextInternal(sessionKey);
return;
}
clearAiGatewayStreamingTextInternal(sessionKey);
upsertTaskThreadInternal(
sessionKey,
@ -420,7 +424,6 @@ extension AppControllerDesktopThreadActions on AppController {
lastResultCode: result.success ? 'success' : 'error',
updatedAtMs: DateTime.now().millisecondsSinceEpoch.toDouble(),
);
await persistGoTaskArtifactsForSessionInternal(sessionKey, result);
if (!result.success) {
appendLocalSessionMessageInternal(
sessionKey,
@ -468,7 +471,18 @@ extension AppControllerDesktopThreadActions on AppController {
),
persistInThreadContext: true,
);
recomputeTasksInternal();
notifyIfActiveInternal();
await persistGoTaskArtifactsForSessionInternal(sessionKey, result);
} catch (error) {
if (!aiGatewayPendingSessionKeysInternal.contains(sessionKey) &&
taskThreadForSessionInternal(
sessionKey,
)?.lifecycleState.lastResultCode ==
'aborted') {
clearAiGatewayStreamingTextInternal(sessionKey);
return;
}
clearAiGatewayStreamingTextInternal(sessionKey);
final recoverableTransportCode =
recoverableAcpHttpTransportCodeInternal(error);

View File

@ -17,6 +17,7 @@ import '../../app/ui_feature_manifest.dart';
import '../../i18n/app_language.dart';
import '../../models/app_models.dart';
import '../../runtime/multi_agent_orchestrator.dart';
import '../../runtime/gateway_acp_client.dart';
import '../../runtime/runtime_models.dart';
import '../../theme/app_palette.dart';
import '../../theme/app_theme.dart';
@ -122,6 +123,12 @@ extension AssistantPageStateClosureInternal on AssistantPageStateInternal {
lifecycleStatus: thread?.lifecycleState.status ?? '',
lastResultCode: thread?.lifecycleState.lastResultCode ?? '',
artifactSyncStatus: thread?.lastArtifactSyncStatus ?? '',
runtimeBudgetMinutes:
gatewayAcpTaskRuntimeBudgetMinutesForParams(<String, dynamic>{
'taskPrompt': currentTask.preview,
'requestedExecutionTarget':
currentTask.executionTarget.promptValue,
}),
);
return SurfaceCard(
@ -161,7 +168,21 @@ extension AssistantPageStateClosureInternal on AssistantPageStateInternal {
),
),
),
AssistantTaskProgressBar(state: progressState),
AssistantTaskProgressBar(
state: progressState,
onStop: progressState.running
? () {
unawaited(controller.abortRun());
}
: null,
onContinue: progressState.interrupted
? () {
unawaited(
controller.sendChatMessage(appText('继续', 'Continue')),
);
}
: null,
),
ColoredBox(
color: palette.canvas,
child: SizedBox(

View File

@ -633,7 +633,11 @@ class GatewayAcpClient {
),
);
final response = await httpRequest.close().timeout(
gatewayAcpHttpResponseTimeoutFor(endpoint, request.method),
gatewayAcpHttpResponseTimeoutFor(
endpoint,
request.method,
request.params,
),
);
statusCode = response.statusCode;
contentType =
@ -1170,14 +1174,6 @@ class GatewayAcpClient {
Uri? _resolveHttpRpcEndpoint([Uri? endpointOverride, String method = '']) {
final endpoint = endpointOverride ?? endpointResolver();
if (_isOpenClawTaskSubmitEndpoint(endpoint) &&
_isOpenClawTaskSubmitMethod(method)) {
return endpoint?.replace(
path: '/gateway/openclaw',
query: null,
fragment: null,
);
}
return resolveAcpHttpRpcEndpoint(endpoint);
}
@ -1236,26 +1232,116 @@ class GatewayAcpClient {
}
}
bool _isOpenClawTaskSubmitEndpoint(Uri? endpoint) {
var path = endpoint?.path.trim() ?? '';
if (!path.startsWith('/')) {
path = '/$path';
}
path = path.replaceFirst(RegExp(r'/+$'), '');
return path == '/gateway/openclaw';
}
bool _isOpenClawTaskSubmitMethod(String method) {
final normalized = method.trim();
return normalized == 'session.start' || normalized == 'session.message';
}
Duration gatewayAcpHttpResponseTimeoutFor(Uri endpoint, String method) {
if (_isOpenClawTaskSubmitEndpoint(endpoint) &&
_isOpenClawTaskSubmitMethod(method)) {
return const Duration(minutes: 10);
Duration gatewayAcpHttpResponseTimeoutFor(
Uri endpoint,
String method, [
Map<String, dynamic> params = const <String, dynamic>{},
]) {
if (!_isOpenClawTaskSubmitMethod(method)) {
return const Duration(seconds: 120);
}
return const Duration(seconds: 120);
return Duration(minutes: gatewayAcpTaskRuntimeBudgetMinutesForParams(params));
}
int gatewayAcpTaskRuntimeBudgetMinutesForParams(Map<String, dynamic> params) {
if (_looksLikeLongArtifactTask(params)) {
return 30;
}
if (_looksLikeGatewayTask(params)) {
return 10;
}
return 2;
}
bool _looksLikeGatewayTask(Map<String, dynamic> params) {
final target = _paramText(params, const <String>[
'requestedExecutionTarget',
'executionTarget',
]).toLowerCase();
if (target == AssistantExecutionTarget.gateway.promptValue) {
return true;
}
final providerText = _paramText(params, const <String>[
'provider',
'gatewayProvider',
'preferredGatewayProviderId',
]).toLowerCase();
if (providerText.contains('openclaw')) {
return true;
}
final routing = params['routing'];
if (routing is Map) {
final preferred = routing['preferredGatewayProviderId']
?.toString()
.trim()
.toLowerCase();
return preferred == kCanonicalGatewayProviderId ||
preferred?.contains('openclaw') == true;
}
return false;
}
bool _looksLikeLongArtifactTask(Map<String, dynamic> params) {
final prompt = _paramText(params, const <String>[
'taskPrompt',
'prompt',
'message',
]);
final lower = prompt.toLowerCase();
final attachments =
_paramListLength(params['attachments']) +
_paramListLength(params['inlineAttachments']);
if (attachments >= 2 || prompt.length >= 1200) {
return true;
}
const markers = <String>[
'生成文件',
'同步生成文件',
'产物',
'附件',
'图片提示词',
'完整调研ppt',
'markdown格式',
'输出markdown',
'输出 完整',
'ppt',
'pptx',
'powerpoint',
'markdown',
'.md',
'javascript',
'.js',
'image prompt',
'artifacts',
'downloadurl',
];
return markers.any(lower.contains);
}
String _paramText(Map<String, dynamic> params, List<String> keys) {
for (final key in keys) {
final value = params[key];
if (value == null) {
continue;
}
final text = value.toString().trim();
if (text.isNotEmpty) {
return text;
}
}
return '';
}
int _paramListLength(Object? value) {
if (value is List) {
return value.length;
}
return 0;
}
class _GatewayAcpRpcRequest {

View File

@ -16,25 +16,40 @@ class AssistantTaskProgressState {
required this.phase,
required this.label,
this.value,
this.runtimeBudgetMinutes,
});
const AssistantTaskProgressState.idle()
: phase = AssistantTaskProgressPhase.idle,
label = '',
value = null;
value = null,
runtimeBudgetMinutes = null;
final AssistantTaskProgressPhase phase;
final String label;
final double? value;
final int? runtimeBudgetMinutes;
bool get visible => phase != AssistantTaskProgressPhase.idle;
bool get interrupted => phase == AssistantTaskProgressPhase.interrupted;
bool get running =>
phase == AssistantTaskProgressPhase.running ||
phase == AssistantTaskProgressPhase.retrying ||
phase == AssistantTaskProgressPhase.continuing ||
phase == AssistantTaskProgressPhase.syncingArtifacts;
}
class AssistantTaskProgressBar extends StatelessWidget {
const AssistantTaskProgressBar({super.key, required this.state});
const AssistantTaskProgressBar({
super.key,
required this.state,
this.onStop,
this.onContinue,
});
final AssistantTaskProgressState state;
final VoidCallback? onStop;
final VoidCallback? onContinue;
@override
Widget build(BuildContext context) {
@ -60,68 +75,135 @@ class AssistantTaskProgressBar extends StatelessWidget {
bottom: BorderSide(color: theme.dividerColor.withValues(alpha: 0.42)),
),
),
child: Column(
mainAxisSize: MainAxisSize.min,
crossAxisAlignment: CrossAxisAlignment.stretch,
child: Row(
children: [
Text(
state.label,
key: const Key('assistant-task-progress-label'),
maxLines: 1,
overflow: TextOverflow.ellipsis,
style: theme.textTheme.labelSmall?.copyWith(
color: color,
fontWeight: FontWeight.w700,
Expanded(
child: Column(
mainAxisSize: MainAxisSize.min,
crossAxisAlignment: CrossAxisAlignment.stretch,
children: [
Text(
state.label,
key: const Key('assistant-task-progress-label'),
maxLines: 1,
overflow: TextOverflow.ellipsis,
style: theme.textTheme.labelSmall?.copyWith(
color: color,
fontWeight: FontWeight.w700,
),
),
const SizedBox(height: 5),
LinearProgressIndicator(
key: const Key('assistant-task-progress-indicator'),
value: state.value,
minHeight: 3,
color: color,
backgroundColor: color.withValues(alpha: 0.16),
borderRadius: BorderRadius.circular(999),
),
],
),
),
const SizedBox(height: 5),
LinearProgressIndicator(
key: const Key('assistant-task-progress-indicator'),
value: state.value,
minHeight: 3,
color: color,
backgroundColor: color.withValues(alpha: 0.16),
borderRadius: BorderRadius.circular(999),
),
if (state.running && onStop != null) ...[
const SizedBox(width: 8),
_AssistantTaskProgressActionButton(
key: const Key('assistant-task-progress-stop-button'),
icon: Icons.stop_rounded,
label: appText('停止', 'Stop'),
color: color,
onPressed: onStop,
),
] else if (state.interrupted && onContinue != null) ...[
const SizedBox(width: 8),
_AssistantTaskProgressActionButton(
key: const Key('assistant-task-progress-continue-button'),
icon: Icons.play_arrow_rounded,
label: appText('继续', 'Continue'),
color: color,
onPressed: onContinue,
),
],
],
),
);
}
}
class _AssistantTaskProgressActionButton extends StatelessWidget {
const _AssistantTaskProgressActionButton({
super.key,
required this.icon,
required this.label,
required this.color,
required this.onPressed,
});
final IconData icon;
final String label;
final Color color;
final VoidCallback? onPressed;
@override
Widget build(BuildContext context) {
return TextButton.icon(
onPressed: onPressed,
icon: Icon(icon, size: 16),
label: Text(label),
style: TextButton.styleFrom(
foregroundColor: color,
minimumSize: const Size(0, 28),
padding: const EdgeInsets.symmetric(horizontal: 8, vertical: 4),
tapTargetSize: MaterialTapTargetSize.shrinkWrap,
visualDensity: VisualDensity.compact,
),
);
}
}
AssistantTaskProgressState assistantTaskProgressState({
required bool pending,
required String lifecycleStatus,
required String lastResultCode,
required String artifactSyncStatus,
int? runtimeBudgetMinutes,
}) {
final syncStatus = artifactSyncStatus.trim().toLowerCase();
final status = lifecycleStatus.trim().toLowerCase();
final budget = runtimeBudgetMinutes == null || runtimeBudgetMinutes <= 0
? null
: runtimeBudgetMinutes;
if (pending && syncStatus == 'syncing') {
return AssistantTaskProgressState(
phase: AssistantTaskProgressPhase.syncingArtifacts,
label: appText('正在同步生成文件...', 'Syncing generated files...'),
value: 0.82,
runtimeBudgetMinutes: budget,
);
}
if (pending && status == 'continuing') {
return AssistantTaskProgressState(
phase: AssistantTaskProgressPhase.continuing,
label: appText('任务继续中...', 'Continuing task...'),
label: _budgetedProgressLabel(
appText('任务继续中', 'Continuing task'),
budget,
),
value: 0.62,
runtimeBudgetMinutes: budget,
);
}
if (pending && status == 'retrying') {
return AssistantTaskProgressState(
phase: AssistantTaskProgressPhase.retrying,
label: appText('任务重试中...', 'Retrying task...'),
label: _budgetedProgressLabel(appText('任务重试中', 'Retrying task'), budget),
value: 0.38,
runtimeBudgetMinutes: budget,
);
}
if (pending) {
return AssistantTaskProgressState(
phase: AssistantTaskProgressPhase.running,
label: appText('任务运行中...', 'Task running...'),
label: _budgetedProgressLabel(appText('任务运行中', 'Task running'), budget),
runtimeBudgetMinutes: budget,
);
}
final result = lastResultCode.trim().toUpperCase();
@ -135,6 +217,13 @@ AssistantTaskProgressState assistantTaskProgressState({
return const AssistantTaskProgressState.idle();
}
String _budgetedProgressLabel(String base, int? minutes) {
if (minutes == null || minutes <= 0) {
return '$base...';
}
return appText('$base,预计最长 $minutes 分钟...', '$base, up to $minutes min...');
}
String _interruptedTaskProgressLabel(String result) {
if (result == 'ACP_HTTP_HANDSHAKE_INTERRUPTED') {
return appText(

View File

@ -15,7 +15,9 @@ void main() {
lifecycleStatus: 'running',
lastResultCode: 'running',
artifactSyncStatus: '',
runtimeBudgetMinutes: 30,
),
onStop: () {},
),
);
@ -23,7 +25,11 @@ void main() {
find.byKey(const Key('assistant-task-progress-bar')),
findsOneWidget,
);
expect(find.text('任务运行中...'), findsOneWidget);
expect(find.text('任务运行中,预计最长 30 分钟...'), findsOneWidget);
expect(
find.byKey(const Key('assistant-task-progress-stop-button')),
findsOneWidget,
);
final indicator = tester.widget<LinearProgressIndicator>(
find.byKey(const Key('assistant-task-progress-indicator')),
);
@ -104,10 +110,15 @@ void main() {
lastResultCode: 'ACP_HTTP_CONNECTION_CLOSED',
artifactSyncStatus: 'interrupted',
),
onContinue: () {},
),
);
expect(find.text('Bridge 响应中断,等待下一次发送续写同一会话。'), findsOneWidget);
expect(
find.byKey(const Key('assistant-task-progress-continue-button')),
findsOneWidget,
);
final indicator = tester.widget<LinearProgressIndicator>(
find.byKey(const Key('assistant-task-progress-indicator')),
);
@ -188,13 +199,21 @@ void main() {
});
}
Widget _buildTestApp(AssistantTaskProgressState state) {
Widget _buildTestApp(
AssistantTaskProgressState state, {
VoidCallback? onStop,
VoidCallback? onContinue,
}) {
return MaterialApp(
theme: AppTheme.light(),
home: Material(
child: SizedBox(
width: 420,
child: AssistantTaskProgressBar(state: state),
child: AssistantTaskProgressBar(
state: state,
onStop: onStop,
onContinue: onContinue,
),
),
),
);

View File

@ -539,6 +539,105 @@ void main() {
},
);
test(
'retries bridge artifact downloads up to five weak network attempts',
() async {
const body = 'download after retries';
var requestCount = 0;
final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0);
addTearDown(() => server.close());
server.listen((socket) async {
requestCount += 1;
final requestBytes = <int>[];
await for (final chunk in socket) {
requestBytes.addAll(chunk);
if (String.fromCharCodes(requestBytes).contains('\r\n\r\n')) {
break;
}
}
if (requestCount < 5) {
socket.destroy();
return;
}
socket.add(
'HTTP/1.1 200 OK\r\n'
'Content-Type: text/plain\r\n'
'Content-Length: ${body.length}\r\n'
'\r\n'
.codeUnits,
);
socket.add(body.codeUnits);
await socket.flush();
await socket.close();
});
final controller = AppController(
environmentOverride: const <String, String>{
'BRIDGE_AUTH_TOKEN': 'bridge-token',
},
);
addTearDown(controller.dispose);
final localWorkspace = await Directory.systemTemp.createTemp(
'xworkmate-retry-artifact-workspace-',
);
addTearDown(() async {
if (await localWorkspace.exists()) {
await localWorkspace.delete(recursive: true);
}
});
controller.upsertTaskThreadInternal(
'session-1',
workspaceBinding: WorkspaceBinding(
workspaceId: 'session-1',
workspaceKind: WorkspaceKind.localFs,
workspacePath: localWorkspace.path,
displayPath: localWorkspace.path,
writable: true,
),
);
final result = GoTaskServiceResult(
success: true,
message: 'hello',
turnId: 'turn-1',
raw: <String, dynamic>{
'artifacts': <Map<String, dynamic>>[
<String, dynamic>{
'relativePath': 'reports/retry.txt',
'downloadUrl':
'http://xworkmate-bridge.svc.plus:${server.port}/retry.txt',
'contentType': 'text/plain',
},
],
},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
);
final clientFactory = _proxiedClientFactory(server.port);
await HttpOverrides.runZoned(() async {
await controller.persistGoTaskArtifactsForSessionInternal(
'session-1',
result,
);
}, createHttpClient: clientFactory);
expect(requestCount, 5);
expect(
await File('${localWorkspace.path}/reports/retry.txt').readAsString(),
body,
);
expect(
controller
.requireTaskThreadForSessionInternal('session-1')
.lastArtifactSyncStatus,
'synced',
);
},
);
test('keeps syncing later artifacts when one download fails', () async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
addTearDown(() => server.close(force: true));

View File

@ -706,6 +706,7 @@ void main() {
expect(fakeGoTaskService.requests, hasLength(2));
expect(fakeGoTaskService.requests.last.resumeSession, isTrue);
await _waitForLastChatMessageText(controller, '全部 6 个文件已生成 ✅');
expect(controller.chatMessages.last.text, '全部 6 个文件已生成 ✅');
final thread = controller.taskThreadForSessionInternal('session-1');
expect(thread?.lifecycleState.status, 'ready');
@ -799,6 +800,7 @@ void main() {
expect(fakeGoTaskService.requests, hasLength(2));
expect(fakeGoTaskService.requests.last.resumeSession, isTrue);
await _waitForLastChatMessageText(controller, '全部 6 个文件已生成 ✅');
expect(controller.chatMessages.last.text, '全部 6 个文件已生成 ✅');
final thread = controller.taskThreadForSessionInternal('session-1');
expect(thread?.lifecycleState.status, 'ready');
@ -957,6 +959,80 @@ void main() {
);
});
test('abortRun cancels only the current pending session', () async {
final fakeGoTaskService = _BlockingGoTaskServiceClient();
final controller = _connectedController(fakeGoTaskService);
addTearDown(controller.dispose);
await controller.switchSession('task-a');
final taskAFuture = controller.sendChatMessage('task A');
await fakeGoTaskService.waitForRequestCount(1);
await controller.switchSession('task-b');
final taskBFuture = controller.sendChatMessage('task B');
await fakeGoTaskService.waitForRequestCount(2);
fakeGoTaskService.emitDelta('task-b', 'streaming text');
expect(controller.assistantSessionHasPendingRun('task-a'), isTrue);
expect(controller.assistantSessionHasPendingRun('task-b'), isTrue);
await controller.abortRun();
expect(fakeGoTaskService.cancelledSessionIds, <String>['task-b']);
expect(controller.assistantSessionHasPendingRun('task-a'), isTrue);
expect(controller.assistantSessionHasPendingRun('task-b'), isFalse);
expect(
controller
.requireTaskThreadForSessionInternal('task-b')
.lifecycleState
.lastResultCode,
'aborted',
);
expect(
controller.aiGatewayStreamingTextBySessionInternal['task-b'],
isNull,
);
fakeGoTaskService.complete(
'task-b',
const GoTaskServiceResult(
success: true,
message: 'late result B',
turnId: 'turn-b',
raw: <String, dynamic>{},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
),
);
await taskBFuture;
expect(
controller.localSessionMessagesInternal['task-b']!.map(
(message) => message.text,
),
isNot(contains('late result B')),
);
fakeGoTaskService.complete(
'task-a',
const GoTaskServiceResult(
success: true,
message: 'result A',
turnId: 'turn-a',
raw: <String, dynamic>{},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
),
);
await taskAFuture;
expect(
controller.localSessionMessagesInternal['task-a']!.map(
(message) => message.text,
),
contains('result A'),
);
});
test(
'sendChatMessage exposes continuing and retrying lifecycle states',
() async {
@ -1126,6 +1202,24 @@ Future<void> _waitForThreadLifecycleStatus(
);
}
Future<void> _waitForLastChatMessageText(
AppController controller,
String expectedText,
) async {
final deadline = DateTime.now().add(const Duration(seconds: 2));
while (DateTime.now().isBefore(deadline)) {
if (controller.chatMessages.isNotEmpty &&
controller.chatMessages.last.text == expectedText) {
return;
}
await Future<void>.delayed(const Duration(milliseconds: 10));
}
expect(
controller.chatMessages.isEmpty ? '' : controller.chatMessages.last.text,
expectedText,
);
}
Future<_CapabilityServerCapture> _startEmptyCapabilityServer() async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
final capture = _CapabilityServerCapture._(
@ -1291,8 +1385,11 @@ class _RecordingGoTaskServiceClient implements GoTaskServiceClient {
class _BlockingGoTaskServiceClient implements GoTaskServiceClient {
final List<GoTaskServiceRequest> requests = <GoTaskServiceRequest>[];
final List<String> cancelledSessionIds = <String>[];
final Map<String, Completer<GoTaskServiceResult>> _pending =
<String, Completer<GoTaskServiceResult>>{};
final Map<String, void Function(GoTaskServiceUpdate)> _updates =
<String, void Function(GoTaskServiceUpdate)>{};
@override
Future<ExternalCodeAgentAcpCapabilities> loadExternalAcpCapabilities({
@ -1314,6 +1411,7 @@ class _BlockingGoTaskServiceClient implements GoTaskServiceClient {
required void Function(GoTaskServiceUpdate update) onUpdate,
}) {
requests.add(request);
_updates[request.sessionId] = onUpdate;
final completer = Completer<GoTaskServiceResult>();
_pending[request.sessionId] = completer;
return completer.future;
@ -1331,19 +1429,43 @@ class _BlockingGoTaskServiceClient implements GoTaskServiceClient {
void complete(String sessionId, GoTaskServiceResult result) {
final completer = _pending.remove(sessionId);
_updates.remove(sessionId);
if (completer == null) {
throw StateError('No pending task for $sessionId.');
}
completer.complete(result);
}
void emitDelta(String sessionId, String text) {
final onUpdate = _updates[sessionId];
if (onUpdate == null) {
throw StateError('No pending update sink for $sessionId.');
}
onUpdate(
GoTaskServiceUpdate(
sessionId: sessionId,
threadId: sessionId,
turnId: 'turn-$sessionId',
type: 'delta',
text: text,
message: '',
pending: true,
error: false,
route: GoTaskServiceRoute.externalAcpSingle,
payload: const <String, dynamic>{},
),
);
}
@override
Future<void> cancelTask({
required GoTaskServiceRoute route,
required AssistantExecutionTarget target,
required String sessionId,
required String threadId,
}) async {}
}) async {
cancelledSessionIds.add(sessionId);
}
@override
Future<void> closeTask({

View File

@ -832,7 +832,7 @@ void main() {
);
test(
'desktop task execution routes OpenClaw through dedicated bridge gateway path',
'desktop task execution routes OpenClaw through managed bridge RPC',
() async {
final capture = await _startAcpHttpServer(
streamResponse: true,
@ -865,8 +865,7 @@ void main() {
final transport = ExternalCodeAgentAcpDesktopTransport(
client: client,
endpointResolver: (_) => capture.baseEndpoint,
taskEndpointResolver: (_) =>
capture.baseEndpoint.replace(path: '/gateway/openclaw'),
taskEndpointResolver: (_) => capture.baseEndpoint,
);
final result = await transport.executeTask(
@ -879,7 +878,8 @@ void main() {
expect(capture.authorizationHeader, 'Bearer bridge-token');
expect(capture.acceptHeader, 'text/event-stream, application/json');
expect(capture.requestPath, '/gateway/openclaw');
expect(capture.requestPath, '/acp/rpc');
expect(capture.requestPath, isNot(contains('/gateway/openclaw')));
expect(capture.requestPath, isNot(contains('/acp-server')));
expect(capture.requestPath, isNot(contains('/acp-server/gateway')));
final params = _lastRequestParams(capture);
@ -915,7 +915,7 @@ void main() {
);
test(
'desktop OpenClaw follow-up routes through dedicated bridge gateway path',
'desktop OpenClaw follow-up routes through managed bridge RPC',
() async {
final capture = await _startAcpHttpServer();
addTearDown(capture.close);
@ -927,8 +927,7 @@ void main() {
final transport = ExternalCodeAgentAcpDesktopTransport(
client: client,
endpointResolver: (_) => capture.baseEndpoint,
taskEndpointResolver: (_) =>
capture.baseEndpoint.replace(path: '/gateway/openclaw'),
taskEndpointResolver: (_) => capture.baseEndpoint,
);
await transport.executeTask(
@ -941,30 +940,42 @@ void main() {
);
expect(capture.acceptHeader, 'text/event-stream, application/json');
expect(capture.requestPath, '/gateway/openclaw');
expect(capture.requestPath, '/acp/rpc');
expect(capture.requestPath, isNot(contains('/gateway/openclaw')));
expect(capture.requestBody, contains('"method":"session.message"'));
},
);
test('OpenClaw task submit uses extended HTTP response timeout', () {
test('task submit uses dynamic HTTP response timeout budgets', () {
final openClawEndpoint = Uri.parse(
'https://xworkmate-bridge.svc.plus/gateway/openclaw',
'https://xworkmate-bridge.svc.plus/acp/rpc',
);
final acpEndpoint = Uri.parse(
'https://xworkmate-bridge.svc.plus/acp/rpc',
);
expect(
gatewayAcpHttpResponseTimeoutFor(openClawEndpoint, 'session.start'),
gatewayAcpHttpResponseTimeoutFor(
openClawEndpoint,
'session.start',
const <String, dynamic>{'requestedExecutionTarget': 'gateway'},
),
const Duration(minutes: 10),
);
expect(
gatewayAcpHttpResponseTimeoutFor(openClawEndpoint, 'session.message'),
const Duration(minutes: 10),
gatewayAcpHttpResponseTimeoutFor(
openClawEndpoint,
'session.message',
const <String, dynamic>{
'taskPrompt': '输出 完整调研PPT 和 Markdown格式 文件',
'requestedExecutionTarget': 'gateway',
},
),
const Duration(minutes: 30),
);
expect(
gatewayAcpHttpResponseTimeoutFor(acpEndpoint, 'session.start'),
const Duration(seconds: 120),
const Duration(minutes: 2),
);
expect(
gatewayAcpHttpResponseTimeoutFor(openClawEndpoint, 'acp.capabilities'),
@ -972,62 +983,58 @@ void main() {
);
});
test(
'desktop controller only uses gateway path for OpenClaw task submit',
() {
final controller = AppController(
environmentOverride: const <String, String>{},
);
addTearDown(controller.dispose);
test('desktop controller keeps task submit on managed bridge RPC', () {
final controller = AppController(
environmentOverride: const <String, String>{},
);
addTearDown(controller.dispose);
final openClawStart = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
),
);
final openClawFollowUp = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
resumeSession: true,
),
);
final unspecifiedGateway = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.unspecified,
),
);
final multiAgentGateway = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
multiAgent: true,
),
);
final agentTask = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.agent,
provider: SingleAgentProvider.codex,
),
);
final openClawStart = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
),
);
final openClawFollowUp = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
resumeSession: true,
),
);
final unspecifiedGateway = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.unspecified,
),
);
final multiAgentGateway = controller
.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
multiAgent: true,
),
);
final agentTask = controller.resolveExternalAcpEndpointForRequestInternal(
_taskRequest(
target: AssistantExecutionTarget.agent,
provider: SingleAgentProvider.codex,
),
);
expect(openClawStart?.path, '/gateway/openclaw');
expect(openClawFollowUp?.path, '/gateway/openclaw');
expect(unspecifiedGateway?.path, '/gateway/openclaw');
expect(multiAgentGateway?.path, '');
expect(agentTask?.path, '');
},
);
expect(openClawStart?.path, '/acp/rpc');
expect(openClawFollowUp?.path, '/acp/rpc');
expect(unspecifiedGateway?.path, '/acp/rpc');
expect(multiAgentGateway?.path, '/acp/rpc');
expect(agentTask?.path, '/acp/rpc');
});
test(
'desktop controller resolves OpenClaw gateway submit on managed bridge origin',
'desktop controller resolves OpenClaw gateway submit to managed bridge RPC',
() {
final controller = AppController(
environmentOverride: const <String, String>{},
@ -1044,10 +1051,10 @@ void main() {
expect(
endpoint.toString(),
'https://xworkmate-bridge.svc.plus/gateway/openclaw',
'https://xworkmate-bridge.svc.plus/acp/rpc',
);
expect(endpoint, isNotNull);
expect(endpoint!.path, isNot('/acp/rpc'));
expect(endpoint!.path, isNot('/gateway/openclaw'));
},
);