Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
151 changes: 92 additions & 59 deletions lib/data/repositories/mostro_storage.dart
Original file line number Diff line number Diff line change
@@ -1,26 +1,91 @@
import 'dart:async';
import 'package:flutter/foundation.dart';
import 'package:mostro_mobile/services/logger_service.dart';
import 'package:mostro_mobile/data/models/payload.dart';
import 'package:sembast/sembast.dart';
import 'package:mostro_mobile/data/models/mostro_message.dart';
import 'package:mostro_mobile/data/repositories/base_storage.dart';

class MostroStorage extends BaseStorage<MostroMessage> {


MostroStorage({required Database db})
: super(db, stringMapStoreFactory.store('orders'));

/// In-memory index by order id (newest first). Every write used to wake
/// one Sembast query listener per OrderNotifier and per visible trade row,
/// each re-filtering the whole unindexed store on the UI isolate. The
/// index is warmed from disk once; watchers are served from memory and
/// demultiplexed per order.
final Map<String, List<MostroMessage>> _byOrder = {};
final StreamController<String> _orderChanges = StreamController.broadcast();
Future<void>? _warmup;

@visibleForTesting
int get debugIndexSize => _byOrder.length;

Future<void> _ensureIndex() {
return _warmup ??= () async {
final all = await getAll();
for (final message in all) {
_indexAdd(message, notify: false);
}
logger.i('Mostro message index warmed: ${_byOrder.length} orders');
}();
}

void _indexAdd(MostroMessage message, {bool notify = true}) {
final orderId = message.id;
if (orderId == null) return;
final list = _byOrder.putIfAbsent(orderId, () => <MostroMessage>[]);
list.add(message);
list.sort((a, b) => (b.timestamp ?? 0).compareTo(a.timestamp ?? 0));
if (notify) _notifyOrder(orderId);
}

void _notifyOrder(String orderId) {
if (!_orderChanges.isClosed) _orderChanges.add(orderId);
}

MostroMessage? _latestFor(String orderId) {
final list = _byOrder[orderId];
return (list == null || list.isEmpty) ? null : list.first;
}

List<MostroMessage> _historyFor(String orderId) =>
List.unmodifiable(_byOrder[orderId] ?? const <MostroMessage>[]);

Stream<R> _watchOrder<R>(String orderId, R Function() read) {
late StreamController<R> controller;
StreamSubscription<String>? changes;
controller = StreamController<R>(
onListen: () async {
await _ensureIndex();
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
if (controller.isClosed) return;
controller.add(read());
changes = _orderChanges.stream
.where((changed) => changed == orderId)
.listen((_) => controller.add(read()));
},
onCancel: () async {
await changes?.cancel();
await controller.close();
},
);
return controller.stream;
}

/// Save or update any MostroMessage
Future<void> addMessage(String key, MostroMessage message) async {
final id = key;
try {
await _ensureIndex();
if (await hasItem(id)) return;
// Add metadata for easier querying
final Map<String, dynamic> dbMap = message.toJson();
message.timestamp ??= DateTime.now().millisecondsSinceEpoch;
dbMap['timestamp'] = message.timestamp;

await store.record(id).put(db, dbMap);
_indexAdd(message);
logger.i(
'Saved message of type ${message.action} with order id ${message.id}',
);
Expand All @@ -47,7 +112,11 @@ class MostroStorage extends BaseStorage<MostroMessage> {
/// Delete all messages
Future<void> deleteAllMessages() async {
try {
await _ensureIndex();
await deleteAll();
Comment thread
grunch marked this conversation as resolved.
Outdated
final orderIds = _byOrder.keys.toList();
_byOrder.clear();
orderIds.forEach(_notifyOrder);
logger.i('All messages deleted');
} catch (e, stack) {
logger.e('deleteAllMessages failed', error: e, stackTrace: stack);
Expand All @@ -57,9 +126,12 @@ class MostroStorage extends BaseStorage<MostroMessage> {

/// Delete all messages by Id regardless of type
Future<void> deleteAllMessagesByOrderId(String orderId) async {
await _ensureIndex();
await deleteWhere(
Filter.equals('id', orderId),
);
_byOrder.remove(orderId);
_notifyOrder(orderId);
}

/// Filter messages by payload type
Expand Down Expand Up @@ -106,63 +178,30 @@ class MostroStorage extends BaseStorage<MostroMessage> {

/// Get the latest message for an order, regardless of type
Future<MostroMessage?> getLatestMessageById(String orderId) async {
final finder = Finder(
filter: Filter.equals('id', orderId),
sortOrders: _getDefaultSort(),
limit: 1,
);

final snapshot = await store.findFirst(db, finder: finder);
if (snapshot != null) {
return MostroMessage.fromJson(snapshot.value);
}
return null;
await _ensureIndex();
return _latestFor(orderId);
}

/// Stream of the latest message for an order
Stream<MostroMessage?> watchLatestMessage(String orderId) {
final query = store.query(
finder: Finder(
filter: Filter.equals('id', orderId),
sortOrders: _getDefaultSort(),
limit: 1,
),
);

return query.onSnapshots(db).map((snaps) =>
snaps.isEmpty ? null : MostroMessage.fromJson(snaps.first.value));
}
Stream<MostroMessage?> watchLatestMessage(String orderId) =>
_watchOrder(orderId, () => _latestFor(orderId));

/// Stream of the latest message for an order whose payload is of type T
Stream<MostroMessage?> watchLatestMessageOfType<T>(String orderId) {
// Watch all messages for the orderId, sorted by timestamp descending
final query = store.query(
finder: Finder(
filter: Filter.equals('id', orderId),
sortOrders: _getDefaultSort(),
),
);
return query.onSnapshots(db).map((snaps) {
for (final snap in snaps) {
final msg = MostroMessage.fromJson(snap.value);
if (msg.payload is T) {
return msg;
Stream<MostroMessage?> watchLatestMessageOfType<T>(String orderId) =>
_watchOrder(orderId, () {
for (final message in _byOrder[orderId] ?? const <MostroMessage>[]) {
if (message.payload is T) return message;
}
}
return null;
});
}
return null;
});

// Use the same sorting across all methods that return lists of messages
List<SortOrder> _getDefaultSort() => [SortOrder('timestamp', false, true)];
/// Stream of all messages for an order (newest first)
Stream<List<MostroMessage>> watchAllMessages(String orderId) =>
_watchOrder(orderId, () => _historyFor(orderId));

/// Stream of all messages for an order
Stream<List<MostroMessage>> watchAllMessages(String orderId) {
return watch(
filter: Filter.equals('id', orderId),
sort: _getDefaultSort(),
);
}
// Sorting for the remaining DB-backed query (request-id lookups are
// transient one-offs during order creation and stay on Sembast).
List<SortOrder> _getDefaultSort() => [SortOrder('timestamp', false, true)];

Stream<MostroMessage?> watchByRequestId(int requestId) {
final query = store.query(
Expand All @@ -178,13 +217,7 @@ class MostroStorage extends BaseStorage<MostroMessage> {
}

Future<List<MostroMessage>> getAllMessagesForOrderId(String orderId) async {
final finder = Finder(
filter: Filter.equals('id', orderId),
sortOrders: [SortOrder('timestamp', false)]);

final snapshots = await store.find(db, finder: finder);
return snapshots
.map((snapshot) => MostroMessage.fromJson(snapshot.value))
.toList();
await _ensureIndex();
return _historyFor(orderId);
}
}
109 changes: 109 additions & 0 deletions test/data/repositories/mostro_storage_index_test.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
import 'package:flutter_test/flutter_test.dart';
import 'package:mostro_mobile/data/models/enums/action.dart';
import 'package:mostro_mobile/data/models/mostro_message.dart';
import 'package:mostro_mobile/data/repositories/mostro_storage.dart';
import 'package:sembast/sembast_memory.dart';

/// Every `orders`-store write used to wake one Sembast query listener per
/// OrderNotifier and per visible trade row, each re-filtering the whole
/// unindexed store on the UI isolate: O((notifiers + rows) × messages) per
/// incoming message. The storage now keeps a single in-memory index by order
/// id; watchers are served from it and Sembast remains the persistence
/// layer, warmed once per cold start.
void main() {
late MostroStorage storage;

MostroMessage message(String orderId, Action action, int timestamp) {
final m = MostroMessage(action: action, id: orderId);
m.timestamp = timestamp;
return m;
}

setUp(() async {
final db = await newDatabaseFactoryMemory().openDatabase('index_test.db');
storage = MostroStorage(db: db);
});

test('the latest-message watcher tracks adds for its order', () async {
final emissions = <MostroMessage?>[];
final sub = storage.watchLatestMessage('a').listen(emissions.add);
addTearDown(sub.cancel);
await Future<void>.delayed(const Duration(milliseconds: 20));

await storage.addMessage('k1', message('a', Action.newOrder, 1000));
await storage.addMessage('k2', message('a', Action.payInvoice, 2000));
await Future<void>.delayed(const Duration(milliseconds: 20));

expect(emissions.last!.action, Action.payInvoice);
});

test('a write for another order does not notify this watcher', () async {
await storage.addMessage('k1', message('a', Action.newOrder, 1000));
var emissions = 0;
final sub = storage.watchLatestMessage('a').skip(1).listen((_) {
emissions++;
});
addTearDown(sub.cancel);
await Future<void>.delayed(const Duration(milliseconds: 20));

await storage.addMessage('k2', message('b', Action.newOrder, 2000));
await Future<void>.delayed(const Duration(milliseconds: 20));

expect(emissions, 0,
reason: 'per-order streams must be demultiplexed in memory');
});

test('history is served newest-first and per order', () async {
await storage.addMessage('k1', message('a', Action.newOrder, 1000));
await storage.addMessage('k2', message('a', Action.payInvoice, 3000));
await storage.addMessage('k3', message('b', Action.newOrder, 2000));

final history = await storage.getAllMessagesForOrderId('a');

expect(history.map((m) => m.action),
[Action.payInvoice, Action.newOrder]);
expect(await storage.getLatestMessageById('b'),
isA<MostroMessage>().having((m) => m.id, 'id', 'b'));
});

test('a cold start warms the index from disk', () async {
await storage.addMessage('k1', message('a', Action.newOrder, 1000));

// New storage over the same database: what a restart looks like.
final restarted = MostroStorage(db: storage.db);
final latest = await restarted.getLatestMessageById('a');

expect(latest!.action, Action.newOrder);
expect(restarted.debugIndexSize, greaterThan(0));
});

test('deleting an order clears it from index and watchers', () async {
await storage.addMessage('k1', message('a', Action.newOrder, 1000));
final emissions = <MostroMessage?>[];
final sub = storage.watchLatestMessage('a').listen(emissions.add);
addTearDown(sub.cancel);
await Future<void>.delayed(const Duration(milliseconds: 20));

await storage.deleteAllMessagesByOrderId('a');
await Future<void>.delayed(const Duration(milliseconds: 20));

expect(emissions.last, isNull);
expect(await storage.getAllMessagesForOrderId('a'), isEmpty);
});

test('duplicate keys are ignored without notifying watchers twice',
() async {
await storage.addMessage('k1', message('a', Action.newOrder, 1000));
var emissions = 0;
final sub = storage.watchLatestMessage('a').skip(1).listen((_) {
emissions++;
});
addTearDown(sub.cancel);
await Future<void>.delayed(const Duration(milliseconds: 20));

await storage.addMessage('k1', message('a', Action.newOrder, 1000));
await Future<void>.delayed(const Duration(milliseconds: 20));

expect(emissions, 0);
});
}
Loading