From a15509f246ea924cdb1cbfcd4c1402b1dda04d39 Mon Sep 17 00:00:00 2001 From: franzap <_@franzap.com> Date: Sun, 5 Jan 2025 13:42:24 -0300 Subject: [PATCH] Buffer and load zappers --- lib/models/nostr_adapter.dart | 7 ++++++ lib/utils/extensions.dart | 31 ++++++++++++++++++++++++++ lib/widgets/zap_receipts.dart | 42 +++++++++++++++++++++++++++-------- 3 files changed, 71 insertions(+), 9 deletions(-) diff --git a/lib/models/nostr_adapter.dart b/lib/models/nostr_adapter.dart index 8057da3..729dbe2 100644 --- a/lib/models/nostr_adapter.dart +++ b/lib/models/nostr_adapter.dart @@ -105,6 +105,13 @@ mixin NostrAdapter> on Adapter { return r.isNotEmpty; } + Iterable existingIds(Iterable ids) { + final r = db.select( + 'SELECT id FROM _keys WHERE id in (${ids.map((e) => '?').join(',')})', + ids.toList()); + return r.map((e) => e['id']); + } + @override Future> findAll( {bool? remote, diff --git a/lib/utils/extensions.dart b/lib/utils/extensions.dart index ab5185f..0049e5d 100644 --- a/lib/utils/extensions.dart +++ b/lib/utils/extensions.dart @@ -1,3 +1,4 @@ +import 'dart:async'; import 'dart:math'; import 'package:dart_emoji/dart_emoji.dart'; @@ -150,3 +151,33 @@ extension TextExt on Text { } const kZapstoreAppIdentifier = 'dev.zapstore.app'; + +// stream utils + +Stream> bufferByTime(Stream source, Duration duration) { + final controller = StreamController>(); + + final buffer = []; + + Timer? timer; + + // Listen to the source stream + source.listen((data) { + buffer.add(data); + + // If there's no active timer, start one + timer ??= Timer(duration, () { + controller.add(List.from(buffer)); // Emit buffered data + buffer.clear(); // Clear the buffer + timer = null; // Reset the timer + }); + }, onDone: () { + // Emit remaining items in the buffer when the source is done + if (buffer.isNotEmpty) { + controller.add(List.from(buffer)); + } + controller.close(); // Close the controller + }); + + return controller.stream; // Return the buffered stream +} diff --git a/lib/widgets/zap_receipts.dart b/lib/widgets/zap_receipts.dart index 6243824..cbb9243 100644 --- a/lib/widgets/zap_receipts.dart +++ b/lib/widgets/zap_receipts.dart @@ -63,7 +63,7 @@ class ZapReceipts extends HookConsumerWidget { class ZapReceiptsNotifier extends StateNotifier>> { Ref ref; final FileMetadata fileMetadata; - StreamSubscription? sub; + StreamSubscription>? sub; ZapReceiptsNotifier(this.ref, this.fileMetadata) : super(AsyncLoading()) { final adapter = ref.users.nostrAdapter; @@ -78,15 +78,15 @@ class ZapReceiptsNotifier extends StateNotifier>> { return; } - final localReceipts = ref.zapReceipts.zapReceiptAdapter + final localZapReceipts = ref.zapReceipts.zapReceiptAdapter .findByRecipient( pubkey: developerPubkey, eventId: fileMetadata.id!.toString()) .toSet(); - state = AsyncData(localReceipts); + _loadZappers(localZapReceipts); + state = AsyncData(localZapReceipts); - print('querying for $developerPubkey and id ${fileMetadata.event.id}'); - - final latestReceiptTimestamp = localReceipts + // NOTE: ideally this caching stuff should be handled by purplebase + final latestReceiptTimestamp = localZapReceipts .sortedBy((z) => z.event.createdAt) .lastOrNull ?.event @@ -106,12 +106,36 @@ class ZapReceiptsNotifier extends StateNotifier>> { ], ); - sub = receiptsResponse.stream.listen((data) { - final model = ZapReceipt.fromJson(data.toJson()).init().saveLocal(); - state = AsyncData({if (state.hasValue) ...state.value!, model}); + sub = bufferByTime(receiptsResponse.stream, Duration(seconds: 1)) + .listen((data) async { + final zapReceipts = + data.map((r) => ZapReceipt.fromJson(r.toJson()).init().saveLocal()); + await _loadZappers(zapReceipts); + state = AsyncData({if (state.hasValue) ...state.value!, ...zapReceipts}); }); } + // Load senders - this tedious work will be handled by purplebase at some point + Future _loadZappers(Iterable receipts) async { + // If we do not have the users locally, trigger a remote fetch + final receiptsPubkeys = receipts.map((r) => r.senderPubkey).toSet(); + final existingPubkeysLocally = + ref.users.nostrAdapter.existingIds(receiptsPubkeys).toSet(); + final missingPubkeysLocally = + receiptsPubkeys.difference(existingPubkeysLocally); + if (missingPubkeysLocally.isNotEmpty) { + final fetchedSenders = + await ref.users.findAll(params: {'authors': missingPubkeysLocally}); + final stillMissingPubkeys = missingPubkeysLocally + .toSet() + .difference(fetchedSenders.map((s) => s.id!.toString()).toSet()); + for (final pubkey in stillMissingPubkeys) { + // If relays do not have it, create a dummy local user + User.fromPubkey(pubkey).init().saveLocal(); + } + } + } + @override void dispose() { sub?.cancel();