Files
mobile-app/lib/features/messaging/data/socket_service.dart

238 lines
7.8 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/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;
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();
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;
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;
// 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(5)
.setReconnectionDelay(1000)
.setReconnectionDelayMax(5000)
.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') {
_log.w('Server rejected connection - token likely expired');
_reconnectWithFreshToken();
}
});
_socket!.onConnectError((error) {
_log.e('Socket connect error: $error');
_isConnected = false;
_connectionController.add(false);
});
// 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));
});
_socket!.connect();
}
Future<void> _reconnectWithFreshToken() async {
final token = await SecureStorage.getAccessToken();
if (token != null && token != _currentToken) {
_currentToken = token;
_socket?.io.options?['auth'] = {'token': token};
_socket?.connect();
}
}
void disconnect() {
_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();
}
}