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 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() 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) } hubConnection?.start()?.blockingAwait() } 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() } }