Skip to content

Commit 7a69a80

Browse files
authored
Merge pull request #254 from relaystr/feat-concurrent-event-verifier-stream
Feat concurrent event verifier stream
2 parents 0c83055 + 16d7ed7 commit 7a69a80

5 files changed

Lines changed: 415 additions & 16 deletions

File tree

packages/ndk/lib/domain_layer/usecases/requests/verify_event_stream.dart

Lines changed: 43 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5,25 +5,54 @@ import '../../repositories/event_verifier.dart';
55
class VerifyEventStream {
66
final Stream<Nip01Event> unverifiedStreamInput;
77
final EventVerifier eventVerifier;
8+
final int maxConcurrent;
9+
810
VerifyEventStream({
911
required this.unverifiedStreamInput,
1012
required this.eventVerifier,
13+
this.maxConcurrent = 100,
1114
});
1215

1316
Stream<Nip01Event> call() {
14-
return unverifiedStreamInput
15-
.asyncMap<Nip01Event>((data) async {
16-
final valid = await eventVerifier.verify(data);
17-
data.validSig = valid; // assign validity
18-
19-
if (!valid) {
20-
Logger.log
21-
.w('WARNING: Event with id ${data.id} has invalid signature');
22-
}
23-
24-
return data;
25-
})
26-
.where((event) => event.validSig == true) // Filter out invalid events
27-
.asBroadcastStream();
17+
return _verifyInParallel().asBroadcastStream();
18+
}
19+
20+
Stream<Nip01Event> _verifyInParallel() async* {
21+
final buffer = <Future<Nip01Event?>>[];
22+
23+
await for (final event in unverifiedStreamInput) {
24+
// Start verification without waiting
25+
final future = _verifyEvent(event);
26+
buffer.add(future);
27+
28+
// Once we hit max concurrent, wait for the first one to complete
29+
if (buffer.length >= maxConcurrent) {
30+
final verified = await buffer.first;
31+
buffer.removeAt(0);
32+
33+
if (verified != null && verified.validSig == true) {
34+
yield verified;
35+
}
36+
}
37+
}
38+
39+
// Process remaining events in buffer
40+
final remaining = await Future.wait(buffer);
41+
for (final event in remaining) {
42+
if (event != null && event.validSig == true) {
43+
yield event;
44+
}
45+
}
46+
}
47+
48+
Future<Nip01Event?> _verifyEvent(Nip01Event data) async {
49+
final valid = await eventVerifier.verify(data);
50+
data.validSig = valid;
51+
52+
if (!valid) {
53+
Logger.log.w('WARNING: Event with id ${data.id} has invalid signature');
54+
}
55+
56+
return data;
2857
}
2958
}
Lines changed: 207 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,207 @@
1+
import 'package:ndk/domain_layer/entities/nip_01_event.dart';
2+
import 'package:ndk/domain_layer/usecases/requests/verify_event_stream.dart';
3+
import 'package:test/test.dart';
4+
5+
import '../mocks/mock_event_verifier.dart';
6+
7+
void main() {
8+
group('VerifyEventStream', () {
9+
late MockEventVerifier mockVerifier;
10+
11+
setUp(() {
12+
mockVerifier = MockEventVerifier(result: true);
13+
});
14+
15+
Nip01Event createMockEvent(String id) {
16+
final event = Nip01Event(
17+
pubKey: 'pubkey$id',
18+
createdAt: DateTime.now().millisecondsSinceEpoch ~/ 1000,
19+
kind: 1,
20+
tags: [],
21+
content: 'content$id',
22+
);
23+
return event;
24+
}
25+
26+
test('should verify and yield valid events', () async {
27+
final events = [
28+
createMockEvent('1'),
29+
createMockEvent('2'),
30+
createMockEvent('3'),
31+
];
32+
final inputStream = Stream.fromIterable(events);
33+
34+
final verifyStream = VerifyEventStream(
35+
unverifiedStreamInput: inputStream,
36+
eventVerifier: mockVerifier,
37+
);
38+
39+
final results = await verifyStream().toList();
40+
41+
expect(results.length, equals(3));
42+
expect(results.every((e) => e.validSig == true), isTrue);
43+
expect(results[0].content, equals('content1'));
44+
expect(results[1].content, equals('content2'));
45+
expect(results[2].content, equals('content3'));
46+
});
47+
48+
test('should verify and yield valid events with small buffer', () async {
49+
final events = [
50+
createMockEvent('1'),
51+
createMockEvent('2'),
52+
createMockEvent('3'),
53+
createMockEvent('4'),
54+
createMockEvent('5'),
55+
createMockEvent('6'),
56+
];
57+
final inputStream = Stream.fromIterable(events);
58+
59+
final verifyStream = VerifyEventStream(
60+
unverifiedStreamInput: inputStream,
61+
eventVerifier: mockVerifier,
62+
maxConcurrent: 2,
63+
);
64+
65+
final results = await verifyStream().toList();
66+
67+
expect(results.length, equals(6));
68+
expect(results.every((e) => e.validSig == true), isTrue);
69+
expect(results[0].content, equals('content1'));
70+
expect(results[1].content, equals('content2'));
71+
expect(results[2].content, equals('content3'));
72+
expect(results[3].content, equals('content4'));
73+
expect(results[4].content, equals('content5'));
74+
expect(results[5].content, equals('content6'));
75+
});
76+
77+
test('should filter out invalid events', () async {
78+
mockVerifier = MockEventVerifier(result: false);
79+
final events = [
80+
createMockEvent('1'),
81+
createMockEvent('2'),
82+
];
83+
final inputStream = Stream.fromIterable(events);
84+
85+
final verifyStream = VerifyEventStream(
86+
unverifiedStreamInput: inputStream,
87+
eventVerifier: mockVerifier,
88+
);
89+
90+
final results = await verifyStream().toList();
91+
92+
expect(results.length, equals(0));
93+
});
94+
95+
test('should handle empty stream', () async {
96+
final inputStream = Stream<Nip01Event>.fromIterable([]);
97+
98+
final verifyStream = VerifyEventStream(
99+
unverifiedStreamInput: inputStream,
100+
eventVerifier: mockVerifier,
101+
);
102+
103+
final results = await verifyStream().toList();
104+
105+
expect(results.length, equals(0));
106+
});
107+
108+
test('should handle stream with small maxConcurrent', () async {
109+
final events = List.generate(5, (i) => createMockEvent('$i'));
110+
final inputStream = Stream.fromIterable(events);
111+
112+
final verifyStream = VerifyEventStream(
113+
unverifiedStreamInput: inputStream,
114+
eventVerifier: mockVerifier,
115+
maxConcurrent: 3,
116+
);
117+
118+
final results = await verifyStream().toList();
119+
120+
expect(results.length, equals(5));
121+
});
122+
123+
test('should return broadcast stream', () async {
124+
final events = [createMockEvent('1')];
125+
final inputStream = Stream.fromIterable(events);
126+
127+
final verifyStream = VerifyEventStream(
128+
unverifiedStreamInput: inputStream,
129+
eventVerifier: mockVerifier,
130+
);
131+
132+
final stream = verifyStream();
133+
134+
expect(stream.isBroadcast, isTrue);
135+
});
136+
137+
test('should allow multiple listeners on broadcast stream', () async {
138+
final events = [
139+
createMockEvent('1'),
140+
createMockEvent('2'),
141+
];
142+
final inputStream = Stream.fromIterable(events);
143+
144+
final verifyStream = VerifyEventStream(
145+
unverifiedStreamInput: inputStream,
146+
eventVerifier: mockVerifier,
147+
);
148+
149+
final stream = verifyStream();
150+
151+
final results1Future = stream.toList();
152+
final results2Future = stream.toList();
153+
154+
final results1 = await results1Future;
155+
final results2 = await results2Future;
156+
157+
expect(results1.length, equals(2));
158+
expect(results2.length, equals(2));
159+
});
160+
161+
test('should process remaining buffer after stream ends', () async {
162+
final events = List.generate(5, (i) => createMockEvent('$i'));
163+
final inputStream = Stream.fromIterable(events);
164+
165+
final verifyStream = VerifyEventStream(
166+
unverifiedStreamInput: inputStream,
167+
eventVerifier: mockVerifier,
168+
maxConcurrent: 10, // Higher than event count
169+
);
170+
171+
final results = await verifyStream().toList();
172+
173+
expect(results.length, equals(5));
174+
expect(results.every((e) => e.validSig == true), isTrue);
175+
});
176+
177+
test('should set validSig property correctly', () async {
178+
final events = [createMockEvent('1')];
179+
final inputStream = Stream.fromIterable(events);
180+
181+
final verifyStream = VerifyEventStream(
182+
unverifiedStreamInput: inputStream,
183+
eventVerifier: mockVerifier,
184+
);
185+
186+
final results = await verifyStream().toList();
187+
188+
expect(results.length, equals(1));
189+
expect(results[0].validSig, equals(true));
190+
});
191+
192+
test('should handle large number of events', () async {
193+
final events = List.generate(100, (i) => createMockEvent('$i'));
194+
final inputStream = Stream.fromIterable(events);
195+
196+
final verifyStream = VerifyEventStream(
197+
unverifiedStreamInput: inputStream,
198+
eventVerifier: mockVerifier,
199+
maxConcurrent: 50,
200+
);
201+
202+
final results = await verifyStream().toList();
203+
204+
expect(results.length, equals(100));
205+
});
206+
});
207+
}

packages/rust_verifier/README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,7 @@ you can read more about it in the [flutter docs](https://docs.flutter.dev/platfo
4141

4242
When enabled the verification is done in a background thread/worker.
4343

44-
## How to build
44+
## How to build the rust_verifier from source [library development]
4545

4646
### normal build
4747

packages/sample-app/lib/main.dart

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import 'package:ndk/ndk.dart';
55
import 'package:ndk_demo/accounts_page.dart';
66
import 'package:ndk_demo/blossom_page.dart';
77
import 'package:ndk_demo/nwc_page.dart';
8+
import 'package:ndk_demo/query_performance.dart';
89
import 'package:ndk_demo/relays_page.dart';
910
import 'package:ndk_demo/verifiers_performance.dart';
1011
import 'package:ndk_demo/zaps_page.dart';
@@ -135,6 +136,7 @@ class _MyHomePageState extends State<MyHomePage>
135136
const Tab(text: nwcTabName),
136137
const Tab(text: "Blossom"),
137138
const Tab(text: 'Verifiers'),
139+
const Tab(text: 'Query Performance'),
138140
// Conditionally add Amber tab if it's part of the design
139141
// For a fixed length of 6, ensure this list matches.
140142
// Example: if Amber is the 6th tab:
@@ -158,7 +160,7 @@ class _MyHomePageState extends State<MyHomePage>
158160
// The main change is how _tabPages is constructed in build() to pass the callback.
159161

160162
_tabController = TabController(
161-
length: 6,
163+
length: 7,
162164
vsync:
163165
this); // Fixed length to 5 (Accounts, Metadata, Relays, NWC, Blossom)
164166
_tabController.addListener(() {
@@ -246,6 +248,7 @@ class _MyHomePageState extends State<MyHomePage>
246248
const Tab(text: nwcTabName),
247249
const Tab(text: "Blossom"),
248250
const Tab(text: 'Verifiers'),
251+
const Tab(text: 'Query Performance'),
249252
// Amber tab removed
250253
];
251254

@@ -256,6 +259,7 @@ class _MyHomePageState extends State<MyHomePage>
256259
const NwcPage(),
257260
BlossomMediaPage(ndk: ndk),
258261
VerifiersPerformancePage(ndk: ndk),
262+
QueryPerformancePage(ndk: ndk),
259263
// AmberPage removed
260264
];
261265

0 commit comments

Comments
 (0)