import 'dart:async'; import 'dart:developer' as dev; import 'package:flutter/foundation.dart'; import 'package:pubnub/pubnub.dart'; import '../Provider/public_keys_provider.dart'; import '../network/kafkahandler.dart'; class DriverPubNubManager extends ChangeNotifier { DriverPubNubManager() { KafkaHandler.instance.driverLocationStream .listen(_handleKafkaLocationUpdate); } PubNub? _pn; late String _uuid; final _updateCtrl = StreamController>.broadcast(); Stream> get rideUpdateStream => _updateCtrl.stream; Subscription? _subscription; final Set _baseChannels = {}; Set _rideChannels = {}; bool get isInitialized => _pn != null; Future connect({required String uuid}) async { _uuid = uuid; if (!publickeys.guardKey(publickeys.pubnubPublishKey, 'pubnubPublishKey') || !publickeys.guardKey(publickeys.pubnubSubscribeKey, 'pubnubSubscribeKey')) { dev.log('🔴 PubNub keys not configured — DriverPubNubManager not connected'); return; } dev.log('🔌 DRIVER PubNub Manager: Connecting for $_uuid'); KafkaHandler.instance.connect(); final keyset = Keyset( publishKey: publickeys.pubnubPublishKey!, subscribeKey: publickeys.pubnubSubscribeKey!, userId: UserId(_uuid), ); _pn = PubNub(defaultKeyset: keyset); _baseChannels.add('ride-offers-public'); _baseChannels.add('driver-offer-$_uuid'); _resubscribe(); } void _resubscribe() { if (_pn == null) { dev.log('⚠️ PubNub not initialized — skipping resubscribe'); return; } _subscription?.dispose(); 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); _subscription!.messages.listen(_handleIncomingMessage); } void subscribeToRide({required int rideId, required String riderUsername}) { dev.log('Adding ride-specific channels for ride $rideId'); _rideChannels = { 'ride-updates-$rideId', 'user-ride_updates-$riderUsername', 'user-ride_updates-$_uuid', }; _resubscribe(); } void unsubscribeFromRide() { dev.log('Removing ride-specific channels.'); KafkaHandler.instance.unsubscribeFromDriverLocations(); if (_rideChannels.isEmpty) return; _rideChannels.clear(); _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 (_pn == null) { dev.log('⚠️ Cannot publish — PubNub not initialized'); return Future.value(); } return _pn! .publish(channel, payload) .catchError((e) => dev.log('🔴 Driver publish error: $e')); } 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'); }