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; 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(); 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; 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'); _isConnected = false; _connectionController.add(false); if (reason == 'io server disconnect' && !_intentionalDisconnect) { _log.w('Server rejected connection - token likely expired'); _reconnectWithFreshToken(); } }); _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)); }); _socket!.connect(); } Future _reconnectWithFreshToken() async { // Wait briefly for any in-flight REST token refresh to complete await Future.delayed(const Duration(milliseconds: 500)); final token = await SecureStorage.getAccessToken(); if (token == null) return; // If token hasn't changed, try triggering a refresh via a lightweight API call if (token == _currentToken) { try { // Make a lightweight call that triggers the API client's 401 interceptor // which auto-refreshes the token final dio = ApiClient.instance.dio; await dio.get('/auth/me'); final refreshedToken = await SecureStorage.getAccessToken(); if (refreshedToken == null || refreshedToken == _currentToken) return; _currentToken = refreshedToken; } catch (_) { return; // Token refresh failed — user will need to re-login } } else { _currentToken = token; } _socket?.io.options?['auth'] = {'token': _currentToken}; _socket?.connect(); } /// Ensure socket is connected — reconnect if needed. Future ensureConnected() async { if (_socket?.connected == true) return; await connect(); } void disconnect() { _intentionalDisconnect = true; _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: 10), 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(); } }