From e4bfeeafcd3649865d3d9612e43e7227bdffeafe Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Mon, 25 May 2026 10:35:30 +0800 Subject: [PATCH] Fix assistant continue task requeue --- ...app_controller_desktop_thread_actions.dart | 170 +++++++++++-- .../app_controller_openclaw_task_queue.dart | 2 + .../assistant_page_state_actions.dart | 8 + .../assistant_page_state_closure.dart | 10 +- .../assistant_task_progress_bar_test.dart | 24 ++ .../assistant_execution_target_test.dart | 238 ++++++++++++++++++ 6 files changed, 432 insertions(+), 20 deletions(-) diff --git a/lib/app/app_controller_desktop_thread_actions.dart b/lib/app/app_controller_desktop_thread_actions.dart index a0f5011b..bd74db6a 100644 --- a/lib/app/app_controller_desktop_thread_actions.dart +++ b/lib/app/app_controller_desktop_thread_actions.dart @@ -249,11 +249,100 @@ extension AppControllerDesktopThreadActions on AppController { sessionsControllerInternal.currentSessionKey, ); } - final currentTarget = assistantExecutionTargetForSession(sessionKey); final resumeSessionHint = shouldResumeGatewaySessionForNextSendInternal( sessionKey, ); - var connectionState = assistantConnectionStateForSession(sessionKey); + await dispatchGatewayChatTurnInternal( + sessionKey: sessionKey, + message: message, + thinking: thinking, + attachments: attachments, + localAttachments: localAttachments, + selectedSkillLabels: selectedSkillLabels, + resumeSessionHint: resumeSessionHint, + ); + } + + Future continueAssistantTaskInternal(String sessionKey) async { + final normalizedSessionKey = normalizedAssistantSessionKeyInternal( + sessionKey.trim().isEmpty + ? sessionsControllerInternal.currentSessionKey + : sessionKey, + ); + final thread = taskThreadForSessionInternal(normalizedSessionKey); + final lifecycleStatus = thread?.lifecycleState.status ?? ''; + final lastResultCode = thread?.lifecycleState.lastResultCode ?? ''; + final artifactSyncStatus = thread?.lastArtifactSyncStatus ?? ''; + if (!isRecoverableAssistantTaskStateInternal( + lifecycleStatus: lifecycleStatus, + lastResultCode: lastResultCode, + artifactSyncStatus: artifactSyncStatus, + )) { + final error = StateError( + appText('当前任务状态不可继续执行。', 'The current task state cannot be continued.'), + ); + appendAssistantThreadMessageInternal( + normalizedSessionKey, + assistantErrorMessageInternal(error.message), + ); + await flushAssistantThreadPersistenceInternal(); + recomputeTasksInternal(); + notifyIfActiveInternal(); + throw error; + } + final lastUserTurn = lastCommittedUserTurnForGatewaySessionInternal( + normalizedSessionKey, + ); + final message = lastUserTurn?.text.trim() ?? ''; + if (message.isEmpty) { + final error = StateError( + appText( + '当前任务没有可恢复的用户请求,请输入需求后重新提交。', + 'This task has no recoverable user request. Enter a request and submit it again.', + ), + ); + appendAssistantThreadMessageInternal( + normalizedSessionKey, + assistantErrorMessageInternal(error.message), + ); + await flushAssistantThreadPersistenceInternal(); + recomputeTasksInternal(); + notifyIfActiveInternal(); + throw error; + } + await dispatchGatewayChatTurnInternal( + sessionKey: normalizedSessionKey, + message: message, + thinking: 'off', + attachments: const [], + localAttachments: const [], + selectedSkillLabels: const [], + resumeSessionHint: + lastResultCode.trim().toUpperCase() != 'ABORTED' && + shouldResumeGatewaySessionForNextSendInternal(normalizedSessionKey), + appendUserTurn: false, + ); + } + + Future dispatchGatewayChatTurnInternal({ + required String sessionKey, + required String message, + required String thinking, + required List attachments, + required List localAttachments, + required List selectedSkillLabels, + required bool resumeSessionHint, + bool appendUserTurn = true, + }) async { + final normalizedSessionKey = normalizedAssistantSessionKeyInternal( + sessionKey, + ); + final currentTarget = assistantExecutionTargetForSession( + normalizedSessionKey, + ); + var connectionState = assistantConnectionStateForSession( + normalizedSessionKey, + ); if (!connectionState.connected && isBridgeAcpRuntimeConfiguredInternal() && bridgeCapabilityRefreshNeededForAssistantTargetInternal( @@ -261,7 +350,9 @@ extension AppControllerDesktopThreadActions on AppController { )) { try { await refreshAcpCapabilitiesInternal(forceRefresh: true); - connectionState = assistantConnectionStateForSession(sessionKey); + connectionState = assistantConnectionStateForSession( + normalizedSessionKey, + ); } catch (_) { // Fallback to existing connection state if refresh fails. } @@ -269,7 +360,7 @@ extension AppControllerDesktopThreadActions on AppController { if (!connectionState.connected) { final error = StateError(connectionState.detailLabel); appendAssistantThreadMessageInternal( - sessionKey, + normalizedSessionKey, assistantErrorMessageInternal(error.message), ); await flushAssistantThreadPersistenceInternal(); @@ -278,14 +369,17 @@ extension AppControllerDesktopThreadActions on AppController { throw error; } await ensureDesktopTaskThreadBindingInternal( - sessionKey, + normalizedSessionKey, executionTarget: currentTarget, ); final workingDirectory = - assistantWorkingDirectoryForSessionInternal(sessionKey)?.trim() ?? ''; + assistantWorkingDirectoryForSessionInternal( + normalizedSessionKey, + )?.trim() ?? + ''; final remoteWorkingDirectoryHint = assistantRemoteWorkingDirectoryHintForSessionInternal( - sessionKey, + normalizedSessionKey, )?.trim() ?? ''; if (workingDirectory.isEmpty) { @@ -296,7 +390,7 @@ extension AppControllerDesktopThreadActions on AppController { ), ); appendAssistantThreadMessageInternal( - sessionKey, + normalizedSessionKey, assistantErrorMessageInternal(error.message), ); await flushAssistantThreadPersistenceInternal(); @@ -311,7 +405,7 @@ extension AppControllerDesktopThreadActions on AppController { } if (providerCatalogForExecutionTarget(currentTarget).isEmpty) { upsertTaskThreadInternal( - sessionKey, + normalizedSessionKey, selectedProvider: SingleAgentProvider.unspecified, selectedProviderSource: ThreadSelectionSource.inherited, latestResolvedProviderId: '', @@ -329,7 +423,7 @@ extension AppControllerDesktopThreadActions on AppController { ), ); appendAssistantThreadMessageInternal( - sessionKey, + normalizedSessionKey, assistantErrorMessageInternal(error.message), ); await flushAssistantThreadPersistenceInternal(); @@ -338,11 +432,13 @@ extension AppControllerDesktopThreadActions on AppController { throw error; } } - final provider = assistantProviderForSession(sessionKey); + final provider = assistantProviderForSession(normalizedSessionKey); final model = currentTarget.isGateway ? '' - : assistantModelForSession(sessionKey); - final routing = buildExternalAcpRoutingForSessionInternal(sessionKey); + : assistantModelForSession(normalizedSessionKey); + final routing = buildExternalAcpRoutingForSessionInternal( + normalizedSessionKey, + ); final dispatch = await codeAgentNodeOrchestratorInternal .buildGatewayDispatch( buildCodeAgentNodeStateInternal(executionTarget: currentTarget), @@ -361,7 +457,7 @@ extension AppControllerDesktopThreadActions on AppController { OpenClawGatewayQueuedTurnInternal( queueId: 'openclaw-${DateTime.now().microsecondsSinceEpoch}-$localMessageCounterInternal', - sessionKey: sessionKey, + sessionKey: normalizedSessionKey, target: currentTarget, provider: provider, message: message, @@ -376,14 +472,15 @@ extension AppControllerDesktopThreadActions on AppController { agentId: dispatch.agentId ?? '', metadata: Map.unmodifiable(dispatch.metadata), resumeSessionHint: resumeSessionHint, + appendUserTurn: appendUserTurn, ), ); return; } await enqueueThreadTurnInternal( - sessionKey, + normalizedSessionKey, () => runGatewayChatTurnInternal( - sessionKey: sessionKey, + sessionKey: normalizedSessionKey, target: currentTarget, provider: provider, message: message, @@ -398,6 +495,7 @@ extension AppControllerDesktopThreadActions on AppController { agentId: dispatch.agentId ?? '', metadata: Map.unmodifiable(dispatch.metadata), resumeSessionHint: resumeSessionHint, + appendUserTurn: appendUserTurn, ), ); recomputeTasksInternal(); @@ -570,7 +668,9 @@ extension AppControllerDesktopThreadActions on AppController { Future enqueueOpenClawGatewayTurnInternal( OpenClawGatewayQueuedTurnInternal turn, ) async { - appendGatewayUserTurnInternal(turn.sessionKey, turn.message); + if (turn.appendUserTurn) { + appendGatewayUserTurnInternal(turn.sessionKey, turn.message); + } if (openClawGatewayActiveTasksInternal >= openClawGatewayMaxActiveTasksInternal && openClawGatewayQueuedTurnsInternal.length >= @@ -974,6 +1074,42 @@ extension AppControllerDesktopThreadActions on AppController { }); } + GatewayChatMessage? lastCommittedUserTurnForGatewaySessionInternal( + String sessionKey, + ) { + final normalizedSessionKey = normalizedAssistantSessionKeyInternal( + sessionKey, + ); + final messages = [ + ...?assistantThreadRecordsInternal[normalizedSessionKey]?.messages, + ...?assistantThreadMessagesInternal[normalizedSessionKey], + ...?localSessionMessagesInternal[normalizedSessionKey], + ]; + for (final message in messages.reversed) { + final role = message.role.trim().toLowerCase(); + if (role == 'user' && !message.pending) { + return message; + } + } + return null; + } + + bool isRecoverableAssistantTaskStateInternal({ + required String lifecycleStatus, + required String lastResultCode, + required String artifactSyncStatus, + }) { + final status = lifecycleStatus.trim().toLowerCase(); + final syncStatus = artifactSyncStatus.trim().toLowerCase(); + final result = lastResultCode.trim().toUpperCase(); + return status == 'interrupted' || + syncStatus == 'interrupted' || + result == 'ABORTED' || + result == 'ERROR' || + result == 'ACP_HTTP_CONNECTION_CLOSED' || + result == 'SESSION_CONTINUATION_UNAVAILABLE'; + } + bool shouldResumeGatewaySessionForNextSendInternal(String sessionKey) { final normalizedSessionKey = normalizedAssistantSessionKeyInternal( sessionKey, diff --git a/lib/app/app_controller_openclaw_task_queue.dart b/lib/app/app_controller_openclaw_task_queue.dart index 61edc6fb..66c5caf4 100644 --- a/lib/app/app_controller_openclaw_task_queue.dart +++ b/lib/app/app_controller_openclaw_task_queue.dart @@ -22,6 +22,7 @@ class OpenClawGatewayQueuedTurnInternal { required this.agentId, required this.metadata, required this.resumeSessionHint, + this.appendUserTurn = true, }); final String queueId; @@ -40,6 +41,7 @@ class OpenClawGatewayQueuedTurnInternal { final String agentId; final Map metadata; final bool resumeSessionHint; + final bool appendUserTurn; bool cancelled = false; } diff --git a/lib/features/assistant/assistant_page_state_actions.dart b/lib/features/assistant/assistant_page_state_actions.dart index ba6acdd9..7b563b5d 100644 --- a/lib/features/assistant/assistant_page_state_actions.dart +++ b/lib/features/assistant/assistant_page_state_actions.dart @@ -355,6 +355,14 @@ extension AssistantPageStateActionsInternal on AssistantPageStateInternal { widget.controller.openSettings(tab: SettingsTab.gateway); } + Future continueCurrentTaskInternal(String sessionKey) async { + try { + await widget.controller.continueAssistantTaskInternal(sessionKey); + } catch (_) { + focusComposerInternal(); + } + } + void focusComposerInternal() { if (!mounted) { return; diff --git a/lib/features/assistant/assistant_page_state_closure.dart b/lib/features/assistant/assistant_page_state_closure.dart index 588943b7..ba1aa2d4 100644 --- a/lib/features/assistant/assistant_page_state_closure.dart +++ b/lib/features/assistant/assistant_page_state_closure.dart @@ -178,9 +178,13 @@ extension AssistantPageStateClosureInternal on AssistantPageStateInternal { } : null, onContinue: progressState.recoverable - ? AssistantPageStateActionsInternal( - this, - ).focusComposerInternal + ? () { + unawaited( + AssistantPageStateActionsInternal( + this, + ).continueCurrentTaskInternal(activeSessionKey), + ); + } : null, ), ColoredBox( diff --git a/test/features/assistant/assistant_task_progress_bar_test.dart b/test/features/assistant/assistant_task_progress_bar_test.dart index 5242d732..64976d68 100644 --- a/test/features/assistant/assistant_task_progress_bar_test.dart +++ b/test/features/assistant/assistant_task_progress_bar_test.dart @@ -129,6 +129,30 @@ void main() { expect(indicator.value, 0.48); }); + testWidgets('invokes continue action for a recoverable interrupted task', ( + tester, + ) async { + var continued = false; + await tester.pumpWidget( + _buildTestApp( + assistantTaskProgressState( + pending: false, + lifecycleStatus: 'interrupted', + lastResultCode: 'ACP_HTTP_CONNECTION_CLOSED', + artifactSyncStatus: 'interrupted', + ), + onContinue: () { + continued = true; + }, + ), + ); + + await tester.tap( + find.byKey(const Key('assistant-task-progress-continue-button')), + ); + expect(continued, isTrue); + }); + testWidgets('shows continue action for a stopped task', (tester) async { var continued = false; await tester.pumpWidget( diff --git a/test/runtime/assistant_execution_target_test.dart b/test/runtime/assistant_execution_target_test.dart index a3f75cb7..b3279a0f 100644 --- a/test/runtime/assistant_execution_target_test.dart +++ b/test/runtime/assistant_execution_target_test.dart @@ -2790,6 +2790,142 @@ void main() { }, ); + test( + 'continueAssistantTaskInternal requeues a stopped OpenClaw task without clearing queued work', + () async { + final fakeGoTaskService = _BlockingGoTaskServiceClient(); + final controller = _connectedGatewayController(fakeGoTaskService); + addTearDown(() { + fakeGoTaskService.completeAll(); + controller.dispose(); + }); + + await _selectGatewaySession(controller, 'continue-active-openclaw'); + await controller.sendChatMessage('active task'); + await fakeGoTaskService.waitForRequestCount(1); + + await _selectGatewaySession(controller, 'continue-queued-openclaw'); + await controller.sendChatMessage('queued before continue'); + await _waitForThreadLifecycleStatus( + controller, + 'continue-queued-openclaw', + 'queued', + ); + + await _selectGatewaySession(controller, 'continue-stopped-openclaw'); + controller.appendLocalSessionMessageInternal( + 'continue-stopped-openclaw', + GatewayChatMessage( + id: 'user-continue-stopped-openclaw', + role: 'user', + text: 'resume stopped openclaw task', + timestampMs: DateTime.now().millisecondsSinceEpoch.toDouble(), + toolCallId: null, + toolName: null, + stopReason: null, + pending: false, + error: false, + ), + persistInThreadContext: true, + ); + controller.upsertTaskThreadInternal( + 'continue-stopped-openclaw', + executionTarget: AssistantExecutionTarget.gateway, + selectedProvider: SingleAgentProvider.openclaw, + selectedProviderSource: ThreadSelectionSource.explicit, + lifecycleStatus: 'ready', + lastResultCode: 'aborted', + ); + + await controller.continueAssistantTaskInternal( + 'continue-stopped-openclaw', + ); + + await _waitForThreadLifecycleStatus( + controller, + 'continue-stopped-openclaw', + 'queued', + ); + expect(fakeGoTaskService.requests, hasLength(1)); + expect( + controller.assistantSessionHasPendingRun('continue-queued-openclaw'), + isTrue, + ); + expect( + controller.assistantSessionHasPendingRun('continue-stopped-openclaw'), + isTrue, + ); + expect( + controller.localSessionMessagesInternal['continue-stopped-openclaw']! + .where( + (message) => + message.role == 'user' && + message.text == 'resume stopped openclaw task', + ) + .length, + 1, + ); + + fakeGoTaskService.complete( + 'continue-active-openclaw', + const GoTaskServiceResult( + success: true, + message: 'active done', + turnId: 'turn-active', + raw: {}, + errorMessage: '', + resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, + ), + ); + await fakeGoTaskService.waitForRequestCount(2); + expect( + fakeGoTaskService.requests.last.sessionId, + 'continue-queued-openclaw', + ); + + fakeGoTaskService.complete( + 'continue-queued-openclaw', + const GoTaskServiceResult( + success: true, + message: 'queued done', + turnId: 'turn-queued', + raw: {}, + errorMessage: '', + resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, + ), + ); + await fakeGoTaskService.waitForRequestCount(3); + + final continuedRequest = fakeGoTaskService.requests.last; + expect(continuedRequest.sessionId, 'continue-stopped-openclaw'); + expect(continuedRequest.resumeSession, isFalse); + expect( + continuedRequest.prompt, + contains('User request:\nresume stopped openclaw task'), + ); + + fakeGoTaskService.complete( + 'continue-stopped-openclaw', + const GoTaskServiceResult( + success: true, + message: 'continued stopped', + turnId: 'turn-continued-stopped', + raw: {}, + errorMessage: '', + resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, + ), + ); + await _waitForThreadLifecycleStatus( + controller, + 'continue-stopped-openclaw', + 'ready', + ); + }, + ); + test( 'stale queued lifecycle without a real queue entry is not pending', () { @@ -3045,6 +3181,108 @@ void main() { }, ); + test( + 'continueAssistantTaskInternal resumes interrupted task without duplicating user turn', + () async { + final fakeGoTaskService = _BlockingGoTaskServiceClient(); + final controller = _connectedController(fakeGoTaskService); + addTearDown(() { + fakeGoTaskService.completeAll(); + controller.dispose(); + }); + + await controller.switchSession('continue-interrupted-task'); + controller.appendLocalSessionMessageInternal( + 'continue-interrupted-task', + GatewayChatMessage( + id: 'user-continue-interrupted', + role: 'user', + text: 'previous interrupted request', + timestampMs: DateTime.now().millisecondsSinceEpoch.toDouble(), + toolCallId: null, + toolName: null, + stopReason: null, + pending: false, + error: false, + ), + persistInThreadContext: true, + ); + controller.upsertTaskThreadInternal( + 'continue-interrupted-task', + lifecycleStatus: 'interrupted', + lastResultCode: 'ACP_HTTP_CONNECTION_CLOSED', + lastArtifactSyncStatus: 'interrupted', + ); + + final continueFuture = controller.continueAssistantTaskInternal( + 'continue-interrupted-task', + ); + await fakeGoTaskService.waitForRequestCount(1); + + final request = fakeGoTaskService.requests.single; + expect(request.sessionId, 'continue-interrupted-task'); + expect(request.resumeSession, isTrue); + expect( + request.prompt, + contains('User request:\nprevious interrupted request'), + ); + expect( + controller.localSessionMessagesInternal['continue-interrupted-task']! + .where( + (message) => + message.role == 'user' && + message.text == 'previous interrupted request', + ) + .length, + 1, + ); + + fakeGoTaskService.complete( + 'continue-interrupted-task', + const GoTaskServiceResult( + success: true, + message: 'continued interrupted', + turnId: 'turn-continued-interrupted', + raw: {}, + errorMessage: '', + resolvedModel: '', + route: GoTaskServiceRoute.externalAcpSingle, + ), + ); + await continueFuture; + }, + ); + + test( + 'continueAssistantTaskInternal fails locally without a committed user turn', + () async { + final fakeGoTaskService = _RecordingGoTaskServiceClient(); + final controller = _connectedController(fakeGoTaskService); + addTearDown(controller.dispose); + + await controller.switchSession('continue-empty-task'); + controller.upsertTaskThreadInternal( + 'continue-empty-task', + lifecycleStatus: 'ready', + lastResultCode: 'aborted', + ); + + await expectLater( + controller.continueAssistantTaskInternal('continue-empty-task'), + throwsA(isA()), + ); + + expect(fakeGoTaskService.requests, isEmpty); + expect( + controller + .assistantThreadMessagesInternal['continue-empty-task']! + .last + .text, + contains('没有可恢复的用户请求'), + ); + }, + ); + test('sendChatMessage resumes after confirmed session activity', () async { final fakeGoTaskService = _RecordingGoTaskServiceClient() ..outcomes.add(