Files
sodalive-ios/SodaLive/Sources/V2/Main/Chat/UserCreatorChat/UserCreatorChatRoomViewModel.swift

563 lines
21 KiB
Swift

import Foundation
import Combine
import UIKit
final class UserCreatorChatRoomViewModel: ObservableObject {
@Published var isLoading = false
@Published var isLoadingNextPage = false
@Published var isShowPopup = false
@Published var errorMessage = ""
@Published private(set) var roomId = 0
@Published private(set) var opponentNickname = ""
@Published private(set) var opponentProfileImageUrl: String?
@Published private(set) var messages = [UserCreatorChatDisplayMessage]()
@Published private(set) var hasMore = false
@Published private(set) var nextCursor: Int?
@Published private(set) var socketState: UserCreatorChatSocketState = .disconnected
private let repository = UserCreatorChatRepository()
private let socketClient = WebSocketChatClient()
private var subscription = Set<AnyCancellable>()
private var ackTimeouts: [String: DispatchWorkItem] = [:]
private var reconnectWorkItem: DispatchWorkItem?
private var isScreenAlive = true
private var shouldReconnect = true
private var isAppForeground = UIApplication.shared.applicationState != .background
private var hasOpenedRoom = false
private var joinRetryCount = 0
private let pageLimit = 20
private let pendingMatchClockSkewToleranceMilliseconds: Int64 = 5 * 60 * 1000
init() {
bindSocket()
observeLifecycle()
}
deinit {
leaveAndClose()
NotificationCenter.default.removeObserver(self)
}
func enter(roomId: Int) {
guard roomId > 0 else { return }
prepareEnter()
openRoom(roomId: roomId)
}
func enter(creatorId: Int) {
guard creatorId > 0 else { return }
prepareEnter()
repository.createRoom(creatorId: creatorId)
.receive(on: DispatchQueue.main)
.sink { [weak self] result in
guard let self else { return }
switch result {
case .finished:
DEBUG_LOG("finish")
case .failure(let error):
ERROR_LOG(error.localizedDescription)
self.applyRestFailure(message: nil)
}
} receiveValue: { [weak self] response in
guard let self else { return }
do {
let decoded = try JSONDecoder().decode(ApiResponse<UserCreatorCreateRoomResponse>.self, from: response.data)
if let data = decoded.data, decoded.success, data.roomId > 0 {
self.openRoom(roomId: data.roomId)
} else {
self.applyRestFailure(message: decoded.message)
}
} catch {
ERROR_LOG(error.localizedDescription)
self.applyRestFailure(message: nil)
}
}
.store(in: &subscription)
}
func sendText(_ text: String) {
let trimmedText = text.trimmingCharacters(in: .whitespacesAndNewlines)
guard trimmedText.isEmpty == false else { return }
guard socketState == .joined, roomId > 0 else { return }
guard let requestId = socketClient.sendText(trimmedText) else { return }
let pendingMessage = UserCreatorChatDisplayMessage(
pendingText: trimmedText,
requestId: requestId,
createdAt: Int64(Date().timeIntervalSince1970 * 1000)
)
messages.append(pendingMessage)
startAckTimeout(requestId: requestId)
}
func loadMore() {
guard roomId > 0,
hasMore,
isLoading == false,
isLoadingNextPage == false,
socketState == .joined else { return }
isLoadingNextPage = true
repository.getMessages(roomId: roomId, cursor: nextCursor, limit: pageLimit)
.receive(on: DispatchQueue.main)
.sink { [weak self] result in
guard let self else { return }
switch result {
case .finished:
DEBUG_LOG("finish")
case .failure(let error):
ERROR_LOG(error.localizedDescription)
self.isLoadingNextPage = false
self.errorMessage = I18n.Common.commonError
self.isShowPopup = true
}
} receiveValue: { [weak self] response in
guard let self else { return }
self.isLoadingNextPage = false
do {
let decoded = try JSONDecoder().decode(ApiResponse<UserCreatorChatMessagesResponse>.self, from: response.data)
if let data = decoded.data, decoded.success {
let previousCursor = self.nextCursor
self.hasMore = data.hasMore
self.nextCursor = data.nextCursor
self.mergeServerMessages(data.messages)
let shouldContinuePastVoiceOnlyPage = data.hasMore
&& data.nextCursor != previousCursor
&& data.messages.contains(where: { $0.messageType == "TEXT" }) == false
if shouldContinuePastVoiceOnlyPage {
DispatchQueue.main.async { [weak self] in
self?.loadMore()
}
}
} else {
self.errorMessage = decoded.message ?? I18n.Common.commonError
self.isShowPopup = true
}
} catch {
ERROR_LOG(error.localizedDescription)
self.errorMessage = I18n.Common.commonError
self.isShowPopup = true
}
}
.store(in: &subscription)
}
func sendVoiceMessage(soundData: Data) {
guard roomId > 0 else { return }
repository.sendVoiceMessage(roomId: roomId, voiceData: soundData)
.receive(on: DispatchQueue.main)
.sink { result in
switch result {
case .finished:
DEBUG_LOG("finish")
case .failure(let error):
ERROR_LOG(error.localizedDescription)
}
} receiveValue: { response in
do {
let decoded = try JSONDecoder().decode(ApiResponse<UserCreatorVoiceMessageResponse>.self, from: response.data)
if let data = decoded.data, decoded.success {
DEBUG_LOG("deliveredRealtime=\(data.deliveredRealtime), pushSent=\(data.pushSent)")
} else if let message = decoded.message {
ERROR_LOG(message)
}
} catch {
ERROR_LOG(error.localizedDescription)
}
}
.store(in: &subscription)
}
func leaveAndClose() {
isScreenAlive = false
shouldReconnect = false
cancelReconnect()
cancelAckTimeouts()
subscription.removeAll()
socketClient.close()
}
private func prepareEnter() {
cancelReconnect()
isScreenAlive = true
shouldReconnect = true
isAppForeground = UIApplication.shared.applicationState != .background
hasOpenedRoom = false
joinRetryCount = 0
isLoading = true
isShowPopup = false
errorMessage = ""
messages = []
hasMore = false
nextCursor = nil
roomId = 0
}
private func openRoom(roomId: Int) {
guard roomId > 0 else { return }
repository.openRoom(roomId: roomId, limit: pageLimit)
.receive(on: DispatchQueue.main)
.sink { [weak self] result in
guard let self else { return }
guard self.isScreenAlive, self.shouldReconnect else { return }
switch result {
case .finished:
DEBUG_LOG("finish")
case .failure(let error):
ERROR_LOG(error.localizedDescription)
self.applyRestFailure(message: nil)
}
} receiveValue: { [weak self] response in
guard let self else { return }
guard self.isScreenAlive, self.shouldReconnect else { return }
do {
let decoded = try JSONDecoder().decode(ApiResponse<UserCreatorOpenRoomResponse>.self, from: response.data)
if let data = decoded.data, decoded.success, data.roomId > 0 {
self.applyOpenRoom(data)
if self.isAppForeground {
self.socketClient.connect(roomId: data.roomId)
}
} else {
self.applyRestFailure(message: decoded.message)
}
} catch {
ERROR_LOG(error.localizedDescription)
self.applyRestFailure(message: nil)
}
}
.store(in: &subscription)
}
private func applyOpenRoom(_ data: UserCreatorOpenRoomResponse) {
roomId = data.roomId
hasOpenedRoom = true
opponentNickname = data.opponentNickname
opponentProfileImageUrl = data.opponentProfileImageUrl
hasMore = data.hasMore
nextCursor = data.nextCursor
messages = data.messages
.sorted { $0.createdAt < $1.createdAt }
.map { UserCreatorChatDisplayMessage(message: $0) }
}
private func bindSocket() {
socketClient.onStateChange = { [weak self] state in
DispatchQueue.main.async {
self?.socketState = state
}
}
socketClient.onEvent = { [weak self] event in
DispatchQueue.main.async {
self?.handleSocketEvent(event)
}
}
}
private func handleSocketEvent(_ event: UserCreatorChatSocketEvent) {
switch event {
case .joined:
cancelReconnect()
joinRetryCount = 0
syncLatestMessages(updatePagination: true) { [weak self] isSuccess in
guard isSuccess else { return }
self?.loadMoreIfLatestPageHasNoText()
}
case .joinFailed(let message):
retryJoinOrFail(messageKey: message)
case .sendAck(let requestId, let message):
handleSendAck(requestId: requestId, message: message)
case .message(let message):
mergeServerMessages([message])
case .pong:
DEBUG_LOG("PONG handled")
case .error(let messageKey):
if socketState != .joined {
retryJoinOrFail(messageKey: messageKey)
} else if let messageKey {
ERROR_LOG(messageKey)
}
case .closed:
reconnectIfNeeded()
}
}
private func handleSendAck(requestId: String, message: UserCreatorChatMessageItem?) {
ackTimeouts[requestId]?.cancel()
ackTimeouts[requestId] = nil
guard let message else {
let pendingMessage = messages.first(where: { $0.requestId == requestId && $0.status == .pending })
syncLatestMessages { [weak self] _ in
guard let self else { return }
if let index = self.messages.firstIndex(where: { $0.requestId == requestId && $0.status == .pending }) {
if self.hasMatchingOwnServerMessage(for: pendingMessage) {
self.messages.remove(at: index)
} else {
self.messages[index].status = .failed
}
}
}
return
}
if let pendingIndex = messages.firstIndex(where: { $0.requestId == requestId }) {
messages.remove(at: pendingIndex)
}
mergeServerMessages([message])
}
private func startAckTimeout(requestId: String) {
let workItem = DispatchWorkItem { [weak self] in
guard let self else { return }
if let index = self.messages.firstIndex(where: { $0.requestId == requestId && $0.status == .pending }) {
self.messages[index].status = .failed
}
self.ackTimeouts[requestId] = nil
}
ackTimeouts[requestId] = workItem
DispatchQueue.main.asyncAfter(deadline: .now() + 15, execute: workItem)
}
private func hasMatchingOwnServerMessage(for pendingMessage: UserCreatorChatDisplayMessage?) -> Bool {
guard let pendingMessage else { return false }
return messages.contains { message in
return message.messageId != nil
&& message.mine
&& message.textMessage == pendingMessage.textMessage
&& isWithinPendingMatchWindow(message.createdAt, pendingMessage.createdAt)
}
}
private func isWithinPendingMatchWindow(_ lhsCreatedAt: Int64, _ rhsCreatedAt: Int64) -> Bool {
abs(lhsCreatedAt - rhsCreatedAt) <= pendingMatchClockSkewToleranceMilliseconds
}
private func mergeServerMessages(_ serverMessages: [UserCreatorChatMessageItem]) {
var merged = messages
let visibleServerMessages = serverMessages
.filter { $0.messageType != "VOICE" }
.sorted { $0.createdAt < $1.createdAt }
for serverMessage in visibleServerMessages {
let serverDisplayMessage = UserCreatorChatDisplayMessage(message: serverMessage)
if let index = merged.firstIndex(where: { $0.messageId == serverMessage.messageId }) {
merged[index] = serverDisplayMessage
} else {
removeMatchingLocalMessage(for: serverMessage, from: &merged)
merged.append(serverDisplayMessage)
}
}
messages = merged.sorted { lhs, rhs in
if lhs.createdAt == rhs.createdAt {
return lhs.id < rhs.id
}
return lhs.createdAt < rhs.createdAt
}
}
private func removeMatchingLocalMessage(
for serverMessage: UserCreatorChatMessageItem,
from messages: inout [UserCreatorChatDisplayMessage]
) {
guard serverMessage.messageType == "TEXT", serverMessage.mine else { return }
let matchingIndexes = messages.indices.filter { index in
let message = messages[index]
return message.messageId == nil
&& message.requestId != nil
&& message.mine
&& message.messageType == "TEXT"
&& message.textMessage == serverMessage.textMessage
&& isWithinPendingMatchWindow(serverMessage.createdAt, message.createdAt)
}
guard let index = matchingIndexes.min(by: { lhs, rhs in
let lhsDistance = abs(serverMessage.createdAt - messages[lhs].createdAt)
let rhsDistance = abs(serverMessage.createdAt - messages[rhs].createdAt)
if lhsDistance == rhsDistance {
return messages[lhs].createdAt > messages[rhs].createdAt
}
return lhsDistance < rhsDistance
}) else { return }
if let requestId = messages[index].requestId {
ackTimeouts[requestId]?.cancel()
ackTimeouts[requestId] = nil
}
messages.remove(at: index)
}
private func syncLatestMessages(
updatePagination: Bool = false,
completion: ((Bool) -> Void)? = nil
) {
guard roomId > 0 else {
completion?(false)
return
}
isLoading = true
repository.openRoom(roomId: roomId, limit: pageLimit)
.receive(on: DispatchQueue.main)
.sink { [weak self] result in
switch result {
case .finished:
DEBUG_LOG("finish")
case .failure(let error):
ERROR_LOG(error.localizedDescription)
self?.applyRestFailure(message: nil)
completion?(false)
}
} receiveValue: { [weak self] response in
guard let self else {
completion?(false)
return
}
do {
let decoded = try JSONDecoder().decode(ApiResponse<UserCreatorOpenRoomResponse>.self, from: response.data)
if let data = decoded.data, decoded.success {
if updatePagination {
self.hasMore = data.hasMore
self.nextCursor = data.nextCursor
}
self.mergeServerMessages(data.messages)
self.isLoading = false
completion?(true)
} else {
self.applyRestFailure(message: decoded.message)
completion?(false)
}
} catch {
ERROR_LOG(error.localizedDescription)
self.applyRestFailure(message: nil)
completion?(false)
}
}
.store(in: &subscription)
}
private func retryJoinOrFail(messageKey: String?) {
if let messageKey {
ERROR_LOG(messageKey)
}
guard joinRetryCount < 3, roomId > 0 else {
cancelReconnect()
isLoading = false
errorMessage = I18n.UserCreatorChat.joinFailureToast
isShowPopup = true
shouldReconnect = false
socketClient.closeWithoutLeave(notify: false)
DispatchQueue.main.asyncAfter(deadline: .now() + 2) { [weak self] in
guard let self, self.isScreenAlive, self.shouldReconnect == false else { return }
AppState.shared.back()
}
return
}
joinRetryCount += 1
isLoading = true
socketClient.closeWithoutLeave(notify: false)
scheduleReconnect()
}
private func reconnectIfNeeded() {
guard isScreenAlive, shouldReconnect, isAppForeground, hasOpenedRoom, roomId > 0 else { return }
isLoading = true
scheduleReconnect()
}
private func observeLifecycle() {
NotificationCenter.default.addObserver(
self,
selector: #selector(handleDidEnterBackground),
name: UIApplication.didEnterBackgroundNotification,
object: nil
)
NotificationCenter.default.addObserver(
self,
selector: #selector(handleWillEnterForeground),
name: UIApplication.willEnterForegroundNotification,
object: nil
)
NotificationCenter.default.addObserver(
self,
selector: #selector(handleCloseNotification),
name: .userCreatorChatSessionClose,
object: nil
)
}
@objc private func handleDidEnterBackground() {
isAppForeground = false
cancelReconnect()
failPendingMessages()
cancelAckTimeouts()
socketClient.close()
}
@objc private func handleWillEnterForeground() {
isAppForeground = true
guard isScreenAlive, shouldReconnect, hasOpenedRoom, roomId > 0 else { return }
cancelReconnect()
isLoading = true
socketClient.connect(roomId: roomId)
}
@objc private func handleCloseNotification() {
DispatchQueue.main.async { [weak self] in
self?.leaveAndClose()
}
}
private func applyRestFailure(message: String?) {
isLoading = false
errorMessage = message ?? I18n.Common.commonError
isShowPopup = true
}
private func cancelAckTimeouts() {
ackTimeouts.values.forEach { $0.cancel() }
ackTimeouts.removeAll()
}
private func failPendingMessages() {
for index in messages.indices where messages[index].status == .pending {
messages[index].status = .failed
}
}
private func loadMoreIfLatestPageHasNoText() {
guard messages.contains(where: { $0.messageType == "TEXT" }) == false else { return }
loadMore()
}
private func scheduleReconnect() {
cancelReconnect()
let workItem = DispatchWorkItem { [weak self] in
guard let self,
self.isScreenAlive,
self.shouldReconnect,
self.isAppForeground,
self.hasOpenedRoom,
self.roomId > 0 else { return }
self.reconnectWorkItem = nil
self.socketClient.connect(roomId: self.roomId)
}
reconnectWorkItem = workItem
DispatchQueue.main.asyncAfter(deadline: .now() + 1, execute: workItem)
}
private func cancelReconnect() {
reconnectWorkItem?.cancel()
reconnectWorkItem = nil
}
}