From 26e51a07c561aba469b6947a437d9d5584f92372 Mon Sep 17 00:00:00 2001 From: Evgeny Poberezkin <2769109+epoberezkin@users.noreply.github.com> Date: Thu, 21 Jul 2022 14:26:46 +0100 Subject: [PATCH] ios: improve concurrency of NSE, process multiple messages (#827) --- .../ios/SimpleX NSE/NotificationService.swift | 181 ++++++++++++------ apps/ios/SimpleXChat/APITypes.swift | 13 -- apps/ios/SimpleXChat/ChatTypes.swift | 51 ++++- apps/ios/SimpleXChat/Notifications.swift | 4 +- 4 files changed, 168 insertions(+), 81 deletions(-) diff --git a/apps/ios/SimpleX NSE/NotificationService.swift b/apps/ios/SimpleX NSE/NotificationService.swift index 32cef1a4bc..1b3cdc8fd9 100644 --- a/apps/ios/SimpleX NSE/NotificationService.swift +++ b/apps/ios/SimpleX NSE/NotificationService.swift @@ -14,17 +14,58 @@ let logger = Logger() let suspendingDelay: UInt64 = 2_000_000_000 +typealias NtfStream = AsyncStream + +actor PendingNtfs { + static let shared = PendingNtfs() + private var ntfStreams: [String: NtfStream] = [:] + private var ntfConts: [String: NtfStream.Continuation] = [:] + + func createStream(_ id: String) { + logger.debug("PendingNtfs.createStream: \(id, privacy: .public)") + if ntfStreams.index(forKey: id) == nil { + ntfStreams[id] = AsyncStream { cont in + ntfConts[id] = cont + logger.debug("PendingNtfs.createStream: store continuation") + } + } + } + + func readStream(_ id: String, for nse: NotificationService, msgCount: Int = 1) async { + logger.debug("PendingNtfs.readStream: \(id, privacy: .public) \(msgCount, privacy: .public)") + if let s = ntfStreams[id] { + logger.debug("PendingNtfs.readStream: has stream") + var rcvCount = max(1, msgCount) + for await ntf in s { + nse.setBestAttemptNtf(ntf) + rcvCount -= 1 + if rcvCount == 0 { break } + } + logger.debug("PendingNtfs.readStream: exiting") + } + } + + func writeStream(_ id: String, _ ntf: UNMutableNotificationContent) { + logger.debug("PendingNtfs.writeStream: \(id, privacy: .public)") + if let cont = ntfConts[id] { + logger.debug("PendingNtfs.writeStream: writing ntf") + cont.yield(ntf) + } + } +} + + + class NotificationService: UNNotificationServiceExtension { var contentHandler: ((UNNotificationContent) -> Void)? - var bestAttemptContent: UNMutableNotificationContent? + var bestAttemptNtf: UNMutableNotificationContent? var badgeCount: Int = 0 override func didReceive(_ request: UNNotificationRequest, withContentHandler contentHandler: @escaping (UNNotificationContent) -> Void) { logger.debug("NotificationService.didReceive") - bestAttemptContent = request.content.mutableCopy() as? UNMutableNotificationContent badgeCount = ntfBadgeCountGroupDefault.get() + 1 ntfBadgeCountGroupDefault.set(badgeCount) - bestAttemptContent?.badge = badgeCount as NSNumber + setBestAttemptNtf(request.content.mutableCopy() as? UNMutableNotificationContent) self.contentHandler = contentHandler let appState = appStateGroupDefault.get() switch appState { @@ -43,22 +84,20 @@ class NotificationService: UNNotificationServiceExtension { logger.debug("NotificationService: app state is \(state.rawValue, privacy: .public)") if state.inactive { receiveNtfMessages(request, contentHandler) - } else if let content = bestAttemptContent { - contentHandler(content) + } else { + deliverBestAttemptNtf() } } default: logger.debug("NotificationService: app state is \(appState.rawValue, privacy: .public)") - if let content = bestAttemptContent { - contentHandler(content) - } + deliverBestAttemptNtf() } } func receiveNtfMessages(_ request: UNNotificationRequest, _ contentHandler: @escaping (UNNotificationContent) -> Void) { logger.debug("NotificationService: receiveNtfMessages") if case .documents = dbContainerGroupDefault.get() { - if let content = bestAttemptContent { contentHandler(content) } + deliverBestAttemptNtf() return } let userInfo = request.content.userInfo @@ -70,27 +109,38 @@ class NotificationService: UNNotificationServiceExtension { if let ntfMsgInfo = apiGetNtfMessage(nonce: nonce, encNtfInfo: encNtfInfo) { logger.debug("NotificationService: receiveNtfMessages: apiGetNtfMessage \(String(describing: ntfMsgInfo), privacy: .public)") if let connEntity = ntfMsgInfo.connEntity { - bestAttemptContent = createConnectionEventNtf(connEntity) - bestAttemptContent?.badge = badgeCount as NSNumber - } - if let content = receiveMessageForNotification() { - logger.debug("NotificationService: receiveMessageForNotification: has message") - content.badge = badgeCount as NSNumber - contentHandler(content) - } else if let content = bestAttemptContent { - logger.debug("NotificationService: receiveMessageForNotification: no message") - contentHandler(content) + setBestAttemptNtf(createConnectionEventNtf(connEntity)) + if let id = connEntity.id { + Task { + logger.debug("NotificationService: receiveNtfMessages: in Task, connEntity id \(id, privacy: .public)") + await PendingNtfs.shared.createStream(id) + await PendingNtfs.shared.readStream(id, for: self, msgCount: ntfMsgInfo.ntfMessages.count) + deliverBestAttemptNtf() + } + } } + return } } + deliverBestAttemptNtf() } override func serviceExtensionTimeWillExpire() { logger.debug("NotificationService.serviceExtensionTimeWillExpire") - // Called just before the extension will be terminated by the system. - // Use this as an opportunity to deliver your "best attempt" at modified content, otherwise the original push payload will be used. - if let contentHandler = self.contentHandler, let content = bestAttemptContent { - contentHandler(content) + deliverBestAttemptNtf() + } + + func setBestAttemptNtf(_ ntf: UNMutableNotificationContent?) { + logger.debug("NotificationService.setBestAttemptNtf") + bestAttemptNtf = ntf + bestAttemptNtf?.badge = badgeCount as NSNumber + } + + private func deliverBestAttemptNtf() { + logger.debug("NotificationService.deliverBestAttemptNtf") + if let handler = contentHandler, let content = bestAttemptNtf { + handler(content) + bestAttemptNtf = nil } } } @@ -100,9 +150,12 @@ func startChat() -> User? { if let user = apiGetActiveUser() { logger.debug("active user \(String(describing: user))") do { - try apiStartChat() - try apiSetFilesFolder(filesFolder: getAppFilesDirectory().path) - chatLastStartGroupDefault.set(Date.now) + let justStarted = try apiStartChat() + if justStarted { + try apiSetFilesFolder(filesFolder: getAppFilesDirectory().path) + chatLastStartGroupDefault.set(Date.now) + Task { await receiveMessages() } + } return user } catch { logger.error("NotificationService startChat error: \(responseError(error), privacy: .public)") @@ -113,43 +166,51 @@ func startChat() -> User? { return nil } -func receiveMessageForNotification() -> UNMutableNotificationContent? { - logger.debug("NotificationService receiveMessages started") +func receiveMessages() async { + logger.debug("NotificationService receiveMessages") while true { - if let res = recvSimpleXMsg() { - logger.debug("NotificationService receiveMessages: \(res.responseType)") - switch res { - case let .contactConnected(contact): - return createContactConnectedNtf(contact) - // case let .contactConnecting(contact): - // TODO profile update - case let .receivedContactRequest(contactRequest): - return createContactRequestNtf(contactRequest) - case let .newChatItem(aChatItem): - let cInfo = aChatItem.chatInfo - var cItem = aChatItem.chatItem - if case .image = cItem.content.msgContent { - if let file = cItem.file, - file.fileSize <= maxImageSize, - privacyAcceptImagesGroupDefault.get() { - cItem = apiReceiveFile(fileId: file.fileId)?.chatItem ?? cItem - } - } - return createMessageReceivedNtf(cInfo, cItem) - // case let .rcvFileComplete(aChatItem): - // TODO file received? - // let cInfo = aChatItem.chatInfo - // let cItem = aChatItem.chatItem - // NtfManager.shared.notifyMessageReceived(cInfo, cItem) - default: - logger.debug("NotificationService ignored event: \(res.responseType)") + if let msg = await chatRecvMsg() { + if let (id, ntf) = await receivedMsgNtf(msg) { + await PendingNtfs.shared.createStream(id) + await PendingNtfs.shared.writeStream(id, ntf) } - } else { - return nil } } } +func chatRecvMsg() async -> ChatResponse? { + await withCheckedContinuation { cont in + let resp = recvSimpleXMsg() + cont.resume(returning: resp) + } +} + +func receivedMsgNtf(_ res: ChatResponse) async -> (String, UNMutableNotificationContent)? { + logger.debug("NotificationService processReceivedMsg: \(res.responseType)") + switch res { + case let .contactConnected(contact): + return (contact.id, createContactConnectedNtf(contact)) +// case let .contactConnecting(contact): +// TODO profile update + case let .receivedContactRequest(contactRequest): + return (UserContact(contactRequest: contactRequest).id, createContactRequestNtf(contactRequest)) + case let .newChatItem(aChatItem): + let cInfo = aChatItem.chatInfo + var cItem = aChatItem.chatItem + if case .image = cItem.content.msgContent { + if let file = cItem.file, + file.fileSize <= maxImageSize, + privacyAcceptImagesGroupDefault.get() { + cItem = apiReceiveFile(fileId: file.fileId)?.chatItem ?? cItem + } + } + return (aChatItem.chatId, createMessageReceivedNtf(cInfo, cItem)) + default: + logger.debug("NotificationService processReceivedMsg ignored event: \(res.responseType)") + return nil + } +} + func apiGetActiveUser() -> User? { let _ = getChatCtrl() let r = sendSimpleXCmd(.showActiveUser) @@ -163,11 +224,11 @@ func apiGetActiveUser() -> User? { } } -func apiStartChat() throws { +func apiStartChat() throws -> Bool { let r = sendSimpleXCmd(.startChat(subscribe: false)) switch r { - case .chatStarted: return - case .chatRunning: return + case .chatStarted: return true + case .chatRunning: return false default: throw r } } diff --git a/apps/ios/SimpleXChat/APITypes.swift b/apps/ios/SimpleXChat/APITypes.swift index 2a7847d7e1..8a0b71aa73 100644 --- a/apps/ios/SimpleXChat/APITypes.swift +++ b/apps/ios/SimpleXChat/APITypes.swift @@ -546,19 +546,6 @@ public func decodeJSON(_ json: String) -> T? { return nil } -func decodeCJSON(_ cjson: UnsafePointer) -> T? { - // TODO is there a way to do it without copying the data? e.g: - // let p = UnsafeMutableRawPointer.init(mutating: UnsafeRawPointer(cjson)) - // let d = Data.init(bytesNoCopy: p, count: strlen(cjson), deallocator: .free) - decodeJSON(String.init(cString: cjson)) -} - -private func getJSONObject(_ cjson: UnsafePointer) -> NSDictionary? { - let s = String.init(cString: cjson) - let d = s.data(using: .utf8)! - return try? JSONSerialization.jsonObject(with: d) as? NSDictionary -} - public func encodeJSON(_ value: T) -> String { let data = try! jsonEncoder.encode(value) return String(decoding: data, as: UTF8.self) diff --git a/apps/ios/SimpleXChat/ChatTypes.swift b/apps/ios/SimpleXChat/ChatTypes.swift index bfdc2b6a62..6f740008b5 100644 --- a/apps/ios/SimpleXChat/ChatTypes.swift +++ b/apps/ios/SimpleXChat/ChatTypes.swift @@ -287,10 +287,24 @@ public struct Connection: Decodable { } public struct UserContact: Decodable { + public var userContactLinkId: Int64 + + public init(userContactLinkId: Int64) { + self.userContactLinkId = userContactLinkId + } + + public init(contactRequest: UserContactRequest) { + self.userContactLinkId = contactRequest.userContactLinkId + } + + public var id: String { + "@>\(userContactLinkId)" + } } public struct UserContactRequest: Decodable, NamedChat { var contactRequestId: Int64 + public var userContactLinkId: Int64 var localDisplayName: ContactName var profile: Profile var createdAt: Date @@ -306,6 +320,7 @@ public struct UserContactRequest: Decodable, NamedChat { public static let sampleData = UserContactRequest( contactRequestId: 1, + userContactLinkId: 1, localDisplayName: "alice", profile: Profile.sampleData, createdAt: .now, @@ -408,7 +423,7 @@ public struct GroupInfo: Identifiable, Decodable, NamedChat { var updatedAt: Date public var id: ChatId { get { "#\(groupId)" } } - var apiId: Int64 { get { groupId } } + public var apiId: Int64 { get { groupId } } public var ready: Bool { get { true } } public var sendMsgEnabled: Bool { get { membership.memberActive } } public var displayName: String { get { groupProfile.displayName } } @@ -470,6 +485,10 @@ public struct GroupMember: Decodable { } } + public var id: String { + "#\(groupId) @\(groupMemberId)" + } + public var chatViewName: String { get { let p = memberProfile @@ -550,11 +569,24 @@ public struct MemberSubError: Decodable { } public enum ConnectionEntity: Decodable { - case rcvDirectMsgConnection(entityConnection: Connection, contact: Contact?) - case rcvGroupMsgConnection(entityConnection: Connection, groupInfo: GroupInfo, groupMember: GroupMember) - case sndFileConnection(entityConnection: Connection, sndFileTransfer: SndFileTransfer) - case rcvFileConnection(entityConnection: Connection, rcvFileTransfer: RcvFileTransfer) - case userContactConnection(entityConnection: Connection, userContact: UserContact) + case rcvDirectMsgConnection(contact: Contact?) + case rcvGroupMsgConnection(groupInfo: GroupInfo, groupMember: GroupMember) + case sndFileConnection(sndFileTransfer: SndFileTransfer) + case rcvFileConnection(rcvFileTransfer: RcvFileTransfer) + case userContactConnection(userContact: UserContact) + + public var id: String? { + switch self { + case let .rcvDirectMsgConnection(contact): + return contact?.id ?? nil + case let .rcvGroupMsgConnection(_, groupMember): + return groupMember.id + case let .userContactConnection(userContact): + return userContact.id + default: + return nil + } + } } public struct NtfMsgInfo: Decodable { @@ -564,6 +596,13 @@ public struct NtfMsgInfo: Decodable { public struct AChatItem: Decodable { public var chatInfo: ChatInfo public var chatItem: ChatItem + + public var chatId: String { + if case let .groupRcv(groupMember) = chatItem.chatDir { + return groupMember.id + } + return chatInfo.id + } } public struct ChatItem: Identifiable, Decodable { diff --git a/apps/ios/SimpleXChat/Notifications.swift b/apps/ios/SimpleXChat/Notifications.swift index fc68c79b37..2a4a1f6b75 100644 --- a/apps/ios/SimpleXChat/Notifications.swift +++ b/apps/ios/SimpleXChat/Notifications.swift @@ -92,7 +92,7 @@ public func createConnectionEventNtf(_ connEntity: ConnectionEntity) -> UNMutabl var body: String? = nil var targetContentIdentifier: String? = nil switch connEntity { - case let .rcvDirectMsgConnection(_, contact): + case let .rcvDirectMsgConnection(contact): if let contact = contact { title = hideContent ? contactHidden : "\(contact.chatViewName):" targetContentIdentifier = contact.id @@ -100,7 +100,7 @@ public func createConnectionEventNtf(_ connEntity: ConnectionEntity) -> UNMutabl title = NSLocalizedString("New contact:", comment: "notification") } body = NSLocalizedString("message received", comment: "notification") - case let .rcvGroupMsgConnection(_, groupInfo, groupMember): + case let .rcvGroupMsgConnection(groupInfo, groupMember): title = groupMsgNtfTitle(groupInfo, groupMember, hideContent: hideContent) body = NSLocalizedString("message received", comment: "notification") targetContentIdentifier = groupInfo.id