import 'dart:async'; import 'dart:developer' as dev; import 'package:pubnub/pubnub.dart'; import 'package:connectivity_plus/connectivity_plus.dart'; import '../../config/constants/app_constant.dart'; import '../../network/kafkahandler.dart'; import '../../providers/utils_singleton.dart'; class UserContext { final String username; UserContext({required this.username}); } class PubNubManager { // CRITICAL FIX: True singleton pattern - NEVER dispose during app lifecycle static PubNubManager? _instance; static PubNubManager get instance { _instance ??= PubNubManager._(); return _instance!; } PubNubManager._(); StreamSubscription? _kafkaLocationSubscription; // Core PubNub components PubNub? _pubnub; Subscription? _subscription; // Timers and subscriptions Timer? _heartbeatTimer; Timer? _reconnectTimer; StreamSubscription? _connectivitySubscription; Timer? _cleanupTimer; Timer? _subscriptionWatchdog; // CRITICAL FIX: Stream controller that NEVER gets disposed during app lifecycle StreamController>? _rideUpdateController; bool _streamControllerCreated = false; Stream> get rideUpdateStream { _ensureStreamController(); return _rideUpdateController!.stream; } // State management - REMOVED _isDisposed completely bool _isInitialized = false; bool _isConnected = false; String? _currentUsername; int _reconnectAttempts = 0; // More conservative settings static const int _maxReconnectAttempts = 3; static const Duration _reconnectBaseDelay = Duration(seconds: 5); static const Duration _heartbeatInterval = Duration(seconds: 45); static const Duration _subscriptionWatchdogInterval = Duration(seconds: 60); // Message deduplication final Set _processedMessages = {}; static const int _maxProcessedMessages = 100; // Track subscription health DateTime? _lastMessageReceived; int _messageCount = 0; bool _subscriptionHealthy = true; // CRITICAL FIX: Stream controller created once and NEVER disposed void _ensureStreamController() { if (!_streamControllerCreated) { _rideUpdateController = StreamController>.broadcast( onListen: () { dev.log( '๐ŸŽง PubNub stream listener added (total listeners: ${_rideUpdateController?.hasListener})'); _startSubscriptionWatchdog(); }, onCancel: () => dev.log('๐ŸŽง PubNub stream listener cancelled'), ); _streamControllerCreated = true; _startMessageCleanup(); dev.log('โœ… Stream controller created permanently for app lifecycle'); } } void _startSubscriptionWatchdog() { _subscriptionWatchdog?.cancel(); _subscriptionWatchdog = Timer.periodic(_subscriptionWatchdogInterval, (timer) { final now = DateTime.now(); if (_lastMessageReceived != null) { final silentDuration = now.difference(_lastMessageReceived!); if (silentDuration.inMinutes > 5 && _isConnected) { dev.log( 'โš ๏ธ Subscription silent for ${silentDuration.inMinutes} minutes, forcing reconnect...'); _subscriptionHealthy = false; forceReconnect(); } } dev.log( '๐Ÿ• Subscription watchdog: ${_messageCount} messages received, healthy: $_subscriptionHealthy'); }); } void _startMessageCleanup() { _cleanupTimer?.cancel(); _cleanupTimer = Timer.periodic(const Duration(minutes: 2), (timer) { if (_processedMessages.length > _maxProcessedMessages) { _processedMessages.clear(); dev.log('๐Ÿงน Cleaned processed messages cache'); } }); } // CRITICAL FIX: Allow re-initialization even if previously "disposed" Future init(UserContext userContext) async { dev.log('๐Ÿ”„ PubNub Manager init called for ${userContext.username}'); // Always ensure stream controller exists first _ensureStreamController(); if (_isInitialized && _currentUsername == userContext.username && _isConnected && _subscriptionHealthy) { dev.log( 'โœ… PubNub already initialized and healthy for ${userContext.username}'); return; } dev.log('๐Ÿ”„ Initializing PubNub Manager for ${userContext.username}...'); _currentUsername = userContext.username; KafkaHandler.instance.connect(); _reconnectAttempts = 0; _subscriptionHealthy = true; try { // Always create fresh PubNub instance await _createFreshPubNubInstance(userContext); _setupConnectivityListener(); await _subscribe(userContext.username); _startHeartbeat(); _isInitialized = true; dev.log('โœ… PubNub Manager initialized successfully'); } catch (e, stackTrace) { dev.log('๐Ÿ”ด Error initializing PubNub: $e', error: e, stackTrace: stackTrace); _handleInitializationError(e); } } Future _createFreshPubNubInstance(UserContext userContext) async { dev.log('๐Ÿ—๏ธ Creating fresh PubNub instance...'); // Dispose old instance completely if (_pubnub != null) { try { await _cleanupSubscriptions(); _pubnub = null; } catch (e) { dev.log('โš ๏ธ Error disposing old PubNub: $e'); } } // Create completely new instance final keyset = Keyset( subscribeKey: subKey, publishKey: pubKey, userId: UserId(userContext.username), ); _pubnub = PubNub(defaultKeyset: keyset); // Small delay to ensure proper initialization await Future.delayed(const Duration(milliseconds: 500)); dev.log('โœ… Fresh PubNub instance created'); } void _handleInitializationError(dynamic error) { _isConnected = false; _subscriptionHealthy = false; _safeAddToStream({ 'error': true, 'message': 'Connection error. Retrying...', 'details': error.toString(), 'timestamp': DateTime.now().millisecondsSinceEpoch, }); _scheduleReconnect(); } void _setupConnectivityListener() { _connectivitySubscription?.cancel(); _connectivitySubscription = Connectivity().onConnectivityChanged.listen( (result) async { List results; if (result is List) { results = result; } else if (result is ConnectivityResult) { results = result; } else { results = [ConnectivityResult.none]; } final nowOnline = results.any((r) => r != ConnectivityResult.none); if (nowOnline && !_isConnected) { dev.log('๐Ÿ“ก Network restored, reconnecting PubNub...'); _reconnectAttempts = 0; await _reconnect(); } else if (!nowOnline && _isConnected) { dev.log('๐Ÿ“ก Network lost'); _isConnected = false; _subscriptionHealthy = false; _safeAddToStream({ 'error': true, 'message': 'Network connection lost. Will reconnect automatically...', 'timestamp': DateTime.now().millisecondsSinceEpoch, }); } }, onError: (error) { dev.log('๐Ÿ”ด Connectivity listener error: $error'); }, ); } Future _subscribe(String username) async { if (_pubnub == null) return; await _cleanupSubscriptions(); final channels = { 'user-ride_request-$username', 'user-ride_updates-$username', // 'driver-location', }; await _kafkaLocationSubscription?.cancel(); // Clean up old listener _kafkaLocationSubscription = KafkaHandler.instance.driverLocationStream.listen((locationData) { dev.log('๐Ÿ“จ Message received via [Kafka/WebSocket]: $locationData'); // Add a type or key to identify it as a location update if needed locationData['type'] = 'DRIVER_LOCATION_UPDATE'; _safeAddToStream(locationData); }); KafkaHandler.instance.subscribeToDriverLocations(); dev.log('๐Ÿ“ก Subscribed to Kafka driver-locations topic'); dev.log('๐Ÿ“ก Subscribing to channels: $channels'); try { _subscription = _pubnub!.subscribe(channels: channels, withPresence: true); // CRITICAL FIX: Robust message handling _subscription!.messages.listen( (envelope) { _lastMessageReceived = DateTime.now(); _messageCount++; _isConnected = true; _subscriptionHealthy = true; _reconnectAttempts = 0; final messageId = _generateMessageId(envelope); if (_processedMessages.contains(messageId)) { dev.log('โš ๏ธ Duplicate message ignored: $messageId'); return; } _processedMessages.add(messageId); dev.log( '๐Ÿ“จ Message received on ${envelope.channel}: ${envelope.payload}'); try { if (envelope.payload is Map) { final payload = Map.from( envelope.payload as Map); payload['_channel'] = envelope.channel; payload['_timestamp'] = DateTime.now().millisecondsSinceEpoch; payload['_messageId'] = messageId; payload['_timetoken'] = envelope.timetoken; _safeAddToStream(payload); } else { dev.log('โš ๏ธ Received non-map payload: ${envelope.payload}'); } } catch (e) { dev.log('๐Ÿ”ด Error processing message: $e'); } }, onError: (error, stackTrace) { dev.log('๐Ÿ”ด Message subscription error: $error'); _isConnected = false; _subscriptionHealthy = false; _safeAddToStream({ 'error': true, 'message': 'Message subscription error. Reconnecting...', 'timestamp': DateTime.now().millisecondsSinceEpoch, }); _scheduleReconnect(); }, onDone: () { dev.log('๐Ÿ”ด Message subscription closed unexpectedly'); _isConnected = false; _subscriptionHealthy = false; _scheduleReconnect(); }, cancelOnError: false, ); // Presence handling _subscription!.presence.listen( (event) { dev.log('๐Ÿ‘ฅ Presence event: $event'); _isConnected = true; _subscriptionHealthy = true; _safeAddToStream({ 'presence': true, 'event': event.toString(), 'connected': true, 'timestamp': DateTime.now().millisecondsSinceEpoch, }); }, onError: (error) { dev.log('๐Ÿ”ด Presence subscription error: $error'); _subscriptionHealthy = false; }, cancelOnError: false, ); dev.log('โœ… Successfully subscribed to PubNub channels'); // Send connection confirmation _safeAddToStream({ 'connected': true, 'message': 'Connected to ride updates', 'channels': channels.toList(), 'timestamp': DateTime.now().millisecondsSinceEpoch, }); } catch (e, stackTrace) { dev.log('๐Ÿ”ด Failed to subscribe: $e'); dev.log('StackTrace: $stackTrace'); _isConnected = false; _subscriptionHealthy = false; _safeAddToStream({ 'error': true, 'message': 'Failed to connect. Retrying...', 'timestamp': DateTime.now().millisecondsSinceEpoch, }); _scheduleReconnect(); } } String _generateMessageId(dynamic envelope) { return '${envelope.channel}_${envelope.timetoken}_${envelope.payload.hashCode}_${DateTime.now().millisecondsSinceEpoch}'; } // CRITICAL FIX: Always safe to add to stream - never fails void _safeAddToStream(Map data) { _ensureStreamController(); try { if (_rideUpdateController != null && !_rideUpdateController!.isClosed) { _rideUpdateController!.add(data); // dev.log('โœ… Added data to stream: ${data.keys.join(', ')}'); } else { dev.log('๐Ÿ”ด Stream controller is null or closed - recreating!'); // Force recreate if somehow closed _streamControllerCreated = false; _ensureStreamController(); _rideUpdateController!.add(data); } } catch (e, stackTrace) { dev.log('๐Ÿ”ด Error adding to stream: $e'); dev.log('StackTrace: $stackTrace'); } } void _scheduleReconnect() { if (_reconnectAttempts >= _maxReconnectAttempts) { dev.log('๐Ÿ”ด Max reconnection attempts reached'); _safeAddToStream({ 'error': true, 'message': 'Connection failed. Please restart the app.', 'maxAttemptsReached': true, 'timestamp': DateTime.now().millisecondsSinceEpoch, }); return; } _reconnectTimer?.cancel(); final delay = _reconnectBaseDelay * (_reconnectAttempts + 1); dev.log( 'โฐ Scheduling reconnect in ${delay.inSeconds} seconds (attempt ${_reconnectAttempts + 1})'); _reconnectTimer = Timer(delay, () { _reconnect(); }); } Future _reconnect() async { if (!_isInitialized || _currentUsername == null) return; _reconnectAttempts++; dev.log( '๐Ÿ”„ Reconnecting PubNub (attempt $_reconnectAttempts/$_maxReconnectAttempts)...'); _safeAddToStream({ 'reconnecting': true, 'message': 'Reconnecting... (${_reconnectAttempts}/$_maxReconnectAttempts)', 'attempt': _reconnectAttempts, 'timestamp': DateTime.now().millisecondsSinceEpoch, }); try { await _createFreshPubNubInstance( UserContext(username: _currentUsername!)); await _subscribe(_currentUsername!); } catch (e) { dev.log('๐Ÿ”ด Reconnection failed: $e'); _scheduleReconnect(); } } void _startHeartbeat() { _heartbeatTimer?.cancel(); _heartbeatTimer = Timer.periodic(_heartbeatInterval, (timer) { if (!_isConnected && _isInitialized && _reconnectAttempts < _maxReconnectAttempts) { dev.log('๐Ÿ’” Heartbeat detected disconnection, reconnecting...'); _reconnect(); } else if (_isConnected && _subscriptionHealthy) { dev.log( '๐Ÿ’“ Heartbeat: Connection healthy, ${_messageCount} messages received'); } }); } Future publish(String channel, Map message) async { if (_pubnub == null) { dev.log('โš ๏ธ Cannot publish - PubNub not initialized'); return false; } if (!_isConnected) { dev.log('โš ๏ธ Not connected, cannot publish message'); return false; } dev.log('๐Ÿ“ค Publishing to $channel: $message'); try { await _pubnub!.publish(channel, message); dev.log('โœ… Message published successfully'); return true; } catch (e) { dev.log('๐Ÿ”ด Publish failed: $e'); _isConnected = false; _subscriptionHealthy = false; _scheduleReconnect(); return false; } } Future _cleanupSubscriptions() async { dev.log('๐Ÿงน Cleaning up PubNub & Kafka subscriptions...'); // โœจ 6. CLEAN UP the Kafka subscription as well await _kafkaLocationSubscription?.cancel(); _kafkaLocationSubscription = null; KafkaHandler.instance.unsubscribeFromDriverLocations(); if (_subscription != null && !_subscription!.isCancelled) { try { _subscription!.dispose(); dev.log('โœ… PubNub subscription disposed'); } catch (e) { dev.log('โš ๏ธ Error disposing subscription: $e'); } } _subscription = null; } // CRITICAL FIX: Never fully dispose during app lifecycle void _cleanup() { dev.log('๐Ÿงน Cleaning up PubNub resources (keeping core alive)...'); _cleanupSubscriptions(); _heartbeatTimer?.cancel(); _heartbeatTimer = null; _reconnectTimer?.cancel(); _reconnectTimer = null; _subscriptionWatchdog?.cancel(); _subscriptionWatchdog = null; // Keep connectivity listener and stream controller alive! // _connectivitySubscription?.cancel(); // DON'T cancel this // _rideUpdateController?.close(); // NEVER close this _isConnected = false; _subscriptionHealthy = false; } // IMPROVED: Better force reconnect Future forceReconnect() async { if (!_isInitialized || _currentUsername == null) { dev.log('โš ๏ธ Cannot force reconnect - not initialized or no username'); return; } dev.log('๐Ÿ”„ Force reconnecting PubNub...'); _reconnectAttempts = 0; _isConnected = false; _subscriptionHealthy = false; _safeAddToStream({ 'reconnecting': true, 'message': 'Force reconnecting...', 'forced': true, 'timestamp': DateTime.now().millisecondsSinceEpoch, }); try { await _createFreshPubNubInstance( UserContext(username: _currentUsername!)); await _subscribe(_currentUsername!); } catch (e) { dev.log('๐Ÿ”ด Force reconnect failed: $e'); _scheduleReconnect(); } } bool get isConnected => _isConnected && _subscriptionHealthy; bool get isInitialized => _isInitialized; Map getStatus() { return { 'isConnected': _isConnected, 'isInitialized': _isInitialized, 'subscriptionHealthy': _subscriptionHealthy, 'currentUsername': _currentUsername, 'reconnectAttempts': _reconnectAttempts, 'hasActiveSubscription': _subscription != null && !_subscription!.isCancelled, 'hasStreamController': _rideUpdateController != null && !_rideUpdateController!.isClosed, 'processedMessagesCount': _processedMessages.length, 'messageCount': _messageCount, 'lastMessageReceived': _lastMessageReceived?.toIso8601String(), 'maxReconnectAttempts': _maxReconnectAttempts, 'streamControllerCreated': _streamControllerCreated, }; } // CRITICAL FIX: Only allow disposal on app termination, not screen navigation void disposeOnAppTermination() { dev.log('๐Ÿ’€ PubNub Manager disposing on app termination'); _cleanup(); _connectivitySubscription?.cancel(); _connectivitySubscription = null; _cleanupTimer?.cancel(); _cleanupTimer = null; // Only now close the stream controller if (_rideUpdateController != null && !_rideUpdateController!.isClosed) { _rideUpdateController!.close(); } _rideUpdateController = null; _streamControllerCreated = false; _processedMessages.clear(); _isInitialized = false; _isConnected = false; _subscriptionHealthy = false; _currentUsername = null; _reconnectAttempts = 0; _messageCount = 0; _lastMessageReceived = null; dev.log('โœ… PubNub Manager disposed completely'); } // REMOVED: Regular dispose() method to prevent accidental disposal }