Buffer and load zappers

This commit is contained in:
franzap
2025-01-05 13:42:24 -03:00
parent df747a64ab
commit a15509f246
3 changed files with 71 additions and 9 deletions
+7
View File
@@ -105,6 +105,13 @@ mixin NostrAdapter<T extends DataModelMixin<T>> on Adapter<T> {
return r.isNotEmpty;
}
Iterable<String> existingIds(Iterable<Object> 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<List<T>> findAll(
{bool? remote,
+31
View File
@@ -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<List<T>> bufferByTime<T>(Stream<T> source, Duration duration) {
final controller = StreamController<List<T>>();
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
}
+33 -9
View File
@@ -63,7 +63,7 @@ class ZapReceipts extends HookConsumerWidget {
class ZapReceiptsNotifier extends StateNotifier<AsyncValue<Set<ZapReceipt>>> {
Ref ref;
final FileMetadata fileMetadata;
StreamSubscription<Nip01Event>? sub;
StreamSubscription<List<Nip01Event>>? sub;
ZapReceiptsNotifier(this.ref, this.fileMetadata) : super(AsyncLoading()) {
final adapter = ref.users.nostrAdapter;
@@ -78,15 +78,15 @@ class ZapReceiptsNotifier extends StateNotifier<AsyncValue<Set<ZapReceipt>>> {
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<AsyncValue<Set<ZapReceipt>>> {
],
);
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<void> _loadZappers(Iterable<ZapReceipt> 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();