import 'dart:async'; import 'package:logger/logger.dart'; import 'package:real_estate_mobile/config/app_config.dart'; import 'package:real_estate_mobile/core/network/api_client.dart'; import 'package:real_estate_mobile/core/storage/secure_storage.dart'; import 'package:real_estate_mobile/features/messaging/data/models/message.dart'; import 'package:real_estate_mobile/features/support_chat/data/models/support_chat_models.dart'; // We'll use a simple WebSocket approach for now since socket_io_client needs to be added // This will be a placeholder that uses REST polling until socket_io_client is added import 'package:socket_io_client/socket_io_client.dart' as io; final _log = Logger(printer: PrettyPrinter(methodCount: 0)); class SocketService { static final SocketService _instance = SocketService._(); factory SocketService() => _instance; SocketService._(); io.Socket? _socket; bool _isConnected = false; bool _intentionalDisconnect = false; int _reconnectAttempts = 0; static const _maxReconnectAttempts = 5; String? _currentToken; // Stream controllers for events final _messageController = StreamController.broadcast(); final _typingStartController = StreamController>.broadcast(); final _typingStopController = StreamController>.broadcast(); final _statusController = StreamController>.broadcast(); final _readController = StreamController>.broadcast(); final _connectionController = StreamController.broadcast(); final _supportMessageController = StreamController.broadcast(); final _supportTypingStartController = StreamController>.broadcast(); final _supportTypingStopController = StreamController>.broadcast(); final _connectionRequestController = StreamController>.broadcast(); final _connectionResponseController = StreamController>.broadcast(); final _deliveredController = StreamController>.broadcast(); Stream get onNewMessage => _messageController.stream; Stream> get onTypingStart => _typingStartController.stream; Stream> get onTypingStop => _typingStopController.stream; Stream> get onStatusChange => _statusController.stream; Stream> get onMessagesRead => _readController.stream; Stream get onConnectionChange => _connectionController.stream; Stream get onSupportMessage => _supportMessageController.stream; Stream> get onSupportTypingStart => _supportTypingStartController.stream; Stream> get onSupportTypingStop => _supportTypingStopController.stream; Stream> get onConnectionRequest => _connectionRequestController.stream; Stream> get onConnectionResponse => _connectionResponseController.stream; Stream> get onMessageDelivered => _deliveredController.stream; bool get isConnected => _isConnected; Future connect() async { final token = await SecureStorage.getAccessToken(); if (token == null) { _log.w('No access token for socket connection'); return; } if (_socket?.connected == true && _currentToken == token) return; _currentToken = token; _intentionalDisconnect = false; // Strip /api/v1 from base URL for socket connection final apiUrl = AppConfig.apiBaseUrl; final baseUrl = apiUrl.replaceAll(RegExp(r'/api/v1/?$'), ''); _socket?.disconnect(); _socket?.dispose(); _socket = io.io(baseUrl, io.OptionBuilder() .setTransports(['websocket', 'polling']) .setAuth({'token': token}) .disableAutoConnect() .enableReconnection() .setReconnectionAttempts(double.maxFinite.toInt()) .setReconnectionDelay(1000) .setReconnectionDelayMax(10000) .build(), ); _socket!.onConnect((_) { _log.i('Socket connected: ${_socket!.id}'); _isConnected = true; _connectionController.add(true); }); _socket!.onDisconnect((reason) { _log.w('Socket disconnected: $reason'); final wasConnected = _isConnected; _isConnected = false; _connectionController.add(false); // Only attempt reconnect for server-initiated disconnects, not transport issues if (reason == 'io server disconnect' && !_intentionalDisconnect) { // If we were stably connected (not a connect-then-immediate-disconnect), // reset the counter — this is a new disconnection event if (wasConnected) { // Check: was the connection stable? (connected for more than 2 seconds) // If not, it's a repeated rejection — increment counter } if (_reconnectAttempts >= _maxReconnectAttempts) { _log.e('Max reconnect attempts ($_maxReconnectAttempts) reached. Server keeps rejecting. Giving up.'); _reconnectAttempts = 0; return; } _reconnectAttempts++; _log.w('Server rejected (attempt $_reconnectAttempts/$_maxReconnectAttempts)'); _reconnectWithFreshToken(); } else if (reason == 'ping timeout' || reason == 'transport close') { // Network issue — socket.io auto-reconnects, reset counter _reconnectAttempts = 0; } }); _socket!.onConnectError((error) { _log.e('Socket connect error: $error'); _isConnected = false; _connectionController.add(false); }); _socket!.onReconnect((_) { _log.i('Socket reconnected successfully'); }); _socket!.onReconnectAttempt((attempt) { _log.w('Socket reconnection attempt #$attempt'); }); _socket!.onReconnectError((error) { _log.e('Socket reconnection error: $error'); }); // Listen for events _socket!.on('new_message', (data) { try { final message = ChatMessage.fromJson(data as Map); _messageController.add(message); } catch (e) { _log.e('Error parsing new_message: $e'); } }); _socket!.on('typing_start', (data) { _typingStartController.add(Map.from(data as Map)); }); _socket!.on('typing_stop', (data) { _typingStopController.add(Map.from(data as Map)); }); _socket!.on('user_status_change', (data) { _statusController.add(Map.from(data as Map)); }); _socket!.on('messages_read', (data) { _readController.add(Map.from(data as Map)); }); // Support chat events _socket!.on('support_new_message', (data) { try { final message = SupportMessage.fromJson(data as Map); _supportMessageController.add(message); } catch (e) { _log.e('Error parsing support_new_message: $e'); } }); _socket!.on('support_typing_start', (data) { _supportTypingStartController.add(Map.from(data as Map)); }); _socket!.on('support_typing_stop', (data) { _supportTypingStopController.add(Map.from(data as Map)); }); // Connection request events (real-time updates) _socket!.on('connection_request', (data) { _connectionRequestController.add(Map.from(data as Map)); }); _socket!.on('connection_response', (data) { _connectionResponseController.add(Map.from(data as Map)); }); _socket!.on('message_delivered', (data) { _deliveredController.add(Map.from(data as Map)); }); _socket!.connect(); } Future _reconnectWithFreshToken() async { // Wait with exponential backoff final delay = Duration(seconds: _reconnectAttempts * 2 + 1); _log.i('Waiting ${delay.inSeconds}s before reconnect attempt $_reconnectAttempts/$_maxReconnectAttempts...'); await Future.delayed(delay); // Try to get a fresh token try { final dio = ApiClient.instance.dio; await dio.get('/auth/me'); } catch (_) {} final freshToken = await SecureStorage.getAccessToken(); if (freshToken == null) { _log.w('No token available, giving up'); return; } _log.i('Reconnecting socket with token...'); _currentToken = freshToken; // Update auth on existing socket and reconnect (don't recreate) if (_socket != null) { _socket!.io.options?['auth'] = {'token': freshToken}; _socket!.connect(); } else { await connect(); } } /// Ensure socket is connected — reconnect if needed. Future ensureConnected() async { if (_socket?.connected == true) return; await connect(); } void disconnect() { _intentionalDisconnect = true; _reconnectAttempts = 0; _socket?.disconnect(); _socket?.dispose(); _socket = null; _isConnected = false; _currentToken = null; } Future joinConversation(String conversationId) async { if (_socket == null || !_isConnected) return; final completer = Completer(); _socket!.emitWithAck('join_conversation', {'conversationId': conversationId}, ack: (data) => completer.complete(), ); await completer.future.timeout( const Duration(seconds: 10), onTimeout: () => _log.w('join_conversation timeout'), ); } void leaveConversation(String conversationId) { _socket?.emit('leave_conversation', {'conversationId': conversationId}); } Future sendMessage( String conversationId, Map messageData, ) async { if (_socket == null || !_isConnected) return null; final completer = Completer(); _socket!.emitWithAck('send_message', { 'conversationId': conversationId, 'message': messageData, }, ack: (data) { try { final response = data as Map; if (response['success'] == true && response['message'] != null) { completer.complete( ChatMessage.fromJson(response['message'] as Map), ); } else { completer.complete(null); } } catch (e) { completer.complete(null); } }); return completer.future.timeout( const Duration(seconds: 5), onTimeout: () => null, ); } void startTyping(String conversationId) { _socket?.emit('typing_start', {'conversationId': conversationId}); } void stopTyping(String conversationId) { _socket?.emit('typing_stop', {'conversationId': conversationId}); } /// Emit a raw event (used by support chat) void emitRaw(String event, Map data) { _socket?.emit(event, data); } void markAsRead(String conversationId) { _socket?.emit('mark_read', {'conversationId': conversationId}); } void dispose() { disconnect(); _messageController.close(); _typingStartController.close(); _typingStopController.close(); _statusController.close(); _readController.close(); _connectionController.close(); _supportMessageController.close(); _supportTypingStartController.close(); _supportTypingStopController.close(); _connectionRequestController.close(); _connectionResponseController.close(); _deliveredController.close(); } }