surface agent stream audit IDs

This commit is contained in:
Hamza Ayed
2026-10-04 19:17:28 +03:00
parent 8266de07c5
commit 77cad71380
3 changed files with 100 additions and 2 deletions
@@ -567,6 +567,11 @@ class ApiRepository {
final streamed = await client
.send(request)
.timeout(const Duration(minutes: 10));
final auditId = streamed.headers['x-agent-audit-id'];
String withAuditId(String message) =>
auditId == null || auditId.isEmpty
? message
: '$message\nمعرّف تتبع الطلب: $auditId';
if (streamed.statusCode < 200 || streamed.statusCode >= 300) {
final response = await http.Response.fromStream(streamed);
_checkStatus(response);
@@ -601,9 +606,11 @@ class ApiRepository {
eventName = null;
}
}
if (streamError != null) throw Exception(streamError);
if (streamError != null) throw Exception(withAuditId(streamError));
final data = resultData;
if (data == null) throw Exception('انقطع اتصال البث قبل اكتمال الوكيل.');
if (data == null) {
throw Exception(withAuditId('انقطع اتصال البث قبل اكتمال الوكيل.'));
}
return AgentRunResult.fromApiData(data);
} finally {
if (identical(_activeRequestClient, client)) _activeRequestClient = null;
@@ -8,6 +8,96 @@ import 'package:flutter_app/core/network/api_repository.dart';
void main() {
group('ApiRepository.runAgent SSE stream', () {
test('includes the audit ID when the server closes before done', () async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
server.listen((request) async {
await utf8.decoder.bind(request).join();
request.response.headers.contentType = ContentType.json;
if (request.uri.path == '/v1/auth/local-session') {
request.response.write(
jsonEncode({
'access_token': 'test-session',
'user': {'mode': 'local'},
}),
);
} else if (request.uri.path == '/v1/agent/run/stream') {
request.response.statusCode = HttpStatus.ok;
request.response.headers.set('x-agent-audit-id', 'trace-closed-42');
request.response.headers.contentType = ContentType(
'text',
'event-stream',
charset: 'utf-8',
);
request.response.write(
'event: progress\ndata: {"message":"بدأ الوكيل"}\n\n',
);
} else {
request.response.statusCode = HttpStatus.notFound;
}
await request.response.close();
});
try {
final repository = ApiRepository(
baseUrl: 'http://${server.address.address}:${server.port}',
sessionTokenStore: EphemeralSessionTokenStore(),
);
await expectLater(
repository.runAgent('اختبار انقطاع البث'),
throwsA(
predicate(
(error) => error.toString().contains('trace-closed-42'),
),
),
);
} finally {
await server.close(force: true);
}
});
test('includes the audit ID for an SSE error event', () async {
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
server.listen((request) async {
await utf8.decoder.bind(request).join();
request.response.statusCode = HttpStatus.ok;
if (request.uri.path == '/v1/auth/local-session') {
request.response.headers.contentType = ContentType.json;
request.response.write(
jsonEncode({
'access_token': 'test-session',
'user': {'mode': 'local'},
}),
);
} else if (request.uri.path == '/v1/agent/run/stream') {
request.response.headers.set('x-agent-audit-id', 'trace-error-17');
request.response.headers.contentType = ContentType(
'text',
'event-stream',
charset: 'utf-8',
);
request.response.write(
'event: error\ndata: {"message":"تعذر إكمال المهمة"}\n\n',
);
}
await request.response.close();
});
try {
final repository = ApiRepository(
baseUrl: 'http://${server.address.address}:${server.port}',
sessionTokenStore: EphemeralSessionTokenStore(),
);
await expectLater(
repository.runAgent('اختبار خطأ البث'),
throwsA(
predicate((error) => error.toString().contains('trace-error-17')),
),
);
} finally {
await server.close(force: true);
}
});
test(
'keeps progress through heartbeat and returns the done result',
() async {