fix: unblock OpenClaw gateway task queue

This commit is contained in:
Haitao Pan 2026-05-28 17:05:07 +08:00
parent 8469537060
commit 191ddc6ca4
6 changed files with 318 additions and 151 deletions

View File

@ -203,8 +203,12 @@ class AppController extends ChangeNotifier {
for (final turn in openClawGatewayQueuedTurnsInternal) { for (final turn in openClawGatewayQueuedTurnsInternal) {
turn.cancelled = true; turn.cancelled = true;
} }
for (final turn in openClawGatewayActiveTurnsInternal.values) {
turn.cancelled = true;
}
openClawGatewayQueuedTurnsInternal.clear(); openClawGatewayQueuedTurnsInternal.clear();
openClawGatewayQueuedTurnsBySessionInternal.clear(); openClawGatewayQueuedTurnsBySessionInternal.clear();
openClawGatewayActiveTurnsInternal.clear();
runtimeEventsSubscriptionInternal?.cancel(); runtimeEventsSubscriptionInternal?.cancel();
detachChildListenersInternal(); detachChildListenersInternal();
runtimeCoordinatorInternal.dispose(); runtimeCoordinatorInternal.dispose();
@ -282,7 +286,11 @@ class AppController extends ChangeNotifier {
final Map<String, List<OpenClawGatewayQueuedTurnInternal>> final Map<String, List<OpenClawGatewayQueuedTurnInternal>>
openClawGatewayQueuedTurnsBySessionInternal = openClawGatewayQueuedTurnsBySessionInternal =
<String, List<OpenClawGatewayQueuedTurnInternal>>{}; <String, List<OpenClawGatewayQueuedTurnInternal>>{};
int openClawGatewayActiveTasksInternal = 0; final Map<String, OpenClawGatewayQueuedTurnInternal>
openClawGatewayActiveTurnsInternal =
<String, OpenClawGatewayQueuedTurnInternal>{};
int get openClawGatewayActiveTasksInternal =>
openClawGatewayActiveTurnsInternal.length;
bool multiAgentRunPendingInternal = false; bool multiAgentRunPendingInternal = false;
int localMessageCounterInternal = 0; int localMessageCounterInternal = 0;
int assistantDraftSessionCounterInternal = 0; int assistantDraftSessionCounterInternal = 0;

View File

@ -298,9 +298,12 @@ extension AppControllerDesktopSettings on AppController {
for (final turn in openClawGatewayQueuedTurnsInternal) { for (final turn in openClawGatewayQueuedTurnsInternal) {
turn.cancelled = true; turn.cancelled = true;
} }
for (final turn in openClawGatewayActiveTurnsInternal.values) {
turn.cancelled = true;
}
openClawGatewayQueuedTurnsInternal.clear(); openClawGatewayQueuedTurnsInternal.clear();
openClawGatewayQueuedTurnsBySessionInternal.clear(); openClawGatewayQueuedTurnsBySessionInternal.clear();
openClawGatewayActiveTasksInternal = 0; openClawGatewayActiveTurnsInternal.clear();
multiAgentRunPendingInternal = false; multiAgentRunPendingInternal = false;
final sessionKey = createAssistantDraftSessionKeyInternal(); final sessionKey = createAssistantDraftSessionKeyInternal();
initializeAssistantThreadContext( initializeAssistantThreadContext(

View File

@ -3,7 +3,6 @@
import 'dart:async'; import 'dart:async';
import 'dart:convert'; import 'dart:convert';
import 'dart:io'; import 'dart:io';
import 'dart:math' as math;
import 'package:flutter/material.dart'; import 'package:flutter/material.dart';
import 'app_metadata.dart'; import 'app_metadata.dart';
import 'app_capabilities.dart'; import 'app_capabilities.dart';
@ -704,7 +703,7 @@ extension AppControllerDesktopThreadActions on AppController {
if (turn.appendUserTurn) { if (turn.appendUserTurn) {
appendGatewayUserTurnInternal(turn.sessionKey, turn.message); appendGatewayUserTurnInternal(turn.sessionKey, turn.message);
} }
if (openClawGatewayActiveTasksInternal >= if (openClawGatewayActiveTurnsInternal.length >=
openClawGatewayMaxActiveTasksInternal && openClawGatewayMaxActiveTasksInternal &&
openClawGatewayQueuedTurnsInternal.length >= openClawGatewayQueuedTurnsInternal.length >=
openClawGatewayMaxQueuedTasksInternal) { openClawGatewayMaxQueuedTasksInternal) {
@ -800,6 +799,24 @@ extension AppControllerDesktopThreadActions on AppController {
return true; return true;
} }
bool removeActiveOpenClawGatewayTurnsForSessionInternal(String sessionKey) {
final normalizedSessionKey = normalizedAssistantSessionKeyInternal(
sessionKey,
);
var removed = false;
openClawGatewayActiveTurnsInternal.removeWhere((_, turn) {
final matches =
normalizedAssistantSessionKeyInternal(turn.sessionKey) ==
normalizedSessionKey;
if (matches) {
turn.cancelled = true;
removed = true;
}
return matches;
});
return removed;
}
void markOpenClawGatewayTurnAbortedInternal(String sessionKey) { void markOpenClawGatewayTurnAbortedInternal(String sessionKey) {
final normalizedSessionKey = normalizedAssistantSessionKeyInternal( final normalizedSessionKey = normalizedAssistantSessionKeyInternal(
sessionKey, sessionKey,
@ -834,7 +851,7 @@ extension AppControllerDesktopThreadActions on AppController {
} }
void drainOpenClawGatewayQueueInternal() { void drainOpenClawGatewayQueueInternal() {
while (openClawGatewayActiveTasksInternal < while (openClawGatewayActiveTurnsInternal.length <
openClawGatewayMaxActiveTasksInternal && openClawGatewayMaxActiveTasksInternal &&
openClawGatewayQueuedTurnsInternal.isNotEmpty) { openClawGatewayQueuedTurnsInternal.isNotEmpty) {
final turn = openClawGatewayQueuedTurnsInternal.removeAt(0); final turn = openClawGatewayQueuedTurnsInternal.removeAt(0);
@ -842,11 +859,29 @@ extension AppControllerDesktopThreadActions on AppController {
if (turn.cancelled) { if (turn.cancelled) {
continue; continue;
} }
openClawGatewayActiveTasksInternal += 1; openClawGatewayActiveTurnsInternal[turn.queueId] = turn;
markOpenClawGatewayRunningTurnInternal(turn.sessionKey);
unawaited(runOpenClawGatewayQueuedTurnInternal(turn)); unawaited(runOpenClawGatewayQueuedTurnInternal(turn));
} }
} }
void markOpenClawGatewayRunningTurnInternal(String sessionKey) {
final startedAtMs = DateTime.now().millisecondsSinceEpoch.toDouble();
aiGatewayPendingSessionKeysInternal.add(sessionKey);
upsertTaskThreadInternal(
sessionKey,
lifecycleStatus: 'running',
lastRunAtMs: startedAtMs,
lastResultCode: 'running',
lastArtifactSyncAtMs: startedAtMs,
lastArtifactSyncStatus: 'running',
lastTaskArtifactRelativePaths: const <String>[],
updatedAtMs: startedAtMs,
);
recomputeTasksInternal();
notifyIfActiveInternal();
}
Future<void> runOpenClawGatewayQueuedTurnInternal( Future<void> runOpenClawGatewayQueuedTurnInternal(
OpenClawGatewayQueuedTurnInternal turn, OpenClawGatewayQueuedTurnInternal turn,
) async { ) async {
@ -883,10 +918,7 @@ extension AppControllerDesktopThreadActions on AppController {
); );
} }
} finally { } finally {
openClawGatewayActiveTasksInternal = math.max( openClawGatewayActiveTurnsInternal.remove(turn.queueId);
0,
openClawGatewayActiveTasksInternal - 1,
);
if (!disposedInternal) { if (!disposedInternal) {
drainOpenClawGatewayQueueInternal(); drainOpenClawGatewayQueueInternal();
recomputeTasksInternal(); recomputeTasksInternal();
@ -1227,7 +1259,9 @@ extension AppControllerDesktopThreadActions on AppController {
// Best effort cancellation only. Local state must still leave pending. // Best effort cancellation only. Local state must still leave pending.
} }
removeQueuedOpenClawGatewayTurnsForSessionInternal(sessionKey); removeQueuedOpenClawGatewayTurnsForSessionInternal(sessionKey);
removeActiveOpenClawGatewayTurnsForSessionInternal(sessionKey);
markOpenClawGatewayTurnAbortedInternal(sessionKey); markOpenClawGatewayTurnAbortedInternal(sessionKey);
drainOpenClawGatewayQueueInternal();
return; return;
} }
} }

View File

@ -502,6 +502,15 @@ extension AppControllerDesktopWorkspaceExecution on AppController {
normalizedAssistantSessionKeyInternal(turn.sessionKey) == normalizedAssistantSessionKeyInternal(turn.sessionKey) ==
normalizedSessionKey, normalizedSessionKey,
); );
openClawGatewayActiveTurnsInternal.removeWhere((_, turn) {
final matches =
normalizedAssistantSessionKeyInternal(turn.sessionKey) ==
normalizedSessionKey;
if (matches) {
turn.cancelled = true;
}
return matches;
});
recomputeTasksInternal(); recomputeTasksInternal();
if (wasCurrent) { if (wasCurrent) {
await ensureActiveAssistantThreadInternal(); await ensureActiveAssistantThreadInternal();

View File

@ -1,7 +1,7 @@
import '../runtime/go_task_service_client.dart'; import '../runtime/go_task_service_client.dart';
import '../runtime/runtime_models.dart'; import '../runtime/runtime_models.dart';
const int openClawGatewayMaxActiveTasksInternal = 1; const int openClawGatewayMaxActiveTasksInternal = 5;
const int openClawGatewayMaxQueuedTasksInternal = 20; const int openClawGatewayMaxQueuedTasksInternal = 20;
class OpenClawGatewayQueuedTurnInternal { class OpenClawGatewayQueuedTurnInternal {

View File

@ -2500,88 +2500,78 @@ void main() {
controller.dispose(); controller.dispose();
}); });
await _selectGatewaySession(controller, 'queue-task-a');
final taskAFuture = controller.sendChatMessage('same prompt');
await fakeGoTaskService.waitForRequestCount(1);
await expectLater(
taskAFuture.timeout(const Duration(milliseconds: 250)),
completes,
);
await _selectGatewaySession(controller, 'queue-task-b');
final queuedAttachment = GatewayChatAttachmentPayload( final queuedAttachment = GatewayChatAttachmentPayload(
type: 'file', type: 'file',
mimeType: 'text/plain', mimeType: 'text/plain',
fileName: 'queued.txt', fileName: 'queued.txt',
content: base64Encode(utf8.encode('queued content')), content: base64Encode(utf8.encode('queued content')),
); );
final taskBFuture = controller.sendChatMessage( for (
'same prompt', var index = 0;
attachments: <GatewayChatAttachmentPayload>[queuedAttachment], index < openClawGatewayMaxActiveTasksInternal;
); index += 1
await _waitForThreadLifecycleStatus( ) {
controller, final sessionKey = 'queue-task-$index';
'queue-task-b', await _selectGatewaySession(controller, sessionKey);
'queued', await expectLater(
); controller
await _selectGatewaySession(controller, 'queue-task-c'); .sendChatMessage('active prompt $index')
final taskCFuture = controller.sendChatMessage('different prompt'); .timeout(const Duration(milliseconds: 250)),
await _waitForThreadLifecycleStatus( completes,
controller, );
'queue-task-c', await fakeGoTaskService.waitForRequestCount(index + 1);
'queued', expect(
); controller
await expectLater( .requireTaskThreadForSessionInternal(sessionKey)
taskBFuture.timeout(const Duration(milliseconds: 250)), .lifecycleState
completes, .status,
); 'running',
await expectLater( );
taskCFuture.timeout(const Duration(milliseconds: 250)), }
completes,
);
expect(fakeGoTaskService.requests, hasLength(1)); await _selectGatewaySession(controller, 'queue-task-waiting');
await expectLater(
controller
.sendChatMessage(
'queued prompt',
attachments: <GatewayChatAttachmentPayload>[queuedAttachment],
)
.timeout(const Duration(milliseconds: 250)),
completes,
);
await _waitForThreadLifecycleStatus(
controller,
'queue-task-waiting',
'queued',
);
expect(
fakeGoTaskService.requests,
hasLength(openClawGatewayMaxActiveTasksInternal),
);
expect( expect(
controller controller
.requireTaskThreadForSessionInternal('queue-task-b') .requireTaskThreadForSessionInternal('queue-task-waiting')
.lifecycleState .lifecycleState
.status, .status,
'queued', 'queued',
); );
expect( expect(
controller controller.assistantSessionHasPendingRun('queue-task-waiting'),
.requireTaskThreadForSessionInternal('queue-task-c')
.lifecycleState
.status,
'queued',
);
expect(
controller.assistantSessionHasPendingRun('queue-task-b'),
isTrue, isTrue,
); );
expect( expect(
controller.assistantSessionHasPendingRun('queue-task-c'), controller.localSessionMessagesInternal['queue-task-waiting']!.map(
isTrue,
);
expect(
controller.localSessionMessagesInternal['queue-task-b']!.map(
(message) => message.text, (message) => message.text,
), ),
contains('same prompt'), contains('queued prompt'),
);
expect(
controller.localSessionMessagesInternal['queue-task-c']!.map(
(message) => message.text,
),
contains('different prompt'),
); );
fakeGoTaskService.complete( fakeGoTaskService.complete(
'queue-task-a', 'queue-task-0',
const GoTaskServiceResult( const GoTaskServiceResult(
success: true, success: true,
message: 'result A', message: 'result 0',
turnId: 'turn-a', turnId: 'turn-0',
raw: <String, dynamic>{}, raw: <String, dynamic>{},
errorMessage: '', errorMessage: '',
resolvedModel: '', resolvedModel: '',
@ -2590,36 +2580,41 @@ void main() {
); );
await _waitForThreadLifecycleStatus( await _waitForThreadLifecycleStatus(
controller, controller,
'queue-task-a', 'queue-task-0',
'ready', 'ready',
); );
await fakeGoTaskService.waitForRequestCount(2); await fakeGoTaskService.waitForRequestCount(
openClawGatewayMaxActiveTasksInternal + 1,
);
final taskBRequest = fakeGoTaskService.requests[1]; final queuedRequest = fakeGoTaskService.requests.last;
expect(taskBRequest.sessionId, 'queue-task-b'); expect(queuedRequest.sessionId, 'queue-task-waiting');
expect(taskBRequest.prompt, contains('TaskThread workspace context:')); expect(queuedRequest.prompt, contains('TaskThread workspace context:'));
expect(taskBRequest.prompt, contains('- sessionKey: queue-task-b'));
expect(taskBRequest.prompt, contains('User request:\nsame prompt'));
expect(taskBRequest.resumeSession, isFalse);
expect(taskBRequest.inlineAttachments, hasLength(1));
expect(taskBRequest.inlineAttachments.single.fileName, 'queued.txt');
expect( expect(
taskBRequest.inlineAttachments.single.content, queuedRequest.prompt,
contains('- sessionKey: queue-task-waiting'),
);
expect(queuedRequest.prompt, contains('User request:\nqueued prompt'));
expect(queuedRequest.resumeSession, isFalse);
expect(queuedRequest.inlineAttachments, hasLength(1));
expect(queuedRequest.inlineAttachments.single.fileName, 'queued.txt');
expect(
queuedRequest.inlineAttachments.single.content,
queuedAttachment.content, queuedAttachment.content,
); );
expect(taskBRequest.workingDirectory, endsWith('/queue-task-b')); expect(queuedRequest.workingDirectory, endsWith('/queue-task-waiting'));
expect(taskBRequest.prompt, contains(taskBRequest.workingDirectory)); expect(queuedRequest.prompt, contains(queuedRequest.workingDirectory));
expect( expect(
taskBRequest.remoteWorkingDirectoryHint, queuedRequest.remoteWorkingDirectoryHint,
endsWith('/threads/queue-task-b'), endsWith('/threads/queue-task-waiting'),
); );
fakeGoTaskService.complete( fakeGoTaskService.complete(
'queue-task-b', 'queue-task-waiting',
const GoTaskServiceResult( const GoTaskServiceResult(
success: true, success: true,
message: 'result B', message: 'queued done',
turnId: 'turn-b', turnId: 'turn-waiting',
raw: <String, dynamic>{}, raw: <String, dynamic>{},
errorMessage: '', errorMessage: '',
resolvedModel: '', resolvedModel: '',
@ -2628,45 +2623,18 @@ void main() {
); );
await _waitForThreadLifecycleStatus( await _waitForThreadLifecycleStatus(
controller, controller,
'queue-task-b', 'queue-task-waiting',
'ready', 'ready',
); );
expect( expect(
controller.localSessionMessagesInternal['queue-task-b']! controller.localSessionMessagesInternal['queue-task-waiting']!
.where( .where(
(message) => (message) =>
message.role == 'user' && message.text == 'same prompt', message.role == 'user' && message.text == 'queued prompt',
) )
.length, .length,
1, 1,
); );
await fakeGoTaskService.waitForRequestCount(3);
final taskCRequest = fakeGoTaskService.requests[2];
expect(taskCRequest.sessionId, 'queue-task-c');
expect(
taskCRequest.prompt,
contains('User request:\ndifferent prompt'),
);
expect(taskCRequest.workingDirectory, endsWith('/queue-task-c'));
expect(taskCRequest.prompt, contains(taskCRequest.workingDirectory));
fakeGoTaskService.complete(
'queue-task-c',
const GoTaskServiceResult(
success: true,
message: 'result C',
turnId: 'turn-c',
raw: <String, dynamic>{},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
),
);
await _waitForThreadLifecycleStatus(
controller,
'queue-task-c',
'ready',
);
}, },
); );
@ -2732,9 +2700,11 @@ void main() {
controller.dispose(); controller.dispose();
}); });
await _selectGatewaySession(controller, 'running-openclaw-task'); final activeSessionKeys = await _startOpenClawActiveTasks(
final runningFuture = controller.sendChatMessage('running'); controller,
await fakeGoTaskService.waitForRequestCount(1); fakeGoTaskService,
prefix: 'running-openclaw-task',
);
await _selectGatewaySession(controller, 'queued-openclaw-task'); await _selectGatewaySession(controller, 'queued-openclaw-task');
final queuedFuture = controller.sendChatMessage('queued'); final queuedFuture = controller.sendChatMessage('queued');
@ -2743,7 +2713,10 @@ void main() {
'queued-openclaw-task', 'queued-openclaw-task',
'queued', 'queued',
); );
expect(fakeGoTaskService.requests, hasLength(1)); expect(
fakeGoTaskService.requests,
hasLength(openClawGatewayMaxActiveTasksInternal),
);
await controller.abortRun(); await controller.abortRun();
@ -2758,7 +2731,7 @@ void main() {
await queuedFuture; await queuedFuture;
fakeGoTaskService.complete( fakeGoTaskService.complete(
'running-openclaw-task', activeSessionKeys.first,
const GoTaskServiceResult( const GoTaskServiceResult(
success: true, success: true,
message: 'running done', message: 'running done',
@ -2769,14 +2742,16 @@ void main() {
route: GoTaskServiceRoute.externalAcpSingle, route: GoTaskServiceRoute.externalAcpSingle,
), ),
); );
await runningFuture;
await _waitForThreadLifecycleStatus( await _waitForThreadLifecycleStatus(
controller, controller,
'running-openclaw-task', activeSessionKeys.first,
'ready', 'ready',
); );
await Future<void>.delayed(const Duration(milliseconds: 50)); await Future<void>.delayed(const Duration(milliseconds: 50));
expect(fakeGoTaskService.requests, hasLength(1)); expect(
fakeGoTaskService.requests,
hasLength(openClawGatewayMaxActiveTasksInternal),
);
}, },
); );
@ -2790,9 +2765,12 @@ void main() {
controller.dispose(); controller.dispose();
}); });
await _selectGatewaySession(controller, 'running-openclaw-stop-task'); final activeSessionKeys = await _startOpenClawActiveTasks(
final runningFuture = controller.sendChatMessage('running'); controller,
await fakeGoTaskService.waitForRequestCount(1); fakeGoTaskService,
prefix: 'running-openclaw-stop-task',
);
final runningSessionKey = activeSessionKeys.first;
await _selectGatewaySession(controller, 'queued-openclaw-after-stop'); await _selectGatewaySession(controller, 'queued-openclaw-after-stop');
final queuedFuture = controller.sendChatMessage('queued'); final queuedFuture = controller.sendChatMessage('queued');
@ -2802,25 +2780,30 @@ void main() {
'queued', 'queued',
); );
await _selectGatewaySession(controller, 'running-openclaw-stop-task'); await _selectGatewaySession(controller, runningSessionKey);
await controller.abortRun(); await controller.abortRun();
expect(fakeGoTaskService.cancelledSessionIds, <String>[ expect(fakeGoTaskService.cancelledSessionIds, <String>[
'running-openclaw-stop-task', runningSessionKey,
]); ]);
expect( expect(
controller.assistantSessionHasPendingRun( controller.assistantSessionHasPendingRun(runningSessionKey),
'running-openclaw-stop-task',
),
isFalse, isFalse,
); );
expect( expect(
controller controller
.requireTaskThreadForSessionInternal('running-openclaw-stop-task') .requireTaskThreadForSessionInternal(runningSessionKey)
.lifecycleState .lifecycleState
.lastResultCode, .lastResultCode,
'aborted', 'aborted',
); );
await fakeGoTaskService.waitForRequestCount(
openClawGatewayMaxActiveTasksInternal + 1,
);
expect(
fakeGoTaskService.requests.last.sessionId,
'queued-openclaw-after-stop',
);
expect( expect(
controller.assistantSessionHasPendingRun( controller.assistantSessionHasPendingRun(
'queued-openclaw-after-stop', 'queued-openclaw-after-stop',
@ -2829,7 +2812,7 @@ void main() {
); );
fakeGoTaskService.complete( fakeGoTaskService.complete(
'running-openclaw-stop-task', runningSessionKey,
const GoTaskServiceResult( const GoTaskServiceResult(
success: true, success: true,
message: 'late stopped result', message: 'late stopped result',
@ -2840,12 +2823,6 @@ void main() {
route: GoTaskServiceRoute.externalAcpSingle, route: GoTaskServiceRoute.externalAcpSingle,
), ),
); );
await runningFuture;
await fakeGoTaskService.waitForRequestCount(2);
expect(
fakeGoTaskService.requests.last.sessionId,
'queued-openclaw-after-stop',
);
fakeGoTaskService.complete( fakeGoTaskService.complete(
'queued-openclaw-after-stop', 'queued-openclaw-after-stop',
@ -2865,7 +2842,10 @@ void main() {
'queued-openclaw-after-stop', 'queued-openclaw-after-stop',
'ready', 'ready',
); );
expect(fakeGoTaskService.requests, hasLength(2)); expect(
fakeGoTaskService.requests,
hasLength(openClawGatewayMaxActiveTasksInternal + 1),
);
}, },
); );
@ -2879,9 +2859,11 @@ void main() {
controller.dispose(); controller.dispose();
}); });
await _selectGatewaySession(controller, 'continue-active-openclaw'); final activeSessionKeys = await _startOpenClawActiveTasks(
await controller.sendChatMessage('active task'); controller,
await fakeGoTaskService.waitForRequestCount(1); fakeGoTaskService,
prefix: 'continue-active-openclaw',
);
await _selectGatewaySession(controller, 'continue-queued-openclaw'); await _selectGatewaySession(controller, 'continue-queued-openclaw');
await controller.sendChatMessage('queued before continue'); await controller.sendChatMessage('queued before continue');
@ -2925,7 +2907,10 @@ void main() {
'continue-stopped-openclaw', 'continue-stopped-openclaw',
'queued', 'queued',
); );
expect(fakeGoTaskService.requests, hasLength(1)); expect(
fakeGoTaskService.requests,
hasLength(openClawGatewayMaxActiveTasksInternal),
);
expect( expect(
controller.assistantSessionHasPendingRun('continue-queued-openclaw'), controller.assistantSessionHasPendingRun('continue-queued-openclaw'),
isTrue, isTrue,
@ -2946,7 +2931,7 @@ void main() {
); );
fakeGoTaskService.complete( fakeGoTaskService.complete(
'continue-active-openclaw', activeSessionKeys.first,
const GoTaskServiceResult( const GoTaskServiceResult(
success: true, success: true,
message: 'active done', message: 'active done',
@ -2957,7 +2942,9 @@ void main() {
route: GoTaskServiceRoute.externalAcpSingle, route: GoTaskServiceRoute.externalAcpSingle,
), ),
); );
await fakeGoTaskService.waitForRequestCount(2); await fakeGoTaskService.waitForRequestCount(
openClawGatewayMaxActiveTasksInternal + 1,
);
expect( expect(
fakeGoTaskService.requests.last.sessionId, fakeGoTaskService.requests.last.sessionId,
'continue-queued-openclaw', 'continue-queued-openclaw',
@ -2975,7 +2962,21 @@ void main() {
route: GoTaskServiceRoute.externalAcpSingle, route: GoTaskServiceRoute.externalAcpSingle,
), ),
); );
await fakeGoTaskService.waitForRequestCount(3); fakeGoTaskService.complete(
activeSessionKeys[1],
const GoTaskServiceResult(
success: true,
message: 'second active done',
turnId: 'turn-second-active',
raw: <String, dynamic>{},
errorMessage: '',
resolvedModel: '',
route: GoTaskServiceRoute.externalAcpSingle,
),
);
await fakeGoTaskService.waitForRequestCount(
openClawGatewayMaxActiveTasksInternal + 2,
);
final continuedRequest = fakeGoTaskService.requests.last; final continuedRequest = fakeGoTaskService.requests.last;
expect(continuedRequest.sessionId, 'continue-stopped-openclaw'); expect(continuedRequest.sessionId, 'continue-stopped-openclaw');
@ -3027,13 +3028,93 @@ void main() {
}, },
); );
test(
'OpenClaw drain starts queued work when no active turns remain',
() async {
final fakeGoTaskService = _BlockingGoTaskServiceClient();
final controller = _connectedGatewayController(fakeGoTaskService);
addTearDown(() {
fakeGoTaskService.completeAll();
controller.dispose();
});
const sessionKey = 'stale-slot-queued-task';
await _selectGatewaySession(controller, sessionKey);
final turn = OpenClawGatewayQueuedTurnInternal(
queueId: 'stale-slot-queued-turn',
sessionKey: sessionKey,
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
message: 'recover queued work',
thinking: 'off',
selectedSkillLabels: const <String>[],
attachments: const <GatewayChatAttachmentPayload>[],
localAttachments: const <CollaborationAttachment>[],
workingDirectory: '/tmp/$sessionKey',
localWorkingDirectory: '/tmp/$sessionKey-local',
remoteWorkingDirectoryHint: '/threads/$sessionKey',
model: '',
routing: const ExternalCodeAgentAcpRoutingConfig.auto(
preferredGatewayTarget: kCanonicalGatewayProviderId,
),
agentId: '',
metadata: const <String, dynamic>{},
resumeSessionHint: false,
);
controller.openClawGatewayQueuedTurnsInternal.add(turn);
controller.openClawGatewayQueuedTurnsBySessionInternal[sessionKey] =
<OpenClawGatewayQueuedTurnInternal>[turn];
controller.markOpenClawGatewayQueuedTurnInternal(sessionKey);
controller.drainOpenClawGatewayQueueInternal();
await fakeGoTaskService.waitForRequestCount(1);
expect(fakeGoTaskService.requests.single.sessionId, sessionKey);
expect(
controller
.requireTaskThreadForSessionInternal(sessionKey)
.lifecycleState
.status,
'running',
);
expect(controller.openClawGatewayActiveTasksInternal, 1);
},
);
test('OpenClaw queue overflow fails without artifact sync', () async { test('OpenClaw queue overflow fails without artifact sync', () async {
final fakeGoTaskService = _BlockingGoTaskServiceClient(); final fakeGoTaskService = _BlockingGoTaskServiceClient();
final controller = _connectedGatewayController(fakeGoTaskService); final controller = _connectedGatewayController(fakeGoTaskService);
addTearDown(controller.dispose); addTearDown(controller.dispose);
controller.openClawGatewayActiveTasksInternal = for (
openClawGatewayMaxActiveTasksInternal; var index = 0;
index < openClawGatewayMaxActiveTasksInternal;
index += 1
) {
final sessionKey = 'queue-full-active-$index';
final turn = OpenClawGatewayQueuedTurnInternal(
queueId: 'queue-full-active-$index',
sessionKey: sessionKey,
target: AssistantExecutionTarget.gateway,
provider: SingleAgentProvider.openclaw,
message: 'active $index',
thinking: 'off',
selectedSkillLabels: const <String>[],
attachments: const <GatewayChatAttachmentPayload>[],
localAttachments: const <CollaborationAttachment>[],
workingDirectory: '/tmp/$sessionKey',
localWorkingDirectory: '/tmp/$sessionKey-local',
remoteWorkingDirectoryHint: '/threads/$sessionKey',
model: '',
routing: const ExternalCodeAgentAcpRoutingConfig.auto(
preferredGatewayTarget: kCanonicalGatewayProviderId,
),
agentId: '',
metadata: const <String, dynamic>{},
resumeSessionHint: false,
);
controller.openClawGatewayActiveTurnsInternal[turn.queueId] = turn;
}
for ( for (
var index = 0; var index = 0;
index < openClawGatewayMaxQueuedTasksInternal; index < openClawGatewayMaxQueuedTasksInternal;
@ -3609,6 +3690,38 @@ Future<void> _selectGatewaySession(
); );
} }
Future<List<String>> _startOpenClawActiveTasks(
AppController controller,
_BlockingGoTaskServiceClient fakeGoTaskService, {
required String prefix,
}) async {
final sessionKeys = <String>[];
for (
var index = 0;
index < openClawGatewayMaxActiveTasksInternal;
index += 1
) {
final sessionKey = '$prefix-$index';
sessionKeys.add(sessionKey);
await _selectGatewaySession(controller, sessionKey);
await expectLater(
controller
.sendChatMessage('active task $index')
.timeout(const Duration(milliseconds: 250)),
completes,
);
await fakeGoTaskService.waitForRequestCount(index + 1);
expect(
controller
.requireTaskThreadForSessionInternal(sessionKey)
.lifecycleState
.status,
'running',
);
}
return sessionKeys;
}
Future<void> _waitForThreadLifecycleStatus( Future<void> _waitForThreadLifecycleStatus(
AppController controller, AppController controller,
String sessionKey, String sessionKey,