mirror of
https://github.com/immich-app/drift.git
synced 2026-09-30 13:22:57 +08:00
Close isolate connection on exit
Fix for https://github.com/simolus3/drift/issues/3449#issuecomment-3045322812 # Conflicts: # drift/CHANGELOG.md # drift/lib/native.dart
This commit is contained in:
committed by
shenlong-tanwen
parent
f1b58a71c3
commit
53ef7e9f19
@@ -1,3 +1,18 @@
|
|||||||
|
## 2.28.0-dev
|
||||||
|
|
||||||
|
- Query delegates: Support nested transactions in `SupportedTransactionDelegate`.
|
||||||
|
- Fix drift server isolates leaking when their only client isolate exits without closing the database.
|
||||||
|
|
||||||
|
## 2.27.0
|
||||||
|
|
||||||
|
- Allow passing `sqlite3` callback to `NativeDatabase` to customize how SQLite
|
||||||
|
bindings are obtained.
|
||||||
|
|
||||||
|
## 2.26.1
|
||||||
|
|
||||||
|
- Add `isNotNull()` column filter for the manager APIs.
|
||||||
|
- Add optional `orderBy` parameter to more aggregate function extension.
|
||||||
|
|
||||||
## 2.26.0
|
## 2.26.0
|
||||||
|
|
||||||
- Add support for window functions with `WindowFunctionExpression`.
|
- Add support for window functions with `WindowFunctionExpression`.
|
||||||
|
|||||||
+13
-9
@@ -431,15 +431,19 @@ class _NativeIsolateStartup {
|
|||||||
|
|
||||||
static Future<void> start(_NativeIsolateStartup startup) async {
|
static Future<void> start(_NativeIsolateStartup startup) async {
|
||||||
await startup.isolateSetup?.call();
|
await startup.isolateSetup?.call();
|
||||||
final isolate = DriftIsolate.inCurrent(() {
|
final isolate = DriftIsolate.inCurrent(
|
||||||
return DatabaseConnection(NativeDatabase(
|
() {
|
||||||
File(startup.path),
|
return DatabaseConnection(NativeDatabase(
|
||||||
logStatements: startup.enableLogs,
|
File(startup.path),
|
||||||
cachePreparedStatements: startup.cachePreparedStatements,
|
logStatements: startup.enableLogs,
|
||||||
enableMigrations: startup.enableMigrations,
|
cachePreparedStatements: startup.cachePreparedStatements,
|
||||||
setup: startup.setup,
|
enableMigrations: startup.enableMigrations,
|
||||||
));
|
setup: startup.setup,
|
||||||
});
|
));
|
||||||
|
},
|
||||||
|
shutdownAfterLastDisconnect: true,
|
||||||
|
killIsolateWhenDone: true,
|
||||||
|
);
|
||||||
|
|
||||||
startup.sendServer.send(isolate);
|
startup.sendServer.send(isolate);
|
||||||
}
|
}
|
||||||
|
|||||||
+19
-10
@@ -36,7 +36,13 @@ Future<(StreamChannel, bool)> connectToServer(
|
|||||||
// If the isolate accepts the connection, it sends us a send port back which
|
// If the isolate accepts the connection, it sends us a send port back which
|
||||||
// is then used for the rest of the communication.
|
// is then used for the rest of the communication.
|
||||||
final receive = ReceivePort('drift client receive');
|
final receive = ReceivePort('drift client receive');
|
||||||
serverConnectPort.send([receive.sendPort, serialize]);
|
serverConnectPort.send([
|
||||||
|
receive.sendPort,
|
||||||
|
serialize,
|
||||||
|
// The server isolate will use addOnExitListener to mark the connection as
|
||||||
|
// closed when the isolate shuts down.
|
||||||
|
Isolate.current.controlPort
|
||||||
|
]);
|
||||||
|
|
||||||
final controller =
|
final controller =
|
||||||
StreamChannelController<Object?>(allowForeignErrors: false, sync: true);
|
StreamChannelController<Object?>(allowForeignErrors: false, sync: true);
|
||||||
@@ -105,16 +111,21 @@ class RunningDriftServer {
|
|||||||
closeConnectionAfterShutdown: closeConnectionAfterShutdown,
|
closeConnectionAfterShutdown: closeConnectionAfterShutdown,
|
||||||
) {
|
) {
|
||||||
final subscription = connectPort.listen((message) {
|
final subscription = connectPort.listen((message) {
|
||||||
if (message is List && message.length == 2) {
|
if (message
|
||||||
|
case [
|
||||||
|
final sendPort as SendPort,
|
||||||
|
final serialize as bool,
|
||||||
|
final closeOnExit as SendPort
|
||||||
|
]) {
|
||||||
if (onlyAcceptSingleConnection) {
|
if (onlyAcceptSingleConnection) {
|
||||||
connectPort.close();
|
connectPort.close();
|
||||||
}
|
}
|
||||||
|
|
||||||
final sendPort = message[0]! as SendPort;
|
final clientIsolate = Isolate(closeOnExit);
|
||||||
final serialize = message[1]! as bool;
|
|
||||||
final receiveForConnection =
|
final receiveForConnection =
|
||||||
ReceivePort('drift channel #${_counter++}');
|
ReceivePort('drift channel #${_counter++}');
|
||||||
sendPort.send(receiveForConnection.sendPort);
|
final sendPortForRemote = receiveForConnection.sendPort;
|
||||||
|
sendPort.send(sendPortForRemote);
|
||||||
|
|
||||||
final controller = StreamChannelController<Object?>(
|
final controller = StreamChannelController<Object?>(
|
||||||
allowForeignErrors: false, sync: true);
|
allowForeignErrors: false, sync: true);
|
||||||
@@ -122,11 +133,6 @@ class RunningDriftServer {
|
|||||||
if (message == disconnectMessage) {
|
if (message == disconnectMessage) {
|
||||||
// Client closed the connection
|
// Client closed the connection
|
||||||
controller.local.sink.close();
|
controller.local.sink.close();
|
||||||
|
|
||||||
if (onlyAcceptSingleConnection) {
|
|
||||||
// The only connection was closed, so shut down the server.
|
|
||||||
server.shutdown();
|
|
||||||
}
|
|
||||||
} else {
|
} else {
|
||||||
controller.local.sink.add(message);
|
controller.local.sink.add(message);
|
||||||
}
|
}
|
||||||
@@ -137,9 +143,12 @@ class RunningDriftServer {
|
|||||||
sendPort.send(disconnectMessage);
|
sendPort.send(disconnectMessage);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
clientIsolate.addOnExitListener(sendPortForRemote,
|
||||||
|
response: disconnectMessage);
|
||||||
_activeConnections++;
|
_activeConnections++;
|
||||||
server.serve(controller.foreign, serialize: serialize).whenComplete(() {
|
server.serve(controller.foreign, serialize: serialize).whenComplete(() {
|
||||||
_activeConnections--;
|
_activeConnections--;
|
||||||
|
clientIsolate.removeErrorListener(sendPortForRemote);
|
||||||
|
|
||||||
if (_activeConnections == 0 && shutDownAfterLastDisconnect) {
|
if (_activeConnections == 0 && shutDownAfterLastDisconnect) {
|
||||||
server.shutdown();
|
server.shutdown();
|
||||||
|
|||||||
@@ -174,6 +174,26 @@ void main() {
|
|||||||
await db.close();
|
await db.close();
|
||||||
}, tags: 'background_isolate');
|
}, tags: 'background_isolate');
|
||||||
|
|
||||||
|
test('can close server when client isolate exits', () async {
|
||||||
|
final shutdownCalled = Completer<void>();
|
||||||
|
final isolate = DriftIsolate.inCurrent(
|
||||||
|
() => NativeDatabase.memory(),
|
||||||
|
shutdownAfterLastDisconnect: true,
|
||||||
|
beforeShutdown: shutdownCalled.complete,
|
||||||
|
);
|
||||||
|
|
||||||
|
await Isolate.run(() async {
|
||||||
|
final connection = await isolate.connect();
|
||||||
|
// Ensure the statement works
|
||||||
|
await TodoDb(connection).customStatement('SELECT 1');
|
||||||
|
|
||||||
|
// Close the isolate without closing the database
|
||||||
|
});
|
||||||
|
|
||||||
|
// This should still shut down the server.
|
||||||
|
await shutdownCalled.future;
|
||||||
|
});
|
||||||
|
|
||||||
test('shutting down will close the underlying executor', () async {
|
test('shutting down will close the underlying executor', () async {
|
||||||
final mockExecutor = MockExecutor();
|
final mockExecutor = MockExecutor();
|
||||||
final isolate =
|
final isolate =
|
||||||
|
|||||||
Reference in New Issue
Block a user