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() 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.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.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.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.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.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 } }