package chats.data.remote.signalr import android.util.Log import io.reactivex.rxjava3.core.Single import com.microsoft.signalr.HubConnection import com.microsoft.signalr.HubConnectionBuilder import com.microsoft.signalr.HubConnectionState import chats.data.remote.dto.ChatDto import chats.data.remote.dto.MessageDto import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.SharedFlow import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.launch import javax.inject.Inject import javax.inject.Singleton @Singleton class ChatHubClient @Inject constructor() { private var hubConnection: HubConnection? = null private val _events = MutableSharedFlow(extraBufferCapacity = 64) val events: SharedFlow = _events.asSharedFlow() private val scope = CoroutineScope(Dispatchers.IO) fun connect(baseUrl: String, accessToken: String) { if (hubConnection?.connectionState == HubConnectionState.CONNECTED) return hubConnection = HubConnectionBuilder.create("${baseUrl}/chatHub") .withAccessTokenProvider(Single.just(accessToken)) .build() setupHandlers() hubConnection?.onClosed { exception -> Log.e("ChatHubClient", "Connection closed", exception) // Optional: Reconnect logic } scope.launch { try { hubConnection?.start()?.blockingAwait() Log.d("ChatHubClient", "SignalR Connected") } catch (e: Exception) { Log.e("ChatHubClient", "SignalR Connection Error", e) } } } private fun setupHandlers() { hubConnection?.let { conn -> conn.on("new_message", { message: MessageDto -> _events.tryEmit(ChatEvent.NewMessage(message)) }, MessageDto::class.java) conn.on("messages_read", { chatId: String, userId: String, lastReadSequenceId: Int -> _events.tryEmit(ChatEvent.MessagesRead(chatId, userId, lastReadSequenceId)) }, String::class.java, String::class.java, Int::class.java) conn.on("user_typing", { chatId: String, userId: String -> _events.tryEmit(ChatEvent.UserTyping(chatId, userId)) }, String::class.java, String::class.java) conn.on("user_online", { userId: String -> _events.tryEmit(ChatEvent.UserOnline(userId)) }, String::class.java) conn.on("new_chat", { chat: ChatDto -> _events.tryEmit(ChatEvent.NewChat(chat)) }, ChatDto::class.java) conn.on("reaction_updated", { messageId: String, chatId: String, userId: String, emoji: String -> _events.tryEmit(ChatEvent.ReactionUpdated(messageId, chatId, userId, emoji)) }, String::class.java, String::class.java, String::class.java, String::class.java) // WebRTC Signaling Handlers conn.on("call_incoming", { chatId: String, from: String, offer: String, callType: String -> _events.tryEmit(ChatEvent.CallIncoming(chatId, from, offer, callType)) }, String::class.java, String::class.java, String::class.java, String::class.java) conn.on("call_answered", { chatId: String, answer: String -> _events.tryEmit(ChatEvent.CallAnswered(chatId, answer)) }, String::class.java, String::class.java) conn.on("ice_candidate", { chatId: String, candidate: String -> _events.tryEmit(ChatEvent.IceCandidateReceived(chatId, candidate)) }, String::class.java, String::class.java) conn.on("call_ended", { chatId: String -> _events.tryEmit(ChatEvent.CallEnded(chatId)) }, String::class.java) } } fun disconnect() { hubConnection?.stop() } }