// lib/network/kafkahandler.dart import 'dart:async'; import 'dart:convert'; import 'dart:developer'; import 'package:stomp_dart_client/stomp.dart'; import 'package:stomp_dart_client/stomp_config.dart'; import 'package:stomp_dart_client/stomp_frame.dart'; import '../config/constants/app_constant.dart'; import '../providers/utils_singleton.dart'; class KafkaHandler { KafkaHandler._(); static final KafkaHandler instance = KafkaHandler._(); StompClient? _client; final Map _subscriptions = {}; // expose a broadcast stream of geojson FeatureCollections // final _ride_heat_sc = StreamController>.broadcast(); // Stream> get rideHeatStream => _ride_heat_sc.stream; final _driver_location_sc = StreamController>.broadcast(); Stream> get driverLocationStream => _driver_location_sc.stream; bool get connected => _client?.connected ?? false; void connect() { if (connected) return; final t = utils.token ?? 'Invalid_token'; // SockJS requires http/https. Add token in query to satisfy your HttpHandshakeInterceptor. final url = '$socketUrlNative2?token=${Uri.encodeQueryComponent('Bearer $t')}'; log("socket url : ${url}"); _client = StompClient( config: StompConfig.SockJS( url: url, onConnect: _onConnect, onWebSocketError: (e) => log('[stomp] ws error: $e'), onDisconnect: (_) => log('[stomp] disconnected'), onStompError: (f) => log('[stomp] stomp error: ${f.body}'), stompConnectHeaders: {'Authorization': 'Bearer $t'}, webSocketConnectHeaders: {'Authorization': 'Bearer $t'}, heartbeatIncoming: const Duration(seconds: 10), heartbeatOutgoing: const Duration(seconds: 10), ), )..activate(); } void _onConnect(StompFrame _) { subscribeToDriverLocations(); } void subscribeToDriverLocations() { const destination = '/topic/public/driver-locations'; // ✨ FIX: Always clean up any previous subscription for this topic first. // This makes the function safe to call multiple times. unsubscribeFromDriverLocations(); if (_client == null || !connected) { log('[stomp] Cannot subscribe, client is not connected. ${_client == null} ou ${!connected}'); return; } final unsubscribeFn = _client!.subscribe( destination: destination, callback: (frame) { try { if (frame.body == null) return; final data = jsonDecode(frame.body!); if (data is Map && data.containsKey('data')) { final payload = data['data'] as Map; _driver_location_sc.add(payload); } } catch (e) { log('[stomp] Error processing driver-location message: $e'); } }, ); _subscriptions[destination] = unsubscribeFn; log('[stomp] Subscribed to $destination'); } void unsubscribeFromDriverLocations() { const destination = '/topic/public/driver-locations'; final unsubscribeFn = _subscriptions.remove(destination); if (unsubscribeFn != null) { unsubscribeFn(); log('[stomp] Unsubscribed from $destination'); } } void dispose() { // ✨ NEW: Unsubscribe from all dynamic topics before deactivating _subscriptions.forEach((key, unsubscribeFn) => unsubscribeFn()); _subscriptions.clear(); _client?.deactivate(); // _ride_heat_sc.close(); _driver_location_sc.close(); // ✨ NEW: Close the new stream controller } }