Skip to content

Commit 11a5e78

Browse files
committed
Week 3 Day 20: Add robust sync engine with batching, deduplication, and termination guarantees
1 parent c8b4e70 commit 11a5e78

2 files changed

Lines changed: 349 additions & 13 deletions

File tree

mobile_app/lib/services/p2p_service.dart

Lines changed: 128 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,19 @@ import '../models/network_envelope.dart';
88
import 'message_cache.dart';
99
import 'gossip_log_service.dart';
1010

11+
/// Represents an active sync session with a specific peer (Day-20)
12+
class SyncSession {
13+
final String peerId;
14+
final Set<String> requestedIds = {};
15+
final Set<String> receivedIds = {};
16+
final Set<String> sentIds = {};
17+
int iterationCount = 0;
18+
bool syncComplete = false;
19+
DateTime? lastRequestTime;
20+
21+
SyncSession(this.peerId);
22+
}
23+
1124
/// Maximum number of incidents per sync_response batch to avoid network flooding.
1225
const int syncBatchSize = 50;
1326

@@ -56,7 +69,16 @@ class P2PService {
5669
/// peer_id → sync_completed (boolean)
5770
final Map<String, bool> _peerSyncState = {};
5871

59-
P2PService({required this.hostUrl}) {
72+
/// Active sync sessions with peers for Day-20 HEAD-based sync
73+
final Map<String, SyncSession> _syncSessions = {};
74+
75+
/// Timer for Day-20 retry strategy
76+
Timer? _retryTimer;
77+
78+
final int port;
79+
80+
P2PService({required this.hostUrl, this.port = 7000}) {
81+
_retryTimer = Timer.periodic(const Duration(milliseconds: 1000), (_) => _checkRetries());
6082
// Day-18: Subscribe to GossipLog valid messages and route them
6183
_gossipLogSub = _gossipLog.validMessages.listen((envelope) {
6284
_routeValidEnvelope(envelope);
@@ -88,7 +110,7 @@ class P2PService {
88110
/// Connects to the local daemon's WebSocket endpoint for receiving messages
89111
void connect() {
90112
final wsBase = hostUrl.replaceFirst('http', 'ws');
91-
final wsUrl = Uri.parse('$wsBase:7000/events');
113+
final wsUrl = Uri.parse('$wsBase:$port/events');
92114

93115
try {
94116
_channel = WebSocketChannel.connect(wsUrl);
@@ -173,7 +195,7 @@ class P2PService {
173195
final missingDeps = _gossipLog.receive(envelope);
174196
if (missingDeps.isNotEmpty) {
175197
debugPrint('[P2P] MISSING_DETECTED: requesting missing deps $missingDeps');
176-
_sendMessageRequest(missingDeps);
198+
_processMissingDeps(envelope.originPeer, missingDeps);
177199
}
178200
} catch (e) {
179201
debugPrint('[P2P] Failed to parse incoming message: $e');
@@ -188,8 +210,78 @@ class P2PService {
188210

189211
final missing = _gossipLog.findMissingMessages(peerHeads);
190212
if (missing.isNotEmpty) {
191-
debugPrint('[P2P] MISSING_DETECTED: requesting ${missing.length} messages');
192-
_sendMessageRequest(missing);
213+
debugPrint('[P2P] MISSING_DETECTED: requesting ${missing.length} messages (from heads_exchange)');
214+
_processMissingDeps(envelope.originPeer, missing);
215+
} else {
216+
final session = _syncSessions.putIfAbsent(envelope.originPeer, () => SyncSession(envelope.originPeer));
217+
if (!session.syncComplete) {
218+
session.syncComplete = true;
219+
debugPrint('[P2P] SYNC_COMPLETED: No missing messages for ${envelope.originPeer}');
220+
}
221+
}
222+
}
223+
224+
void _processMissingDeps(String peerId, List<String> missing) {
225+
if (missing.isEmpty) return;
226+
227+
final session = _syncSessions.putIfAbsent(peerId, () => SyncSession(peerId));
228+
229+
if (session.syncComplete) {
230+
// Re-activate if there are new missing messages discovered
231+
session.syncComplete = false;
232+
}
233+
234+
if (session.iterationCount >= 10) {
235+
debugPrint('[P2P] SYNC_TERMINATED: Max iterations reached for $peerId');
236+
session.syncComplete = true;
237+
return;
238+
}
239+
240+
final toRequest = <String>[];
241+
for (final id in missing) {
242+
if (session.requestedIds.contains(id)) {
243+
debugPrint('[P2P] REQUEST_SKIPPED_DUPLICATE: $id already requested from $peerId');
244+
} else {
245+
toRequest.add(id);
246+
session.requestedIds.add(id);
247+
}
248+
}
249+
250+
if (toRequest.isEmpty) {
251+
final pendingToRequest = session.requestedIds.difference(session.receivedIds);
252+
if (pendingToRequest.isEmpty) {
253+
session.syncComplete = true;
254+
debugPrint('[P2P] SYNC_COMPLETED: No new missing messages for $peerId');
255+
}
256+
return;
257+
}
258+
259+
session.iterationCount++;
260+
session.lastRequestTime = DateTime.now();
261+
debugPrint('[P2P] SYNC_PROGRESS: Iteration ${session.iterationCount} for $peerId');
262+
263+
_sendMessageRequest(toRequest, peerId);
264+
}
265+
266+
void _checkRetries() {
267+
final now = DateTime.now();
268+
for (final session in _syncSessions.values) {
269+
if (!session.syncComplete && session.lastRequestTime != null) {
270+
if (now.difference(session.lastRequestTime!) > const Duration(milliseconds: 1000)) {
271+
final pendingToRequest = session.requestedIds.difference(session.receivedIds).toList();
272+
if (pendingToRequest.isNotEmpty) {
273+
if (session.iterationCount >= 10) {
274+
debugPrint('[P2P] SYNC_TERMINATED: Max iterations reached for ${session.peerId}');
275+
session.syncComplete = true;
276+
continue;
277+
}
278+
session.iterationCount++;
279+
session.lastRequestTime = now;
280+
debugPrint('[P2P] SYNC_PROGRESS: Iteration ${session.iterationCount} for ${session.peerId} (Retry)');
281+
_sendMessageRequest(pendingToRequest, session.peerId);
282+
}
283+
}
284+
}
193285
}
194286
}
195287

@@ -198,16 +290,31 @@ class P2PService {
198290
.map((e) => e.toString())
199291
.toList();
200292

201-
debugPrint('[P2P] MESSAGE_REQUEST_RECEIVED: for ${requestedIds.length} messages');
293+
final peerId = envelope.originPeer;
294+
final session = _syncSessions.putIfAbsent(peerId, () => SyncSession(peerId));
295+
296+
debugPrint('[P2P] MESSAGE_REQUEST_RECEIVED: for ${requestedIds.length} messages from $peerId');
202297

298+
// Day-20 Part 7: Topological sort happens inside fetchAndSortMessages
203299
final messagesToSend = _gossipLog.fetchAndSortMessages(requestedIds);
204-
if (messagesToSend.isNotEmpty) {
205-
_sendMessageResponse(messagesToSend);
300+
301+
// Day-20 Part 6: Deduplicate responses
302+
final toSend = <NetworkEnvelope>[];
303+
for (final msg in messagesToSend) {
304+
if (!session.sentIds.contains(msg.msgId)) {
305+
toSend.add(msg);
306+
session.sentIds.add(msg.msgId);
307+
}
308+
}
309+
310+
if (toSend.isNotEmpty) {
311+
_sendMessageResponse(toSend, peerId);
206312
}
207313
}
208314

209315
void _handleMessageResponse(NetworkEnvelope envelope) {
210316
debugPrint('[P2P] MESSAGE_RESPONSE_RECEIVED: from ${envelope.originPeer}');
317+
final session = _syncSessions.putIfAbsent(envelope.originPeer, () => SyncSession(envelope.originPeer));
211318
final messagesList = envelope.payload['messages'] as List<dynamic>? ?? [];
212319

213320
for (final msgData in messagesList) {
@@ -218,6 +325,7 @@ class P2PService {
218325
: Map<String, dynamic>.from(msgData as Map);
219326

220327
final msgEnvelope = NetworkEnvelope.fromJson(msgMap);
328+
session.receivedIds.add(msgEnvelope.msgId);
221329

222330
// Add to deduplication cache to prevent re-processing
223331
if (_messageCache.isDuplicate(msgEnvelope.msgId)) {
@@ -229,15 +337,21 @@ class P2PService {
229337

230338
if (missingDeps.isNotEmpty) {
231339
debugPrint('[P2P] MISSING_DETECTED: requesting missing deps $missingDeps');
232-
_sendMessageRequest(missingDeps);
340+
_processMissingDeps(envelope.originPeer, missingDeps);
233341
}
234342
} catch (e) {
235343
debugPrint('[P2P] Failed to parse message in response: $e');
236344
}
237345
}
346+
347+
final pendingToRequest = session.requestedIds.difference(session.receivedIds);
348+
if (pendingToRequest.isEmpty && !session.syncComplete) {
349+
session.syncComplete = true;
350+
debugPrint('[P2P] SYNC_COMPLETED: No new missing messages for ${envelope.originPeer}');
351+
}
238352
}
239353

240-
Future<void> _sendMessageRequest(List<String> requestedIds) async {
354+
Future<void> _sendMessageRequest(List<String> requestedIds, [String? targetPeer]) async {
241355
if (requestedIds.isEmpty) return;
242356

243357
final envelope = NetworkEnvelope(
@@ -250,15 +364,15 @@ class P2PService {
250364
},
251365
);
252366

253-
debugPrint('[P2P] MESSAGE_REQUEST_SENT: requesting ids $requestedIds');
367+
debugPrint('[P2P] BATCH_REQUEST_SENT: requesting ${requestedIds.length} ids');
254368

255369
// Add to our own dedup cache to prevent self-echo
256370
_messageCache.isDuplicate(envelope.msgId);
257371

258372
await _sendEnvelopeHttp(envelope);
259373
}
260374

261-
Future<void> _sendMessageResponse(List<NetworkEnvelope> messages) async {
375+
Future<void> _sendMessageResponse(List<NetworkEnvelope> messages, [String? targetPeer]) async {
262376
final payloadMessages = messages.map((m) => m.toJson()).toList();
263377

264378
final envelope = NetworkEnvelope(
@@ -419,7 +533,7 @@ class P2PService {
419533

420534
/// Sends a NetworkEnvelope via HTTP POST to the daemon
421535
Future<bool> _sendEnvelopeHttp(NetworkEnvelope envelope) async {
422-
final uri = Uri.parse('$hostUrl:7000/broadcast');
536+
final uri = Uri.parse('$hostUrl:$port/broadcast');
423537

424538
try {
425539
final response = await http.post(
@@ -485,6 +599,7 @@ class P2PService {
485599
}
486600

487601
void dispose() {
602+
_retryTimer?.cancel();
488603
_gossipLogSub?.cancel();
489604
_gossipLog.dispose();
490605
_channel?.sink.close();

0 commit comments

Comments
 (0)