package core.database.data import androidx.paging.PagingSource import androidx.room.* import kotlinx.coroutines.flow.Flow /** * Статус синхронизации сообщения с сервером */ enum class SyncStatus { SYNCED, // Сообщение успешно синхронизировано SYNCING, // Сообщение отправляется на сервер FAILED // Ошибка синхронизации } /** * Entity для хранения чатов в локальной базе данных */ @Entity(tableName = "chats") data class ChatEntity( @PrimaryKey val id: String, val name: String, val avatar: String?, val type: String = "personal", // personal, group, saved val lastMessageId: String? = null, val lastMessageText: String? = null, val lastMessageAt: Long = 0, val unreadCount: Int = 0, val isPinned: Boolean = false, val lastUpdated: Long = System.currentTimeMillis() ) /** * Entity для хранения сообщений в локальной базе данных * Поддерживает офлайн-работу и фоновую синхронизацию */ @Entity(tableName = "messages") data class MessageEntity( @PrimaryKey val id: String, val chatId: String, val senderId: String, val senderName: String, val senderAvatar: String?, val content: String?, val sequenceId: Int, val createdAt: String, val mediaType: String, val mediaJson: String, val reactionsJson: String, val isRead: Boolean, val replyToId: String? = null, // Поля для офлайн-синхронизации @ColumnInfo(defaultValue = "SYNCED") val syncStatus: SyncStatus = SyncStatus.SYNCED, @ColumnInfo(defaultValue = "0") val isDeletedLocally: Boolean = false, @ColumnInfo(defaultValue = "0") val isEditedLocally: Boolean = false, val editedContent: String? = null, @ColumnInfo(defaultValue = "0") val lastUpdated: Long = 0 ) @Dao interface MessageDao { // ==================== Flow для UI ==================== @Query("SELECT * FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0 ORDER BY sequenceId ASC") fun getMessages(chatId: String): Flow> // ==================== Paging 3 ==================== @Query("SELECT * FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0 ORDER BY sequenceId DESC") fun getMessagesPagingSource(chatId: String): PagingSource // Загружаем сообщения С МЕНЬШИМ sequenceId (старые), порядок DESC для PagingSource @Query("SELECT * FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0 AND sequenceId < :sequenceId ORDER BY sequenceId DESC LIMIT :limit") suspend fun getMessagesBefore(chatId: String, sequenceId: Int, limit: Int): List // Загружаем сообщения до и ВКЛЮЧАЯ sequenceId, порядок DESC @Query("SELECT * FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0 AND sequenceId <= :sequenceId ORDER BY sequenceId DESC LIMIT :limit") suspend fun getMessagesUpToAndIncluding(chatId: String, sequenceId: Int, limit: Int): List // Загружаем сообщения С БОЛЬШИМ sequenceId (новые), порядок ASC @Query("SELECT * FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0 AND sequenceId > :sequenceId ORDER BY sequenceId ASC LIMIT :limit") suspend fun getMessagesAfter(chatId: String, sequenceId: Int, limit: Int): List @Query("SELECT MAX(sequenceId) FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0") suspend fun getMaxSequenceId(chatId: String): Int? @Query("SELECT MIN(sequenceId) FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0") suspend fun getMinSequenceId(chatId: String): Int? // ==================== Основные операции ==================== @Insert(onConflict = OnConflictStrategy.REPLACE) suspend fun insertMessages(messages: List) @Insert(onConflict = OnConflictStrategy.REPLACE) suspend fun insertMessage(message: MessageEntity) /** * Upsert - вставляет или обновляет сообщение * Приоритет: серверные данные > локальные (кроме сообщений в процессе отправки) */ @Transaction suspend fun upsertMessage(message: MessageEntity) { val existing = getMessageById(message.id) if (existing != null) { // Сохраняем локальные изменения если сообщение в процессе отправки val syncedMessage = when { existing.syncStatus == SyncStatus.SYNCING || existing.isEditedLocally -> { message.copy( syncStatus = existing.syncStatus, isEditedLocally = existing.isEditedLocally, editedContent = existing.editedContent, lastUpdated = existing.lastUpdated ) } else -> message } insertMessage(syncedMessage) } else { insertMessage(message) } } @Transaction suspend fun upsertMessages(messages: List) { messages.forEach { upsertMessage(it) } } @Query("SELECT * FROM messages WHERE id = :id LIMIT 1") suspend fun getMessageById(id: String): MessageEntity? @Query("SELECT * FROM messages WHERE id = :id LIMIT 1") fun getMessageByIdFlow(id: String): Flow @Query("DELETE FROM messages WHERE chatId = :chatId") suspend fun clearChat(chatId: String) @Query("DELETE FROM messages WHERE id = :messageId") suspend fun deleteMessage(messageId: String) @Query("UPDATE messages SET isRead = 1 WHERE chatId = :chatId AND sequenceId <= :lastReadSequenceId AND isRead = 0") suspend fun markMessagesAsRead(chatId: String, lastReadSequenceId: Int) // ==================== Офлайн операции ==================== /** * Помечает сообщение как удалённое локально * Фактическое удаление произойдёт после синхронизации с сервером */ @Query("UPDATE messages SET isDeletedLocally = 1, syncStatus = 'SYNCING', lastUpdated = :timestamp WHERE id = :messageId") suspend fun markAsDeletedLocally(messageId: String, timestamp: Long = System.currentTimeMillis()) /** * Помечает сообщение как отредактированное локально */ @Query("UPDATE messages SET isEditedLocally = 1, editedContent = :newContent, syncStatus = 'SYNCING', lastUpdated = :timestamp WHERE id = :messageId") suspend fun markAsEditedLocally(messageId: String, newContent: String, timestamp: Long = System.currentTimeMillis()) /** * Сбрасывает локальные изменения после успешной синхронизации */ @Query("UPDATE messages SET syncStatus = 'SYNCED', isDeletedLocally = 0, isEditedLocally = 0, editedContent = NULL WHERE id = :messageId") suspend fun markAsSynced(messageId: String) /** * Устанавливает статус FAILED для сообщений с ошибкой синхронизации */ @Query("UPDATE messages SET syncStatus = 'FAILED' WHERE id = :messageId") suspend fun markAsSyncFailed(messageId: String) /** * Получает все сообщения, требующие синхронизации */ @Query("SELECT * FROM messages WHERE syncStatus = 'SYNCING' OR isDeletedLocally = 1 ORDER BY lastUpdated ASC") suspend fun getPendingSyncMessages(): List @Query("SELECT * FROM messages WHERE syncStatus = 'SYNCING' OR isDeletedLocally = 1 ORDER BY lastUpdated ASC") fun getPendingSyncMessagesFlow(): Flow> /** * Получает сообщения с_failed статусом для повторной отправки */ @Query("SELECT * FROM messages WHERE syncStatus = 'FAILED' ORDER BY lastUpdated ASC") suspend fun getFailedSyncMessages(): List // ==================== Утилиты ==================== @Query("SELECT COUNT(*) FROM messages WHERE chatId = :chatId AND isDeletedLocally = 0") suspend fun getMessagesCount(chatId: String): Int @Query("SELECT COUNT(*) FROM messages WHERE chatId = :chatId AND isRead = 0 AND isDeletedLocally = 0") suspend fun getUnreadCount(chatId: String): Int @Query("SELECT EXISTS(SELECT 1 FROM messages WHERE id = :id)") suspend fun exists(id: String): Boolean } @Dao interface ChatDao { @Query("SELECT * FROM chats ORDER BY isPinned DESC, lastMessageAt DESC") fun getAllChatsFlow(): Flow> @Query("SELECT * FROM chats ORDER BY isPinned DESC, lastMessageAt DESC") suspend fun getAllChats(): List @Query("SELECT * FROM chats WHERE id = :chatId LIMIT 1") suspend fun getChatById(chatId: String): ChatEntity? @Query("SELECT * FROM chats WHERE id = :chatId LIMIT 1") fun getChatByIdFlow(chatId: String): Flow @Insert(onConflict = OnConflictStrategy.REPLACE) suspend fun insertChat(chat: ChatEntity) @Insert(onConflict = OnConflictStrategy.REPLACE) suspend fun insertChats(chats: List) @Query("DELETE FROM chats WHERE id = :chatId") suspend fun deleteChat(chatId: String) @Query("DELETE FROM chats") suspend fun clearAll() @Query("UPDATE chats SET unreadCount = :count WHERE id = :chatId") suspend fun updateUnreadCount(chatId: String, count: Int) @Query("UPDATE chats SET lastMessageId = :lastMessageId, lastMessageText = :lastMessageText, lastMessageAt = :lastMessageAt WHERE id = :chatId") suspend fun updateLastMessage(chatId: String, lastMessageId: String?, lastMessageText: String?, lastMessageAt: Long) } @Database(entities = [MessageEntity::class, ChatEntity::class], version = 3) @TypeConverters(SyncStatusConverter::class) abstract class ChatDatabase : RoomDatabase() { abstract fun messageDao(): MessageDao abstract fun chatDao(): ChatDao companion object { const val DATABASE_NAME = "knot_chat_database" } } /** * Конвертер для Enum SyncStatus */ class SyncStatusConverter { @TypeConverter fun fromSyncStatus(status: SyncStatus): String = status.name @TypeConverter fun toSyncStatus(value: String): SyncStatus = SyncStatus.valueOf(value) }