fix: stabilize ACP closed flow handling
This commit is contained in:
parent
f657272cea
commit
5fa7c86f9c
@ -95,19 +95,28 @@ extension AppControllerDesktopThreadBinding on AppController {
|
||||
final normalizedSessionKey = normalizedAssistantSessionKeyInternal(
|
||||
sessionKey,
|
||||
);
|
||||
final baseWorkspace = settings.workspacePath.trim().isNotEmpty
|
||||
? settings.workspacePath.trim()
|
||||
: resolvedUserHomeDirectoryInternal.trim();
|
||||
if (baseWorkspace.isEmpty) {
|
||||
final homeDirectory = resolvedUserHomeDirectoryInternal.trim();
|
||||
if (homeDirectory.isEmpty) {
|
||||
return '';
|
||||
}
|
||||
final threadWorkspace =
|
||||
'${trimTrailingPathSeparatorInternal(baseWorkspace)}/.xworkmate/threads/${threadWorkspaceDirectoryNameInternal(normalizedSessionKey)}';
|
||||
'${trimTrailingPathSeparatorInternal(homeDirectory)}/.xworkmate/threads/${threadWorkspaceDirectoryNameInternal(normalizedSessionKey)}';
|
||||
return ensureLocalWorkspaceDirectoryInternal(threadWorkspace)
|
||||
? threadWorkspace
|
||||
: '';
|
||||
}
|
||||
|
||||
String localThreadWorkspaceDisplayPathInternal(String sessionKey) {
|
||||
final homeDirectory = resolvedUserHomeDirectoryInternal.trim();
|
||||
if (homeDirectory.isEmpty) {
|
||||
return '';
|
||||
}
|
||||
final normalizedSessionKey = normalizedAssistantSessionKeyInternal(
|
||||
sessionKey,
|
||||
);
|
||||
return '\$HOME/.xworkmate/threads/${threadWorkspaceDirectoryNameInternal(normalizedSessionKey)}';
|
||||
}
|
||||
|
||||
String remoteThreadWorkspacePathInternal(
|
||||
String sessionKey,
|
||||
ThreadOwnerScope ownerScope,
|
||||
@ -190,11 +199,14 @@ extension AppControllerDesktopThreadBinding on AppController {
|
||||
WorkspaceBinding? existingBinding,
|
||||
}) {
|
||||
final localPath = localThreadWorkspacePathInternal(sessionKey);
|
||||
final displayPath = localPath.isEmpty
|
||||
? ''
|
||||
: localThreadWorkspaceDisplayPathInternal(sessionKey);
|
||||
return WorkspaceBinding(
|
||||
workspaceId: normalizedAssistantSessionKeyInternal(sessionKey),
|
||||
workspaceKind: WorkspaceKind.localFs,
|
||||
workspacePath: localPath,
|
||||
displayPath: localPath,
|
||||
displayPath: displayPath,
|
||||
writable: existingBinding?.writable ?? true,
|
||||
);
|
||||
}
|
||||
@ -277,11 +289,13 @@ extension AppControllerDesktopThreadBinding on AppController {
|
||||
normalizedSessionKey,
|
||||
ownerScope: ownerScope,
|
||||
workspaceBinding: workspaceBinding,
|
||||
lastRemoteWorkingDirectory: remoteThreadWorkspacePathInternal(
|
||||
normalizedSessionKey,
|
||||
ownerScope,
|
||||
),
|
||||
lastRemoteWorkspaceRefKind: WorkspaceRefKind.remotePath,
|
||||
lastRemoteWorkingDirectory:
|
||||
snapshot.record?.lastRemoteWorkingDirectory?.trim().isNotEmpty == true
|
||||
? snapshot.record?.lastRemoteWorkingDirectory
|
||||
: remoteThreadWorkspacePathInternal(normalizedSessionKey, ownerScope),
|
||||
lastRemoteWorkspaceRefKind:
|
||||
snapshot.record?.lastRemoteWorkspaceRefKind ??
|
||||
WorkspaceRefKind.remotePath,
|
||||
executionBinding: buildDesktopExecutionBindingInternal(
|
||||
executionTarget: snapshot.executionTarget,
|
||||
existingBinding: snapshot.record?.executionBinding,
|
||||
|
||||
@ -863,7 +863,7 @@ class GatewayAcpClient {
|
||||
required String requestId,
|
||||
required void Function(Map<String, dynamic>) onNotification,
|
||||
}) async {
|
||||
final completer = Completer<Map<String, dynamic>>();
|
||||
Map<String, dynamic>? resolvedResponse;
|
||||
final eventLines = <String>[];
|
||||
|
||||
void consumeEventPayload(String payload) {
|
||||
@ -874,9 +874,7 @@ class GatewayAcpClient {
|
||||
final json = _decodeMap(trimmed);
|
||||
if (stringValue(json['id']) == requestId &&
|
||||
(json.containsKey('result') || json.containsKey('error'))) {
|
||||
if (!completer.isCompleted) {
|
||||
completer.complete(json);
|
||||
}
|
||||
resolvedResponse ??= json;
|
||||
return;
|
||||
}
|
||||
if ((stringValue(json['method']) ?? '').isNotEmpty) {
|
||||
@ -890,6 +888,9 @@ class GatewayAcpClient {
|
||||
if (eventLines.isNotEmpty) {
|
||||
consumeEventPayload(eventLines.join('\n'));
|
||||
eventLines.clear();
|
||||
if (resolvedResponse != null) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
@ -898,16 +899,16 @@ class GatewayAcpClient {
|
||||
}
|
||||
}
|
||||
|
||||
if (eventLines.isNotEmpty) {
|
||||
if (eventLines.isNotEmpty && resolvedResponse == null) {
|
||||
consumeEventPayload(eventLines.join('\n'));
|
||||
}
|
||||
if (!completer.isCompleted) {
|
||||
final resolved = resolvedResponse;
|
||||
if (resolved == null) {
|
||||
throw GatewayAcpException(
|
||||
'ACP SSE ended without JSON-RPC response for request $requestId',
|
||||
code: 'ACP_SSE_NO_RESULT',
|
||||
);
|
||||
}
|
||||
final resolved = await completer.future;
|
||||
_throwIfJsonRpcError(resolved);
|
||||
return resolved;
|
||||
}
|
||||
|
||||
@ -4,10 +4,82 @@ import 'package:crypto/crypto.dart' as crypto;
|
||||
import 'package:flutter_test/flutter_test.dart';
|
||||
import 'package:xworkmate/app/app_controller.dart';
|
||||
import 'package:xworkmate/app/app_controller_desktop_runtime_coordination_impl.dart';
|
||||
import 'package:xworkmate/app/app_controller_desktop_thread_binding.dart';
|
||||
import 'package:xworkmate/runtime/go_task_service_client.dart';
|
||||
import 'package:xworkmate/runtime/runtime_models.dart';
|
||||
|
||||
void main() {
|
||||
test(
|
||||
'converges managed local thread workspaces to the user home root',
|
||||
() async {
|
||||
final controller = AppController(
|
||||
environmentOverride: const <String, String>{},
|
||||
);
|
||||
addTearDown(controller.dispose);
|
||||
|
||||
final home = await Directory.systemTemp.createTemp(
|
||||
'xworkmate-home-thread-root-',
|
||||
);
|
||||
final oldWorkspace = await Directory.systemTemp.createTemp(
|
||||
'xworkmate-app-worktree-thread-root-',
|
||||
);
|
||||
addTearDown(() async {
|
||||
if (await home.exists()) {
|
||||
await home.delete(recursive: true);
|
||||
}
|
||||
if (await oldWorkspace.exists()) {
|
||||
await oldWorkspace.delete(recursive: true);
|
||||
}
|
||||
});
|
||||
controller.resolvedUserHomeDirectoryInternal = home.path;
|
||||
|
||||
const sessionKey = 'draft-1778207741322';
|
||||
final oldThreadWorkspace = Directory(
|
||||
'${oldWorkspace.path}/.xworkmate/threads/$sessionKey',
|
||||
);
|
||||
await oldThreadWorkspace.create(recursive: true);
|
||||
|
||||
controller.upsertTaskThreadInternal(
|
||||
sessionKey,
|
||||
workspaceBinding: WorkspaceBinding(
|
||||
workspaceId: sessionKey,
|
||||
workspaceKind: WorkspaceKind.localFs,
|
||||
workspacePath: oldThreadWorkspace.path,
|
||||
displayPath: oldThreadWorkspace.path,
|
||||
writable: true,
|
||||
),
|
||||
messages: const <GatewayChatMessage>[
|
||||
GatewayChatMessage(
|
||||
id: 'assistant-1',
|
||||
role: 'assistant',
|
||||
text: 'kept message',
|
||||
timestampMs: 1,
|
||||
toolCallId: null,
|
||||
toolName: null,
|
||||
stopReason: null,
|
||||
pending: false,
|
||||
error: false,
|
||||
),
|
||||
],
|
||||
lastRemoteWorkingDirectory: '/remote/thread/workspace',
|
||||
lastRemoteWorkspaceRefKind: WorkspaceRefKind.remotePath,
|
||||
);
|
||||
|
||||
await controller.ensureDesktopTaskThreadBindingInternal(sessionKey);
|
||||
|
||||
final expectedWorkspace = '${home.path}/.xworkmate/threads/$sessionKey';
|
||||
final thread = controller.requireTaskThreadForSessionInternal(sessionKey);
|
||||
expect(thread.workspaceBinding.workspacePath, expectedWorkspace);
|
||||
expect(
|
||||
thread.workspaceBinding.displayPath,
|
||||
'\$HOME/.xworkmate/threads/$sessionKey',
|
||||
);
|
||||
expect(Directory(expectedWorkspace).existsSync(), isTrue);
|
||||
expect(thread.lastRemoteWorkingDirectory, '/remote/thread/workspace');
|
||||
expect(thread.messages.single.text, 'kept message');
|
||||
},
|
||||
);
|
||||
|
||||
test(
|
||||
'keeps local workspace binding separate from remote execution workspace',
|
||||
() {
|
||||
|
||||
@ -249,6 +249,67 @@ void main() {
|
||||
expect((response['result'] as Map)['ok'], true);
|
||||
});
|
||||
|
||||
test(
|
||||
'returns SSE final response before a truncated chunked close is reported',
|
||||
() async {
|
||||
final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0);
|
||||
addTearDown(() => server.close());
|
||||
server.listen((socket) async {
|
||||
final requestBytes = <int>[];
|
||||
var headerEnd = -1;
|
||||
await for (final chunk in socket) {
|
||||
requestBytes.addAll(chunk);
|
||||
final raw = utf8.decode(requestBytes, allowMalformed: true);
|
||||
headerEnd = raw.indexOf('\r\n\r\n');
|
||||
if (headerEnd < 0) {
|
||||
continue;
|
||||
}
|
||||
if (raw.contains('"id"') && raw.contains('"method"')) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
final rawRequest = utf8.decode(requestBytes, allowMalformed: true);
|
||||
final id =
|
||||
RegExp(
|
||||
r'"id"\s*:\s*"([^"]+)"',
|
||||
).firstMatch(rawRequest)?.group(1) ??
|
||||
'request-id';
|
||||
final event = utf8.encode(
|
||||
'data: ${jsonEncode(<String, dynamic>{
|
||||
'jsonrpc': '2.0',
|
||||
'id': id,
|
||||
'result': <String, dynamic>{'ok': true},
|
||||
})}\n\n',
|
||||
);
|
||||
socket
|
||||
..add(
|
||||
ascii.encode(
|
||||
'HTTP/1.1 200 OK\r\n'
|
||||
'Content-Type: text/event-stream\r\n'
|
||||
'Transfer-Encoding: chunked\r\n'
|
||||
'Connection: keep-alive\r\n'
|
||||
'\r\n'
|
||||
'${event.length.toRadixString(16)}\r\n',
|
||||
),
|
||||
)
|
||||
..add(event)
|
||||
..add(ascii.encode('\r\n'));
|
||||
await socket.flush();
|
||||
socket.destroy();
|
||||
});
|
||||
|
||||
final endpoint = Uri.parse('http://127.0.0.1:${server.port}');
|
||||
final client = GatewayAcpClient(endpointResolver: () => endpoint);
|
||||
|
||||
final response = await client.request(
|
||||
method: 'acp.capabilities',
|
||||
params: const <String, dynamic>{},
|
||||
);
|
||||
|
||||
expect((response['result'] as Map)['ok'], true);
|
||||
},
|
||||
);
|
||||
|
||||
test(
|
||||
'normalizes raw authorization override into bearer header for HTTP',
|
||||
() async {
|
||||
|
||||
Loading…
Reference in New Issue
Block a user