feat: Dart based drop multi threading

This commit is contained in:
Tien Do Nam
2026-07-16 03:46:08 +02:00
parent bf70a5647a
commit ba7ee985ad
6 changed files with 55 additions and 146 deletions
@@ -122,7 +122,6 @@ class IsolateHttpUploadActionResult {
}
class IsolateHttpUploadAction extends ReduxActionWithResult<IsolateController, ParentIsolateState, IsolateHttpUploadActionResult> {
final int isolateIndex;
final String? remoteSessionId;
final String remoteFileToken;
final String fileId;
@@ -132,7 +131,6 @@ class IsolateHttpUploadAction extends ReduxActionWithResult<IsolateController, P
final Device device;
IsolateHttpUploadAction({
required this.isolateIndex,
required this.remoteSessionId,
required this.remoteFileToken,
required this.fileId,
@@ -144,7 +142,10 @@ class IsolateHttpUploadAction extends ReduxActionWithResult<IsolateController, P
@override
(ParentIsolateState, IsolateHttpUploadActionResult) reduce() {
final connection = state.httpUpload[isolateIndex];
final connection = state.httpUpload;
if (connection == null) {
throw StateError('httpUpload is not initialized');
}
final task = HttpUploadTask(
remoteSessionId: remoteSessionId,
@@ -173,13 +174,11 @@ class IsolateHttpUploadAction extends ReduxActionWithResult<IsolateController, P
}
class IsolateHttpUploadFilesAction extends ReduxActionWithResult<IsolateController, ParentIsolateState, IsolateHttpUploadActionResult> {
final int isolateIndex;
final String? remoteSessionId;
final List<HttpUploadFile> files;
final Device device;
IsolateHttpUploadFilesAction({
required this.isolateIndex,
required this.remoteSessionId,
required this.files,
required this.device,
@@ -187,7 +186,10 @@ class IsolateHttpUploadFilesAction extends ReduxActionWithResult<IsolateControll
@override
(ParentIsolateState, IsolateHttpUploadActionResult) reduce() {
final connection = state.httpUpload[isolateIndex];
final connection = state.httpUpload;
if (connection == null) {
throw StateError('httpUpload is not initialized');
}
final taskId = IdProvider.instance.getNextId();
final progress = connection.sendWrappedTaskAndListenStream(
task: HttpUploadFilesTask(
@@ -209,17 +211,18 @@ class IsolateHttpUploadFilesAction extends ReduxActionWithResult<IsolateControll
}
class IsolateHttpUploadCancelAction extends ReduxAction<IsolateController, ParentIsolateState> {
final int isolateIndex;
final int taskId;
IsolateHttpUploadCancelAction({
required this.isolateIndex,
required this.taskId,
});
@override
ParentIsolateState reduce() {
final connection = state.httpUpload[isolateIndex];
final connection = state.httpUpload;
if (connection == null) {
throw StateError('httpUpload is not initialized');
}
connection.sendToIsolate(
SendToIsolateData(
@@ -12,8 +12,6 @@ import 'package:typed_isolates/typed_isolates.dart';
part 'parent_isolate_provider.mapper.dart';
const _uploadIsolateCount = 2;
/// Holds the state of the parent isolate that is visible in the main Flutter isolate.
/// The [ParentIsolateState.syncState] is synchronized with all child isolates.
/// Additionally, holds the objects to communicate with the child isolates.
@@ -22,8 +20,7 @@ class ParentIsolateState with ParentIsolateStateMappable {
final SyncState syncState;
final IsolateConnector<IsolateTaskStreamResult<Device>, SendToIsolateData<IsolateTask<HttpScanTask>>>? httpScanDiscovery;
final IsolateConnector<Device, SendToIsolateData<MulticastTask>>? multicastDiscovery;
final List<IsolateConnector<IsolateTaskStreamResult<double>, SendToIsolateData<IsolateTask<BaseHttpUploadTask>>>> httpUpload;
int get uploadIsolateCount => httpUpload.length;
final IsolateConnector<IsolateTaskStreamResult<double>, SendToIsolateData<IsolateTask<BaseHttpUploadTask>>>? httpUpload;
ParentIsolateState({
required this.syncState,
@@ -36,7 +33,7 @@ class ParentIsolateState with ParentIsolateStateMappable {
syncState: syncState,
httpScanDiscovery: null,
multicastDiscovery: null,
httpUpload: [],
httpUpload: null,
);
@override
@@ -89,39 +86,31 @@ class IsolateSetupAction extends AsyncReduxAction<IsolateController, ParentIsola
),
);
final httpUploadIsolates = List.generate(
_uploadIsolateCount,
(index) async {
final httpUpload =
await TypedIsolates.startIsolate<IsolateTaskStreamResult<double>, SendToIsolateData<IsolateTask<BaseHttpUploadTask>>, InitialData>(
task: setupHttpUploadIsolate,
param: InitialData(
syncState: state.syncState,
logLevel: Logger.root.level,
),
);
final httpUpload =
await TypedIsolates.startIsolate<IsolateTaskStreamResult<double>, SendToIsolateData<IsolateTask<BaseHttpUploadTask>>, InitialData>(
task: setupHttpUploadIsolate,
param: InitialData(
syncState: state.syncState,
logLevel: Logger.root.level,
),
);
if (uriContentStreamResolver != null) {
httpUpload.sendToIsolate(
SendToIsolateData(
syncState: null,
data: IsolateTask(
id: -1,
data: HttpUploadSetContentStreamResolverTask(resolver: uriContentStreamResolver!),
),
),
);
}
return httpUpload;
},
growable: false,
);
if (uriContentStreamResolver != null) {
httpUpload.sendToIsolate(
SendToIsolateData(
syncState: null,
data: IsolateTask(
id: -1,
data: HttpUploadSetContentStreamResolverTask(resolver: uriContentStreamResolver!),
),
),
);
}
return state.copyWith(
httpScanDiscovery: httpScanDiscovery,
multicastDiscovery: multicastDiscovery,
httpUpload: await Future.wait(httpUploadIsolates),
httpUpload: httpUpload,
);
}
}
@@ -131,9 +120,7 @@ class IsolateDisposeAction extends ReduxAction<IsolateController, ParentIsolateS
ParentIsolateState reduce() {
state.httpScanDiscovery?.isolate.kill();
state.multicastDiscovery?.isolate.kill();
for (final httpUpload in state.httpUpload) {
httpUpload.isolate.kill();
}
state.httpUpload?.isolate.kill();
return state;
}
}
@@ -49,20 +49,16 @@ class ParentIsolateStateMapper extends ClassMapperBase<ParentIsolateState> {
IsolateConnector<Device, SendToIsolateData<MulticastTask>>
>
_f$multicastDiscovery = Field('multicastDiscovery', _$multicastDiscovery);
static List<
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
>
static IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>?
_$httpUpload(ParentIsolateState v) => v.httpUpload;
static const Field<
ParentIsolateState,
List<
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
>
_f$httpUpload = Field('httpUpload', _$httpUpload);
@@ -156,25 +152,6 @@ abstract class ParentIsolateStateCopyWith<
>
implements ClassCopyWith<$R, $In, $Out> {
SyncStateCopyWith<$R, SyncState, SyncState> get syncState;
ListCopyWith<
$R,
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>,
ObjectCopyWith<
$R,
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>,
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
>
>
get httpUpload;
$R call({
SyncState? syncState,
IsolateConnector<
@@ -184,11 +161,9 @@ abstract class ParentIsolateStateCopyWith<
httpScanDiscovery,
IsolateConnector<Device, SendToIsolateData<MulticastTask>>?
multicastDiscovery,
List<
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>?
httpUpload,
});
@@ -209,47 +184,17 @@ class _ParentIsolateStateCopyWithImpl<$R, $Out>
SyncStateCopyWith<$R, SyncState, SyncState> get syncState =>
$value.syncState.copyWith.$chain((v) => call(syncState: v));
@override
ListCopyWith<
$R,
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>,
ObjectCopyWith<
$R,
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>,
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
>
>
get httpUpload => ListCopyWith(
$value.httpUpload,
(v, t) => ObjectCopyWith(v, $identity, t),
(v) => call(httpUpload: v),
);
@override
$R call({
SyncState? syncState,
Object? httpScanDiscovery = $none,
Object? multicastDiscovery = $none,
List<
IsolateConnector<
IsolateTaskStreamResult<double>,
SendToIsolateData<IsolateTask<BaseHttpUploadTask>>
>
>?
httpUpload,
Object? httpUpload = $none,
}) => $apply(
FieldCopyWithData({
if (syncState != null) #syncState: syncState,
if (httpScanDiscovery != $none) #httpScanDiscovery: httpScanDiscovery,
if (multicastDiscovery != $none) #multicastDiscovery: multicastDiscovery,
if (httpUpload != null) #httpUpload: httpUpload,
if (httpUpload != $none) #httpUpload: httpUpload,
}),
);
@override
@@ -50,11 +50,9 @@ class SendSessionState with SendSessionStateMappable implements SessionState {
}
class SendingTask {
final int isolateIndex;
final int taskId;
SendingTask({
required this.isolateIndex,
required this.taskId,
});
}
-1
View File
@@ -418,7 +418,6 @@ class _ProgressPageState extends State<ProgressPage> with Refena {
.notifier(sendProvider)
.sendFile(
sessionId: widget.sessionId,
isolateIndex: 0,
file: sendSession.files[file.id]!,
isRetry: true,
);
+8 -31
View File
@@ -1,5 +1,4 @@
import 'dart:async';
import 'dart:collection';
import 'dart:convert';
import 'package:flutter/material.dart';
@@ -316,31 +315,13 @@ class SendNotifier extends Notifier<Map<String, SendSessionState>> {
state: (s) => s?.copyWith(startTime: DateTime.now().millisecondsSinceEpoch),
);
final queue = Queue<SendingFile>()..addAll(files.values);
final concurrency = ref.read(parentIsolateProvider).uploadIsolateCount;
_logger.info('Sending files using $concurrency concurrent isolates');
final futures = List.generate(concurrency, (index) async {
while (true) {
final file = switch (queue.isEmpty) {
true => null,
false => queue.removeFirst(),
};
if (file == null) {
break;
}
await sendFile(
sessionId: sessionId,
isolateIndex: index,
file: file,
isRetry: false,
);
}
});
await Future.wait(futures);
for (final file in files.values) {
await sendFile(
sessionId: sessionId,
file: file,
isRetry: false,
);
}
_finish(sessionId: sessionId);
}
@@ -384,7 +365,6 @@ class SendNotifier extends Notifier<Map<String, SendSessionState>> {
/// Returns true, if the next file should be sent.
Future<bool> sendFile({
required String sessionId,
required int isolateIndex,
required SendingFile file,
required bool isRetry,
}) async {
@@ -430,7 +410,6 @@ class SendNotifier extends Notifier<Map<String, SendSessionState>> {
.redux(parentIsolateProvider)
.dispatchTakeResult(
IsolateHttpUploadAction(
isolateIndex: isolateIndex,
remoteSessionId: remoteSessionId,
remoteFileToken: token,
fileId: file.file.id,
@@ -449,7 +428,6 @@ class SendNotifier extends Notifier<Map<String, SendSessionState>> {
sendingTasks: [
...?s.sendingTasks,
SendingTask(
isolateIndex: isolateIndex,
taskId: taskResult.taskId,
),
],
@@ -481,7 +459,7 @@ class SendNotifier extends Notifier<Map<String, SendSessionState>> {
state = state.updateSession(
sessionId: sessionId,
state: (s) => s?.copyWith(
sendingTasks: s.sendingTasks?.where((task) => !(task.isolateIndex == isolateIndex && task.taskId == taskResult.taskId)).toList(),
sendingTasks: s.sendingTasks?.where((task) => task.taskId != taskResult.taskId).toList(),
),
);
}
@@ -560,7 +538,6 @@ class SendNotifier extends Notifier<Map<String, SendSessionState>> {
.redux(parentIsolateProvider)
.dispatch(
IsolateHttpUploadCancelAction(
isolateIndex: task.isolateIndex,
taskId: task.taskId,
),
);