You cannot select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
cloudsolutions-atoms/lib/modules/cx_module/chat/services/signalr_service.dart

339 lines
9.9 KiB
Dart

import 'dart:async';
import 'dart:developer';
import 'package:flutter/foundation.dart';
import 'package:signalr_netcore/hub_connection.dart';
import 'package:signalr_netcore/signalr_client.dart';
import 'package:test_sa/controllers/api_routes/urls.dart';
import 'package:test_sa/core/storage/auth_storage.dart';
class SignalRService {
// Singleton instance - DO NOT instantiate directly, use DI (Provider)
static final SignalRService _instance = SignalRService._internal();
factory SignalRService() => _instance;
SignalRService._internal();
// Single HubConnection instance - ONLY created and managed by this service
HubConnection? _hubConnection;
// Connection state management
bool _isInitializing = false;
Completer<bool>? _initializationCompleter;
// Authentication data
String? _userId;
String? _authToken;
String? _currentConversationId;
// Event handler registrations for re-registration after reconnect
final Map<String, List<Function(List<Object?>?)>> _eventHandlers = {};
// Connection state stream for reactive updates
final StreamController<HubConnectionState> _connectionStateController =
StreamController<HubConnectionState>.broadcast();
/// Stream of connection state changes for consumers to listen
Stream<HubConnectionState> get connectionStateStream => _connectionStateController.stream;
/// Get the current HubConnection (read-only access)
HubConnection? get hubConnection => _hubConnection;
/// Check if connected
bool get isConnected => _hubConnection?.state == HubConnectionState.Connected;
/// Get current connection state
HubConnectionState? get connectionState => _hubConnection?.state;
/// Get connection ID (for debugging and verification)
String? get connectionId => _hubConnection?.connectionId;
/// Initialize SignalR connection from stored credentials
/// Used for background/terminated app states when credentials aren't in memory
/// Returns true if successfully connected, false otherwise
Future<bool> initializeFromStorage() async {
try {
final credentials = await AuthStorage.getCredentials();
if (credentials == null) {
return false;
}
final success = await initialize(
userId: credentials.userId,
authToken: credentials.accessToken,
conversationId: null,
);
return success;
} catch (e, stackTrace) {
log('❌ [SIGNALR] Error initializing from storage: $e',
name: 'SignalRService', error: e, stackTrace: stackTrace);
return false;
}
}
/// Initialize SignalR connection with authentication
/// Thread-safe: Multiple callers will wait for the same connection task
Future<bool> initialize({
required String userId,
required String authToken,
String? conversationId,
}) async {
if (_isInitializing && _initializationCompleter != null) {
return await _initializationCompleter!.future;
}
if (isConnected && _userId == userId && _authToken == authToken) {
if (conversationId != null && conversationId != _currentConversationId) {
await _joinConversation(conversationId);
}
return true;
}
_isInitializing = true;
_initializationCompleter = Completer<bool>();
try {
_userId = userId;
_authToken = authToken;
_currentConversationId = conversationId;
await _disposeConnection();
final httpOp = HttpConnectionOptions(
skipNegotiation: false,
logMessageContent: kDebugMode,
transport: HttpTransportType.WebSockets,
requestTimeout: 30000,
);
_hubConnection = HubConnectionBuilder()
.withUrl(
"${URLs.chatHubUrlChat}?UserId=$userId&source=Desktop&access_token=$authToken",
options: httpOp,
)
.withAutomaticReconnect(retryDelays: <int>[2000, 5000, 10000, 20000])
.build();
_setupReconnectionHandlers();
try {
final startFuture = _hubConnection!.start();
if (startFuture != null) {
await startFuture.timeout(
const Duration(seconds: 30),
onTimeout: () {
throw TimeoutException('SignalR connection timeout after 30 seconds');
},
);
}
} catch (e) {
rethrow;
}
if (_hubConnection!.state != HubConnectionState.Connected) {
throw Exception('SignalR failed to connect. State: ${_hubConnection!.state}');
}
_connectionStateController.add(HubConnectionState.Connected);
if (conversationId != null) {
await _joinConversation(conversationId);
}
_reregisterAllHandlers();
_isInitializing = false;
_initializationCompleter?.complete(true);
return true;
} catch (e, stackTrace) {
log('❌ [SIGNALR] Error initializing connection: $e',
name: 'SignalRService', error: e, stackTrace: stackTrace);
_isInitializing = false;
_initializationCompleter?.complete(false);
return false;
}
}
/// Setup reconnection handlers
void _setupReconnectionHandlers() {
if (_hubConnection == null) return;
_hubConnection!.onclose(({Exception? error}) {
_connectionStateController.add(HubConnectionState.Disconnected);
});
_hubConnection!.onreconnecting(({Exception? error}) {
_connectionStateController.add(HubConnectionState.Reconnecting);
});
_hubConnection!.onreconnected(({String? connectionId}) async {
_connectionStateController.add(HubConnectionState.Connected);
if (_currentConversationId != null) {
await _joinConversation(_currentConversationId!);
}
_reregisterAllHandlers();
});
}
/// Join a conversation
Future<void> _joinConversation(String conversationId) async {
try {
if (_hubConnection?.state != HubConnectionState.Connected) {
return;
}
await _hubConnection!.invoke("JoinConversation", args: [conversationId]);
_currentConversationId = conversationId;
} catch (e) {
log('❌ [SIGNALR] Error joining conversation: $e', name: 'SignalRService');
}
}
/// Register an event handler
/// Events are stored and automatically re-registered after reconnect
void on(String eventName, Function(List<Object?>?) handler) {
if (!_eventHandlers.containsKey(eventName)) {
_eventHandlers[eventName] = [];
}
if (_eventHandlers[eventName]!.contains(handler)) {
return;
}
_eventHandlers[eventName]!.add(handler);
if (_hubConnection != null) {
_hubConnection!.on(eventName, handler);
}
}
/// Unregister an event handler
void off(String eventName, [Function(List<Object?>?)? handler]) {
if (handler != null) {
_eventHandlers[eventName]?.remove(handler);
if (_eventHandlers[eventName]?.isEmpty ?? false) {
_eventHandlers.remove(eventName);
}
} else {
_eventHandlers.remove(eventName);
}
if (_hubConnection != null) {
_hubConnection!.off(eventName, method: handler);
}
}
/// Re-register all event handlers (after reconnection)
void _reregisterAllHandlers() {
if (_hubConnection == null) return;
for (final entry in _eventHandlers.entries) {
final eventName = entry.key;
final handlers = entry.value;
for (final handler in handlers) {
_hubConnection!.on(eventName, handler);
}
}
}
/// Invoke a SignalR method
Future<Object?> invoke(String methodName, {List<Object>? args}) async {
if (_hubConnection?.state != HubConnectionState.Connected) {
throw Exception('SignalR not connected. Current state: ${_hubConnection?.state}');
}
try {
final result = await _hubConnection!.invoke(methodName, args: args);
return result;
} catch (e, stackTrace) {
log('❌ [SIGNALR] Error invoking $methodName', name: 'SignalRService');
log(' Error type: ${e.runtimeType}', name: 'SignalRService');
log(' Error message: $e', name: 'SignalRService');
log(' Connection state: ${_hubConnection?.state}', name: 'SignalRService');
log(' Connection ID: ${_hubConnection?.connectionId}', name: 'SignalRService');
if (args != null && args.isNotEmpty) {
log(' Args that failed: $args', name: 'SignalRService');
}
log(' Stack trace: $stackTrace', name: 'SignalRService');
throw Exception('SignalR invoke failed: $methodName - $e');
}
}
/// Ensure connection is ready (connect if needed)
/// Thread-safe: Multiple callers will wait for the same connection task
Future<bool> ensureConnected() async {
if (isConnected) {
return true;
}
if (_hubConnection?.state == HubConnectionState.Reconnecting) {
final startTime = DateTime.now();
while (_hubConnection?.state == HubConnectionState.Reconnecting) {
await Future.delayed(const Duration(milliseconds: 500));
if (DateTime.now().difference(startTime).inSeconds > 10) {
break;
}
}
if (isConnected) {
return true;
}
}
if (_userId != null && _authToken != null) {
return await initialize(
userId: _userId!,
authToken: _authToken!,
conversationId: _currentConversationId,
);
}
return false;
}
/// Dispose connection
Future<void> _disposeConnection() async {
try {
if (_hubConnection != null) {
await _hubConnection!.stop();
_hubConnection = null;
}
} catch (e) {
log('⚠️ [SIGNALR] Error disposing connection: $e', name: 'SignalRService');
}
}
/// Reset service (for logout)
Future<void> reset() async {
await _disposeConnection();
_eventHandlers.clear();
_userId = null;
_authToken = null;
_currentConversationId = null;
_isInitializing = false;
_initializationCompleter = null;
}
/// Disconnect from SignalR (alias for reset)
Future<void> disconnect() async {
await reset();
}
/// Dispose stream controller (called when app is terminating)
void dispose() {
_connectionStateController.close();
}
}