Прочтение
This commit is contained in:
@@ -16,6 +16,12 @@ import javax.inject.Singleton
|
||||
|
||||
import kotlinx.coroutines.delay
|
||||
|
||||
data class ReadMessagesRequest(
|
||||
val chatId: String,
|
||||
val lastReadMessageId: String,
|
||||
val lastReadSequenceId: Int
|
||||
)
|
||||
|
||||
enum class ConnectionStatus { CONNECTED, CONNECTING, DISCONNECTED }
|
||||
|
||||
@Singleton
|
||||
@@ -38,7 +44,7 @@ class ChatHubClient @Inject constructor() {
|
||||
lastToken = accessToken
|
||||
_status.value = ConnectionStatus.CONNECTING
|
||||
|
||||
hubConnection = HubConnectionBuilder.create("${baseUrl}/chatHub")
|
||||
hubConnection = HubConnectionBuilder.create("${baseUrl}/hubs/chat")
|
||||
.withAccessTokenProvider(Single.just(accessToken))
|
||||
.build()
|
||||
|
||||
@@ -87,8 +93,13 @@ class ChatHubClient @Inject constructor() {
|
||||
_events.tryEmit(ChatEvent.NewChat(chat))
|
||||
}, ChatDto::class.java)
|
||||
|
||||
conn.on("reaction_updated", { messageId: String, chatId: String, userId: String, emoji: String ->
|
||||
conn.on("reaction_added", { messageId: String, chatId: String, userId: String, username: String, emoji: String ->
|
||||
_events.tryEmit(ChatEvent.ReactionUpdated(messageId, chatId, userId, emoji))
|
||||
}, String::class.java, String::class.java, String::class.java, String::class.java, String::class.java)
|
||||
|
||||
conn.on("reaction_removed", { messageId: String, chatId: String, userId: String, emoji: String ->
|
||||
// Using ReactionUpdated with empty emoji to signal removal or just a specific removal event
|
||||
_events.tryEmit(ChatEvent.ReactionUpdated(messageId, chatId, userId, ""))
|
||||
}, String::class.java, String::class.java, String::class.java, String::class.java)
|
||||
|
||||
// WebRTC Signaling Handlers
|
||||
@@ -114,4 +125,10 @@ class ChatHubClient @Inject constructor() {
|
||||
hubConnection?.stop()
|
||||
_status.value = ConnectionStatus.DISCONNECTED
|
||||
}
|
||||
|
||||
fun readMessages(request: ReadMessagesRequest) {
|
||||
if (hubConnection?.connectionState == HubConnectionState.CONNECTED) {
|
||||
hubConnection?.invoke("read_messages", request)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,8 @@ import chats.data.remote.api.ChatApi
|
||||
import chats.data.remote.api.SendMessageRequest
|
||||
import chats.data.remote.dto.ChatDto
|
||||
import chats.data.remote.dto.MessageDto
|
||||
import chats.data.remote.dto.MediaItemDto
|
||||
import chats.data.remote.dto.ReactionDto
|
||||
import chats.domain.model.Chat
|
||||
import chats.domain.model.Message
|
||||
import chats.domain.model.MediaType
|
||||
@@ -14,12 +16,16 @@ import okhttp3.MediaType.Companion.toMediaTypeOrNull
|
||||
import okhttp3.MultipartBody
|
||||
import okhttp3.RequestBody.Companion.asRequestBody
|
||||
import javax.inject.Inject
|
||||
import kotlinx.coroutines.flow.map
|
||||
|
||||
class ChatRepositoryImpl @Inject constructor(
|
||||
private val api: ChatApi,
|
||||
private val tokenManager: TokenManager,
|
||||
private val serverConfig: ServerConfig
|
||||
private val serverConfig: ServerConfig,
|
||||
private val messageDao: core.database.data.MessageDao,
|
||||
private val hubClient: chats.data.remote.signalr.ChatHubClient
|
||||
) : ChatRepository {
|
||||
private val gson = com.google.gson.Gson()
|
||||
|
||||
override suspend fun getChats(): List<Chat> {
|
||||
val currentUserId = tokenManager.getUserId() ?: ""
|
||||
@@ -27,9 +33,24 @@ class ChatRepositoryImpl @Inject constructor(
|
||||
return api.getChats().map { it.toDomain(currentUserId, baseUrl) }
|
||||
}
|
||||
|
||||
override fun getMessagesFlow(chatId: String): kotlinx.coroutines.flow.Flow<List<Message>> {
|
||||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||||
return messageDao.getMessages(chatId).map { entities: List<core.database.data.MessageEntity> ->
|
||||
entities.map { it.toDomain(baseUrl, gson) }
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun getMessages(chatId: String): List<Message> {
|
||||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||||
return api.getMessages(chatId).map { it.toDomain(baseUrl) }
|
||||
return try {
|
||||
val messages = api.getMessages(chatId)
|
||||
// Save to DB
|
||||
messageDao.insertMessages(messages.map { msg: MessageDto -> msg.toEntity(baseUrl, gson) })
|
||||
messages.map { it.toDomain(baseUrl) }
|
||||
} catch (e: Exception) {
|
||||
// If network fails, caller should ideally use the Flow from DB
|
||||
emptyList()
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun sendMessage(
|
||||
@@ -55,8 +76,21 @@ class ChatRepositoryImpl @Inject constructor(
|
||||
api.sendTypingStatus(chatId)
|
||||
}
|
||||
|
||||
override suspend fun markMessagesAsRead(chatId: String, lastMessageId: String) {
|
||||
api.markMessagesAsRead(chatId, lastMessageId)
|
||||
override suspend fun markMessagesAsRead(chatId: String, lastMessageId: String, lastReadSequenceId: Int) {
|
||||
val request = chats.data.remote.signalr.ReadMessagesRequest(
|
||||
chatId = chatId,
|
||||
lastReadMessageId = lastMessageId,
|
||||
lastReadSequenceId = lastReadSequenceId
|
||||
)
|
||||
hubClient.readMessages(request)
|
||||
}
|
||||
|
||||
override suspend fun saveMessage(message: Message) {
|
||||
messageDao.insertMessages(listOf(message.toEntity(gson)))
|
||||
}
|
||||
|
||||
override suspend fun deleteLocalMessage(messageId: String) {
|
||||
messageDao.deleteMessage(messageId)
|
||||
}
|
||||
|
||||
override suspend fun uploadMedia(file: java.io.File): String {
|
||||
@@ -98,6 +132,24 @@ fun ChatDto.toDomain(currentUserId: String, baseUrl: String): Chat {
|
||||
)
|
||||
}
|
||||
|
||||
fun Message.toEntity(gson: com.google.gson.Gson): core.database.data.MessageEntity {
|
||||
return core.database.data.MessageEntity(
|
||||
id = id,
|
||||
chatId = chatId,
|
||||
senderId = senderId,
|
||||
senderName = senderName,
|
||||
senderAvatar = senderAvatar,
|
||||
content = content,
|
||||
sequenceId = sequenceId,
|
||||
createdAt = createdAt,
|
||||
mediaType = mediaType.name.lowercase(),
|
||||
mediaJson = gson.toJson(media),
|
||||
reactionsJson = gson.toJson(reactions),
|
||||
isRead = isRead,
|
||||
replyToId = replyTo?.id
|
||||
)
|
||||
}
|
||||
|
||||
fun MessageDto.toDomain(baseUrl: String): Message {
|
||||
val domainMediaType = when (type) {
|
||||
"gif" -> MediaType.GIF
|
||||
@@ -135,10 +187,70 @@ fun MessageDto.toDomain(baseUrl: String): Message {
|
||||
},
|
||||
mediaType = domainMediaType,
|
||||
reactions = reactions?.associate { it.emoji to it.count } ?: emptyMap(),
|
||||
isRead = true, // Network messages are usually considered read when fetched or handled by server
|
||||
replyTo = replyTo?.toDomain(baseUrl)
|
||||
)
|
||||
}
|
||||
|
||||
fun MessageDto.toEntity(baseUrl: String, gson: com.google.gson.Gson): core.database.data.MessageEntity {
|
||||
return core.database.data.MessageEntity(
|
||||
id = id,
|
||||
chatId = chatId ?: "",
|
||||
senderId = senderId ?: "",
|
||||
senderName = sender?.displayName ?: "Unknown",
|
||||
senderAvatar = sender?.avatarUrl?.ensureAbsoluteUrl(baseUrl),
|
||||
content = content,
|
||||
sequenceId = sequenceId ?: 0,
|
||||
createdAt = createdAt ?: "",
|
||||
mediaType = type ?: "text",
|
||||
mediaJson = gson.toJson(media),
|
||||
reactionsJson = gson.toJson(reactions),
|
||||
isRead = true,
|
||||
replyToId = replyTo?.id
|
||||
)
|
||||
}
|
||||
|
||||
fun core.database.data.MessageEntity.toDomain(baseUrl: String, gson: com.google.gson.Gson): Message {
|
||||
val mediaTypeEnum = when (mediaType) {
|
||||
"gif" -> MediaType.GIF
|
||||
"image", "photo" -> MediaType.IMAGE
|
||||
"video" -> MediaType.VIDEO
|
||||
"audio", "voice" -> MediaType.AUDIO
|
||||
"file" -> MediaType.FILE
|
||||
else -> MediaType.TEXT
|
||||
}
|
||||
|
||||
val mediaTypeToken = object : com.google.gson.reflect.TypeToken<List<chats.data.remote.dto.MediaItemDto>>() {}.type
|
||||
val mediaDtos: List<chats.data.remote.dto.MediaItemDto> = gson.fromJson(mediaJson, mediaTypeToken) ?: emptyList()
|
||||
|
||||
val reactionsTypeToken = object : com.google.gson.reflect.TypeToken<List<chats.data.remote.dto.ReactionDto>>() {}.type
|
||||
val reactionDtos: List<chats.data.remote.dto.ReactionDto> = gson.fromJson(reactionsJson, reactionsTypeToken) ?: emptyList()
|
||||
|
||||
return Message(
|
||||
id = id,
|
||||
chatId = chatId,
|
||||
senderId = senderId,
|
||||
senderName = senderName,
|
||||
senderAvatar = senderAvatar,
|
||||
content = content,
|
||||
sequenceId = sequenceId,
|
||||
createdAt = createdAt,
|
||||
media = mediaDtos.map {
|
||||
chats.domain.model.Media(
|
||||
id = it.id,
|
||||
type = it.type,
|
||||
url = it.url.ensureAbsoluteUrl(baseUrl),
|
||||
filename = it.filename,
|
||||
size = it.size,
|
||||
duration = it.duration
|
||||
)
|
||||
},
|
||||
mediaType = mediaTypeEnum,
|
||||
reactions = reactionDtos.associate { it.emoji to it.count },
|
||||
isRead = isRead
|
||||
)
|
||||
}
|
||||
|
||||
fun String.ensureAbsoluteUrl(baseUrl: String): String {
|
||||
return if (this.startsWith("http")) {
|
||||
this
|
||||
|
||||
Reference in New Issue
Block a user