import 'dart:async'; import 'dart:developer' as dev; import 'package:flutter/foundation.dart'; // REQUIRED import 'package:pubnub/pubnub.dart'; import '../config/constants/app_constant.dart'; import '../Provider/utils_singleton.dart'; import '../network/kafkahandler.dart'; class DriverPubNubManager extends ChangeNotifier { DriverPubNubManager() { // 👈 --- LISTEN TO THE NEW STREAM FROM KAFKA HANDLER --- KafkaHandler.instance.driverLocationStream.listen(_handleKafkaLocationUpdate); } late PubNub _pn; late String _uuid; final _updateCtrl = StreamController>.broadcast(); Stream> get rideUpdateStream => _updateCtrl.stream; Subscription? _subscription; // bool _isInitialized = false; // ✅ Store the current channels to easily add/remove ride-specific ones final Set _baseChannels = {}; Set _rideChannels = {}; Future connect({required String uuid}) async { // if (_isInitialized) return; // _isInitialized = true; _uuid = uuid; dev.log('🔌 DRIVER PubNub Manager v5: Connecting for $_uuid'); KafkaHandler.instance.connect(); final keyset = Keyset( publishKey: pubKey, subscribeKey: subKey, userId: UserId(_uuid), ); _pn = PubNub(defaultKeyset: keyset); // Define the channels the driver always listens to _baseChannels.add('ride-offers-public'); _baseChannels.add('driver-offer-$_uuid'); _resubscribe(); } /// ✅ This is the single, central method for managing subscriptions. void _resubscribe() { _subscription?.dispose(); // Clean up the old subscription first // Combine the base channels with any ride-specific channels final allChannels = {..._baseChannels, ..._rideChannels}; if (allChannels.isEmpty) { dev.log('No channels to subscribe to.'); return; } dev.log('✉️ Updating subscription. Subscribing to: $allChannels'); _subscription = _pn.subscribe(channels: allChannels); // The listener is attached to the new subscription object _subscription!.messages.listen(_handleIncomingMessage); } void subscribeToRide({required int rideId, required String riderUsername}) { dev.log('Adding ride-specific channels for ride $rideId'); // KafkaHandler.instance.subscribeToDriverLocations(); // Define the channels for this specific ride _rideChannels = { 'ride-updates-$rideId', // 'driver-location', 'user-ride_updates-$riderUsername', 'user-ride_updates-$_uuid' }; // Trigger a resubscribe to add the new channels _resubscribe(); } /// ✅ This is the method your UI will call when a ride is over. void unsubscribeFromRide() { dev.log('Removing ride-specific channels.'); KafkaHandler.instance.unsubscribeFromDriverLocations(); if (_rideChannels.isEmpty) return; // Nothing to do _rideChannels.clear(); // Trigger a resubscribe to remove the old ride channels _resubscribe(); } void _handleIncomingMessage(Envelope envelope) { if (envelope.payload is! Map) return; final data = envelope.payload as Map; dev.log('📨 Driver received on [${envelope.channel}]: $data'); _updateCtrl.add(data); } Future _publish(String channel, Map payload) { // if (!_isInitialized) return Future.value(); return _pn.publish(channel, payload) .catchError((e) => dev.log('🔴 Driver publish error: $e')); } // --- No changes needed to your `send...` methods --- // Future sendAcceptRide({required int rideId,required String username}) async { // subscribeToRide(rideId: rideId, riderUsername: username); // await _publish('commands-ingress', { // 'type': 'ACCEPT_RIDE', // 'payload': { // 'rideId': rideId, // 'accepted': true, // 'driverid': utils.user!.driver.id // } // }); // } // Future sendDeclineRide({required int rideId}) => // _publish('commands-ingress', { // 'type': 'ACCEPT_RIDE', // 'payload': { // 'rideId': rideId, // 'accepted': false, // 'driverid': utils.user!.driver.id // } // }); // Future sendRideUpdate(String type, String stop, // {required int rideId, required String riderusername}) => // _publish('commands-ingress', { // 'type': 'UPDATE_RIDE', // 'payload': { // 'rideId': rideId, // 'type': type, // 'message': "new ride update", // 'riderusername': riderusername, // 'stoplocation': stop // } // }); // Future sendRideCancel(String type, // {required int rideId, required String riderusername}) => // _publish('commands-ingress', { // 'type': 'UPDATE_RIDE', // 'payload': { // 'rideId': rideId, // 'type': type, // 'message': "RIDE_CANCEL_USER", // 'riderusername': riderusername // } // }); void _handleKafkaLocationUpdate(Map data) { dev.log('📨 Driver received via [Kafka/WebSocket]: $data'); _updateCtrl.add(data); } void _err(Object e) => dev.log('🔴 PubNub stream error: $e'); }