feat(#367): consolidate all found message stores into the pinned DB
#363 pins the DB to one directory and migrates ONE prior store in, but a user who ran differently-built copies can have data split across several stores. On startup (native only, after the prefs->drift migration) this discovers every store on the machine and unions its bulk blobs (messages by id, contacts by public key) into the current one, so nothing shows as a gap. Sources are read-only and never deleted. Not a one-shot: instead of a permanent flag it tracks each store's signature (mtime+size), so a store that a stray older build later creates or grows is re-merged rather than stranded. A store over a 200 MB guard, or one that cannot be read, is skipped and NOT recorded as done, so it retries later. Identity dedup is key-order independent. sqlite3 (dart:ffi) stays out of the web build via a conditional import. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>release/1.2.2
parent
60e865b900
commit
61b0eab423
@ -0,0 +1,161 @@
|
||||
import 'dart:convert';
|
||||
import 'dart:io';
|
||||
|
||||
import 'package:flutter/foundation.dart';
|
||||
|
||||
import '../storage/drift/blob_store.dart';
|
||||
import '../storage/drift/offband_database.dart';
|
||||
import '../storage/prefs_manager.dart';
|
||||
import '../utils/app_logger.dart';
|
||||
// sqlite3 (dart:ffi) is native-only; the conditional import keeps it out of the
|
||||
// web build (#363/#367).
|
||||
import '../storage/drift/db_snapshot_web.dart'
|
||||
if (dart.library.io) '../storage/drift/db_snapshot_io.dart';
|
||||
|
||||
/// Consolidation of message stores left in other locations (#367).
|
||||
///
|
||||
/// #363 pins the DB to one directory and migrates ONE prior store into it, but
|
||||
/// a user who ran differently-built copies can have data split across several
|
||||
/// stores. This unions the bulk blobs (messages by id, contacts by public key)
|
||||
/// from every discovered store into the current one, so nothing shows as a gap.
|
||||
///
|
||||
/// Native-only, and it NEVER deletes a source. Runs in `main()` after the
|
||||
/// prefs→drift migration and before the stores are read into memory, so the UI
|
||||
/// sees the consolidated data.
|
||||
///
|
||||
/// It is NOT a one-shot: instead of a permanent done-flag, it records each
|
||||
/// consolidated store's signature (mtime + size). A store that is new, or that
|
||||
/// a stray older build has since GROWN, has a different signature and is merged
|
||||
/// again on the next launch, so data added later is never stranded. A store
|
||||
/// whose signature is unchanged is not re-read, keeping startup cheap.
|
||||
class StoreConsolidationService {
|
||||
StoreConsolidationService._();
|
||||
|
||||
static const _sigKey = 'store_consolidation_sigs_v1';
|
||||
|
||||
/// Stores larger than this are skipped rather than read into memory, so a
|
||||
/// pathologically large file cannot stall startup. Well above any realistic
|
||||
/// message history; aging-out/archiving is a separate future feature.
|
||||
static const _defaultMaxStoreBytes = 200 * 1024 * 1024;
|
||||
|
||||
static Future<void> run() async {
|
||||
// Web has a single OPFS/IndexedDB-backed store; there is nothing to scan,
|
||||
// and sqlite3 is unavailable there.
|
||||
if (kIsWeb) return;
|
||||
|
||||
final prefs = PrefsManager.instance;
|
||||
try {
|
||||
final current = await OffbandDatabase.pinnedDatabasePath();
|
||||
final others = (await OffbandDatabase.discoverStorePaths())
|
||||
.where((path) => path != current && File(path).existsSync())
|
||||
.toSet();
|
||||
|
||||
final prior = _readSignatures(prefs.getString(_sigKey));
|
||||
final currentSigs = {for (final path in others) path: _signatureOf(path)};
|
||||
final toMerge = pathsNeedingMerge(prior, currentSigs);
|
||||
if (toMerge.isEmpty) return; // nothing new or changed since last launch
|
||||
|
||||
final result = await consolidateStores(toMerge);
|
||||
|
||||
// Persist signatures ONLY for stores that were actually consolidated. A
|
||||
// store that was skipped (unreadable, or over the size guard) keeps no
|
||||
// signature, so it is retried on a later launch rather than stranded -
|
||||
// its data lands the moment it becomes readable or drops under the guard.
|
||||
final mergedSigs = <String, String>{
|
||||
for (final path in toMerge)
|
||||
if (!result.unmerged.contains(path)) path: currentSigs[path]!,
|
||||
};
|
||||
await prefs.setString(_sigKey, jsonEncode({...prior, ...mergedSigs}));
|
||||
appLogger.info(
|
||||
'Store consolidation: merged from ${result.stores} store(s), '
|
||||
'${result.changed} key(s) updated.',
|
||||
tag: 'Storage',
|
||||
);
|
||||
} catch (e) {
|
||||
// Signatures are NOT persisted on failure, so it retries next launch.
|
||||
// Sources are untouched, so no data is lost.
|
||||
appLogger.error(
|
||||
'Store consolidation failed: $e. Will retry next launch; no data lost.',
|
||||
tag: 'Storage',
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// The subset of [current] paths whose signature is new or differs from
|
||||
/// [prior] - i.e. stores that appeared or changed and must be (re)merged.
|
||||
/// Visible for testing.
|
||||
@visibleForTesting
|
||||
static List<String> pathsNeedingMerge(
|
||||
Map<String, String> prior,
|
||||
Map<String, String> current,
|
||||
) => [
|
||||
for (final e in current.entries)
|
||||
if (prior[e.key] != e.value) e.key,
|
||||
];
|
||||
|
||||
/// Reads each store in [otherPaths] and unions its bulk blobs into the current
|
||||
/// [BlobStore]. A store over [maxStoreBytes], or one that cannot be read, is
|
||||
/// skipped and reported in `unmerged` (so the caller does not record it as
|
||||
/// done) - never fatal. Visible for testing.
|
||||
@visibleForTesting
|
||||
static Future<({int stores, int changed, Set<String> unmerged})>
|
||||
consolidateStores(
|
||||
Iterable<String> otherPaths, {
|
||||
int maxStoreBytes = _defaultMaxStoreBytes,
|
||||
}) async {
|
||||
var stores = 0, changed = 0;
|
||||
final unmerged = <String>{};
|
||||
for (final path in otherPaths) {
|
||||
final file = File(path);
|
||||
if (file.existsSync() && file.lengthSync() > maxStoreBytes) {
|
||||
appLogger.warn(
|
||||
'Consolidation: store $path is '
|
||||
'${(file.lengthSync() / 1024 / 1024).round()}MB, over the '
|
||||
'${(maxStoreBytes / 1024 / 1024).round()}MB guard; skipped to keep '
|
||||
'startup responsive.',
|
||||
tag: 'Storage',
|
||||
);
|
||||
unmerged.add(path);
|
||||
continue;
|
||||
}
|
||||
|
||||
Map<String, String> blobs;
|
||||
try {
|
||||
blobs = readStoredBlobs(path);
|
||||
} catch (e) {
|
||||
appLogger.warn(
|
||||
'Consolidation: could not read store $path: $e (skipped)',
|
||||
tag: 'Storage',
|
||||
);
|
||||
unmerged.add(path);
|
||||
continue;
|
||||
}
|
||||
stores++;
|
||||
for (final entry in blobs.entries) {
|
||||
if (!BlobStore.isBulkKey(entry.key)) continue;
|
||||
if (await BlobStore.instance.mergeBlob(entry.key, entry.value)) {
|
||||
changed++;
|
||||
}
|
||||
}
|
||||
}
|
||||
return (stores: stores, changed: changed, unmerged: unmerged);
|
||||
}
|
||||
|
||||
static Map<String, String> _readSignatures(String? raw) {
|
||||
if (raw == null || raw.isEmpty) return {};
|
||||
try {
|
||||
final decoded = jsonDecode(raw);
|
||||
if (decoded is Map) {
|
||||
return {
|
||||
for (final e in decoded.entries) e.key as String: e.value as String,
|
||||
};
|
||||
}
|
||||
} catch (_) {}
|
||||
return {};
|
||||
}
|
||||
|
||||
static String _signatureOf(String path) {
|
||||
final f = File(path);
|
||||
return '${f.lastModifiedSync().millisecondsSinceEpoch}:${f.lengthSync()}';
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,266 @@
|
||||
import 'dart:io';
|
||||
|
||||
import 'package:drift/native.dart';
|
||||
import 'package:flutter_test/flutter_test.dart';
|
||||
import 'package:meshcore_open/services/store_consolidation_service.dart';
|
||||
import 'package:meshcore_open/storage/drift/blob_store.dart';
|
||||
import 'package:meshcore_open/storage/drift/offband_database.dart';
|
||||
import 'package:sqlite3/sqlite3.dart';
|
||||
|
||||
/// #367: union message stores left in other locations into the current store,
|
||||
/// so a user with data split across differently-built copies loses nothing.
|
||||
void main() {
|
||||
late OffbandDatabase db;
|
||||
late BlobStore store;
|
||||
late Directory tmp;
|
||||
|
||||
setUp(() {
|
||||
db = OffbandDatabase(NativeDatabase.memory());
|
||||
store = BlobStore(db);
|
||||
BlobStore.overrideForTest(store);
|
||||
tmp = Directory.systemTemp.createTempSync('consolidate_test');
|
||||
});
|
||||
|
||||
tearDown(() async {
|
||||
BlobStore.clearTestOverride();
|
||||
await db.close();
|
||||
tmp.deleteSync(recursive: true);
|
||||
});
|
||||
|
||||
// Writes a real drift-shaped store file with the given blobs.
|
||||
String makeStore(String name, Map<String, String> blobs) {
|
||||
final path = '${tmp.path}/$name/offband_store.sqlite';
|
||||
Directory('${tmp.path}/$name').createSync(recursive: true);
|
||||
final sdb = sqlite3.open(path);
|
||||
sdb.execute('CREATE TABLE stored_blobs(key TEXT PRIMARY KEY, value TEXT)');
|
||||
final stmt = sdb.prepare(
|
||||
'INSERT INTO stored_blobs(key, value) VALUES(?, ?)',
|
||||
);
|
||||
blobs.forEach((k, v) => stmt.execute([k, v]));
|
||||
stmt.close();
|
||||
sdb.close();
|
||||
return path;
|
||||
}
|
||||
|
||||
group('mergeBlob', () {
|
||||
const key = 'channel_messages_devpsk_a';
|
||||
|
||||
test('writes when the key is absent', () async {
|
||||
expect(await store.mergeBlob(key, '[{"messageId":"a"}]'), isTrue);
|
||||
expect(await store.read(key), '[{"messageId":"a"}]');
|
||||
});
|
||||
|
||||
test('unions by id, keeping existing and adding missing', () async {
|
||||
await store.write(key, '[{"messageId":"a"},{"messageId":"b"}]');
|
||||
expect(
|
||||
await store.mergeBlob(key, '[{"messageId":"b"},{"messageId":"c"}]'),
|
||||
isTrue,
|
||||
);
|
||||
expect(
|
||||
await store.read(key),
|
||||
'[{"messageId":"a"},{"messageId":"b"},{"messageId":"c"}]',
|
||||
);
|
||||
});
|
||||
|
||||
test('reports no change when the incoming is a subset', () async {
|
||||
await store.write(key, '[{"messageId":"a"},{"messageId":"b"}]');
|
||||
expect(await store.mergeBlob(key, '[{"messageId":"a"}]'), isFalse);
|
||||
expect(await store.read(key), '[{"messageId":"a"},{"messageId":"b"}]');
|
||||
});
|
||||
|
||||
// Identity-edge coverage for the merge logic Gemini flagged as unverified.
|
||||
// _mergeBulk resolves identity in this order: messageId, then publicKey,
|
||||
// then the object's canonical JSON. On a collision the CURRENT copy is
|
||||
// kept and only unseen incoming items are appended. These pin each branch.
|
||||
|
||||
test(
|
||||
'identity falls back to publicKey (contacts) when no messageId',
|
||||
() async {
|
||||
const ck = 'contacts_dev';
|
||||
await store.write(ck, '[{"publicKey":"A","name":"current"}]');
|
||||
// Same pk A (kept as-is), new pk B (added).
|
||||
expect(
|
||||
await store.mergeBlob(
|
||||
ck,
|
||||
'[{"publicKey":"A","name":"stale"},{"publicKey":"B"}]',
|
||||
),
|
||||
isTrue,
|
||||
);
|
||||
expect(
|
||||
await store.read(ck),
|
||||
'[{"publicKey":"A","name":"current"},{"publicKey":"B"}]',
|
||||
);
|
||||
},
|
||||
);
|
||||
|
||||
test(
|
||||
'keeps the current copy on an id collision (an edit does not win)',
|
||||
() async {
|
||||
await store.write(key, '[{"messageId":"a","text":"orig"}]');
|
||||
// Same messageId, different text: the current copy must survive, and the
|
||||
// merge reports no change.
|
||||
expect(
|
||||
await store.mergeBlob(key, '[{"messageId":"a","text":"edited"}]'),
|
||||
isFalse,
|
||||
);
|
||||
expect(await store.read(key), '[{"messageId":"a","text":"orig"}]');
|
||||
},
|
||||
);
|
||||
|
||||
test('id-less objects dedupe by whole-object identity', () async {
|
||||
await store.write(key, '[{"text":"hi","ts":1}]');
|
||||
// Identical object deduped; a differing one appended.
|
||||
expect(
|
||||
await store.mergeBlob(
|
||||
key,
|
||||
'[{"text":"hi","ts":1},{"text":"yo","ts":2}]',
|
||||
),
|
||||
isTrue,
|
||||
);
|
||||
expect(
|
||||
await store.read(key),
|
||||
'[{"text":"hi","ts":1},{"text":"yo","ts":2}]',
|
||||
);
|
||||
});
|
||||
|
||||
test('id-less dedupe is independent of key order', () async {
|
||||
// Two stores from different builds can serialise the same object with
|
||||
// keys in a different order; that must still dedupe, not duplicate.
|
||||
await store.write(key, '[{"ts":1,"text":"a"}]');
|
||||
expect(
|
||||
await store.mergeBlob(key, '[{"text":"a","ts":1},{"text":"b","ts":2}]'),
|
||||
isTrue,
|
||||
);
|
||||
// The reordered duplicate collapsed; only the genuinely new object added.
|
||||
expect(
|
||||
await store.read(key),
|
||||
'[{"ts":1,"text":"a"},{"text":"b","ts":2}]',
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
group('pathsNeedingMerge (re-runs when stores appear or grow)', () {
|
||||
test('includes new and changed paths, skips unchanged', () {
|
||||
final prior = {'a': '1:100', 'b': '2:200'};
|
||||
// a unchanged, b changed (an old build grew it), c is new.
|
||||
final current = {'a': '1:100', 'b': '9:999', 'c': '3:300'};
|
||||
final need = StoreConsolidationService.pathsNeedingMerge(prior, current)
|
||||
..sort();
|
||||
expect(need, ['b', 'c']);
|
||||
});
|
||||
|
||||
test('is empty when nothing changed (cheap no-op startup)', () {
|
||||
final sigs = {'a': '1:1', 'b': '2:2'};
|
||||
expect(StoreConsolidationService.pathsNeedingMerge(sigs, sigs), isEmpty);
|
||||
});
|
||||
});
|
||||
|
||||
group('consolidateStores', () {
|
||||
test(
|
||||
'unions messages from every other store into the current one',
|
||||
() async {
|
||||
const key = 'channel_messages_devpsk_public';
|
||||
// Current store already holds a, b.
|
||||
await store.write(key, '[{"messageId":"a"},{"messageId":"b"}]');
|
||||
|
||||
// Two stranded stores: one adds c (and re-states b), one adds d, plus a
|
||||
// non-bulk key that must be ignored.
|
||||
final s1 = makeStore('one', {
|
||||
key: '[{"messageId":"b"},{"messageId":"c"}]',
|
||||
});
|
||||
final s2 = makeStore('two', {
|
||||
key: '[{"messageId":"d"}]',
|
||||
'ui_sort_option': 'manual', // not a bulk key
|
||||
});
|
||||
|
||||
final result = await StoreConsolidationService.consolidateStores([
|
||||
s1,
|
||||
s2,
|
||||
]);
|
||||
|
||||
expect(result.stores, 2);
|
||||
expect(
|
||||
await store.read(key),
|
||||
'[{"messageId":"a"},{"messageId":"b"},'
|
||||
'{"messageId":"c"},{"messageId":"d"}]',
|
||||
);
|
||||
// The non-bulk key was not imported.
|
||||
expect(await store.read('ui_sort_option'), isNull);
|
||||
},
|
||||
);
|
||||
|
||||
test('skips an unreadable store without failing', () async {
|
||||
final good = makeStore('good', {'contacts_dev': '[{"publicKey":"A"}]'});
|
||||
final missing = '${tmp.path}/gone/offband_store.sqlite';
|
||||
|
||||
final result = await StoreConsolidationService.consolidateStores([
|
||||
missing,
|
||||
good,
|
||||
]);
|
||||
|
||||
expect(result.stores, 1, reason: 'only the readable store counted');
|
||||
expect(result.unmerged, contains(missing), reason: 'retried next launch');
|
||||
expect(await store.read('contacts_dev'), '[{"publicKey":"A"}]');
|
||||
});
|
||||
|
||||
test(
|
||||
'skips corrupt / zero-byte / schemaless stores, merges the valid one',
|
||||
() async {
|
||||
// Zero-byte file.
|
||||
final empty = '${tmp.path}/empty/offband_store.sqlite';
|
||||
Directory('${tmp.path}/empty').createSync();
|
||||
File(empty).writeAsBytesSync(const []);
|
||||
// A valid SQLite file that lacks the stored_blobs table.
|
||||
final noTable = '${tmp.path}/notable/offband_store.sqlite';
|
||||
Directory('${tmp.path}/notable').createSync();
|
||||
final s = sqlite3.open(noTable);
|
||||
s.execute('CREATE TABLE other(x TEXT)');
|
||||
s.close();
|
||||
// Non-SQLite garbage bytes.
|
||||
final garbage = '${tmp.path}/garbage/offband_store.sqlite';
|
||||
Directory('${tmp.path}/garbage').createSync();
|
||||
File(garbage).writeAsBytesSync(List<int>.filled(64, 7));
|
||||
// One good store.
|
||||
final good = makeStore('good2', {
|
||||
'contacts_dev': '[{"publicKey":"Z"}]',
|
||||
});
|
||||
|
||||
final result = await StoreConsolidationService.consolidateStores([
|
||||
empty,
|
||||
noTable,
|
||||
garbage,
|
||||
good,
|
||||
]);
|
||||
|
||||
expect(result.stores, 1, reason: 'only the valid store was read');
|
||||
expect(
|
||||
result.unmerged,
|
||||
containsAll([empty, noTable, garbage]),
|
||||
reason:
|
||||
'skipped stores are reported so they are retried, not stranded',
|
||||
);
|
||||
expect(await store.read('contacts_dev'), '[{"publicKey":"Z"}]');
|
||||
},
|
||||
);
|
||||
|
||||
test('skips a store larger than the size guard', () async {
|
||||
// A store padded well past a small cap, and a tiny one under it.
|
||||
final bigVal =
|
||||
'[${List.generate(20000, (i) => '{"messageId":"m$i"}').join(',')}]';
|
||||
final big = makeStore('big', {'channel_messages_devpsk_big': bigVal});
|
||||
final small = makeStore('small', {'contacts_dev': '[{"publicKey":"A"}]'});
|
||||
final cap = File(small).lengthSync() + 1024;
|
||||
expect(File(big).lengthSync(), greaterThan(cap));
|
||||
|
||||
final result = await StoreConsolidationService.consolidateStores([
|
||||
big,
|
||||
small,
|
||||
], maxStoreBytes: cap);
|
||||
|
||||
expect(result.stores, 1, reason: 'the oversized store was skipped');
|
||||
expect(result.unmerged, contains(big), reason: 'retried if it shrinks');
|
||||
expect(await store.read('contacts_dev'), '[{"publicKey":"A"}]');
|
||||
expect(await store.read('channel_messages_devpsk_big'), isNull);
|
||||
});
|
||||
});
|
||||
}
|
||||
Loading…
Reference in new issue