Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion lib/database/database.dart
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import 'dart:async';
import 'dart:io';

import 'package:bluebubbles/env.dart';
import 'package:bluebubbles/helpers/backend/startup_tasks.dart';
import 'package:bluebubbles/helpers/helpers.dart';
import 'package:bluebubbles/database/models.dart';
Expand Down Expand Up @@ -75,7 +76,8 @@ class Database {
Logger.info(
"Database init: SettingsSvc.finishedSetup = $setupFinished, PrefsSvc.finishedSetup = $setupFinished2");

if (!setupFinished) {
// Setup only ever runs on main, so only main may act on an unfinished one.
if (!setupFinished && !isIsolate) {
Logger.warn("Clearing database because setup is not finished...");

Database.attachments.removeAll();
Expand Down
57 changes: 27 additions & 30 deletions lib/services/backend/settings/desktop_shared_preferences_store.dart
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,11 @@ import 'dart:async';
import 'dart:convert';
import 'dart:io';

import 'package:bluebubbles/services/backend/filesystem/filesystem_service.dart';
import 'package:bluebubbles/utils/file_utils.dart';
import 'package:bluebubbles/utils/logger/logger.dart';
import 'package:flutter/foundation.dart' show debugPrint;
import 'package:path/path.dart' as p;
import 'package:path_provider_linux/path_provider_linux.dart';
import 'package:path_provider_windows/path_provider_windows.dart';
import 'package:shared_preferences_platform_interface/shared_preferences_async_platform_interface.dart';
import 'package:shared_preferences_platform_interface/types.dart';

Expand All @@ -29,15 +28,18 @@ import 'package:shared_preferences_platform_interface/types.dart';
/// - Writes are serialized across isolates AND processes with an exclusively
/// created lock file. (`RandomAccessFile.lock` is not enough: POSIX fcntl
/// locks are process-owned and do not exclude isolates within one process.)
/// - Writes are atomic: temp file + rename, so the store can never be
/// truncated by a crash mid-write.
/// - Writes are atomic where the filesystem allows it: temp file + rename, so
/// the store can never be truncated by a crash mid-write. Where the rename
/// can't replace the file, it falls back to writing in place and the `.bak`
/// below covers it.
/// - Every successful write also refreshes a `.bak` copy (under the same
/// lock), and an unparseable file is quarantined and restored from that
/// backup — at registration and on mid-session reads — instead of crashing
/// the app or losing all settings.
///
/// Storage location and format are identical to the stock implementation, so
/// existing user data carries over untouched. Custom file names via
/// Format is identical to the stock implementation, and so is the location
/// except on MSIX, where [FilesystemService] has already migrated the whole
/// directory to the real path behind the AppData redirect. Custom file names via
/// platform-specific [SharedPreferencesOptions] subclasses are not supported;
/// the app only ever uses the defaults.
base class DesktopSharedPreferencesStore extends SharedPreferencesAsyncPlatform {
Expand All @@ -49,7 +51,6 @@ base class DesktopSharedPreferencesStore extends SharedPreferencesAsyncPlatform
static const Duration _staleLockTimeout = Duration(seconds: 10);
static const Duration _lockRetryDelay = Duration(milliseconds: 5);

String? _cachedDirectoryPath;
Future<void> _writeQueue = Future.value();

/// Registers this store as the [SharedPreferencesAsyncPlatform] and
Expand Down Expand Up @@ -133,31 +134,21 @@ base class DesktopSharedPreferencesStore extends SharedPreferencesAsyncPlatform
Future<Set<String>> getKeys(GetPreferencesParameters parameters, SharedPreferencesOptions options) async =>
(await getPreferences(parameters, options)).keys.toSet();

Future<String> _getDirectoryPath() async {
if (_cachedDirectoryPath != null) return _cachedDirectoryPath!;
// Instantiated directly (instead of going through path_provider) so this
// works in background isolates without plugin registration, exactly like
// the stock implementations do.
final String? directory = Platform.isWindows
? await PathProviderWindows().getApplicationSupportPath()
: await PathProviderLinux().getApplicationSupportPath();
if (directory == null) {
throw const FileSystemException('Unable to resolve the application support directory for preferences');
}
return _cachedDirectoryPath = directory;
}

Future<File> _getDataFile() async => File(p.join(await _getDirectoryPath(), _fileName));
// Not path_provider's support path directly: on MSIX that one is the AppData
// redirect, where a same-folder rename can fail as cross-device. Every
// isolate's init registers FilesystemService before this store.
Future<File> _getDataFile() async => File(p.join(FilesystemSvc.appDocDir.path, _fileName));

Future<File> _getBackupFile() async => File('${(await _getDataFile()).path}$_backupSuffix');

/// Returns the parsed contents of [file], `{}` for an existing-but-empty
/// file, or null when the file is missing or unparseable.
/// Returns the parsed contents of [file], or null when it is missing, empty
/// or unparseable. Empty means a reader caught an in-place write mid-truncate,
/// never "no preferences set" — that is stored as `{}`.
Map<String, Object>? _parseFile(File file) {
try {
if (!file.existsSync()) return null;
final String contents = file.readAsStringSync();
if (contents.isEmpty) return <String, Object>{};
if (contents.isEmpty) return null;
final Object? decoded = json.decode(contents);
return decoded is Map ? decoded.cast<String, Object>() : null;
} on FormatException catch (e) {
Expand Down Expand Up @@ -239,20 +230,26 @@ base class DesktopSharedPreferencesStore extends SharedPreferencesAsyncPlatform
Future<void> _atomicWrite(Map<String, Object> prefs) async {
final File file = await _getDataFile();
final File tmp = File('${file.path}.tmp');
final String contents = json.encode(prefs);
try {
final RandomAccessFile raf = tmp.openSync(mode: FileMode.write);
try {
raf.writeStringSync(json.encode(prefs));
raf.writeStringSync(contents);
raf.flushSync();
} finally {
raf.closeSync();
}
// Not atomic if it has to fall back to a copy — the `.bak` refreshed
// below and the corrupt-file recovery are what cover that.
await moveFile(tmp, file.path);
} on FileSystemException catch (e) {
_log('Failed to save preferences: $e');
return;
// Neither rename nor copy can replace a file another handle holds open
// (`File.copy` deletes the destination first), but an in-place write can.
_log('Atomic preferences write failed, writing in place: $e');
try {
file.writeAsStringSync(contents, flush: true);
} on FileSystemException catch (e) {
_log('Failed to save preferences: $e');
return;
}
}
try {
file.copySync((await _getBackupFile()).path);
Expand Down
13 changes: 7 additions & 6 deletions lib/services/backend/sync/incremental_sync_manager.dart
Original file line number Diff line number Diff line change
Expand Up @@ -364,7 +364,7 @@ class IncrementalSyncManager extends SyncManager {
syncedChats.addAll(chatCache);

// For each chat, bulk sync the messages
final pageMessageIds = <int>[];
final pageMessageIdsByChat = <String, List<int>>{};
final pageLatestMessageIdPerChat = <String, int>{};

for (var item in messagesToSync.entries) {
Expand All @@ -387,9 +387,10 @@ class IncrementalSyncManager extends SyncManager {
latestMessageIdPerChat[item.key] = latest.id!;
}

// Collect per-page data for the progressive UI update event.
for (final m in syncResult.messages) {
if (m.id != null) pageMessageIds.add(m.id!);
// Grouped by chat so the main thread can skip hydrating messages for closed chats.
final ids = syncResult.messages.where((m) => m.id != null).map((m) => m.id!).toList();
if (ids.isNotEmpty) {
pageMessageIdsByChat[item.key] = ids;
}
if (latest != null) {
pageLatestMessageIdPerChat[item.key] = latest.id!;
Expand All @@ -398,9 +399,9 @@ class IncrementalSyncManager extends SyncManager {

// Emit a per-page event so the main thread can update the UI incrementally
// without waiting for all pages to finish.
if (isIsolate && pageMessageIds.isNotEmpty) {
if (isIsolate && pageMessageIdsByChat.isNotEmpty) {
IsolateEventEmitter.emit(IsolateEvent.incrementalSyncPageComplete, {
'messageIds': pageMessageIds,
'messageIdsByChat': pageMessageIdsByChat,
'latestMessageIdPerChat': pageLatestMessageIdPerChat,
});
}
Expand Down
45 changes: 29 additions & 16 deletions lib/services/backend/sync/sync_service.dart
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
import 'dart:async';
import 'dart:math' as math;

import 'package:bluebubbles/database/database.dart';
import 'package:bluebubbles/database/models.dart';
import 'package:bluebubbles/helpers/ui/async_task.dart';
import 'package:bluebubbles/helpers/ui/ui_helpers.dart';
import 'package:bluebubbles/services/backend/interfaces/contact_v2_interface.dart';
import 'package:bluebubbles/services/backend/interfaces/sync_interface.dart';
Expand All @@ -24,6 +26,7 @@ class SyncService {
final RxBool isIncrementalSyncing = false.obs;

static const Duration _incrementalSyncCooldown = Duration(seconds: 30);
static const int _dispatchChunkSize = 50;
DateTime? _lastIncrementalSyncTimestamp;

FullSyncManager? _manager;
Expand Down Expand Up @@ -73,25 +76,17 @@ class SyncService {
final processedMessageIds = <int>{};
final processedSubtitleByChat = <String, int>{}; // chatGuid → message DB ID

// This runs on the UI thread (isolate event listeners are invoked synchronously),
// so a page of up to `batchSize` messages must never be hydrated in one block.
Future<void> onPageComplete(dynamic data) async {
if (data is! Map<String, dynamic>) return;
final messageIds = (data['messageIds'] as List).cast<int>();
final messageIdsByChat =
(data['messageIdsByChat'] as Map).map((k, v) => MapEntry(k as String, (v as List).cast<int>()));
final latestPerChat = Map<String, int>.from(data['latestMessageIdPerChat'] as Map);

// Hydrate the page's messages and dispatch to any open chat view immediately.
final messages = Database.messages.getMany(messageIds).whereType<Message>().toList();
for (final message in messages) {
if (message.id != null) processedMessageIds.add(message.id!);
final chatGuid = message.chat.target?.guid;
if (chatGuid == null || message.guid == null) continue;
if (Get.isRegistered<MessagesService>(tag: chatGuid)) {
unawaited(Get.find<MessagesService>(tag: chatGuid).addNewMessage(message));
}
}

// Update chat subtitles for the per-page latest message per chat.
// repositionImmediate: false so a page touching many chats rebuilds the list once.
for (final entry in latestPerChat.entries) {
final msg = Database.messages.get(entry.value);
final msg = await runAsync(() => Database.messages.get(entry.value));
if (msg == null) continue;
// If this chat was created for the first time during this sync,
// ChatState doesn't exist yet — register it so updateChatLatestMessage
Expand All @@ -100,9 +95,27 @@ class SyncService {
final chat = msg.chat.target;
if (chat != null) await ChatsSvc.addChat(chat, immediate: true);
}
ChatsSvc.updateChatLatestMessage(entry.key, msg);
ChatsSvc.updateChatLatestMessage(entry.key, msg, repositionImmediate: false);
processedSubtitleByChat[entry.key] = entry.value;
}

// Only an open conversation view needs the individual messages — every other
// chat is served by the subtitle update above.
for (final entry in messageIdsByChat.entries) {
if (maybeFindMessagesSvc(entry.key) == null) continue;
for (int i = 0; i < entry.value.length; i += _dispatchChunkSize) {
// Re-resolve every chunk: the user can leave this chat between chunks
final service = maybeFindMessagesSvc(entry.key);
if (service == null) break;
final chunk = entry.value.sublist(i, math.min(i + _dispatchChunkSize, entry.value.length));
final messages = await runAsync(() => Database.messages.getMany(chunk).whereType<Message>().toList());
for (final message in messages) {
if (message.guid == null) continue;
await service.addNewMessage(message);
if (message.id != null) processedMessageIds.add(message.id!);
}
}
}
}

final syncIsolate = GetIt.I<IncrementalSyncIsolate>();
Expand Down Expand Up @@ -134,7 +147,7 @@ class SyncService {
for (final entry in latestPerChat.entries) {
final message = entry.value;
if (message.id != null && processedSubtitleByChat[entry.key] == message.id) continue;
ChatsSvc.updateChatLatestMessage(entry.key, message);
ChatsSvc.updateChatLatestMessage(entry.key, message, repositionImmediate: false);
}

// Dispatch newly synced messages to any currently active chat view.
Expand Down
2 changes: 1 addition & 1 deletion lib/services/isolates/isolate_event.dart
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ enum IsolateEvent {

/// Emitted after each page of messages is persisted during incremental sync.
/// Payload: `Map<String, dynamic>` with keys:
/// `messageIds` — `List<int>` of all DB IDs saved in this page
/// `messageIdsByChat` — `Map<String, List<int>>` chatGuid → DB IDs saved in this page
/// `latestMessageIdPerChat` — `Map<String, int>` chatGuid → latest message DB ID in this page
incrementalSyncPageComplete,

Expand Down
7 changes: 5 additions & 2 deletions lib/services/ui/chat/chats_service.dart
Original file line number Diff line number Diff line change
Expand Up @@ -1585,7 +1585,10 @@ class ChatsService {
/// older delta message as a chat's latest, which would rewind its sort order.
/// The 2s tolerance allows a temp->real GUID swap. [allowOlder] opts out for the
/// post-deletion recompute, which must fall back to an older surviving message.
void updateChatLatestMessage(String chatGuid, Message message, {bool allowOlder = false}) {
/// [repositionImmediate] controls the chat list rebuild only — the ChatState update is always
/// synchronous. Pass false when updating many chats in a row so the rebuilds coalesce.
void updateChatLatestMessage(String chatGuid, Message message,
{bool allowOlder = false, bool repositionImmediate = true}) {
final state = getChatState(chatGuid);
if (state == null) return;

Expand All @@ -1611,7 +1614,7 @@ class ChatsService {
state.updateSubtitleInternal(
message.getNotificationText(hideContactInfo: hideContactInfo, hideMessageContent: hideMessageContent));
state.chat.setLatestMessage(message);
_repositionChat(state.chat, immediate: true);
_repositionChat(state.chat, immediate: repositionImmediate);
}

/// Set chat text field text
Expand Down
24 changes: 8 additions & 16 deletions lib/utils/file_utils.dart
Original file line number Diff line number Diff line change
Expand Up @@ -28,17 +28,14 @@ Future<void> revealInFileManager(String path) async {
}

/// Moves [source] onto [targetPath], by rename where the filesystem allows it.
Future<void> moveFile(File source, String targetPath, {int renameAttempts = 5}) async {
for (int attempt = 1; attempt <= renameAttempts; attempt++) {
try {
await source.rename(targetPath);
return;
} on PathNotFoundException {
rethrow;
} on FileSystemException catch (e) {
if (isCrossDeviceError(e) || attempt >= renameAttempts) break;
await Future.delayed(const Duration(milliseconds: 10));
}
Future<void> moveFile(File source, String targetPath) async {
try {
await source.rename(targetPath);
return;
} on PathNotFoundException {
rethrow;
} on FileSystemException {
// Cross-device, or the destination can't be replaced.
}
await source.copy(targetPath);
try {
Expand All @@ -48,11 +45,6 @@ Future<void> moveFile(File source, String targetPath, {int renameAttempts = 5})
}
}

/// Win32 ERROR_NOT_SAME_DEVICE (17) on Windows, POSIX EXDEV (18) everywhere
/// else — Dart reports Win32 codes there and errno elsewhere, so the two can't
/// be checked together: 17 is EEXIST on POSIX.
bool isCrossDeviceError(FileSystemException e) => e.osError?.errorCode == (Platform.isWindows ? 17 : 18);

/// Desktop "Save As": asks where to put the file, then puts it there. Returns
/// the path it was saved to, or null if the user cancelled the dialog.
Future<String?> saveFileAs({
Expand Down
1 change: 1 addition & 0 deletions linux/my_application.cc
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,7 @@ static void my_application_activate(GApplication* application) {

g_autoptr(FlDartProject) project = fl_dart_project_new();
fl_dart_project_set_dart_entrypoint_arguments(project, self->dart_entrypoint_arguments);
fl_dart_project_set_ui_thread_policy(project, FL_UI_THREAD_POLICY_RUN_ON_SEPARATE_THREAD);

FlView* view = fl_view_new(project);

Expand Down
2 changes: 1 addition & 1 deletion windows/runner/main.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ int APIENTRY wWinMain(_In_ HINSTANCE instance, _In_opt_ HINSTANCE prev,
::CoInitializeEx(nullptr, COINIT_APARTMENTTHREADED);

flutter::DartProject project(L"data");
// project.set_ui_thread_policy(flutter::UIThreadPolicy::RunOnSeparateThread);
project.set_ui_thread_policy(flutter::UIThreadPolicy::RunOnSeparateThread);

std::vector<std::string> command_line_arguments = GetCommandLineArguments();

Expand Down
Loading