327 lines
11 KiB
Dart
327 lines
11 KiB
Dart
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;
|
|
bool _isReconnecting = false;
|
|
int _reconnectAttempts = 0;
|
|
static const _maxReconnectAttempts = 3;
|
|
String? _currentToken;
|
|
|
|
// Stream controllers for events
|
|
final _messageController = StreamController<ChatMessage>.broadcast();
|
|
final _typingStartController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _typingStopController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _statusController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _readController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _connectionController = StreamController<bool>.broadcast();
|
|
final _supportMessageController = StreamController<SupportMessage>.broadcast();
|
|
final _supportTypingStartController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _supportTypingStopController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _connectionRequestController = StreamController<Map<String, dynamic>>.broadcast();
|
|
final _connectionResponseController = StreamController<Map<String, dynamic>>.broadcast();
|
|
|
|
Stream<ChatMessage> get onNewMessage => _messageController.stream;
|
|
Stream<Map<String, dynamic>> get onTypingStart => _typingStartController.stream;
|
|
Stream<Map<String, dynamic>> get onTypingStop => _typingStopController.stream;
|
|
Stream<Map<String, dynamic>> get onStatusChange => _statusController.stream;
|
|
Stream<Map<String, dynamic>> get onMessagesRead => _readController.stream;
|
|
Stream<bool> get onConnectionChange => _connectionController.stream;
|
|
Stream<SupportMessage> get onSupportMessage => _supportMessageController.stream;
|
|
Stream<Map<String, dynamic>> get onSupportTypingStart => _supportTypingStartController.stream;
|
|
Stream<Map<String, dynamic>> get onSupportTypingStop => _supportTypingStopController.stream;
|
|
Stream<Map<String, dynamic>> get onConnectionRequest => _connectionRequestController.stream;
|
|
Stream<Map<String, dynamic>> get onConnectionResponse => _connectionResponseController.stream;
|
|
bool get isConnected => _isConnected;
|
|
|
|
Future<void> 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;
|
|
_reconnectAttempts = 0; // Reset on successful connect
|
|
_isReconnecting = false;
|
|
_connectionController.add(true);
|
|
});
|
|
|
|
_socket!.onDisconnect((reason) {
|
|
_log.w('Socket disconnected: $reason');
|
|
_isConnected = false;
|
|
_connectionController.add(false);
|
|
|
|
if (reason == 'io server disconnect' && !_intentionalDisconnect && !_isReconnecting) {
|
|
if (_reconnectAttempts < _maxReconnectAttempts) {
|
|
_reconnectAttempts++;
|
|
_log.w('Server rejected connection (attempt $_reconnectAttempts/$_maxReconnectAttempts) - refreshing token...');
|
|
_reconnectWithFreshToken();
|
|
} else {
|
|
_log.e('Max reconnect attempts reached - giving up');
|
|
}
|
|
}
|
|
});
|
|
|
|
_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<String, dynamic>);
|
|
_messageController.add(message);
|
|
} catch (e) {
|
|
_log.e('Error parsing new_message: $e');
|
|
}
|
|
});
|
|
|
|
_socket!.on('typing_start', (data) {
|
|
_typingStartController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
_socket!.on('typing_stop', (data) {
|
|
_typingStopController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
_socket!.on('user_status_change', (data) {
|
|
_statusController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
_socket!.on('messages_read', (data) {
|
|
_readController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
// Support chat events
|
|
_socket!.on('support_new_message', (data) {
|
|
try {
|
|
final message = SupportMessage.fromJson(data as Map<String, dynamic>);
|
|
_supportMessageController.add(message);
|
|
} catch (e) {
|
|
_log.e('Error parsing support_new_message: $e');
|
|
}
|
|
});
|
|
|
|
_socket!.on('support_typing_start', (data) {
|
|
_supportTypingStartController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
_socket!.on('support_typing_stop', (data) {
|
|
_supportTypingStopController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
// Connection request events (real-time updates)
|
|
_socket!.on('connection_request', (data) {
|
|
_connectionRequestController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
_socket!.on('connection_response', (data) {
|
|
_connectionResponseController.add(Map<String, dynamic>.from(data as Map));
|
|
});
|
|
|
|
_socket!.connect();
|
|
}
|
|
|
|
Future<void> _reconnectWithFreshToken() async {
|
|
_isReconnecting = true;
|
|
// Wait with exponential backoff
|
|
final delay = Duration(seconds: _reconnectAttempts * 2);
|
|
_log.i('Waiting ${delay.inSeconds}s before reconnect attempt...');
|
|
await Future.delayed(delay);
|
|
|
|
final token = await SecureStorage.getAccessToken();
|
|
if (token == null) {
|
|
_log.w('No access token available for socket reconnection');
|
|
return;
|
|
}
|
|
|
|
// If token hasn't changed, try triggering a refresh via a lightweight API call
|
|
if (token == _currentToken) {
|
|
_log.i('Token unchanged, triggering refresh via /auth/me...');
|
|
try {
|
|
final dio = ApiClient.instance.dio;
|
|
await dio.get('/auth/me');
|
|
final refreshedToken = await SecureStorage.getAccessToken();
|
|
if (refreshedToken == null || refreshedToken == _currentToken) {
|
|
_log.w('Token refresh did not produce a new token');
|
|
// Still try reconnecting — token might be valid but socket had a glitch
|
|
} else {
|
|
_log.i('Got fresh token after refresh');
|
|
}
|
|
} catch (e) {
|
|
_log.e('Token refresh failed: $e');
|
|
// Still try reconnecting with existing token
|
|
}
|
|
}
|
|
|
|
// Get the latest token (may have been refreshed)
|
|
final freshToken = await SecureStorage.getAccessToken();
|
|
if (freshToken == null) return;
|
|
|
|
_log.i('Reconnecting socket with fresh token...');
|
|
// Full reconnect — dispose old socket and create new one
|
|
_socket?.disconnect();
|
|
_socket?.dispose();
|
|
_socket = null;
|
|
_isConnected = false;
|
|
_currentToken = null;
|
|
|
|
// Small delay before reconnecting
|
|
await Future.delayed(const Duration(milliseconds: 500));
|
|
await connect();
|
|
}
|
|
|
|
/// Ensure socket is connected — reconnect if needed.
|
|
Future<void> ensureConnected() async {
|
|
if (_socket?.connected == true) return;
|
|
await connect();
|
|
}
|
|
|
|
void disconnect() {
|
|
_intentionalDisconnect = true;
|
|
_isReconnecting = false;
|
|
_reconnectAttempts = 0;
|
|
_socket?.disconnect();
|
|
_socket?.dispose();
|
|
_socket = null;
|
|
_isConnected = false;
|
|
_currentToken = null;
|
|
}
|
|
|
|
Future<void> joinConversation(String conversationId) async {
|
|
if (_socket == null || !_isConnected) return;
|
|
final completer = Completer<void>();
|
|
_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<ChatMessage?> sendMessage(
|
|
String conversationId,
|
|
Map<String, dynamic> messageData,
|
|
) async {
|
|
if (_socket == null || !_isConnected) return null;
|
|
|
|
final completer = Completer<ChatMessage?>();
|
|
_socket!.emitWithAck('send_message', {
|
|
'conversationId': conversationId,
|
|
'message': messageData,
|
|
}, ack: (data) {
|
|
try {
|
|
final response = data as Map<String, dynamic>;
|
|
if (response['success'] == true && response['message'] != null) {
|
|
completer.complete(
|
|
ChatMessage.fromJson(response['message'] as Map<String, dynamic>),
|
|
);
|
|
} 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<String, dynamic> 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();
|
|
}
|
|
}
|